Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
/*
* Copyright 2026 Nil MALHOMME (malhomme.nil+oss@icloud.com)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.github.ktestify.exceptions;

import java.util.Collection;

/**
* Thrown when a single logical operation (e.g. a multi-row Cucumber {@code DataTable} driving one producer/consumer
* call) resolves more than one distinct topic, where exactly one is required.
*
* <p>This is a guard-rail exception: a DataTable listing several instructions is only allowed to target a single topic
* per call. Mixing topics in one DataTable is almost always an authoring mistake — split it into separate step
* invocations instead.
*
* @since 0.4.0
*/
public class TopicMismatchException extends RuntimeException {

public TopicMismatchException(String message) {
super(message);
}

/**
* Creates an exception describing the distinct topics found where a single topic was expected.
*
* @param distinctTopics the distinct namespaced topic names encountered
* @return a new TopicMismatchException with a descriptive message
*/
public static TopicMismatchException forTopics(Collection<String> distinctTopics) {
return new TopicMismatchException(
"A DataTable can only reference a single topic per step, but " + distinctTopics.size()
+ " distinct topics were found: " + distinctTopics
+ ". Split this into separate step invocations, one per topic.");
}
}
26 changes: 23 additions & 3 deletions src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ public final class ConsumerContext<K, V> {
private final Long consumerDeltaTime;
private final boolean isBatchConsumer;
private final int batchSize;
private final Long referenceTimestamp;

private ConsumerContext(
Topic topic,
Expand All @@ -50,7 +51,8 @@ private ConsumerContext(
Long readTimeout,
Long consumerDeltaTime,
boolean isBatchConsumer,
int batchSize) {
int batchSize,
Long referenceTimestamp) {
this.topic = topic;
this.properties = properties;
this.consumer = consumer;
Expand All @@ -62,14 +64,15 @@ private ConsumerContext(
this.consumerDeltaTime = consumerDeltaTime;
this.isBatchConsumer = isBatchConsumer;
this.batchSize = batchSize;
this.referenceTimestamp = referenceTimestamp;
}

/**
* Convenience accessor for single-record consumers. Returns the first element of {@link #matchFilePaths}, or
* {@code null} if the list is empty.
*/
public String getMatchFilePath() {
return matchFilePaths != null && !matchFilePaths.isEmpty() ? matchFilePaths.get(0) : null;
return matchFilePaths != null && !matchFilePaths.isEmpty() ? matchFilePaths.getFirst() : null;
}

public static <K, V> Builder<K, V> builder() {
Expand All @@ -89,6 +92,7 @@ public static final class Builder<K, V> {
private Long consumerDeltaTime;
private boolean isBatchConsumer;
private int batchSize;
private Long referenceTimestamp;

public Builder<K, V> topic(Topic topic) {
this.topic = topic;
Expand Down Expand Up @@ -159,6 +163,21 @@ public Builder<K, V> batchSize(int batchSize) {
return this;
}

/**
* Pins the "now" reference used by {@code calculateDeltaTime()} to a fixed epoch-millisecond value instead of
* letting it be recomputed via {@code System.currentTimeMillis()} at fetch time.
*
* <p>Useful when a single Cucumber step orchestrates multiple internal fetches (e.g. a batch consumer, or
* several chained {@code Then} steps in quick succession) and needs a consistent seek offset across all of
* them, avoiding timestamp drift.
*
* @param referenceTimestamp epoch milliseconds to use as "now", or {@code null} to use the real clock
*/
public Builder<K, V> referenceTimestamp(Long referenceTimestamp) {
this.referenceTimestamp = referenceTimestamp;
return this;
}

public ConsumerContext<K, V> build() {
Topic validatedTopic = Topic.validateTopic(topic, Topic.Type.OUTPUT);

Expand All @@ -183,7 +202,8 @@ public ConsumerContext<K, V> build() {
readTimeout,
consumerDeltaTime,
isBatchConsumer,
batchSize);
batchSize,
referenceTimestamp);
}

private static <T> T requireNonNull(T value, String message) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,12 @@ private void subscribeAndAwaitAssignment() {
/**
* Calculates the earliest timestamp to read from.
*
* <p>Priority order:
* <p>The "now" reference used below is either {@link ConsumerContext#getReferenceTimestamp()}, when the caller has
* pinned it — or the live {@code System.currentTimeMillis()} otherwise. Pinning "now" lets a single Cucumber step
* spawn several internal fetches (e.g. a batch consumer, or multiple {@code Then} steps executed in quick
* succession) without the seek offset drifting forward as wall-clock time advances between them.
*
* <p>Priority order for the delta itself:
*
* <ol>
* <li>Explicit {@code consumerDeltaTime} set on the {@link ConsumerContext} (milliseconds)
Expand All @@ -185,9 +190,12 @@ private void subscribeAndAwaitAssignment() {
* </ol>
*/
private long calculateDeltaTime() {
long now =
context.getReferenceTimestamp() != null ? context.getReferenceTimestamp() : System.currentTimeMillis();

// 1. Explicit value from context (already in ms)
if (context.getConsumerDeltaTime() != null) {
long delta = System.currentTimeMillis() - context.getConsumerDeltaTime();
long delta = now - context.getConsumerDeltaTime();
log.debug("Using consumer delta time from context: {}ms", context.getConsumerDeltaTime());
return delta;
}
Expand All @@ -198,7 +206,7 @@ private long calculateDeltaTime() {
if (deltaTimeStr != null && !deltaTimeStr.isEmpty()) {
log.debug(MESSAGE_CONSUMER_DELTA_TIME_FROM_DATATABLE, deltaTimeStr);
try {
long delta = System.currentTimeMillis() - (Long.parseLong(deltaTimeStr) * 1000);
long delta = now - (Long.parseLong(deltaTimeStr) * 1000);
log.debug(MESSAGE_CONSUMER_DELTA_TIME_IN_TIMESTAMP, delta);
return delta;
} catch (NumberFormatException e) {
Expand All @@ -208,7 +216,7 @@ private long calculateDeltaTime() {

// 3. Framework default
log.debug(MESSAGE_CONSUMER_NO_DELTA_TIME_FOUND, defaultDeltaMs);
return System.currentTimeMillis() - defaultDeltaMs;
return now - defaultDeltaMs;
}

/**
Expand Down
73 changes: 73 additions & 0 deletions src/main/java/io/github/ktestify/utils/TopicUtils.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
/*
* Copyright 2026 Nil MALHOMME (malhomme.nil+oss@icloud.com)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.github.ktestify.utils;

import io.github.ktestify.exceptions.TopicMismatchException;
import io.github.ktestify.models.Topic;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import lombok.experimental.UtilityClass;

/**
* Utilities for validating and comparing {@link Topic} instances.
*
* <p>Transport-agnostic on purpose: any adapter (Kafka today, IBM MQ or others in the future) that drives one physical
* operation from a multi-row {@code DataTable} can reuse {@link #assertSingleTopic(List)} to enforce that all rows
* target the same topic.
*
* @since 0.4.0
*/
@UtilityClass
public final class TopicUtils {

/**
* Asserts that every {@link Topic} in the given list resolves to the same physical topic (namespaced topic name +
* type). Returns that single topic if so.
*
* @param topics the topics resolved from each row of a DataTable, in row order
* @return the single common topic
* @throws IllegalArgumentException if {@code topics} is null or empty
* @throws TopicMismatchException if more than one distinct topic is found
*/
public static Topic assertSingleTopic(List<Topic> topics) {
if (topics == null || topics.isEmpty()) {
throw new IllegalArgumentException("At least one topic must be provided.");
}

Set<String> distinct = new LinkedHashSet<>();
for (Topic topic : topics) {
distinct.add(identity(topic));
}

if (distinct.size() > 1) {
throw TopicMismatchException.forTopics(distinct);
}

return topics.getFirst();
}

/**
* Returns a stable identity string for a topic, combining its namespaced name and type. Two {@link Topic} instances
* (e.g. resolved via alias vs. via name) that point to the same physical topic will produce the same identity even
* if they are not the same object reference.
*/
private static String identity(Topic topic) {
String namespacedTopic = topic != null ? topic.getNamespacedTopic() : null;
Topic.Type type = topic != null ? topic.getTopicType() : null;
return namespacedTopic + "#" + type;
}
}
128 changes: 128 additions & 0 deletions src/test/java/io/github/ktestify/io/kafka/RawKafkaConsumerTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -390,4 +390,132 @@ void batchRecordsAreDeduplicated() throws Exception {
.call());
}
}

// =========================================================================
// referenceTimestamp — pinned "now" fixes clock-drift across delayed fetches
// (see https://github.com/ktestify/ktestify-cucumber/issues/38)
// =========================================================================

@Nested
@DisplayName("referenceTimestamp — pinned 'now' avoids clock drift")
class ReferenceTimestamp {

/** Narrow enough that a few seconds of drift pushes the seek window past the seeded record. */
private static final long NARROW_DELTA_TIME_MS = 2_000L;

/** Simulated step-processing delay between producing the record and fetching it. */
private static final long SIMULATED_STEP_DELAY_MS = 4_000L;

@Test
@DisplayName("without referenceTimestamp, a record is missed once the live-clock delta window drifts past it")
void recordIsMissedDueToClockDriftWithoutReferenceTimestamp() throws Exception {
seedRecord("KEY-1", "{\"orderId\":\"ORD-DRIFT\"}");

// Simulate the delay a slow Cucumber step (or a previous DataTable row) would introduce
// before this consumer actually seeks — this is exactly the drift the maintainer described
// in issue #38.
Thread.sleep(SIMULATED_STEP_DELAY_MS);

ConsumerContext<String, String> ctx = ConsumerContext.<String, String>builder()
.topic(outputTopic())
.consumer(KafkaClientFactory.createRawConsumer(
KtestifyConfig.getOrLoad(), "drift-consumer-" + UUID.randomUUID()))
.readTimeout(3_000L)
.consumerDeltaTime(NARROW_DELTA_TIME_MS)
// No referenceTimestamp — "now" is resolved live, at seek time.
.build();

assertThrows(
ConsumerException.class,
() -> new RawKafkaConsumer(ctx).call(),
"Expected the record to fall outside the live-clock seek window after the simulated delay.");
}

@Test
@DisplayName("with referenceTimestamp pinned before the delay, the record is still found")
void recordIsFoundWhenReferenceTimestampIsPinned() throws Exception {
// Pin "now" BEFORE seeding + the simulated delay, exactly like a Cucumber step
// would capture Instant.now() once at the top of the step.
long pinnedNow = System.currentTimeMillis();

seedRecord("KEY-1", "{\"orderId\":\"ORD-PINNED\"}");

Thread.sleep(SIMULATED_STEP_DELAY_MS);

ConsumerContext<String, String> ctx = ConsumerContext.<String, String>builder()
.topic(outputTopic())
.consumer(KafkaClientFactory.createRawConsumer(
KtestifyConfig.getOrLoad(), "pinned-consumer-" + UUID.randomUUID()))
.readTimeout(3_000L)
.consumerDeltaTime(NARROW_DELTA_TIME_MS)
.referenceTimestamp(pinnedNow)
.build();

boolean result = new RawKafkaConsumer(ctx).call();

assertTrue(
result,
"Expected the record to be found because the seek window was pinned before the delay, "
+ "not recomputed against the live (drifted) clock.");
}

@Test
@DisplayName("two sequential fetches sharing the same referenceTimestamp compute identical seek windows")
void sequentialFetchesShareSameSeekWindow() throws Exception {
long pinnedNow = System.currentTimeMillis();

seedRecord("KEY-1", "{\"orderId\":\"ORD-A\"}");
seedRecord("KEY-2", "{\"orderId\":\"ORD-B\"}");

// First "row" — simulate a small delay before it runs.
Thread.sleep(1_500L);
boolean firstResult = new RawKafkaConsumer(ConsumerContext.<String, String>builder()
.topic(outputTopic())
.consumer(KafkaClientFactory.createRawConsumer(
KtestifyConfig.getOrLoad(), "seq-consumer-1-" + UUID.randomUUID()))
.readTimeout(3_000L)
.consumerDeltaTime(NARROW_DELTA_TIME_MS)
.expectedRecordKey("KEY-1")
.referenceTimestamp(pinnedNow)
.build())
.call();

// Second "row" — additional delay elapses before it runs too.
Thread.sleep(1_500L);
boolean secondResult = new RawKafkaConsumer(ConsumerContext.<String, String>builder()
.topic(outputTopic())
.consumer(KafkaClientFactory.createRawConsumer(
KtestifyConfig.getOrLoad(), "seq-consumer-2-" + UUID.randomUUID()))
.readTimeout(3_000L)
.consumerDeltaTime(NARROW_DELTA_TIME_MS)
.expectedRecordKey("KEY-2")
.referenceTimestamp(pinnedNow)
.build())
.call();

assertTrue(firstResult, "First row should find its record using the pinned reference timestamp.");
assertTrue(
secondResult,
"Second row should still find its record using the SAME pinned reference timestamp, "
+ "despite additional wall-clock time having elapsed between rows.");
}

@Test
@DisplayName("referenceTimestamp does not break normal (non-delayed) consumption")
void recordIsFoundImmediatelyWithoutDelay() throws Exception {
long pinnedNow = System.currentTimeMillis();
seedRecord(null, "{\"orderId\":\"ORD-IMMEDIATE\"}");

ConsumerContext<String, String> ctx = ConsumerContext.<String, String>builder()
.topic(outputTopic())
.consumer(KafkaClientFactory.createRawConsumer(
KtestifyConfig.getOrLoad(), "immediate-consumer-" + UUID.randomUUID()))
.readTimeout(5_000L)
.consumerDeltaTime(60_000L)
.referenceTimestamp(pinnedNow)
.build();

assertTrue(new RawKafkaConsumer(ctx).call());
}
}
}
Loading
Loading