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
Expand Up @@ -118,18 +118,21 @@ public MskCluster createCluster(String clusterName, String kafkaVersion) {
public MskCluster createCluster(CreateClusterRequest request) {
validateCreateRequest(request);
String clusterName = request.getClusterName();
if (storage.scan(k -> true).stream().anyMatch(c -> c.getClusterName().equals(clusterName))) {
if (storage.scan(k -> true).stream().anyMatch(c -> isCurrentRegion(c)
&& c.getClusterName().equals(clusterName))) {
throw new AwsException("ConflictException", "Cluster already exists: " + clusterName, 409);
}

String accountId = regionResolver.getAccountId();
String clusterArn = AwsArnUtils.Arn.of("kafka", config.defaultRegion(), accountId, "cluster/" + clusterName + "/" + java.util.UUID.randomUUID()).toString();
String clusterArn = AwsArnUtils.Arn.of("kafka", regionResolver.getRegion(), accountId,
"cluster/" + clusterName + "/" + java.util.UUID.randomUUID()).toString();

String kafkaVersion = request.getKafkaVersion();
String resolvedKafkaVersion = (kafkaVersion == null || kafkaVersion.isBlank()) ? DEFAULT_KAFKA_VERSION : kafkaVersion;
MskCluster cluster = new MskCluster(clusterArn, clusterName, resolvedKafkaVersion);
cluster.setClusterType(PROVISIONED_CLUSTER_TYPE);
cluster.setAccountId(accountId);
cluster.setResourceRegion(regionResolver.getRegion());
cluster.setVolumeId(String.format("%06x", new SecureRandom().nextInt(0xFFFFFF)));

if (request.getNumberOfBrokerNodes() != null) {
Expand Down Expand Up @@ -326,19 +329,21 @@ private MskCluster createServerlessCluster(CreateClusterV2Request request) {
throw badRequest("clusterName",
"clusterName must be between 1 and " + MAX_CLUSTER_NAME_LENGTH + " characters.");
}
if (storage.scan(k -> true).stream().anyMatch(c -> c.getClusterName().equals(clusterName))) {
if (storage.scan(k -> true).stream().anyMatch(c -> isCurrentRegion(c)
&& c.getClusterName().equals(clusterName))) {
throw new AwsException("ConflictException", "Cluster already exists: " + clusterName, 409);
}

String accountId = regionResolver.getAccountId();
String clusterArn = AwsArnUtils.Arn.of("kafka", config.defaultRegion(), accountId,
String clusterArn = AwsArnUtils.Arn.of("kafka", regionResolver.getRegion(), accountId,
"cluster/" + clusterName + "/" + UUID.randomUUID()).toString();

MskCluster cluster = new MskCluster(clusterArn, clusterName, DEFAULT_KAFKA_VERSION);
cluster.setClusterType(SERVERLESS_CLUSTER_TYPE);
cluster.setServerless(request.getServerless());
cluster.setTags(request.getTags());
cluster.setAccountId(accountId);
cluster.setResourceRegion(regionResolver.getRegion());
cluster.setVolumeId(String.format("%06x", new SecureRandom().nextInt(0xFFFFFF)));

// Provisioned-only members must not surface on a serverless cluster.
Expand Down Expand Up @@ -377,21 +382,21 @@ public MskCluster describeClusterV1(String clusterArn) {

/** ListClusters for the v1 API, which likewise cannot represent serverless clusters. */
public List<MskCluster> listProvisionedClusters() {
return storage.scan(k -> true).stream().filter(c -> !isServerless(c)).toList();
return storage.scan(k -> true).stream().filter(this::isCurrentRegion)
.filter(c -> !isServerless(c)).toList();
}

public MskCluster describeCluster(String clusterArn) {
return storage.get(clusterArn)
return storage.get(clusterArn).filter(this::isCurrentRegion)
.orElseThrow(() -> new AwsException("NotFoundException", "Cluster not found: " + clusterArn, 404));
}

public List<MskCluster> listClusters() {
return storage.scan(k -> true);
return storage.scan(k -> true).stream().filter(this::isCurrentRegion).toList();
}

public void deleteCluster(String clusterArn) {
MskCluster cluster = storage.get(clusterArn)
.orElseThrow(() -> new AwsException("NotFoundException", "Cluster not found: " + clusterArn, 404));
MskCluster cluster = describeCluster(clusterArn);

cluster.setState(ClusterState.DELETING);
if (!config.services().msk().mock()) {
Expand Down Expand Up @@ -541,7 +546,7 @@ public ConfigurationRevisionDetail describeConfigurationRevision(String arn, lon
// DescribeConfiguration, whose 400 is a deliberate terraform-provider contract.

public Map<String, String> listTagsForResource(String arn) {
MskCluster cluster = storage.get(arn).orElse(null);
MskCluster cluster = storage.get(arn).filter(this::isCurrentRegion).orElse(null);
if (cluster != null) {
return cluster.getTags() != null ? cluster.getTags() : Map.of();
}
Expand All @@ -553,7 +558,7 @@ public Map<String, String> listTagsForResource(String arn) {
}

public void tagResource(String arn, Map<String, String> tags) {
MskCluster cluster = storage.get(arn).orElse(null);
MskCluster cluster = storage.get(arn).filter(this::isCurrentRegion).orElse(null);
if (cluster != null) {
cluster.setTags(merged(cluster.getTags(), tags));
storage.put(arn, cluster);
Expand All @@ -569,7 +574,7 @@ public void tagResource(String arn, Map<String, String> tags) {
}

public void untagResource(String arn, List<String> tagKeys) {
MskCluster cluster = storage.get(arn).orElse(null);
MskCluster cluster = storage.get(arn).filter(this::isCurrentRegion).orElse(null);
if (cluster != null) {
cluster.setTags(without(cluster.getTags(), tagKeys));
storage.put(arn, cluster);
Expand Down Expand Up @@ -643,7 +648,7 @@ private void putCluster(MskCluster cluster) {
@Override
public List<ExplorerResource> getResources() {
List<ExplorerResource> resources = new ArrayList<>();
for (MskCluster cluster : storage.scan(k -> true)) {
for (MskCluster cluster : storage.scan(k -> true).stream().filter(this::isCurrentRegion).toList()) {
String arn = cluster.getClusterArn();
if (arn == null) {
continue;
Expand All @@ -662,4 +667,14 @@ public List<ExplorerResource> getResources() {
public Set<SupportedResourceType> getSupportedResourceTypes() {
return Set.of(new SupportedResourceType("kafka:cluster", "kafka", true));
}

private boolean isCurrentRegion(MskCluster cluster) {
if (cluster.getResourceRegion() != null) {
return regionResolver.getRegion().equals(cluster.getResourceRegion());
}
// Older records have no request-region marker and their ARN always used the configured
// default. Keep them visible until they are recreated so an upgrade does not hide or
// strand persisted clusters whose original request region was not stored.
return true;
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package io.github.hectorvent.floci.services.msk;

import io.github.hectorvent.floci.config.EmulatorConfig;
import io.github.hectorvent.floci.core.common.AwsArnUtils;
import io.github.hectorvent.floci.core.common.RegionResolver;
import io.github.hectorvent.floci.core.common.docker.ContainerBuilder;
import io.github.hectorvent.floci.core.common.docker.ContainerDetector;
Expand All @@ -19,6 +20,7 @@
import java.io.Closeable;
import java.net.HttpURLConnection;
import java.net.URI;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
Expand Down Expand Up @@ -149,8 +151,9 @@ public void startContainer(MskCluster cluster) {
.withDockerNetwork(config.services().dockerNetwork())
.withLogRotation()
.withLabels(ContainerStorageHelper.resourceIdentityLabels(
"msk", cluster.getClusterName(), regionResolver.getAccountId(),
regionResolver.getDefaultRegion()));
"msk", cluster.getClusterName(),
AwsArnUtils.accountOrDefault(cluster.getClusterArn(), regionResolver.getAccountId()),
AwsArnUtils.regionOrDefault(cluster.getClusterArn(), regionResolver.getDefaultRegion())));

if (!containerDetector.isRunningInContainer()) {
specBuilder.withPortBinding(KAFKA_PORT, kafkaHostPort).withDynamicPort(ADMIN_PORT);
Expand All @@ -165,8 +168,7 @@ public void startContainer(MskCluster cluster) {
"/var/lib/redpanda/data");
} else {
// Legacy host-path mode: host-persistent-path is an absolute path
String hostDataPath = ContainerStorageHelper.hostResourcePath(config, "msk", cluster.getClusterName())
.toAbsolutePath().toString();
String hostDataPath = legacyCompatibleHostPath(cluster).toAbsolutePath().toString();
if (!containerDetector.isRunningInContainer()) {
ContainerStorageHelper.ensureHostDir(hostDataPath);
}
Expand Down Expand Up @@ -200,12 +202,13 @@ public void startContainer(MskCluster cluster) {
: info.containerId();
String logGroup = "/aws/msk/cluster/" + cluster.getClusterName();
String logStream = logStreamer.generateLogStreamName(shortId);
String region = regionResolver.getDefaultRegion();
String region = AwsArnUtils.regionOrDefault(cluster.getClusterArn(), regionResolver.getDefaultRegion());

Closeable logHandle = logStreamer.attach(
info.containerId(), logGroup, logStream, region, "msk:" + cluster.getClusterName());
String account = AwsArnUtils.accountOrDefault(cluster.getClusterArn(), regionResolver.getDefaultAccountId());
Closeable logHandle = logStreamer.attachForAccount(
account, info.containerId(), logGroup, logStream, region, "msk:" + cluster.getClusterName());
if (logHandle != null) {
logStreams.put(cluster.getClusterName(), logHandle);
logStreams.put(clusterIdentityKey(cluster), logHandle);
}
}

Expand Down Expand Up @@ -249,7 +252,7 @@ public void stopContainer(MskCluster cluster) {
}

// Close log stream
Closeable logHandle = logStreams.remove(cluster.getClusterName());
Closeable logHandle = logStreams.remove(clusterIdentityKey(cluster));

lifecycleManager.stopAndRemove(cluster.getContainerId(), logHandle);
LOG.infov("Redpanda container {0} stopped and removed", cluster.getContainerId());
Expand Down Expand Up @@ -280,4 +283,23 @@ public void removeClusterStorage(MskCluster cluster) {
ContainerStorageHelper.removeStorage(config, lifecycleManager,
"msk", cluster.getVolumeId(), cluster.getClusterName());
}

private String clusterIdentityKey(MskCluster cluster) {
return cluster.getClusterArn() != null ? cluster.getClusterArn() : cluster.getClusterName();
}

private String clusterStorageId(MskCluster cluster) {
String account = AwsArnUtils.accountOrDefault(cluster.getClusterArn(), regionResolver.getDefaultAccountId());
String region = AwsArnUtils.regionOrDefault(cluster.getClusterArn(), regionResolver.getDefaultRegion());
return ContainerStorageHelper.dockerName(config,
"msk-" + account + "-" + region + "-" + cluster.getClusterName());
}

private Path legacyCompatibleHostPath(MskCluster cluster) {
Path scopedPath = ContainerStorageHelper.hostResourcePath(config, "msk", clusterStorageId(cluster));
Path legacyPath = ContainerStorageHelper.hostResourcePath(config, "msk", cluster.getClusterName());
return cluster.getResourceRegion() == null && Files.exists(legacyPath)
? legacyPath
: scopedPath;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,10 @@ public class MskCluster {
@JsonProperty("accountId")
private String accountId;

// Region is persisted separately so records created before regional ARNs remain identifiable.
@JsonProperty("resourceRegion")
private String resourceRegion;

// 6-char hex generated once at creation for stable, collision-free volume/container naming
@JsonProperty("volumeId")
private String volumeId;
Expand Down Expand Up @@ -180,6 +184,9 @@ public MskCluster(String clusterArn, String clusterName, String kafkaVersion) {
public String getAccountId() { return accountId; }
public void setAccountId(String accountId) { this.accountId = accountId; }

public String getResourceRegion() { return resourceRegion; }
public void setResourceRegion(String resourceRegion) { this.resourceRegion = resourceRegion; }

public String getVolumeId() { return volumeId; }
public void setVolumeId(String volumeId) { this.volumeId = volumeId; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ class MskServiceTest {
private StorageFactory storageFactory;
private EmulatorConfig config;
private RedpandaManager redpandaManager;
private RegionResolver regionResolver;
// MskService's constructor creates the cluster backend first and the configuration backend
// second; captured here so tests can seed raw entries directly into the configuration store
// (e.g. to simulate a pre-revision-history persisted entry) without exposing it from MskService.
Expand Down Expand Up @@ -87,7 +88,10 @@ void setUp() {
when(config.defaultRegion()).thenReturn("us-east-1");

redpandaManager = Mockito.mock(RedpandaManager.class);
RegionResolver regionResolver = new RegionResolver("us-east-1", "000000000000");
regionResolver = Mockito.mock(RegionResolver.class);
when(regionResolver.getRegion()).thenReturn("us-east-1");
when(regionResolver.getAccountId()).thenReturn("000000000000");
when(regionResolver.getDefaultRegion()).thenReturn("us-east-1");
mskService = new MskService(storageFactory, config, regionResolver, redpandaManager);
}

Expand Down Expand Up @@ -129,6 +133,53 @@ void listClusters() {
assertEquals(2, clusters.size());
}

@Test
void sameClusterNameIsAllowedInDifferentRegionsAndListsOnlyCurrentRegion() {
MskCluster east = mskService.createCluster("shared-name");

when(regionResolver.getRegion()).thenReturn("eu-west-1");
MskCluster west = mskService.createCluster("shared-name");

assertNotEquals(east.getClusterArn(), west.getClusterArn());
assertTrue(east.getClusterArn().contains(":us-east-1:"));
assertTrue(west.getClusterArn().contains(":eu-west-1:"));
assertEquals(List.of(west), mskService.listClusters());

when(regionResolver.getRegion()).thenReturn("us-east-1");
assertEquals(List.of(east), mskService.listClusters());
}

@Test
void legacyClusterWithoutResourceRegionRemainsVisibleAfterRegionalIsolation() {
MskCluster legacy = mskService.createCluster("legacy-cluster");
legacy.setResourceRegion(null);

when(regionResolver.getRegion()).thenReturn("eu-west-1");

assertEquals(List.of(legacy), mskService.listClusters());
assertEquals(legacy, mskService.describeCluster(legacy.getClusterArn()));
}

@Test
void sameClusterNameIsRejectedWithinOneRegion() {
mskService.createCluster("shared-name");

AwsException error = assertThrows(AwsException.class, () -> mskService.createCluster("shared-name"));

assertEquals("ConflictException", error.getErrorCode());
}

@Test
void clusterCannotBeDescribedOrDeletedFromAnotherRegion() {
MskCluster east = mskService.createCluster("regional-cluster");

when(regionResolver.getRegion()).thenReturn("eu-west-1");

assertThrows(AwsException.class, () -> mskService.describeCluster(east.getClusterArn()));
assertThrows(AwsException.class, () -> mskService.deleteCluster(east.getClusterArn()));
assertEquals(0, mskService.listClusters().size());
}

@Test
void deleteCluster() {
MskCluster cluster = mskService.createCluster("test-cluster");
Expand Down Expand Up @@ -355,6 +406,8 @@ void clusterSurvivesTheStorageMapperRoundTripWithItsInternalFields() throws Exce
assertNotNull(reloaded.getVolumeId());
assertEquals(cluster.getAccountId(), reloaded.getAccountId());
assertNotNull(reloaded.getAccountId());
assertEquals(cluster.getResourceRegion(), reloaded.getResourceRegion());
assertNotNull(reloaded.getResourceRegion());

// and the client-facing metadata survives too
assertEquals(3, reloaded.getNumberOfBrokerNodes());
Expand Down
Loading
Loading