发送Rpc请求。
(RaftRpc<TArgument, TResult> rpc, ToLongFunction<Protocol<?>> handle)
| 96 | * 发送Rpc请求。 |
| 97 | */ |
| 98 | public <TArgument extends Serializable, TResult extends Serializable> |
| 99 | void send(RaftRpc<TArgument, TResult> rpc, ToLongFunction<Protocol<?>> handle) { |
| 100 | if (handle == null) |
| 101 | throw new IllegalArgumentException("null handle"); |
| 102 | if (pendingLimit > 0 && pending.size() > pendingLimit) // UrgentPending不限制。 |
| 103 | throw new IllegalStateException("too many pending"); |
| 104 | |
| 105 | // 由于interface不能把setter弄成保护的,实际上外面可以修改。 |
| 106 | // 简单检查一下吧。 |
| 107 | if (rpc.getUnique().getRequestId() != 0) |
| 108 | throw new IllegalStateException("RaftRpc.UniqueRequestId != 0. Need A Fresh RaftRpc"); |
| 109 | |
| 110 | rpc.getUnique().setRequestId(uniqueRequestIdGenerator.next()); |
| 111 | // 外面可以设置clientId,默认使用Generator.getName(); |
| 112 | if (rpc.getUnique().getClientId().isEmpty()) |
| 113 | rpc.getUnique().setClientId(uniqueRequestIdGenerator.getName()); |
| 114 | rpc.setCreateTime(System.currentTimeMillis()); |
| 115 | rpc.setSendTime(rpc.getCreateTime()); |
| 116 | if (rpc.getTimeout() == 0) // set default timeout |
| 117 | rpc.setTimeout(raftConfig.getAgentTimeout()); |
| 118 | |
| 119 | rpc.handle = handle; |
| 120 | if (pending.putIfAbsent(rpc.getUnique().getRequestId(), rpc) != null) |
| 121 | throw new IllegalStateException("duplicate requestId rpc=" + rpc); |
| 122 | |
| 123 | rpc.setResponseHandle(p -> sendHandle(p, rpc)); |
| 124 | ConnectorProxy leader = this.leader; |
| 125 | if (!ProxyAgent.send(client, proxyAgent, rpc, leader, |
| 126 | leader != null ? leader.getConnector().TryGetReadySocket() : null)) |
| 127 | logger.debug("send failed: leader={}, rpc={}", leader, rpc); |
| 128 | } |
| 129 | |
| 130 | private <TArgument extends Serializable, TResult extends Serializable> |
| 131 | long sendHandle(Rpc<TArgument, TResult> p, RaftRpc<TArgument, TResult> rpc) { |
no test coverage detected