Skip to content
Open
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
Expand Up @@ -22,11 +22,8 @@
import io.github.ktestify.match.MatchResult;
import io.github.ktestify.match.RecordMatcher;
import io.github.ktestify.models.ConsumedRecord;
import io.github.ktestify.models.Topic;
import java.util.List;
import java.util.Map;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;

/**
* Thin coordinator that wires a {@link KafkaRecordFetcher} (transport) with a {@link RecordMatcher} (assertion) and
Expand Down Expand Up @@ -66,28 +63,6 @@ protected AbstractKafkaConsumer(ConsumerContext<K, V> context, RecordMatcher<V>
matcher.getClass().getSimpleName());
}

/**
* Legacy convenience constructor for callers that previously passed topic + consumer + properties.
*
* @param topic the topic to consume from
* @param consumer the Kafka consumer instance
* @param properties the consumer properties map
* @param matcher the assertion strategy
* @deprecated Build a {@link ConsumerContext} and use {@link #AbstractKafkaConsumer(ConsumerContext,
* RecordMatcher)} instead.
*/
@Deprecated
protected AbstractKafkaConsumer(
Topic topic, Consumer<K, V> consumer, Map<String, String> properties, RecordMatcher<V> matcher) {
this(
ConsumerContext.<K, V>builder()
.topic(topic)
.consumer(consumer)
.properties(properties)
.build(),
matcher);
}

/**
* Fetches records from Kafka, then asserts them with the configured matcher.
*
Expand Down Expand Up @@ -137,6 +112,7 @@ protected MatchContext buildMatchContext() {
.matchFilePaths(context.getMatchFilePaths())
.excludedFields(context.getExcludedFields())
.strictMatching(false)
.keyMatchStrategy(context.getKeyMatchStrategy())
.build();
}
}
21 changes: 21 additions & 0 deletions src/main/java/io/github/ktestify/io/kafka/ConsumerContext.java
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

import io.github.ktestify.config.KtestifyConfig;
import io.github.ktestify.exceptions.ConsumerException;
import io.github.ktestify.match.KeyMatchStrategy;
import io.github.ktestify.models.Topic;
import java.util.Collections;
import java.util.List;
Expand All @@ -31,6 +32,7 @@ public final class ConsumerContext<K, V> {
private final Map<String, String> properties;
private final Consumer<K, V> consumer;
private final String expectedRecordKey;
private final KeyMatchStrategy keyMatchStrategy;
private final String matchMethod;
private final List<String> matchFilePaths;
private final List<String> excludedFields;
Expand All @@ -45,6 +47,7 @@ private ConsumerContext(
Map<String, String> properties,
Consumer<K, V> consumer,
String expectedRecordKey,
KeyMatchStrategy keyMatchStrategy,
String matchMethod,
List<String> matchFilePaths,
List<String> excludedFields,
Expand All @@ -57,6 +60,7 @@ private ConsumerContext(
this.properties = properties;
this.consumer = consumer;
this.expectedRecordKey = expectedRecordKey;
this.keyMatchStrategy = keyMatchStrategy != null ? keyMatchStrategy : KeyMatchStrategy.EXACT;
this.matchMethod = matchMethod;
this.matchFilePaths = matchFilePaths != null ? matchFilePaths : Collections.emptyList();
this.excludedFields = excludedFields != null ? excludedFields : Collections.emptyList();
Expand Down Expand Up @@ -85,6 +89,7 @@ public static final class Builder<K, V> {
private Map<String, String> properties;
private Consumer<K, V> consumer;
private String expectedRecordKey;
private KeyMatchStrategy keyMatchStrategy;
private String matchMethod;
private List<String> matchFilePaths;
private List<String> excludedFields;
Expand Down Expand Up @@ -114,6 +119,21 @@ public Builder<K, V> expectedRecordKey(String expectedRecordKey) {
return this;
}

/**
* Sets the strategy used to compare {@link #expectedRecordKey} against the actual record key during the
* fetch-time pre-filter in {@code KafkaRecordFetcher.passesKeyFilter()}.
*
* <p>Defaults to {@link KeyMatchStrategy#EXACT} when not set, preserving backward compatibility.
*
* @param keyMatchStrategy the match strategy, or {@code null} to use the default
* @return this builder
* @since 1.1.5
*/
public Builder<K, V> keyMatchStrategy(KeyMatchStrategy keyMatchStrategy) {
this.keyMatchStrategy = keyMatchStrategy;
return this;
}

