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 @@ -17,6 +17,8 @@
*/
package io.aiven.inkless.metadata;

import java.util.Set;
import java.util.function.Supplier;
import java.util.regex.Matcher;
import java.util.regex.Pattern;

Expand All @@ -25,6 +27,18 @@ class ClientAZExtractor {

// Visible for testing
static String getClientAZ(final String clientId) {
return getClientAZ(clientId, Set::of);
}

/**
* Extract the client AZ from the client ID, validating against known rack values if provided.
*
* @param clientId the client ID string, may be null
* @param racksSupplier supplier of valid rack/AZ identifiers, evaluated when needed.
* An empty set disables validation.
* @return the extracted AZ, or null if not found or not matchable
*/
static String getClientAZ(final String clientId, final Supplier<Set<String>> racksSupplier) {
if (clientId == null) {
return null;
}
Expand All @@ -33,6 +47,30 @@ static String getClientAZ(final String clientId) {
if (!matcher.find()) {
return null;
}
return matcher.group("az");

final String rawAZ = matcher.group("az");
if (rawAZ == null || rawAZ.isEmpty()) {
return rawAZ;
}

final Set<String> racks = racksSupplier.get();
if (racks == null || racks.isEmpty()) {
return rawAZ;
}

if (racks.contains(rawAZ)) {
return rawAZ;
}

String azCandidate = rawAZ;
int lastHyphen;
while ((lastHyphen = azCandidate.lastIndexOf('-')) > 0) {
azCandidate = azCandidate.substring(0, lastHyphen);
if (racks.contains(azCandidate)) {
return azCandidate;
}
}

return null;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
Expand Down Expand Up @@ -422,7 +423,8 @@ private static String normalizeAZ(final String az) {
* Resolve the client AZ using based on client ID or listener name.
*/
private String resolveClientAZ(final ListenerName listenerName, final String clientId) {
final String explicitAZ = normalizeAZ(ClientAZExtractor.getClientAZ(clientId));
final String explicitAZ = normalizeAZ(
ClientAZExtractor.getClientAZ(clientId, () -> computeKnownRacks(listenerName)));
if (explicitAZ != null) {
return explicitAZ;
}
Expand All @@ -436,6 +438,18 @@ private String resolveClientAZ(final ListenerName listenerName, final String cli
return null;
}

private Set<String> computeKnownRacks(final ListenerName listenerName) {
final Set<String> racks = new HashSet<>(clientAzListenerMap.values());
if (listenerName != null) {
metadataView.getAliveBrokerNodes(listenerName)
.stream()
.map(Node::rack)
.filter(r -> r != null && !r.isBlank())
.forEach(racks::add);
}
return racks;
}

/**
* Record cross-AZ routing only when we can confirm the leader is in a different AZ.
* If client AZ is unknown or the leader's rack is unset, we cannot confirm cross-AZ.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,15 @@

import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;

import java.util.Set;
import java.util.stream.Stream;

import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.params.provider.Arguments.arguments;

class ClientAZExtractorTest {

Expand Down Expand Up @@ -64,4 +70,66 @@ void getClientAZWhenProvidedEmptyValue(final String clientId) {
void getClientAZWhenOnlySeeminglyCorrect(final String clientId) {
assertThat(ClientAZExtractor.getClientAZ(clientId)).isNull();
}

@Test
void withoutKnownRacksReturnRawValueIncludingSuffix() {
assertThat(ClientAZExtractor.getClientAZ(
"kafka-connect,diskless_az=use1-az2-575969bf-aae2-47c5-98d6-8750ef0536f8"
)).isEqualTo("use1-az2-575969bf-aae2-47c5-98d6-8750ef0536f8");
}

@Test
void withoutKnownRacksReturnRawValueWithBackingStoreSuffix() {
assertThat(ClientAZExtractor.getClientAZ(
"connect-cluster-kafka-connect,diskless_az=use1-az2-configs"
)).isEqualTo("use1-az2-configs");
}

@ParameterizedTest
@MethodSource("knownRacksCases")
void getClientAZWithKnownRacks(final String clientId, final Set<String> knownRacks, final String expected) {
assertThat(ClientAZExtractor.getClientAZ(clientId, () -> knownRacks)).isEqualTo(expected);
}

@ParameterizedTest
@MethodSource("noValidationCases")
void getClientAZFallsBackToRawValueWithoutValidation(final Set<String> knownRacks) {
assertThat(ClientAZExtractor.getClientAZ("my-app,diskless_az=use1-az2-suffix", () -> knownRacks))
.isEqualTo("use1-az2-suffix");
}

static Stream<Arguments> knownRacksCases() {
final Set<String> euWest = Set.of("eu-west-1a", "eu-west-1b", "eu-west-1c");
final Set<String> use1 = Set.of("use1-az2", "use1-az4", "use1-az6");
return Stream.of(
arguments("my-app,diskless_az=eu-west-1a", euWest, "eu-west-1a"),
arguments("diskless_az=az1", Set.of("az0", "az1", "az2"), "az1"),
arguments("kafka-connect,diskless_az=use1-az2-575969bf-aae2-47c5-98d6-8750ef0536f8", use1, "use1-az2"),
arguments("connect-cluster-kafka-connect,diskless_az=use1-az2-configs", use1, "use1-az2"),
arguments("connect-cluster-kafka-connect,diskless_az=eu-west-1a-offsets", euWest, "eu-west-1a"),
arguments("connect-cluster-kafka-connect,diskless_az=eu-west-1a-statuses", euWest, "eu-west-1a"),
arguments("my-app,diskless_az=us-east-1a-extra-suffix", Set.of("us-east-1a", "us-east-1b"), "us-east-1a"),
arguments("kafka-connect,diskless_az=eu-west-1a-deadbeef-1234-5678-abcd-ef0123456789", euWest, "eu-west-1a"),
arguments("kafka-connect,diskless_az=az1-575969bf-aae2-47c5-98d6-8750ef0536f8", Set.of("az1", "az2", "az3"), "az1"),
// Progressive stripping removes from the right, so the longest prefix that matches wins.
arguments("my-app,diskless_az=us-east-1a-suffix", Set.of("us", "us-east", "us-east-1a"), "us-east-1a"),
arguments("kafka-connect,diskless_az=use1-az2-575969bf-aae2-47c5-98d6-8750ef0536f8", Set.of("use1-az2"), "use1-az2"),
arguments("my-app,diskless_az=", euWest, ""),
arguments("my-app,diskless_az=completely-wrong-value", euWest, null),
arguments("my-app,diskless_az=unknown", euWest, null),
arguments("my-app,diskless_az=eu-west-kaboom", euWest, null),
arguments(null, euWest, null),
arguments("connector-producer-orders-0", euWest, null),
arguments("connector-producer-orders-source-0", use1, null),
arguments("connector-consumer-my-sink-0", use1, null),
arguments("source->target|MirrorSourceConnector-0|replication-consumer", euWest, null)
);
}

static Stream<Arguments> noValidationCases() {
return Stream.of(
arguments((Set<String>) null),
arguments(Set.of())
);
}
}
Loading