diff --git a/metadata/src/main/java/org/apache/kafka/controller/AclControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/AclControlManager.java index 7569e88ef0f00..661719df1b56c 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/AclControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/AclControlManager.java @@ -26,6 +26,7 @@ import org.apache.kafka.common.metadata.AccessControlEntryRecord; import org.apache.kafka.common.metadata.RemoveAccessControlEntryRecord; import org.apache.kafka.common.requests.ApiError; +import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.metadata.authorizer.ClusterMetadataAuthorizer; import org.apache.kafka.metadata.authorizer.StandardAcl; import org.apache.kafka.metadata.authorizer.StandardAclWithId; @@ -37,6 +38,7 @@ import org.apache.kafka.timeline.SnapshotRegistry; import org.apache.kafka.timeline.TimelineHashMap; import org.apache.kafka.timeline.TimelineHashSet; +import org.slf4j.Logger; import java.util.ArrayList; import java.util.Collections; @@ -64,12 +66,44 @@ * completed, which is another reason the prepare / complete callbacks are needed. */ public class AclControlManager { + static class Builder { + private LogContext logContext = null; + private SnapshotRegistry snapshotRegistry = null; + private Optional authorizer = Optional.empty(); + + Builder setLogContext(LogContext logContext) { + this.logContext = logContext; + return this; + } + + Builder setSnapshotRegistry(SnapshotRegistry snapshotRegistry) { + this.snapshotRegistry = snapshotRegistry; + return this; + } + + Builder setClusterMetadataAuthorizer(Optional authorizer) { + this.authorizer = authorizer; + return this; + } + + AclControlManager build() { + if (logContext == null) logContext = new LogContext(); + if (snapshotRegistry == null) snapshotRegistry = new SnapshotRegistry(logContext); + return new AclControlManager(logContext, snapshotRegistry, authorizer); + } + } + + private final Logger log; private final TimelineHashMap idToAcl; private final TimelineHashSet existingAcls; private final Optional authorizer; - AclControlManager(SnapshotRegistry snapshotRegistry, - Optional authorizer) { + AclControlManager( + LogContext logContext, + SnapshotRegistry snapshotRegistry, + Optional authorizer + ) { + this.log = logContext.logger(AclControlManager.class); this.idToAcl = new TimelineHashMap<>(snapshotRegistry, 0); this.existingAcls = new TimelineHashSet<>(snapshotRegistry, 0); this.authorizer = authorizer; @@ -184,8 +218,10 @@ static void validateFilter(AclBindingFilter filter) { } } - public void replay(AccessControlEntryRecord record, - Optional snapshotId) { + public void replay( + AccessControlEntryRecord record, + Optional snapshotId + ) { StandardAclWithId aclWithId = StandardAclWithId.fromRecord(record); idToAcl.put(aclWithId.id(), aclWithId.acl()); existingAcls.add(aclWithId.acl()); @@ -194,10 +230,14 @@ public void replay(AccessControlEntryRecord record, a.addAcl(aclWithId.id(), aclWithId.acl()); }); } + log.info("Replayed AccessControlEntryRecord for {}, setting {}", record.id(), + aclWithId.acl()); } - public void replay(RemoveAccessControlEntryRecord record, - Optional snapshotId) { + public void replay( + RemoveAccessControlEntryRecord record, + Optional snapshotId + ) { StandardAcl acl = idToAcl.remove(record.id()); if (acl == null) { throw new RuntimeException("Unable to replay " + record + ": no acl with " + @@ -212,6 +252,7 @@ public void replay(RemoveAccessControlEntryRecord record, a.removeAcl(record.id()); }); } + log.info("Replayed RemoveAccessControlEntryRecord for {}, removing {}", record.id(), acl); } Map idToAcl() { diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java index d969925e98ba6..dd86e530ebcce 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClientQuotaControlManager.java @@ -26,9 +26,11 @@ import org.apache.kafka.common.quota.ClientQuotaAlteration; import org.apache.kafka.common.quota.ClientQuotaEntity; import org.apache.kafka.common.requests.ApiError; +import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.server.common.ApiMessageAndVersion; import org.apache.kafka.timeline.SnapshotRegistry; import org.apache.kafka.timeline.TimelineHashMap; +import org.slf4j.Logger; import java.net.InetAddress; import java.net.UnknownHostException; @@ -45,11 +47,38 @@ public class ClientQuotaControlManager { + static class Builder { + private LogContext logContext = null; + private SnapshotRegistry snapshotRegistry = null; + + Builder setLogContext(LogContext logContext) { + this.logContext = logContext; + return this; + } + + Builder setSnapshotRegistry(SnapshotRegistry snapshotRegistry) { + this.snapshotRegistry = snapshotRegistry; + return this; + } + + ClientQuotaControlManager build() { + if (logContext == null) logContext = new LogContext(); + if (snapshotRegistry == null) snapshotRegistry = new SnapshotRegistry(logContext); + return new ClientQuotaControlManager(logContext, snapshotRegistry); + } + } + + private final Logger log; + private final SnapshotRegistry snapshotRegistry; final TimelineHashMap> clientQuotaData; - ClientQuotaControlManager(SnapshotRegistry snapshotRegistry) { + ClientQuotaControlManager( + LogContext logContext, + SnapshotRegistry snapshotRegistry + ) { + this.log = logContext.logger(ClientQuotaControlManager.class); this.snapshotRegistry = snapshotRegistry; this.clientQuotaData = new TimelineHashMap<>(snapshotRegistry, 0); } @@ -109,8 +138,11 @@ public void replay(ClientQuotaRecord record) { if (quotas.size() == 0) { clientQuotaData.remove(entity); } + log.info("Replayed ClientQuotaRecord for {} removing {}.", entity, record.key()); } else { quotas.put(record.key(), record.value()); + log.info("Replayed ClientQuotaRecord for {} setting {} to {}.", + entity, record.key(), record.value()); } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java index f0986f1d1a253..cd13aa6430771 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -437,11 +437,13 @@ public void replay(RegisterBrokerRecord record, long offset) { heartbeatManager.register(brokerId, record.fenced()); } if (prevRegistration == null) { - log.info("Registered new broker: {}", record); + log.info("Replayed initial RegisterBrokerRecord for broker {}: {}", record.brokerId(), record); } else if (prevRegistration.incarnationId().equals(record.incarnationId())) { - log.info("Re-registered broker incarnation: {}", record); + log.info("Replayed RegisterBrokerRecord modifying the registration for broker {}: {}", + record.brokerId(), record); } else { - log.info("Re-registered broker id {}: {}", brokerId, record); + log.info("Replayed RegisterBrokerRecord establishing a new incarnation of broker {}: {}", + record.brokerId(), record); } } @@ -458,7 +460,7 @@ public void replay(UnregisterBrokerRecord record) { } else { if (heartbeatManager != null) heartbeatManager.remove(brokerId); brokerRegistrations.remove(brokerId); - log.info("Unregistered broker: {}", record); + log.info("Replayed {}", record); } } @@ -520,6 +522,8 @@ private void replayRegistrationChange( inControlledShutdownChange ); if (!curRegistration.equals(nextRegistration)) { + log.info("Replayed {} modifying the registration for broker {}: {}", + record.getClass().getSimpleName(), brokerId, record); brokerRegistrations.put(brokerId, nextRegistration); } else { log.info("Ignoring no-op registration change for {}", curRegistration); diff --git a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java index 2cbf51c9baf80..adc9d195f45af 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java @@ -411,9 +411,11 @@ public void replay(ConfigRecord record) { configData.remove(configResource); } if (configSchema.isSensitive(record)) { - log.info("{}: set configuration {} to {}", configResource, record.name(), Password.HIDDEN); + log.info("Replayed ConfigRecord for {} which set configuration {} to {}", + configResource, record.name(), Password.HIDDEN); } else { - log.info("{}: set configuration {} to {}", configResource, record.name(), record.value()); + log.info("Replayed ConfigRecord for {} which set configuration {} to {}", + configResource, record.name(), record.value()); } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java index b197518391622..93754811c74b5 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/FeatureControlManager.java @@ -325,24 +325,31 @@ public void replay(FeatureLevelRecord record) { } if (record.name().equals(MetadataVersion.FEATURE_NAME)) { MetadataVersion mv = MetadataVersion.fromFeatureLevel(record.featureLevel()); - log.info("Setting metadata version to {}", mv); metadataVersion.set(mv); + log.info("Replayed a FeatureLevelRecord setting metadata version to {}", mv); } else { if (record.featureLevel() == 0) { - log.info("Removing feature {}", record.name()); finalizedVersions.remove(record.name()); + log.info("Replayed a FeatureLevelRecord removing feature {}", record.name()); } else { - log.info("Setting feature {} to {}", record.name(), record.featureLevel()); finalizedVersions.put(record.name(), record.featureLevel()); + log.info("Replayed a FeatureLevelRecord setting feature {} to {}", + record.name(), record.featureLevel()); } } } public void replay(ZkMigrationStateRecord record) { - ZkMigrationState recordState = ZkMigrationState.of(record.zkMigrationState()); - ZkMigrationState currentState = migrationControlState.get(); - log.info("Transitioning ZK migration state from {} to {}", currentState, recordState); - migrationControlState.set(recordState); + ZkMigrationState newState = ZkMigrationState.of(record.zkMigrationState()); + ZkMigrationState previousState = migrationControlState.get(); + if (previousState.equals(newState)) { + log.debug("Replayed a ZkMigrationStateRecord which did not alter the state from {}.", + previousState); + } else { + migrationControlState.set(newState); + log.info("Replayed a ZkMigrationStateRecord changing the migration state from {} to {}.", + previousState, newState); + } } boolean isControllerId(int nodeId) { diff --git a/metadata/src/main/java/org/apache/kafka/controller/ProducerIdControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ProducerIdControlManager.java index 932af0c2e754e..f2a29609edda4 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ProducerIdControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ProducerIdControlManager.java @@ -19,21 +19,62 @@ import org.apache.kafka.common.errors.UnknownServerException; import org.apache.kafka.common.metadata.ProducerIdsRecord; +import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.server.common.ApiMessageAndVersion; import org.apache.kafka.server.common.ProducerIdsBlock; import org.apache.kafka.timeline.SnapshotRegistry; import org.apache.kafka.timeline.TimelineLong; import org.apache.kafka.timeline.TimelineObject; +import org.slf4j.Logger; import java.util.Collections; public class ProducerIdControlManager { + static class Builder { + private LogContext logContext = null; + private SnapshotRegistry snapshotRegistry = null; + private ClusterControlManager clusterControlManager = null; + + Builder setLogContext(LogContext logContext) { + this.logContext = logContext; + return this; + } + + Builder setSnapshotRegistry(SnapshotRegistry snapshotRegistry) { + this.snapshotRegistry = snapshotRegistry; + return this; + } + + Builder setClusterControlManager(ClusterControlManager clusterControlManager) { + this.clusterControlManager = clusterControlManager; + return this; + } + + ProducerIdControlManager build() { + if (logContext == null) logContext = new LogContext(); + if (snapshotRegistry == null) snapshotRegistry = new SnapshotRegistry(logContext); + if (clusterControlManager == null) { + throw new RuntimeException("You must specify ClusterControlManager."); + } + return new ProducerIdControlManager( + logContext, + clusterControlManager, + snapshotRegistry); + } + } + + private final Logger log; private final ClusterControlManager clusterControlManager; private final TimelineObject nextProducerBlock; private final TimelineLong brokerEpoch; - ProducerIdControlManager(ClusterControlManager clusterControlManager, SnapshotRegistry snapshotRegistry) { + private ProducerIdControlManager( + LogContext logContext, + ClusterControlManager clusterControlManager, + SnapshotRegistry snapshotRegistry + ) { + this.log = logContext.logger(ProducerIdControlManager.class); this.clusterControlManager = clusterControlManager; this.nextProducerBlock = new TimelineObject<>(snapshotRegistry, ProducerIdsBlock.EMPTY); this.brokerEpoch = new TimelineLong(snapshotRegistry); @@ -71,7 +112,9 @@ void replay(ProducerIdsRecord record) { throw new RuntimeException("Next Producer ID from replayed record (" + record.nextProducerId() + ")" + " is not greater than current next Producer ID in block (" + nextBlock + ")"); } else { - nextProducerBlock.set(new ProducerIdsBlock(record.brokerId(), record.nextProducerId(), ProducerIdsBlock.PRODUCER_ID_BLOCK_SIZE)); + log.info("Replaying ProducerIdsRecord {}", record); + nextProducerBlock.set(new ProducerIdsBlock(record.brokerId(), record.nextProducerId(), + ProducerIdsBlock.PRODUCER_ID_BLOCK_SIZE)); brokerEpoch.set(record.brokerEpoch()); } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index 4da4298a1461e..7b11c26cd0bc5 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -1829,7 +1829,10 @@ private QuorumController( setStaticConfig(staticConfig). setNodeId(nodeId). build(); - this.clientQuotaControlManager = new ClientQuotaControlManager(snapshotRegistry); + this.clientQuotaControlManager = new ClientQuotaControlManager.Builder(). + setLogContext(logContext). + setSnapshotRegistry(snapshotRegistry). + build(); this.featureControl = new FeatureControlManager.Builder(). setLogContext(logContext). setQuorumFeatures(quorumFeatures). @@ -1851,7 +1854,11 @@ private QuorumController( setFeatureControlManager(featureControl). setZkMigrationEnabled(zkMigrationEnabled). build(); - this.producerIdControlManager = new ProducerIdControlManager(clusterControl, snapshotRegistry); + this.producerIdControlManager = new ProducerIdControlManager.Builder(). + setLogContext(logContext). + setSnapshotRegistry(snapshotRegistry). + setClusterControlManager(clusterControl). + build(); this.leaderImbalanceCheckIntervalNs = leaderImbalanceCheckIntervalNs; this.maxIdleIntervalNs = maxIdleIntervalNs; this.replicationControl = new ReplicationControlManager.Builder(). @@ -1871,10 +1878,14 @@ private QuorumController( build(); this.authorizer = authorizer; authorizer.ifPresent(a -> a.setAclMutator(this)); - this.aclControlManager = new AclControlManager(snapshotRegistry, authorizer); + this.aclControlManager = new AclControlManager.Builder(). + setLogContext(logContext). + setSnapshotRegistry(snapshotRegistry). + setClusterMetadataAuthorizer(authorizer). + build(); this.logReplayTracker = new LogReplayTracker.Builder(). - setLogContext(logContext). - build(); + setLogContext(logContext). + build(); this.raftClient = raftClient; this.bootstrapMetadata = bootstrapMetadata; this.maxRecordsPerBatch = maxRecordsPerBatch; diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 2d416d0eea6e4..2dab250cee50d 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -377,7 +377,19 @@ private ReplicationControlManager( } public void replay(TopicRecord record) { - topicsByName.put(record.name(), record.topicId()); + Uuid existingUuid = topicsByName.put(record.name(), record.topicId()); + if (existingUuid != null) { + // We don't currently support sending a second TopicRecord for the same topic name... + // unless, of course, there is a RemoveTopicRecord in between. + if (existingUuid.equals(record.topicId())) { + throw new RuntimeException("Found duplicate TopicRecord for " + record.name() + + " with topic ID " + record.topicId()); + } else { + throw new RuntimeException("Found duplicate TopicRecord for " + record.name() + + " with a different ID than before. Previous ID was " + existingUuid + + " and new ID is " + record.topicId()); + } + } if (Topic.hasCollisionChars(record.name())) { String normalizedName = Topic.unifyCollisionChars(record.name()); TimelineHashSet topicNames = topicsWithCollisionChars.get(normalizedName); @@ -389,7 +401,7 @@ public void replay(TopicRecord record) { } topics.put(record.topicId(), new TopicControlInfo(record.name(), snapshotRegistry, record.topicId())); - log.info("Created topic {} with topic ID {}.", record.name(), record.topicId()); + log.info("Replayed TopicRecord for topic {} with topic ID {}.", record.name(), record.topicId()); } public void replay(PartitionRecord record) { @@ -403,13 +415,16 @@ public void replay(PartitionRecord record) { String description = topicInfo.name + "-" + record.partitionId() + " with topic ID " + record.topicId(); if (prevPartInfo == null) { - log.info("Created partition {} and {}.", description, newPartInfo); + log.info("Replayed PartitionRecord for new partition {} and {}.", description, + newPartInfo); topicInfo.parts.put(record.partitionId(), newPartInfo); brokersToIsrs.update(record.topicId(), record.partitionId(), null, newPartInfo.isr, NO_LEADER, newPartInfo.leader); updateReassigningTopicsIfNeeded(record.topicId(), record.partitionId(), false, isReassignmentInProgress(newPartInfo)); } else if (!newPartInfo.equals(prevPartInfo)) { + log.info("Replayed PartitionRecord for existing partition {} and {}.", description, + newPartInfo); newPartInfo.maybeLogPartitionChange(log, description, prevPartInfo); topicInfo.parts.put(record.partitionId(), newPartInfo); brokersToIsrs.update(record.topicId(), record.partitionId(), prevPartInfo.isr, @@ -473,8 +488,8 @@ public void replay(PartitionChangeRecord record) { if (record.removingReplicas() != null || record.addingReplicas() != null) { log.info("Replayed partition assignment change {} for topic {}", record, topicInfo.name); - } else if (log.isTraceEnabled()) { - log.trace("Replayed partition change {} for topic {}", record, topicInfo.name); + } else if (log.isDebugEnabled()) { + log.debug("Replayed partition change {} for topic {}", record, topicInfo.name); } } @@ -514,7 +529,7 @@ public void replay(RemoveTopicRecord record) { } brokersToIsrs.removeTopicEntryForBroker(topic.id, NO_LEADER); - log.info("Removed topic {} with ID {}.", topic.name, record.topicId()); + log.info("Replayed RemoveTopicRecord for topic {} with ID {}.", topic.name, record.topicId()); } ControllerResult createTopics( diff --git a/metadata/src/main/java/org/apache/kafka/controller/ScramControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ScramControlManager.java index e9a5e5db43736..2a876053e4021 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ScramControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ScramControlManager.java @@ -322,7 +322,7 @@ public void replay(RemoveUserScramCredentialRecord record) { if (credentials.remove(key) == null) { throw new RuntimeException("Unable to find credential to delete: " + key); } - log.info("Removed SCRAM credential for {} with mechanism {}.", + log.info("Replayed RemoveUserScramCredentialRecord for {} with mechanism {}.", key.username, key.mechanism); } @@ -334,11 +334,11 @@ public void replay(UserScramCredentialRecord record) { record.serverKey(), record.iterations()); if (credentials.put(key, value) == null) { - log.info("Created new SCRAM credential for {} with mechanism {}.", - key.username, key.mechanism); + log.info("Replayed UserScramCredentialRecord creating new entry for {} with " + + "mechanism {}.", key.username, key.mechanism); } else { - log.info("Modified SCRAM credential for {} with mechanism {}.", - key.username, key.mechanism); + log.info("Replayed UserScramCredentialRecord modifying existing entry for {} " + + "with mechanism {}.", key.username, key.mechanism); } } diff --git a/metadata/src/test/java/org/apache/kafka/controller/AclControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/AclControlManagerTest.java index 566fa4acd5471..0806455dcf3f2 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/AclControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/AclControlManagerTest.java @@ -203,7 +203,9 @@ public void configure(Map configs) { public void testLoadSnapshot() { SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); snapshotRegistry.getOrCreateSnapshot(0); - AclControlManager manager = new AclControlManager(snapshotRegistry, Optional.empty()); + AclControlManager manager = new AclControlManager.Builder(). + setSnapshotRegistry(snapshotRegistry). + build(); // Load TEST_ACLS into the AclControlManager. Set loadedAcls = new HashSet<>(); @@ -236,8 +238,7 @@ public void testLoadSnapshot() { @Test public void testAddAndDelete() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - AclControlManager manager = new AclControlManager(snapshotRegistry, Optional.empty()); + AclControlManager manager = new AclControlManager.Builder().build(); MockClusterMetadataAuthorizer authorizer = new MockClusterMetadataAuthorizer(); authorizer.loadSnapshot(manager.idToAcl()); manager.replay(StandardAclWithIdTest.TEST_ACLS.get(0).toRecord(), Optional.empty()); @@ -248,8 +249,7 @@ public void testAddAndDelete() { @Test public void testCreateAclDeleteAcl() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - AclControlManager manager = new AclControlManager(snapshotRegistry, Optional.empty()); + AclControlManager manager = new AclControlManager.Builder().build(); MockClusterMetadataAuthorizer authorizer = new MockClusterMetadataAuthorizer(); authorizer.loadSnapshot(manager.idToAcl()); @@ -311,8 +311,7 @@ public void testCreateAclDeleteAcl() { @Test public void testDeleteDedupe() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - AclControlManager manager = new AclControlManager(snapshotRegistry, Optional.empty()); + AclControlManager manager = new AclControlManager.Builder().build(); MockClusterMetadataAuthorizer authorizer = new MockClusterMetadataAuthorizer(); authorizer.loadSnapshot(manager.idToAcl()); diff --git a/metadata/src/test/java/org/apache/kafka/controller/ClientQuotaControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ClientQuotaControlManagerTest.java index 12f7477a6ca9d..48c985bdfdb78 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ClientQuotaControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ClientQuotaControlManagerTest.java @@ -25,10 +25,8 @@ import org.apache.kafka.common.quota.ClientQuotaAlteration; import org.apache.kafka.common.quota.ClientQuotaEntity; import org.apache.kafka.common.requests.ApiError; -import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.metadata.RecordTestUtils; import org.apache.kafka.server.common.ApiMessageAndVersion; -import org.apache.kafka.timeline.SnapshotRegistry; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -52,8 +50,7 @@ public class ClientQuotaControlManagerTest { @Test public void testInvalidEntityTypes() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - ClientQuotaControlManager manager = new ClientQuotaControlManager(snapshotRegistry); + ClientQuotaControlManager manager = new ClientQuotaControlManager.Builder().build(); // Unknown type "foo" assertInvalidEntity(manager, entity("foo", "bar")); @@ -79,8 +76,7 @@ private void assertInvalidEntity(ClientQuotaControlManager manager, ClientQuotaE @Test public void testInvalidQuotaKeys() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - ClientQuotaControlManager manager = new ClientQuotaControlManager(snapshotRegistry); + ClientQuotaControlManager manager = new ClientQuotaControlManager.Builder().build(); ClientQuotaEntity entity = entity(ClientQuotaEntity.USER, "user-1"); // Invalid + valid keys @@ -103,8 +99,7 @@ private void assertInvalidQuota(ClientQuotaControlManager manager, ClientQuotaEn @Test public void testAlterAndRemove() { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - ClientQuotaControlManager manager = new ClientQuotaControlManager(snapshotRegistry); + ClientQuotaControlManager manager = new ClientQuotaControlManager.Builder().build(); ClientQuotaEntity userEntity = userEntity("user-1"); List alters = new ArrayList<>(); @@ -178,8 +173,7 @@ public void testAlterAndRemove() { @Test public void testEntityTypes() throws Exception { - SnapshotRegistry snapshotRegistry = new SnapshotRegistry(new LogContext()); - ClientQuotaControlManager manager = new ClientQuotaControlManager(snapshotRegistry); + ClientQuotaControlManager manager = new ClientQuotaControlManager.Builder().build(); Map> quotasToTest = new HashMap<>(); quotasToTest.put(userClientEntity("user-1", "client-id-1"), diff --git a/metadata/src/test/java/org/apache/kafka/controller/MockAclControlManager.java b/metadata/src/test/java/org/apache/kafka/controller/MockAclControlManager.java index 16d71e8786cbb..e14c8e2e4f5ba 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/MockAclControlManager.java +++ b/metadata/src/test/java/org/apache/kafka/controller/MockAclControlManager.java @@ -33,7 +33,7 @@ public class MockAclControlManager extends AclControlManager { public MockAclControlManager(LogContext logContext, Optional authorizer) { - super(new SnapshotRegistry(logContext), authorizer); + super(new LogContext(), new SnapshotRegistry(logContext), authorizer); } public List createAndReplayAcls(List acls) { diff --git a/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java index b50083169e6ab..fb02e0ebb4070 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ProducerIdControlManagerTest.java @@ -72,7 +72,10 @@ public void setUp() { clusterControl.replay(brokerRecord, 100L); } - this.producerIdControlManager = new ProducerIdControlManager(clusterControl, snapshotRegistry); + this.producerIdControlManager = new ProducerIdControlManager.Builder(). + setClusterControlManager(clusterControl). + setSnapshotRegistry(snapshotRegistry). + build(); } @Test diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java index f44c410e60f53..40e11a315aa4f 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerTest.java @@ -2527,4 +2527,26 @@ private static List isrWithDefaultEpoch(Integer... isr) { return Arrays.stream(isr).map(brokerId -> brokerState(brokerId, defaultBrokerEpoch(brokerId))) .collect(Collectors.toList()); } + + @Test + public void testDuplicateTopicIdReplay() { + ReplicationControlTestContext ctx = new ReplicationControlTestContext(); + ReplicationControlManager replicationControl = ctx.replicationControl; + replicationControl.replay(new TopicRecord(). + setName("foo"). + setTopicId(Uuid.fromString("Ktv3YkMQRe-MId4VkkrMyw"))); + assertEquals("Found duplicate TopicRecord for foo with topic ID Ktv3YkMQRe-MId4VkkrMyw", + assertThrows(RuntimeException.class, + () -> replicationControl.replay(new TopicRecord(). + setName("foo"). + setTopicId(Uuid.fromString("Ktv3YkMQRe-MId4VkkrMyw")))). + getMessage()); + assertEquals("Found duplicate TopicRecord for foo with a different ID than before. " + + "Previous ID was Ktv3YkMQRe-MId4VkkrMyw and new ID is 8auUWq8zQqe_99H_m2LAmw", + assertThrows(RuntimeException.class, + () -> replicationControl.replay(new TopicRecord(). + setName("foo"). + setTopicId(Uuid.fromString("8auUWq8zQqe_99H_m2LAmw")))). + getMessage()); + } }