diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index f67e0d6ec8..3703ff5039 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -52,7 +52,7 @@ import apache.rocketmq.v1.ReportThreadStackTraceResponse; import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; -import io.netty.util.concurrent.CompleteFuture; +import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.constant.LoggerName; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -65,80 +65,80 @@ public class ClusterGrpcService implements GrpcService { } @Override - public CompleteFuture queryRoute(Context ctx, QueryRouteRequest request) { + public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { return null; } @Override - public CompleteFuture heartbeat(Context ctx, HeartbeatRequest request) { + public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { return null; } @Override - public CompleteFuture healthCheck(Context ctx, HealthCheckRequest request) { + public CompletableFuture healthCheck(Context ctx, HealthCheckRequest request) { return null; } @Override - public CompleteFuture sendMessage(Context ctx, SendMessageRequest request) { + public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { return null; } @Override - public CompleteFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { + public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { return null; } - @Override public CompleteFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { + @Override public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { return null; } - @Override public CompleteFuture ackMessage(Context ctx, AckMessageRequest request) { + @Override public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { return null; } - @Override public CompleteFuture nackMessage(Context ctx, NackMessageRequest request) { + @Override public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { return null; } @Override - public CompleteFuture forwardMessageToDeadLetterQueue(Context ctx, + public CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { return null; } - @Override public CompleteFuture endTransaction(Context ctx, EndTransactionRequest request) { + @Override public CompletableFuture endTransaction(Context ctx, EndTransactionRequest request) { return null; } - @Override public CompleteFuture queryOffset(Context ctx, QueryOffsetRequest request) { + @Override public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { return null; } - @Override public CompleteFuture pullMessage(Context ctx, PullMessageRequest request) { + @Override public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { return null; } - @Override public CompleteFuture pollCommand(Context ctx, PollCommandRequest request) { + @Override public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { return null; } - @Override public CompleteFuture reportThreadStackTrace(Context ctx, + @Override public CompletableFuture reportThreadStackTrace(Context ctx, ReportThreadStackTraceRequest request) { return null; } - @Override public CompleteFuture reportMessageConsumptionResult(Context ctx, + @Override public CompletableFuture reportMessageConsumptionResult(Context ctx, ReportMessageConsumptionResultRequest request) { return null; } - @Override public CompleteFuture notifyClientTermination(Context ctx, + @Override public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { return null; } - @Override public CompleteFuture changeInvisibleDuration(Context ctx, + @Override public CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java index 4a791f3b3c..bd51796922 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java @@ -52,46 +52,46 @@ import apache.rocketmq.v1.ReportThreadStackTraceResponse; import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; -import io.netty.util.concurrent.CompleteFuture; import org.apache.rocketmq.proxy.common.StartAndShutdown; +import java.util.concurrent.CompletableFuture; public interface GrpcService extends StartAndShutdown { - CompleteFuture queryRoute(Context ctx, QueryRouteRequest request); + CompletableFuture queryRoute(Context ctx, QueryRouteRequest request); - CompleteFuture heartbeat(Context ctx, HeartbeatRequest request); + CompletableFuture heartbeat(Context ctx, HeartbeatRequest request); - CompleteFuture healthCheck(Context ctx, HealthCheckRequest request); + CompletableFuture healthCheck(Context ctx, HealthCheckRequest request); - CompleteFuture sendMessage(Context ctx, SendMessageRequest request); + CompletableFuture sendMessage(Context ctx, SendMessageRequest request); - CompleteFuture queryAssignment(Context ctx, QueryAssignmentRequest request); + CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request); - CompleteFuture receiveMessage(Context ctx, ReceiveMessageRequest request); + CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request); - CompleteFuture ackMessage(Context ctx, AckMessageRequest request); + CompletableFuture ackMessage(Context ctx, AckMessageRequest request); - CompleteFuture nackMessage(Context ctx, NackMessageRequest request); + CompletableFuture nackMessage(Context ctx, NackMessageRequest request); - CompleteFuture forwardMessageToDeadLetterQueue(Context ctx, + CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request); - CompleteFuture endTransaction(Context ctx, EndTransactionRequest request); + CompletableFuture endTransaction(Context ctx, EndTransactionRequest request); - CompleteFuture queryOffset(Context ctx, QueryOffsetRequest request); + CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request); - CompleteFuture pullMessage(Context ctx, PullMessageRequest request); + CompletableFuture pullMessage(Context ctx, PullMessageRequest request); - CompleteFuture pollCommand(Context ctx, PollCommandRequest request); + CompletableFuture pollCommand(Context ctx, PollCommandRequest request); - CompleteFuture reportThreadStackTrace(Context ctx, + CompletableFuture reportThreadStackTrace(Context ctx, ReportThreadStackTraceRequest request); - CompleteFuture reportMessageConsumptionResult(Context ctx, + CompletableFuture reportMessageConsumptionResult(Context ctx, ReportMessageConsumptionResultRequest request); - CompleteFuture notifyClientTermination(Context ctx, + CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request); - CompleteFuture changeInvisibleDuration(Context ctx, + CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index c415fa8daf..3433075cec 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -53,7 +53,6 @@ import apache.rocketmq.v1.ReportThreadStackTraceResponse; import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; -import io.netty.util.concurrent.CompleteFuture; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -86,19 +85,19 @@ public class LocalGrpcService implements GrpcService { this.sendChannelManager = new ChannelManager<>(); } - @Override public CompleteFuture queryRoute(Context ctx, QueryRouteRequest request) { + @Override public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { return null; } - @Override public CompleteFuture heartbeat(Context ctx, HeartbeatRequest request) { + @Override public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { return null; } - @Override public CompleteFuture healthCheck(Context ctx, HealthCheckRequest request) { + @Override public CompletableFuture healthCheck(Context ctx, HealthCheckRequest request) { return null; } - @Override public CompleteFuture sendMessage(Context ctx, SendMessageRequest request) { + @Override public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { SendMessageRequestHeader requestHeader = Converter.buildSendMessageRequestHeader(request); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader); Message message = request.getMessage(); @@ -128,60 +127,60 @@ public class LocalGrpcService implements GrpcService { } @Override - public CompleteFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { + public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { return null; } - @Override public CompleteFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { + @Override public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { return null; } - @Override public CompleteFuture ackMessage(Context ctx, AckMessageRequest request) { + @Override public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { return null; } - @Override public CompleteFuture nackMessage(Context ctx, NackMessageRequest request) { + @Override public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { return null; } @Override - public CompleteFuture forwardMessageToDeadLetterQueue(Context ctx, + public CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { return null; } - @Override public CompleteFuture endTransaction(Context ctx, EndTransactionRequest request) { + @Override public CompletableFuture endTransaction(Context ctx, EndTransactionRequest request) { return null; } - @Override public CompleteFuture queryOffset(Context ctx, QueryOffsetRequest request) { + @Override public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { return null; } - @Override public CompleteFuture pullMessage(Context ctx, PullMessageRequest request) { + @Override public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { return null; } - @Override public CompleteFuture pollCommand(Context ctx, PollCommandRequest request) { + @Override public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { return null; } - @Override public CompleteFuture reportThreadStackTrace(Context ctx, + @Override public CompletableFuture reportThreadStackTrace(Context ctx, ReportThreadStackTraceRequest request) { return null; } - @Override public CompleteFuture reportMessageConsumptionResult(Context ctx, + @Override public CompletableFuture reportMessageConsumptionResult(Context ctx, ReportMessageConsumptionResultRequest request) { return null; } - @Override public CompleteFuture notifyClientTermination(Context ctx, + @Override public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { return null; } - @Override public CompleteFuture changeInvisibleDuration(Context ctx, + @Override public CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { return null; }