[ISSUE #3949] Add code in GrpcMessagingProcessor

This commit is contained in:
zhouxiang
2022-07-13 11:29:09 +08:00
parent 115030a80f
commit 9e2dc9329c
2 changed files with 185 additions and 1 deletions
@@ -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<HeartbeatResponse> responseObserver) {
@Override
public void queryRoute(QueryRouteRequest request, StreamObserver<QueryRouteResponse> responseObserver) {
CompletableFuture<QueryRouteResponse> 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<HeartbeatResponse> responseObserver) {
CompletableFuture<HeartbeatResponse> 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<HealthCheckResponse> responseObserver) {
CompletableFuture<HealthCheckResponse> 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<QueryAssignmentResponse> responseObserver) {
CompletableFuture<QueryAssignmentResponse> 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<ReceiveMessageResponse> responseObserver) {
CompletableFuture<ReceiveMessageResponse> 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<AckMessageResponse> responseObserver) {
CompletableFuture<AckMessageResponse> 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<NackMessageResponse> responseObserver) {
CompletableFuture<NackMessageResponse> 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<ForwardMessageToDeadLetterQueueResponse> responseObserver) {
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> 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<EndTransactionResponse> responseObserver) {
CompletableFuture<EndTransactionResponse> 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<QueryOffsetResponse> responseObserver) {
CompletableFuture<QueryOffsetResponse> 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<PullMessageResponse> responseObserver) {
CompletableFuture<PullMessageResponse> 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<PollCommandResponse> responseObserver) {
CompletableFuture<PollCommandResponse> 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<ReportThreadStackTraceResponse> responseObserver) {
CompletableFuture<ReportThreadStackTraceResponse> 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<ReportMessageConsumptionResultResponse> responseObserver) {
CompletableFuture<ReportMessageConsumptionResultResponse> 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<NotifyClientTerminationResponse> responseObserver) {
CompletableFuture<NotifyClientTerminationResponse> 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<ChangeInvisibleDurationResponse> responseObserver) {
CompletableFuture<ChangeInvisibleDurationResponse> future = grpcForwardService.changeInvisibleDuration(Context.current(), request);
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
.exceptionally(e -> {
ResponseWriter.writeException(responseObserver, e);
return null;
});
}
}
@@ -28,6 +28,10 @@ public class ResponseWriter {
public static <T> void write(StreamObserver<T> observer, final T response) {
if (observer instanceof ServerCallStreamObserver) {
if (response == null) {
return;
}
final ServerCallStreamObserver<T> serverCallStreamObserver = (ServerCallStreamObserver<T>) observer;
if (serverCallStreamObserver.isCancelled()) {
LOGGER.warn("client has cancelled the request. response to write: {}", response);