public Builder<K, V> matchMethod(String matchMethod) {
this.matchMethod = matchMethod;
return this;
Expand Down Expand Up @@ -196,6 +216,7 @@ public ConsumerContext<K, V> build() {
validatedProps,
validatedConsumer,
expectedRecordKey,
keyMatchStrategy,
matchMethod,
matchFilePaths,
excludedFields,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,7 @@
import io.github.ktestify.models.ConsumedRecord;
import io.github.ktestify.models.MatchedRecord;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import lombok.extern.slf4j.Slf4j;
Expand Down Expand Up @@ -337,7 +333,7 @@ private void registerAsMatched(ConsumedRecord<V> record) {

/**
* Returns {@code true} if no key-filter is configured, or if the record key matches the expected key from the
* context / properties.
* context / properties using the configured {@link io.github.ktestify.match.KeyMatchStrategy}.
*/
private boolean passesKeyFilter(ConsumerRecord<K, V> record) {
// Context takes priority over properties map
Expand All @@ -351,7 +347,7 @@ private boolean passesKeyFilter(ConsumerRecord<K, V> record) {
}

String recordKey = record.key() != null ? record.key().toString() : null;
if (expectedKey.equals(recordKey)) {
if (context.getKeyMatchStrategy().matches(expectedKey, recordKey)) {
log.info(MESSAGE_CONSUMER_RECORD_MATCHES_EXPECTED_KEY, expectedKey);
return true;
}
Expand Down
140 changes: 140 additions & 0 deletions src/main/java/io/github/ktestify/match/KeyMatchStrategy.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
/*
* 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.match;

/**
* Strategies for comparing a record key against an expected key.
*
* <p>Used in two places:
*
* <ul>
* <li>{@code KafkaRecordFetcher.passesKeyFilter()} as a pre-filter during the Kafka poll loop
* <li>{@code KeyRecordMatcher}, {@code FileKeyRecordMatcher}, {@code AvroKeyRecordMatcher}, and
* {@code AvroFileKeyRecordMatcher} as the post-fetch assertion
* </ul>
*
* <p>The default strategy is {@link #EXACT}, which preserves the original {@code String.equals()} behavior. Other
* strategies allow matching dynamically generated keys by prefix, suffix, substring, or regular expression.
*
* @since 1.1.4
*/
public enum KeyMatchStrategy {

/**
* Exact equality: {@code expected.equals(actual)}.
*
* <p>This is the default and preserves backward compatibility for feature files that do not specify a
* {@code keyMatchStrategy} column.
*
* @since 1.1.4
*/
EXACT {
@Override
public boolean matches(String expected, String actual) {
return expected != null && expected.equals(actual);
}
},

/**
* Substring match: {@code actual.contains(expected)}.
*
* <p>Useful when the record key contains a known fragment embedded in a larger dynamically generated value.
*
* @since 1.1.4
*/
CONTAINS {
@Override
public boolean matches(String expected, String actual) {
return expected != null && actual != null && actual.contains(expected);
}
},

/**
* Prefix match: {@code actual.startsWith(expected)}.
*
* <p>Useful when the record key starts with a known prefix followed by a dynamically generated suffix (e.g.
* {@code ORD-<uuid>}).
*
* @since 1.1.4
*/
STARTS_WITH {
@Override
public boolean matches(String expected, String actual) {
return expected != null && actual != null && actual.startsWith(expected);
}
},

/**
* Suffix match: {@code actual.endsWith(expected)}.
*
* <p>Useful when the record key ends with a known suffix preceded by a dynamically generated prefix.
*
* @since 1.1.4
*/
ENDS_WITH {
@Override
public boolean matches(String expected, String actual) {
return expected != null && actual != null && actual.endsWith(expected);
}
},

/**
* Regular expression match: {@code actual.matches(expected)}.
*
* <p>The {@code expected} string is interpreted as a Java regular expression. Useful for arbitrary patterns such as
* {@code ORD-\d{6}} that cannot be expressed with prefix, suffix, or substring matching.
*
* @since 1.1.4
*/
REGEX {
@Override
public boolean matches(String expected, String actual) {
return expected != null && actual != null && actual.matches(expected);
}
};

/**
* Tests whether the {@code actual} record key satisfies this strategy given the {@code expected} key.
*
* @param expected the expected key value (or pattern for {@link #REGEX})
* @param actual the actual record key, may be {@code null} when the Kafka record has no key
* @return {@code true} if the actual key matches according to this strategy
* @since 1.1.4
*/
public abstract boolean matches(String expected, String actual);

/**
* Parses a strategy name from a DataTable column value.
*
* <p>Matching is case-insensitive and tolerant of hyphens, underscores, and spaces. For example,
* {@code "starts_with"}, {@code "starts-with"}, and {@code "STARTS WITH"} all resolve to {@link #STARTS_WITH}.
*
* @param value the raw column value, may be {@code null} or blank
* @return the parsed strategy, or {@link #EXACT} when the value is {@code null}, blank, or unrecognized
* @since 1.1.4
*/
public static KeyMatchStrategy fromString(String value) {
if (value == null || value.isBlank()) {
return EXACT;
}
String normalized = value.trim().toUpperCase().replace('-', '_').replace(' ', '_');
try {
return KeyMatchStrategy.valueOf(normalized);
} catch (IllegalArgumentException e) {
return EXACT;
}
}
}
12 changes: 12 additions & 0 deletions src/main/java/io/github/ktestify/match/MatchContext.java
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,18 @@ public class MatchContext {
/** Expected value for {@link #matchKey}. */
String matchValue;

/**
* Strategy used to compare {@link #matchKey} against the actual record key in key-related matchers
* ({@code KeyRecordMatcher}, {@code FileKeyRecordMatcher}, {@code AvroKeyRecordMatcher},
* {@code AvroFileKeyRecordMatcher}).
*
* <p>Defaults to {@link KeyMatchStrategy#EXACT}, preserving the original exact-equality behavior.
*
* @since 1.1.5
*/
@Builder.Default
KeyMatchStrategy keyMatchStrategy = KeyMatchStrategy.EXACT;

/**
* Multiple key/value pairs for multi-field inline matching.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,13 +54,13 @@ public MatchResult match(List<ConsumedRecord<GenericRecord>> records, MatchConte
throw new ComparisonException("AvroFileKeyRecordMatcher requires matchFilePath to be set.");
}

ConsumedRecord<GenericRecord> record = records.get(0);
ConsumedRecord<GenericRecord> record = records.getFirst();
String actualKey = record.getKey();
String expectedKey = context.getMatchKey();
String expectedValue = FileUtils.getFileContent(FileUtils.getFile(context.getMatchFilePath()));
String actualValue = toJson(record.getValue());

boolean keyMatches = expectedKey.equals(actualKey);
boolean keyMatches = context.getKeyMatchStrategy().matches(expectedKey, actualKey);
boolean valueMatches =
AvroUtils.doesAvroRecordsSmartMatches(AvroUtils.getPrettyAvroValue(expectedValue), actualValue);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,16 +41,24 @@ public MatchResult match(List<ConsumedRecord<GenericRecord>> records, MatchConte
}

String expectedKey = context.getMatchKey();
String actualKey = records.get(0).getKey();
String actualKey = records.getFirst().getKey();

if (expectedKey.equals(actualKey)) {
log.info("Avro record key matches expected key '{}'.", expectedKey);
if (context.getKeyMatchStrategy().matches(expectedKey, actualKey)) {
log.info(
"Avro record key matches expected key '{}' using {} strategy.",
expectedKey,
context.getKeyMatchStrategy());
return MatchResult.pass(expectedKey, actualKey);
}

log.error("Avro record key mismatch β€” expected: '{}', actual: '{}'", expectedKey, actualKey);
log.error(
"Avro record key mismatch, expected: '{}', actual: '{}', strategy: {}",
expectedKey,
actualKey,
context.getKeyMatchStrategy());
return MatchResult.fail(
"Avro record key does not match β€” expected: '" + expectedKey + "', actual: '" + actualKey + "'.",
"Avro record key does not match, expected: '" + expectedKey + "', actual: '" + actualKey
+ "', strategy: " + context.getKeyMatchStrategy() + ".",
expectedKey,
actualKey);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,17 +50,17 @@ public MatchResult match(List<ConsumedRecord<String>> records, MatchContext cont
throw new ComparisonException("FileKeyRecordMatcher requires matchFilePath to be set.");
}

ConsumedRecord<String> record = records.get(0);
ConsumedRecord<String> record = records.getFirst();
String expectedValue = FileUtils.getFileContent(FileUtils.getFile(context.getMatchFilePath()));
String actualValue = record.getValue();
String expectedKey = context.getMatchKey();
String actualKey = record.getKey();

boolean keyMatches = expectedKey.equals(actualKey);
boolean keyMatches = context.getKeyMatchStrategy().matches(expectedKey, actualKey);
boolean valueMatches = actualValue.equals(expectedValue);

if (!keyMatches) {
log.error("Key mismatch β€” expected: '{}', actual: '{}'", expectedKey, actualKey);
log.error("Key mismatch, expected: '{}', actual: '{}'", expectedKey, actualKey);
}
if (!valueMatches) {
log.error(
Expand Down
Loading
Loading