mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
Rename ProxyOutResult to ProxyRelayResult
This commit is contained in:
+6
-6
@@ -33,7 +33,7 @@ import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.config.ProxyConfig;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyOutResult;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyRelayResult;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyRelayService;
|
||||
|
||||
public class GrpcChannelManager implements StartAndShutdown {
|
||||
@@ -89,13 +89,13 @@ public class GrpcChannelManager implements StartAndShutdown {
|
||||
return channelRef.get();
|
||||
}
|
||||
|
||||
public <T> String addResponseFuture(CompletableFuture<ProxyOutResult<T>> responseFuture) {
|
||||
public <T> String addResponseFuture(CompletableFuture<ProxyRelayResult<T>> responseFuture) {
|
||||
String nonce = this.nextNonce();
|
||||
this.resultNonceFutureMap.put(nonce, new ResultFuture<>(responseFuture));
|
||||
return nonce;
|
||||
}
|
||||
|
||||
public <T> CompletableFuture<ProxyOutResult<T>> getAndRemoveResponseFuture(String nonce) {
|
||||
public <T> CompletableFuture<ProxyRelayResult<T>> getAndRemoveResponseFuture(String nonce) {
|
||||
ResultFuture<T> resultFuture = this.resultNonceFutureMap.remove(nonce);
|
||||
if (resultFuture != null) {
|
||||
return resultFuture.future;
|
||||
@@ -120,7 +120,7 @@ public class GrpcChannelManager implements StartAndShutdown {
|
||||
if (System.currentTimeMillis() - resultFuture.createTime > timeOutMs) {
|
||||
resultFuture = this.resultNonceFutureMap.remove(nonce);
|
||||
if (resultFuture != null) {
|
||||
resultFuture.future.complete(new ProxyOutResult<>(ResponseCode.SYSTEM_BUSY, "call remote timeout", null));
|
||||
resultFuture.future.complete(new ProxyRelayResult<>(ResponseCode.SYSTEM_BUSY, "call remote timeout", null));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -137,10 +137,10 @@ public class GrpcChannelManager implements StartAndShutdown {
|
||||
}
|
||||
|
||||
protected static class ResultFuture<T> {
|
||||
public CompletableFuture<ProxyOutResult<T>> future;
|
||||
public CompletableFuture<ProxyRelayResult<T>> future;
|
||||
public long createTime = System.currentTimeMillis();
|
||||
|
||||
public ResultFuture(CompletableFuture<ProxyOutResult<T>> future) {
|
||||
public ResultFuture(CompletableFuture<ProxyRelayResult<T>> future) {
|
||||
this.future = future;
|
||||
}
|
||||
}
|
||||
|
||||
+3
-3
@@ -36,7 +36,7 @@ import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestH
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyChannel;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyOutResult;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyRelayResult;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyRelayService;
|
||||
import org.apache.rocketmq.proxy.service.transaction.TransactionId;
|
||||
import org.apache.rocketmq.remoting.common.RemotingUtil;
|
||||
@@ -154,7 +154,7 @@ public class GrpcClientChannel extends ProxyChannel {
|
||||
@Override
|
||||
protected CompletableFuture<Void> processGetConsumerRunningInfo(RemotingCommand command,
|
||||
GetConsumerRunningInfoRequestHeader header,
|
||||
CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> responseFuture) {
|
||||
CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> responseFuture) {
|
||||
if (!header.isJstackEnable()) {
|
||||
return CompletableFuture.completedFuture(null);
|
||||
}
|
||||
@@ -169,7 +169,7 @@ public class GrpcClientChannel extends ProxyChannel {
|
||||
@Override
|
||||
protected CompletableFuture<Void> processConsumeMessageDirectly(RemotingCommand command,
|
||||
ConsumeMessageDirectlyResultRequestHeader header,
|
||||
MessageExt messageExt, CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> responseFuture) {
|
||||
MessageExt messageExt, CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> responseFuture) {
|
||||
this.getTelemetryCommandStreamObserver().onNext(TelemetryCommand.newBuilder()
|
||||
.setVerifyMessageCommand(VerifyMessageCommand.newBuilder()
|
||||
.setNonce(this.grpcChannelManager.addResponseFuture(responseFuture))
|
||||
|
||||
@@ -64,7 +64,7 @@ import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.common.GrpcProxyException;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.processor.MessagingProcessor;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyOutResult;
|
||||
import org.apache.rocketmq.proxy.service.relay.ProxyRelayResult;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
|
||||
public class ClientActivity extends AbstractMessingActivity {
|
||||
@@ -247,17 +247,17 @@ public class ClientActivity extends AbstractMessingActivity {
|
||||
protected void reportThreadStackTrace(Context ctx, Status status, ThreadStackTrace request) {
|
||||
String nonce = request.getNonce();
|
||||
String threadStack = request.getThreadStackTrace();
|
||||
CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> responseFuture = this.grpcChannelManager.getAndRemoveResponseFuture(nonce);
|
||||
CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> responseFuture = this.grpcChannelManager.getAndRemoveResponseFuture(nonce);
|
||||
if (responseFuture != null) {
|
||||
try {
|
||||
if (status.getCode().equals(Code.OK)) {
|
||||
ConsumerRunningInfo runningInfo = new ConsumerRunningInfo();
|
||||
runningInfo.setJstack(threadStack);
|
||||
responseFuture.complete(new ProxyOutResult<>(ResponseCode.SUCCESS, "", runningInfo));
|
||||
responseFuture.complete(new ProxyRelayResult<>(ResponseCode.SUCCESS, "", runningInfo));
|
||||
} else if (status.getCode().equals(Code.VERIFY_MESSAGE_FORBIDDEN)) {
|
||||
responseFuture.complete(new ProxyOutResult<>(ResponseCode.NO_PERMISSION, "forbidden to verify message", null));
|
||||
responseFuture.complete(new ProxyRelayResult<>(ResponseCode.NO_PERMISSION, "forbidden to verify message", null));
|
||||
} else {
|
||||
responseFuture.complete(new ProxyOutResult<>(ResponseCode.SYSTEM_ERROR, "verify message failed", null));
|
||||
responseFuture.complete(new ProxyRelayResult<>(ResponseCode.SYSTEM_ERROR, "verify message failed", null));
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
responseFuture.completeExceptionally(t);
|
||||
@@ -267,11 +267,11 @@ public class ClientActivity extends AbstractMessingActivity {
|
||||
|
||||
protected void reportVerifyMessageResult(Context ctx, Status status, VerifyMessageResult request) {
|
||||
String nonce = request.getNonce();
|
||||
CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> responseFuture = this.grpcChannelManager.getAndRemoveResponseFuture(nonce);
|
||||
CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> responseFuture = this.grpcChannelManager.getAndRemoveResponseFuture(nonce);
|
||||
if (responseFuture != null) {
|
||||
try {
|
||||
ConsumeMessageDirectlyResult result = this.buildConsumeMessageDirectlyResult(status, request);
|
||||
responseFuture.complete(new ProxyOutResult<>(ResponseCode.SUCCESS, "", result));
|
||||
responseFuture.complete(new ProxyRelayResult<>(ResponseCode.SUCCESS, "", result));
|
||||
} catch (Throwable t) {
|
||||
responseFuture.completeExceptionally(t);
|
||||
}
|
||||
|
||||
+2
-2
@@ -29,12 +29,12 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
public class ClusterProxyRelayService implements ProxyRelayService {
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(RemotingCommand command,
|
||||
public CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(RemotingCommand command,
|
||||
GetConsumerRunningInfoRequestHeader header) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(
|
||||
@Override public CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(
|
||||
RemotingCommand command, ConsumeMessageDirectlyResultRequestHeader header) {
|
||||
return null;
|
||||
}
|
||||
|
||||
+3
-3
@@ -36,9 +36,9 @@ public class LocalProxyRelayService implements ProxyRelayService {
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(RemotingCommand command,
|
||||
public CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(RemotingCommand command,
|
||||
GetConsumerRunningInfoRequestHeader header) {
|
||||
CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> future = new CompletableFuture<>();
|
||||
CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> future = new CompletableFuture<>();
|
||||
future.thenAccept(proxyOutResult -> {
|
||||
if (proxyOutResult.getCode() == ResponseCode.SUCCESS && proxyOutResult.getResult() != null) {
|
||||
ConsumerRunningInfo consumerRunningInfo = proxyOutResult.getResult();
|
||||
@@ -59,7 +59,7 @@ public class LocalProxyRelayService implements ProxyRelayService {
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(RemotingCommand command,
|
||||
public CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(RemotingCommand command,
|
||||
ConsumeMessageDirectlyResultRequestHeader header) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -116,13 +116,13 @@ public abstract class ProxyChannel extends AbstractChannel {
|
||||
protected abstract CompletableFuture<Void> processGetConsumerRunningInfo(
|
||||
RemotingCommand command,
|
||||
GetConsumerRunningInfoRequestHeader header,
|
||||
CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> responseFuture);
|
||||
CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> responseFuture);
|
||||
|
||||
protected abstract CompletableFuture<Void> processConsumeMessageDirectly(
|
||||
RemotingCommand command,
|
||||
ConsumeMessageDirectlyResultRequestHeader header,
|
||||
MessageExt messageExt,
|
||||
CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> responseFuture);
|
||||
CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> responseFuture);
|
||||
|
||||
@Override
|
||||
public ChannelConfig config() {
|
||||
|
||||
+2
-2
@@ -17,12 +17,12 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.service.relay;
|
||||
|
||||
public class ProxyOutResult<T> {
|
||||
public class ProxyRelayResult<T> {
|
||||
private int code;
|
||||
private String remark;
|
||||
private T result;
|
||||
|
||||
public ProxyOutResult(int code, String remark, T result) {
|
||||
public ProxyRelayResult(int code, String remark, T result) {
|
||||
this.code = code;
|
||||
this.remark = remark;
|
||||
this.result = result;
|
||||
@@ -25,12 +25,12 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public interface ProxyRelayService {
|
||||
|
||||
CompletableFuture<ProxyOutResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(
|
||||
CompletableFuture<ProxyRelayResult<ConsumerRunningInfo>> processGetConsumerRunningInfo(
|
||||
RemotingCommand command,
|
||||
GetConsumerRunningInfoRequestHeader header
|
||||
);
|
||||
|
||||
CompletableFuture<ProxyOutResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(
|
||||
CompletableFuture<ProxyRelayResult<ConsumeMessageDirectlyResult>> processConsumeMessageDirectly(
|
||||
RemotingCommand command,
|
||||
ConsumeMessageDirectlyResultRequestHeader header
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user