diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java index 1aed35126d..f39857ba84 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java @@ -67,7 +67,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc private final org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService clusterGrpcService; public ClusterGrpcService() { - this.clusterGrpcService = new org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService();; + this.clusterGrpcService = new org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService(); this.appendStartAndShutdown(clusterGrpcService); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java index 515039b263..e8caea4203 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java @@ -93,7 +93,8 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc this.consumerService = new ConsumerService(connectorManager, grpcClientManager); this.producerService = new ProducerService(connectorManager); this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager, grpcClientManager); - this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager); + this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, + channelManager, grpcClientManager, pollCommandResponseManager); this.pullMessageService = new PullMessageService(connectorManager); this.transactionService = new TransactionService(connectorManager, channelManager); @@ -108,7 +109,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { - return null; + return clientService.heartbeat(ctx, request); } @Override @@ -160,18 +161,18 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { - return null; + return clientService.notifyClientTermination(ctx, request); } @Override public CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { - return null; + return consumerService.changeInvisibleDuration(ctx, request); } @Override public StreamObserver telemetry(Context ctx, StreamObserver responseObserver) { - return null; + return clientService.telemetry(ctx, responseObserver); } private class ClusterGrpcServiceStartAndShutdown implements StartAndShutdown { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index 7267b78b85..ac744b4b3c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -18,6 +18,8 @@ package org.apache.rocketmq.proxy.grpc.v2.service.cluster; import apache.rocketmq.v2.AckMessageRequest; import apache.rocketmq.v2.AckMessageResponse; +import apache.rocketmq.v2.ChangeInvisibleDurationRequest; +import apache.rocketmq.v2.ChangeInvisibleDurationResponse; import apache.rocketmq.v2.Code; import apache.rocketmq.v2.Message; import apache.rocketmq.v2.NackMessageRequest; @@ -69,7 +71,8 @@ public class ConsumerService extends BaseService { private volatile ResponseHook receiveMessageHook; private volatile ResponseHook ackNoMatchedMessageHook; private volatile ResponseHook ackMessageHook; - private volatile ResponseHook nackMessageResponseResponseHook; + private volatile ResponseHook nackMessageHook; + private volatile ResponseHook changeInvisibleDurationHook; private final GrpcClientManager grpcClientManager; @@ -246,8 +249,8 @@ public class ConsumerService extends BaseService { public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> { - if (nackMessageResponseResponseHook != null) { - nackMessageResponseResponseHook.beforeResponse(ctx, request, response, throwable); + if (nackMessageHook != null) { + nackMessageHook.beforeResponse(ctx, request, response, throwable); } }); try { @@ -323,19 +326,106 @@ public class ConsumerService extends BaseService { .build(); } + public CompletableFuture changeInvisibleDuration(Context ctx, + ChangeInvisibleDurationRequest request) { + CompletableFuture future = new CompletableFuture<>(); + future.whenComplete((response, throwable) -> { + if (changeInvisibleDurationHook != null) { + changeInvisibleDurationHook.beforeResponse(ctx, request, response, throwable); + } + }); + try { + ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); + String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); + + ChangeInvisibleTimeRequestHeader requestHeader = convertToChangeInvisibleTimeRequestHeader(ctx, request); + CompletableFuture resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader); + resultFuture + .thenAccept(result -> { + try { + future.complete(convertToChangeInvisibleDurationResponse(ctx, request, result)); + } catch (Throwable throwable) { + future.completeExceptionally(throwable); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); + } catch (Throwable t) { + future.completeExceptionally(t); + } + return future; + } + + protected ChangeInvisibleTimeRequestHeader convertToChangeInvisibleTimeRequestHeader(Context ctx, + ChangeInvisibleDurationRequest request) { + return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request); + } + + protected ChangeInvisibleDurationResponse convertToChangeInvisibleDurationResponse(Context ctx, + ChangeInvisibleDurationRequest request, AckResult ackResult) { + if (AckStatus.OK.equals(ackResult.getStatus())) { + return ChangeInvisibleDurationResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .setReceiptHandle(ackResult.getExtraInfo()) + .build(); + } + return ChangeInvisibleDurationResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "changeInvisibleDuration failed: status is abnormal")) + .build(); + } + + public ReadQueueSelector getReadQueueSelector() { + return readQueueSelector; + } + public void setReadQueueSelector(ReadQueueSelector readQueueSelector) { this.readQueueSelector = readQueueSelector; } - public void setReceiveMessageHook(ResponseHook receiveMessageHook) { + public ResponseHook getReceiveMessageHook() { + return receiveMessageHook; + } + + public void setReceiveMessageHook( + ResponseHook receiveMessageHook) { this.receiveMessageHook = receiveMessageHook; } - public void setAckNoMatchedMessageHook(ResponseHook ackNoMatchedMessageHook) { + public ResponseHook getAckNoMatchedMessageHook() { + return ackNoMatchedMessageHook; + } + + public void setAckNoMatchedMessageHook( + ResponseHook ackNoMatchedMessageHook) { this.ackNoMatchedMessageHook = ackNoMatchedMessageHook; } - public void setAckMessageHook(ResponseHook ackMessageHook) { + public ResponseHook getAckMessageHook() { + return ackMessageHook; + } + + public void setAckMessageHook( + ResponseHook ackMessageHook) { this.ackMessageHook = ackMessageHook; } + + public ResponseHook getNackMessageHook() { + return nackMessageHook; + } + + public void setNackMessageHook( + ResponseHook nackMessageHook) { + this.nackMessageHook = nackMessageHook; + } + + public ResponseHook getChangeInvisibleDurationHook() { + return changeInvisibleDurationHook; + } + + public void setChangeInvisibleDurationHook( + ResponseHook changeInvisibleDurationHook) { + this.changeInvisibleDurationHook = changeInvisibleDurationHook; + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java index 6df2e87db6..996e58a89c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java @@ -16,14 +16,21 @@ */ package org.apache.rocketmq.proxy.grpc.v2.service.cluster; -import apache.rocketmq.v1.ConsumerData; -import apache.rocketmq.v1.HeartbeatRequest; -import apache.rocketmq.v1.NoopCommand; -import apache.rocketmq.v1.NotifyClientTerminationRequest; -import apache.rocketmq.v1.PollCommandRequest; -import apache.rocketmq.v1.PollCommandResponse; -import apache.rocketmq.v1.Resource; +import apache.rocketmq.v2.ClientOverwrittenSettings; +import apache.rocketmq.v2.ClientSettings; +import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Direction; +import apache.rocketmq.v2.HeartbeatRequest; +import apache.rocketmq.v2.HeartbeatResponse; +import apache.rocketmq.v2.NotifyClientTerminationRequest; +import apache.rocketmq.v2.NotifyClientTerminationResponse; +import apache.rocketmq.v2.Publishing; +import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.Settings; +import apache.rocketmq.v2.Subscription; +import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; +import io.grpc.stub.StreamObserver; import java.time.Duration; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; @@ -34,12 +41,17 @@ import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.MQVersion; +import org.apache.rocketmq.common.consumer.ConsumeFromWhere; +import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.proxy.channel.ChannelManager; -import org.apache.rocketmq.proxy.connector.ConnectorManager; -import org.apache.rocketmq.proxy.grpc.v1.adapter.GrpcConverter; import org.apache.rocketmq.proxy.common.PollResponseManager; -import org.apache.rocketmq.proxy.grpc.v1.adapter.channel.GrpcClientChannel; +import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -51,11 +63,13 @@ public class ForwardClientService extends BaseService { private final ConsumerManager consumerManager; private final ProducerManager producerManager; private final PollResponseManager pollCommandResponseManager; + private final GrpcClientManager grpcClientManager; public ForwardClientService( ConnectorManager connectorManager, ScheduledExecutorService scheduledExecutorService, ChannelManager channelManager, + GrpcClientManager grpcClientManager, PollResponseManager pollCommandResponseManager ) { super(connectorManager); @@ -65,6 +79,7 @@ public class ForwardClientService extends BaseService { Duration.ofSeconds(10).toMillis(), TimeUnit.MILLISECONDS); this.channelManager = channelManager; + this.grpcClientManager = grpcClientManager; this.pollCommandResponseManager = pollCommandResponseManager; this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() { @@ -80,90 +95,150 @@ public class ForwardClientService extends BaseService { this.producerManager.setProducerOfflineListener(connectorManager.getTransactionHeartbeatRegisterService()::onProducerGroupOffline); } - public void heartbeat(Context ctx, HeartbeatRequest request) { - String language = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.LANGUAGE); - LanguageCode languageCode = LanguageCode.valueOf(language); - String clientId = request.getClientId(); + public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { + CompletableFuture future = new CompletableFuture<>(); - if (request.hasProducerData()) { - String producerGroup = GrpcConverter.wrapResourceWithNamespace(request.getProducerData().getGroup()); - GrpcClientChannel channel = GrpcClientChannel.create(ctx, channelManager, producerGroup, clientId, pollCommandResponseManager); - ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); - producerManager.registerProducer(producerGroup, clientChannelInfo); - } + try { + String language = InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE); + String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); + LanguageCode languageCode = LanguageCode.valueOf(language); - if (request.hasConsumerData()) { - ConsumerData consumerData = request.getConsumerData(); - String consumerGroup = GrpcConverter.wrapResourceWithNamespace(consumerData.getGroup()); - GrpcClientChannel channel = GrpcClientChannel.create(ctx, channelManager, consumerGroup, clientId, pollCommandResponseManager); - ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); - - consumerManager.registerConsumer( - consumerGroup, - clientChannelInfo, - GrpcConverter.buildConsumeType(consumerData.getConsumeType()), - GrpcConverter.buildMessageModel(consumerData.getConsumeModel()), - GrpcConverter.buildConsumeFromWhere(consumerData.getConsumePolicy()), - GrpcConverter.buildSubscriptionDataSet(consumerData.getSubscriptionsList()), - false - ); - } - } - - public void unregister(Context ctx, NotifyClientTerminationRequest request) { - String clientId = request.getClientId(); - - if (request.hasProducerGroup()) { - String producerGroup = GrpcConverter.wrapResourceWithNamespace(request.getProducerGroup()); - GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, producerGroup, clientId); - if (channel != null) { - producerManager.doChannelCloseEvent(producerGroup, channel); - } - } - - if (request.hasConsumerGroup()) { - String consumerGroup = GrpcConverter.wrapResourceWithNamespace(request.getConsumerGroup()); - GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, consumerGroup, clientId); - if (channel != null) { - consumerManager.doChannelCloseEvent(consumerGroup, channel); - } - } - } - - public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { - CompletableFuture future = new CompletableFuture<>(); - PollCommandResponse noopCommandResponse = PollCommandResponse.newBuilder().setNoopCommand( - NoopCommand.newBuilder().build() - ).build(); - - String clientId = request.getClientId(); - switch (request.getGroupCase()) { - case PRODUCER_GROUP: - Resource producerGroup = request.getProducerGroup(); - String producerGroupName = GrpcConverter.wrapResourceWithNamespace(producerGroup); - GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(this.channelManager, producerGroupName, clientId); - if (producerChannel == null) { - future.complete(noopCommandResponse); - } else { - producerChannel.setClientObserver(future); + ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); + switch (clientSettings.getClientType()) { + case PRODUCER: { + for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) { + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); + GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager); + ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); + // use topic name as producer group + producerManager.registerProducer(topicName, clientChannelInfo); + } + break; } - break; - case CONSUMER_GROUP: - Resource consumerGroup = request.getConsumerGroup(); - String consumerGroupName = GrpcConverter.wrapResourceWithNamespace(consumerGroup); - GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(this.channelManager, consumerGroupName, clientId); - if (consumerChannel == null) { - future.complete(noopCommandResponse); - } else { - consumerChannel.setClientObserver(future); + case PULL_CONSUMER: + case PUSH_CONSUMER: + case SIMPLE_CONSUMER: { + if (!request.hasGroup()) { + throw new ProxyException(Code.ILLEGAL_CONSUMER_GROUP, "group cannot be empty for consumer"); + } + String consumerGroup = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + GrpcClientChannel channel = GrpcClientChannel.create(ctx, channelManager, consumerGroup, clientId, pollCommandResponseManager); + ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); + + consumerManager.registerConsumer( + consumerGroup, + clientChannelInfo, + GrpcConverter.buildConsumeType(clientSettings.getClientType()), + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + GrpcConverter.buildSubscriptionDataSet(clientSettings.getSettings() + .getSubscription() + .getSubscriptionsList()), + false + ); + break; } - break; - default: - break; + default: { + throw new IllegalArgumentException("ClientType not exist " + clientSettings.getClientType()); + } + } + future.complete(HeartbeatResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .build()); + return future; + } catch (Throwable t) { + future.completeExceptionally(t); } return future; } + public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { + CompletableFuture future = new CompletableFuture<>(); + + try { + String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); + ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); + + switch (clientSettings.getClientType()) { + case PRODUCER: + for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) { + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); + // user topic name as producer group + GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, topicName, clientId); + if (channel != null) { + producerManager.doChannelCloseEvent(topicName, channel); + } + } + break; + case PULL_CONSUMER: + case PUSH_CONSUMER: + case SIMPLE_CONSUMER: + if (!request.hasGroup()) { + throw new ProxyException(Code.ILLEGAL_CONSUMER_GROUP, "group cannot be empty for consumer"); + } + String consumerGroup = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, consumerGroup, clientId); + if (channel != null) { + consumerManager.doChannelCloseEvent(consumerGroup, channel); + } + break; + default: + break; + } + future.complete(NotifyClientTerminationResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .build()); + } catch (Throwable t) { + future.completeExceptionally(t); + } + return future; + } + + public StreamObserver telemetry(Context ctx, StreamObserver responseObserver) { + String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); + return new StreamObserver() { + @Override + public void onNext(TelemetryCommand request) { + if (request.getCommandCase() == TelemetryCommand.CommandCase.CLIENT_SETTINGS) { + ClientSettings clientSettings = request.getClientSettings(); + grpcClientManager.updateClientSettings(clientId, clientSettings); + Settings settings = clientSettings.getSettings(); + if (settings.hasPublishing()) { + Publishing publishing = settings.getPublishing(); + for (Resource topic : publishing.getTopicsList()) { + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); + GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager); + producerChannel.setClientObserver(responseObserver); + } + } + if (settings.hasSubscription()) { + Subscription subscription = settings.getSubscription(); + String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup()); + GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager); + consumerChannel.setClientObserver(responseObserver); + } + responseObserver.onNext(TelemetryCommand.newBuilder() + .setClientOverwrittenSettings(ClientOverwrittenSettings.newBuilder() + .setNonce(clientSettings.getNonce()) + .setDirection(Direction.RESPONSE) + .setSettings(settings) + .build()) + .build()); + } + } + + @Override + public void onError(Throwable t) { + + } + + @Override + public void onCompleted() { + responseObserver.onCompleted(); + } + }; + } + private void scanNotActiveChannel() { try { this.consumerManager.scanNotActiveChannel(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java index 4ea9ccd179..88bd809660 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java @@ -37,7 +37,7 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; -import org.apache.rocketmq.proxy.grpc.v1.adapter.channel.GrpcClientChannel; +import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.remoting.common.RemotingHelper; public class TransactionService extends BaseService implements TransactionStateChecker { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java index fee3efce17..c2d9fdedea 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java @@ -68,7 +68,7 @@ public class ForwardClientServiceTest extends BaseServiceTest { assertNotNull(channel); assertTrue(channel instanceof GrpcClientChannel); - clientService.unregister(ctx, NotifyClientTerminationRequest.newBuilder() + clientService.notifyClientTermination(ctx, NotifyClientTerminationRequest.newBuilder() .setClientId("clientId") .setProducerGroup(Resource.newBuilder() .setName("producerGroup") @@ -127,7 +127,7 @@ public class ForwardClientServiceTest extends BaseServiceTest { assertEquals("*", consumerGroupInfo.getSubscriptionTable().get("topic").getSubString()); - clientService.unregister(ctx, NotifyClientTerminationRequest.newBuilder() + clientService.notifyClientTermination(ctx, NotifyClientTerminationRequest.newBuilder() .setClientId("clientId") .setConsumerGroup(Resource.newBuilder() .setName("consumerGroup")