From 08000137b7f07d705079603cd55672a03f866f85 Mon Sep 17 00:00:00 2001 From: "Colin P. McCabe" Date: Wed, 10 May 2023 10:50:15 -0700 Subject: [PATCH] MINOR: Standardize controller log4j output for replaying records Standardize controller log4j output for replaying important records. The log message should include word "replayed" to make it clear that this is a record replay. Log the replay of records for ACLs, client quotas, and producer IDs, which were previously not logged. Also fix a case where we weren't logging changes to broker registrations. AclControlManager, ClientQuotaControlManager, and ProducerIdControlManager didn't previously have a log4j logger object, so this PR adds one. It also converts them to using Builder objects. This makes junit tests more readable because we don't need to specify paramaters where the test can use the default (like LogContexts). Throw an exception in replay if we get another TopicRecord for a topic which already exists. --- .../kafka/controller/AclControlManager.java | 53 ++++++++++++++++--- .../controller/ClientQuotaControlManager.java | 34 +++++++++++- .../controller/ClusterControlManager.java | 12 +++-- .../ConfigurationControlManager.java | 6 ++- .../controller/FeatureControlManager.java | 21 +++++--- .../controller/ProducerIdControlManager.java | 47 +++++++++++++++- .../kafka/controller/QuorumController.java | 21 ++++++-- .../controller/ReplicationControlManager.java | 27 +++++++--- .../kafka/controller/ScramControlManager.java | 10 ++-- .../controller/AclControlManagerTest.java | 13 +++-- .../ClientQuotaControlManagerTest.java | 14 ++--- .../controller/MockAclControlManager.java | 2 +- .../ProducerIdControlManagerTest.java | 5 +- .../ReplicationControlManagerTest.java | 22 ++++++++ 14 files changed, 230 insertions(+), 57 deletions(-) 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()); + } }