aliehsaeedii commented on code in PR #22165:
URL: https://github.com/apache/kafka/pull/22165#discussion_r3559281841
##########
streams/src/test/java/org/apache/kafka/streams/state/internals/TimeOrderedKeyValueBufferTest.java:
##########
@@ -306,10 +306,74 @@ public void shouldReturnPriorValueForBufferedKey(final
String testName, final Fu
context.setRecordContext(recordContext);
buffer.put(1L, new Record<>("A", new Change<>("new-value",
"old-value"), 0L), recordContext);
buffer.put(1L, new Record<>("B", new Change<>("new-value", null), 0L),
recordContext);
- assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", -1, new
RecordHeaders()))));
+ assertThat(buffer.priorValueForBuffered("A"),
is(Maybe.defined(ValueTimestampHeaders.make("old-value", 0L, new
RecordHeaders()))));
assertThat(buffer.priorValueForBuffered("B"), is(Maybe.defined(null)));
}
+ @ParameterizedTest
+ @MethodSource("parameters")
+ public void shouldPropagateHeadersThroughEviction(final String testName,
final Function<String, B> bufferSupplier) {
+ setup(testName, bufferSupplier);
+ final TimeOrderedKeyValueBuffer<String, String, Change<String>> buffer
= bufferSupplier.apply(testName);
+ final MockInternalProcessorContext<?, ?> context = makeContext();
+ buffer.init(context, buffer);
+
+ final RecordHeaders headers = new RecordHeaders(new Header[]{new
RecordHeader("h1", "v1".getBytes(UTF_8))});
Review Comment:
All other calls are using `utf8` too, so I keep it as-is:)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]