Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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
2 changes: 2 additions & 0 deletions checkstyle/import-control-core.xml
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,8 @@
</subpackage>

<subpackage name="test">
<allow pkg="org.apache.kafka.controller"/>
<allow pkg="org.apache.kafka.metadata"/>
<allow pkg="kafka.test.annotation"/>
<allow pkg="kafka.test.junit"/>
<allow pkg="kafka.network"/>
Expand Down
2 changes: 1 addition & 1 deletion checkstyle/suppressions.xml
Original file line number Diff line number Diff line change
Expand Up @@ -265,7 +265,7 @@

<!-- metadata -->
<suppress checks="ClassDataAbstractionCoupling"
files="(ReplicationControlManager).java"/>
files="(ReplicationControlManager|ReplicationControlManagerTest).java"/>
<suppress checks="ClassFanOutComplexity"
files="(QuorumController|ReplicationControlManager).java"/>
<suppress checks="CyclomaticComplexity"
Expand Down
203 changes: 179 additions & 24 deletions core/src/main/scala/kafka/server/ControllerApis.scala

Large diffs are not rendered by default.

222 changes: 222 additions & 0 deletions core/src/test/java/kafka/test/MockController.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,222 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You 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 kafka.test;

import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.errors.NotControllerException;
import org.apache.kafka.common.message.AlterIsrRequestData;
import org.apache.kafka.common.message.AlterIsrResponseData;
import org.apache.kafka.common.message.BrokerHeartbeatRequestData;
import org.apache.kafka.common.message.BrokerRegistrationRequestData;
import org.apache.kafka.common.message.CreateTopicsRequestData;
import org.apache.kafka.common.message.CreateTopicsResponseData;
import org.apache.kafka.common.message.ElectLeadersRequestData;
import org.apache.kafka.common.message.ElectLeadersResponseData;
import org.apache.kafka.common.protocol.Errors;
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.controller.Controller;
import org.apache.kafka.controller.ResultOrError;
import org.apache.kafka.metadata.BrokerHeartbeatReply;
import org.apache.kafka.metadata.BrokerRegistrationReply;
import org.apache.kafka.metadata.FeatureMapAndEpoch;

import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CompletableFuture;


public class MockController implements Controller {
private final static NotControllerException NOT_CONTROLLER_EXCEPTION =
new NotControllerException("This is not the correct controller for this cluster.");

public static class Builder {
private final Map<String, MockTopic> initialTopics = new HashMap<>();

public Builder newInitialTopic(String name, Uuid id) {
initialTopics.put(name, new MockTopic(name, id));
return this;
}

public MockController build() {
return new MockController(initialTopics.values());
}
}

private volatile boolean active = true;

private MockController(Collection<MockTopic> initialTopics) {
for (MockTopic topic : initialTopics) {
topics.put(topic.id, topic);
topicNameToId.put(topic.name, topic.id);
}
}

@Override
public CompletableFuture<AlterIsrResponseData> alterIsr(AlterIsrRequestData request) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<CreateTopicsResponseData> createTopics(CreateTopicsRequestData request) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<Void> unregisterBroker(int brokerId) {
throw new UnsupportedOperationException();
}

static class MockTopic {
private final String name;
private final Uuid id;

MockTopic(String name, Uuid id) {
this.name = name;
this.id = id;
}
}

private final Map<String, Uuid> topicNameToId = new HashMap<>();

private final Map<Uuid, MockTopic> topics = new HashMap<>();

@Override
synchronized public CompletableFuture<Map<String, ResultOrError<Uuid>>>
findTopicIds(Collection<String> topicNames) {
Map<String, ResultOrError<Uuid>> results = new HashMap<>();
for (String topicName : topicNames) {
if (!topicNameToId.containsKey(topicName)) {
results.put(topicName, new ResultOrError<>(new ApiError(Errors.UNKNOWN_TOPIC_OR_PARTITION)));
} else {
results.put(topicName, new ResultOrError<>(topicNameToId.get(topicName)));
}
}
return CompletableFuture.completedFuture(results);
}

@Override
synchronized public CompletableFuture<Map<Uuid, ResultOrError<String>>>
findTopicNames(Collection<Uuid> topicIds) {
Map<Uuid, ResultOrError<String>> results = new HashMap<>();
for (Uuid topicId : topicIds) {
MockTopic topic = topics.get(topicId);
if (topic == null) {
results.put(topicId, new ResultOrError<>(new ApiError(Errors.UNKNOWN_TOPIC_ID)));
} else {
results.put(topicId, new ResultOrError<>(topic.name));
}
}
return CompletableFuture.completedFuture(results);
}

@Override
synchronized public CompletableFuture<Map<Uuid, ApiError>>
deleteTopics(Collection<Uuid> topicIds) {
if (!active) {
CompletableFuture<Map<Uuid, ApiError>> future = new CompletableFuture<>();
future.completeExceptionally(NOT_CONTROLLER_EXCEPTION);
return future;
}
Map<Uuid, ApiError> results = new HashMap<>();
for (Uuid topicId : topicIds) {
MockTopic topic = topics.remove(topicId);
if (topic == null) {
results.put(topicId, new ApiError(Errors.UNKNOWN_TOPIC_ID));
} else {
topicNameToId.remove(topic.name);
results.put(topicId, ApiError.NONE);
}
}
return CompletableFuture.completedFuture(results);
}

@Override
public CompletableFuture<Map<ConfigResource, ResultOrError<Map<String, String>>>> describeConfigs(Map<ConfigResource, Collection<String>> resources) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<ElectLeadersResponseData> electLeaders(ElectLeadersRequestData request) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<FeatureMapAndEpoch> finalizedFeatures() {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<Map<ConfigResource, ApiError>> incrementalAlterConfigs(
Map<ConfigResource, Map<String, Map.Entry<AlterConfigOp.OpType, String>>> configChanges,
boolean validateOnly) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<Map<ConfigResource, ApiError>> legacyAlterConfigs(
Map<ConfigResource, Map<String, String>> newConfigs, boolean validateOnly) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<BrokerHeartbeatReply>
processBrokerHeartbeat(BrokerHeartbeatRequestData request) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<BrokerRegistrationReply>
registerBroker(BrokerRegistrationRequestData request) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<Void> waitForReadyBrokers(int minBrokers) {
throw new UnsupportedOperationException();
}

@Override
public CompletableFuture<Map<ClientQuotaEntity, ApiError>>
alterClientQuotas(Collection<ClientQuotaAlteration> quotaAlterations, boolean validateOnly) {
throw new UnsupportedOperationException();
}

@Override
public void beginShutdown() {
this.active = false;
}

public void setActive(boolean active) {
this.active = active;
}

@Override
public long curClaimEpoch() {
return active ? 1 : -1;
}

@Override
public void close() {
beginShutdown();
}
}
Loading