MCPcopy Create free account
hub / github.com/e2wugui/zeze / send

Method send

ZezeJava/ZezeJava/src/main/java/Zeze/Raft/ProxyAgent.java:119–165  ·  view source on GitHub ↗

如果启用了代理,则把rpc包装成代理协议,发送出去; 否则按原始raft请求发送出去。 @param proxyAgent 启用代理的实例 @param rpc 发送的rpc @param leader leader连接器,可能是原始的,也可能是伪造的。可能为null。 @param leaderSocket leader.Socket。可能为null。 @return 发送结果,可能失败。 @see ProxyServer send

(Service localService,
							   ProxyAgent proxyAgent,
							   RaftRpc<?, ?> rpc,
							   Agent.ConnectorProxy leader,
							   AsyncSocket leaderSocket)

Source from the content-addressed store, hash-verified

117 * @see ProxyServer send
118 */
119 @SuppressWarnings("unchecked")
120 public static boolean send(Service localService,
121 ProxyAgent proxyAgent,
122 RaftRpc<?, ?> rpc,
123 Agent.ConnectorProxy leader,
124 AsyncSocket leaderSocket) {
125
126 if (null != proxyAgent) {
127 if (null != leader) {
128 var proxyArgument = new ProxyArgument(leader.getName(), rpc);
129 var proxyRpc = new ProxyRequest(proxyArgument);
130 // leaderSocket 就是从leader中获取的,这里是为了在循环中发送的时候不用每次获取,优化!
131 //logger.info("send to {}", leaderSocket.getRemoteAddress());
132 return proxyRpc.Send(leaderSocket, (proxyRpcThis) -> {
133 if (proxyRpc.getResultCode() == 0) {
134 var outFh = new OutObject<Service.ProtocolFactoryHandle<?>>();
135 var resultRpc = (Rpc<?, ?>)Protocol.decode(
136 localService::findProtocolFactoryHandle,
137 ByteBuffer.Wrap(proxyRpc.Result.getData()),
138 outFh);
139 if (null != resultRpc) {
140 if (null != rpc.getResponseHandle()) {
141 @SuppressWarnings("rawtypes")
142 var originHandle = (ProtocolHandle)rpc.getResponseHandle();
143 localService.dispatchRpcResponse(resultRpc, originHandle, outFh.value);
144 } else if (rpc.getFuture() != null){
145 rpc.getFuture().setRawResult(resultRpc.Result);
146 }
147 } else {
148 logger.error("Agent ProxyRequest({}) resultRpc not found.", proxyArgument.getRaftName());
149 }
150 } else {
151 if (Procedure.Timeout != proxyRpc.getResultCode())
152 logger.error("Agent ProxyRequest({}) error={}",
153 proxyArgument.getRaftName(), IModule.getErrorCode(proxyRpc.getResultCode()));
154 else
155 logger.info("Agent ProxyRequest({}) timeout.", proxyArgument.getRaftName());
156 }
157 return 0;
158 }, proxyAgent.rpcTimeout);
159 }
160 // leader 还没有选出。
161 return false;
162 }
163 // 旧的独立的直接的raft访问发送方式。
164 return rpc.Send(leaderSocket);
165 }
166}

Callers 3

sendMethod · 0.95
sendForWaitMethod · 0.95
resendMethod · 0.95

Calls 12

decodeMethod · 0.95
WrapMethod · 0.95
getRaftNameMethod · 0.95
getErrorCodeMethod · 0.95
getResponseHandleMethod · 0.80
setRawResultMethod · 0.80
getNameMethod · 0.65
getDataMethod · 0.65
SendMethod · 0.45
getResultCodeMethod · 0.45
dispatchRpcResponseMethod · 0.45
getFutureMethod · 0.45

Tested by

no test coverage detected