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: * *

    *
  1. 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())); + } +}