diff --git a/src/main/java/io/github/ktestify/exceptions/TopicMismatchException.java b/src/main/java/io/github/ktestify/exceptions/TopicMismatchException.java
new file mode 100644
index 0000000..188352a
--- /dev/null
+++ b/src/main/java/io/github/ktestify/exceptions/TopicMismatchException.java
@@ -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.
+ *
+ *
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 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.");
+ }
+}
diff --git a/src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java b/src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java
index 441bea2..0dec8a0 100644
--- a/src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java
+++ b/src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java
@@ -38,6 +38,7 @@ public final class ConsumerContext {
private final Long consumerDeltaTime;
private final boolean isBatchConsumer;
private final int batchSize;
+ private final Long referenceTimestamp;
private ConsumerContext(
Topic topic,
@@ -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;
@@ -62,6 +64,7 @@ private ConsumerContext(
this.consumerDeltaTime = consumerDeltaTime;
this.isBatchConsumer = isBatchConsumer;
this.batchSize = batchSize;
+ this.referenceTimestamp = referenceTimestamp;
}
/**
@@ -69,7 +72,7 @@ private ConsumerContext(
* {@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 Builder builder() {
@@ -89,6 +92,7 @@ public static final class Builder {
private Long consumerDeltaTime;
private boolean isBatchConsumer;
private int batchSize;
+ private Long referenceTimestamp;
public Builder topic(Topic topic) {
this.topic = topic;
@@ -159,6 +163,21 @@ public Builder 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.
+ *
+ * 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 referenceTimestamp(Long referenceTimestamp) {
+ this.referenceTimestamp = referenceTimestamp;
+ return this;
+ }
+
public ConsumerContext build() {
Topic validatedTopic = Topic.validateTopic(topic, Topic.Type.OUTPUT);
@@ -183,7 +202,8 @@ public ConsumerContext build() {
readTimeout,
consumerDeltaTime,
isBatchConsumer,
- batchSize);
+ batchSize,
+ referenceTimestamp);
}
private static T requireNonNull(T value, String message) {
diff --git a/src/main/java/io/github/ktestify/io/kafka/KafkaRecordFetcher.java b/src/main/java/io/github/ktestify/io/kafka/KafkaRecordFetcher.java
index cb7d465..045209e 100644
--- a/src/main/java/io/github/ktestify/io/kafka/KafkaRecordFetcher.java
+++ b/src/main/java/io/github/ktestify/io/kafka/KafkaRecordFetcher.java
@@ -176,7 +176,12 @@ private void subscribeAndAwaitAssignment() {
/**
* Calculates the earliest timestamp to read from.
*
- * Priority order:
+ *
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.
+ *
+ *
Priority order for the delta itself:
*
*
* - Explicit {@code consumerDeltaTime} set on the {@link ConsumerContext} (milliseconds)
@@ -185,9 +190,12 @@ private void subscribeAndAwaitAssignment() {
*
*/
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;
}
@@ -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) {
@@ -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;
}
/**
diff --git a/src/main/java/io/github/ktestify/utils/TopicUtils.java b/src/main/java/io/github/ktestify/utils/TopicUtils.java
new file mode 100644
index 0000000..56f5ac1
--- /dev/null
+++ b/src/main/java/io/github/ktestify/utils/TopicUtils.java
@@ -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.
+ *
+ * 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 topics) {
+ if (topics == null || topics.isEmpty()) {
+ throw new IllegalArgumentException("At least one topic must be provided.");
+ }
+
+ Set 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;
+ }
+}
diff --git a/src/test/java/io/github/ktestify/io/kafka/RawKafkaConsumerTest.java b/src/test/java/io/github/ktestify/io/kafka/RawKafkaConsumerTest.java
index de085dc..95bc220 100644
--- a/src/test/java/io/github/ktestify/io/kafka/RawKafkaConsumerTest.java
+++ b/src/test/java/io/github/ktestify/io/kafka/RawKafkaConsumerTest.java
@@ -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 ctx = ConsumerContext.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 ctx = ConsumerContext.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.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.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 ctx = ConsumerContext.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());
+ }
+ }
}
diff --git a/src/test/java/io/github/ktestify/utils/TopicUtilsTest.java b/src/test/java/io/github/ktestify/utils/TopicUtilsTest.java
new file mode 100644
index 0000000..e2bce18
--- /dev/null
+++ b/src/test/java/io/github/ktestify/utils/TopicUtilsTest.java
@@ -0,0 +1,70 @@
+/*
+ * 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 static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import io.github.ktestify.exceptions.TopicMismatchException;
+import io.github.ktestify.models.Topic;
+import java.util.List;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+@DisplayName("TopicUtils")
+class TopicUtilsTest {
+
+ private Topic topic(String name, String alias, Topic.Type type) {
+ return Topic.builder().topicName(name).topicAlias(alias).topicType(type).build();
+ }
+
+ @Test
+ @DisplayName("returns the common topic when all rows resolve to the same topic")
+ void returnsCommonTopicWhenAllSame() {
+ Topic byName = topic("orders", null, Topic.Type.INPUT);
+ Topic byAlias = topic("orders", "orders-alias", Topic.Type.INPUT);
+
+ Topic result = TopicUtils.assertSingleTopic(List.of(byName, byAlias));
+
+ assertEquals("orders", result.getTopicName());
+ }
+
+ @Test
+ @DisplayName("throws TopicMismatchException when rows resolve to different topics")
+ void throwsWhenTopicsDiffer() {
+ Topic ordersTopic = topic("orders", null, Topic.Type.INPUT);
+ Topic paymentsTopic = topic("payments", null, Topic.Type.INPUT);
+
+ assertThrows(
+ TopicMismatchException.class, () -> TopicUtils.assertSingleTopic(List.of(ordersTopic, paymentsTopic)));
+ }
+
+ @Test
+ @DisplayName("throws TopicMismatchException when the same topic name has different types")
+ void throwsWhenTopicTypesDiffer() {
+ Topic input = topic("orders", null, Topic.Type.INPUT);
+ Topic output = topic("orders", null, Topic.Type.OUTPUT);
+
+ assertThrows(TopicMismatchException.class, () -> TopicUtils.assertSingleTopic(List.of(input, output)));
+ }
+
+ @Test
+ @DisplayName("throws IllegalArgumentException when the list is null or empty")
+ void throwsOnNullOrEmpty() {
+ assertThrows(IllegalArgumentException.class, () -> TopicUtils.assertSingleTopic(null));
+ assertThrows(IllegalArgumentException.class, () -> TopicUtils.assertSingleTopic(List.of()));
+ }
+}