From 9e2dc9329cdba18555ec3afd8c9e06e6784ad6eb Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Tue, 15 Mar 2022 11:47:42 +0800 Subject: [PATCH] [ISSUE #3949] Add code in GrpcMessagingProcessor --- .../proxy/grpc/GrpcMessagingProcessor.java | 182 +++++++++++++++++- .../proxy/grpc/common/ResponseWriter.java | 4 + 2 files changed, 185 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java index 69edec9c96..f1f026b4f0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java @@ -17,11 +17,39 @@ package org.apache.rocketmq.proxy.grpc; +import apache.rocketmq.v1.AckMessageRequest; +import apache.rocketmq.v1.AckMessageResponse; +import apache.rocketmq.v1.ChangeInvisibleDurationRequest; +import apache.rocketmq.v1.ChangeInvisibleDurationResponse; +import apache.rocketmq.v1.EndTransactionRequest; +import apache.rocketmq.v1.EndTransactionResponse; +import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest; +import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v1.HealthCheckRequest; import apache.rocketmq.v1.HealthCheckResponse; import apache.rocketmq.v1.HeartbeatRequest; import apache.rocketmq.v1.HeartbeatResponse; import apache.rocketmq.v1.MessagingServiceGrpc; +import apache.rocketmq.v1.NackMessageRequest; +import apache.rocketmq.v1.NackMessageResponse; +import apache.rocketmq.v1.NotifyClientTerminationRequest; +import apache.rocketmq.v1.NotifyClientTerminationResponse; +import apache.rocketmq.v1.PollCommandRequest; +import apache.rocketmq.v1.PollCommandResponse; +import apache.rocketmq.v1.PullMessageRequest; +import apache.rocketmq.v1.PullMessageResponse; +import apache.rocketmq.v1.QueryAssignmentRequest; +import apache.rocketmq.v1.QueryAssignmentResponse; +import apache.rocketmq.v1.QueryOffsetRequest; +import apache.rocketmq.v1.QueryOffsetResponse; +import apache.rocketmq.v1.QueryRouteRequest; +import apache.rocketmq.v1.QueryRouteResponse; +import apache.rocketmq.v1.ReceiveMessageRequest; +import apache.rocketmq.v1.ReceiveMessageResponse; +import apache.rocketmq.v1.ReportMessageConsumptionResultRequest; +import apache.rocketmq.v1.ReportMessageConsumptionResultResponse; +import apache.rocketmq.v1.ReportThreadStackTraceRequest; +import apache.rocketmq.v1.ReportThreadStackTraceResponse; import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; @@ -41,12 +69,34 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic this.grpcForwardService = grpcForwardService; } - public void heartbeat(HeartbeatRequest request, StreamObserver responseObserver) { + @Override + public void queryRoute(QueryRouteRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.queryRoute(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + @Override + public void heartbeat(HeartbeatRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.heartbeat(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); } @Override public void healthCheck(HealthCheckRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.healthCheck(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); } @Override @@ -58,4 +108,134 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic return null; }); } + + @Override + public void queryAssignment(QueryAssignmentRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.queryAssignment(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void receiveMessage(ReceiveMessageRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.receiveMessage(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void ackMessage(AckMessageRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.ackMessage(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void nackMessage(NackMessageRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.nackMessage(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void forwardMessageToDeadLetterQueue(ForwardMessageToDeadLetterQueueRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.forwardMessageToDeadLetterQueue(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void endTransaction(EndTransactionRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.endTransaction(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void queryOffset(QueryOffsetRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.queryOffset(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void pullMessage(PullMessageRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.pullMessage(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void pollCommand(PollCommandRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.pollCommand(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void reportThreadStackTrace(ReportThreadStackTraceRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.reportThreadStackTrace(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void reportMessageConsumptionResult(ReportMessageConsumptionResultRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.reportMessageConsumptionResult(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void notifyClientTermination(NotifyClientTerminationRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.notifyClientTermination(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } + + @Override + public void changeInvisibleDuration(ChangeInvisibleDurationRequest request, StreamObserver responseObserver) { + CompletableFuture future = grpcForwardService.changeInvisibleDuration(Context.current(), request); + future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) + .exceptionally(e -> { + ResponseWriter.writeException(responseObserver, e); + return null; + }); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java index d524cc055e..5d8fff2998 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java @@ -28,6 +28,10 @@ public class ResponseWriter { public static void write(StreamObserver observer, final T response) { if (observer instanceof ServerCallStreamObserver) { + if (response == null) { + return; + } + final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver) observer; if (serverCallStreamObserver.isCancelled()) { LOGGER.warn("client has cancelled the request. response to write: {}", response);