| 10 | import Zeze.Util.TaskCompletionSource; |
| 11 | |
| 12 | public abstract class RaftRpc<TArgument extends Serializable, TResult extends Serializable> |
| 13 | extends ProxyableRpc<TArgument, TResult> implements IRaftRpc { |
| 14 | private long createTime; |
| 15 | private UniqueRequestId unique = new UniqueRequestId(); |
| 16 | private long sendTime; |
| 17 | |
| 18 | TaskCompletionSource<RaftRpc<TArgument, TResult>> future; |
| 19 | ToLongFunction<Protocol<?>> handle; |
| 20 | |
| 21 | @Override |
| 22 | public int getFamilyClass() { |
| 23 | return isRequest() ? FamilyClass.RaftRequest : FamilyClass.RaftResponse; |
| 24 | } |
| 25 | |
| 26 | @Override |
| 27 | public long getCreateTime() { |
| 28 | return createTime; |
| 29 | } |
| 30 | |
| 31 | @Override |
| 32 | public void setCreateTime(long value) { |
| 33 | createTime = value; |
| 34 | } |
| 35 | |
| 36 | @Override |
| 37 | public UniqueRequestId getUnique() { |
| 38 | return unique; |
| 39 | } |
| 40 | |
| 41 | @Override |
| 42 | public void setUnique(UniqueRequestId value) { |
| 43 | unique = value; |
| 44 | } |
| 45 | |
| 46 | @Override |
| 47 | public long getSendTime() { |
| 48 | return sendTime; |
| 49 | } |
| 50 | |
| 51 | @Override |
| 52 | public void setSendTime(long value) { |
| 53 | sendTime = value; |
| 54 | } |
| 55 | |
| 56 | @Override |
| 57 | public boolean Send(AsyncSocket socket) { |
| 58 | // 1. |
| 59 | // 通过原始raft连接方式发送请求。 |
| 60 | // 由于每次发送需要新的rpc.sessionId,所以这里新建了一个桥接类发送, |
| 61 | // 并不会真正发送原始的RaftRpc。 |
| 62 | // 2. |
| 63 | // 通过Proxy发送的时候,ProxyRequest已经是一个新的请求,所有的dispatch特别处理, |
| 64 | // 不会使用这个方法发送请求了。 |
| 65 | // 3. |
| 66 | // Rpc.SendResult 针对上面两个发送方式,也有两个路径。 |
| 67 | // 3.1 原始方式通过SessionId找到本地存根并最终找到本RaftRpc.ResponseHandle把结果派发出去。 |
| 68 | // 3.2 Proxy方式由ProxyRequest的嵌套ResponseHandle派发出去。其中SendResult需要拦截。see 本类。 |
| 69 | var bridge = new RaftRpcBridge<>(this); |
nothing calls this directly
no outgoing calls
no test coverage detected