appcore_capabilities/
peer_rpc_invoker.rs1use crate::{
18 CapabilityError, CapabilityRequest, CapabilityResponse, CapabilityResult,
19 RemoteCapabilityInvoker,
20};
21use appcore_core::CapabilityMode;
22use appcore_distributed_contracts::{
23 PeerRecord, PeerRpcCallKind, PeerRpcClientExecutor, PeerRpcOutboundRequest,
24};
25
26pub struct PeerRpcRemoteCapabilityInvoker<C> {
28 client: C,
29}
30
31impl<C> PeerRpcRemoteCapabilityInvoker<C> {
32 pub fn new(client: C) -> Self {
34 Self { client }
35 }
36}
37
38impl<C> RemoteCapabilityInvoker for PeerRpcRemoteCapabilityInvoker<C>
39where
40 C: PeerRpcClientExecutor,
41{
42 fn invoke_remote(
43 &self,
44 peer: &PeerRecord,
45 request: &CapabilityRequest,
46 ) -> CapabilityResult<CapabilityResponse> {
47 let endpoint_url = peer_rpc_endpoint(peer).ok_or_else(|| {
48 CapabilityError::RemoteEndpointUnavailable(request.capability.clone())
49 })?;
50 let kind = peer_rpc_call_kind(request.mode)?;
51 let response = self
52 .client
53 .call_peer(
54 endpoint_url,
55 kind,
56 PeerRpcOutboundRequest::new(
57 request.request_id.clone(),
58 peer.identity.core_id.clone(),
59 request.capability.clone(),
60 request.payload.clone(),
61 request.idempotency_key.clone(),
62 request.trace.clone(),
63 ),
64 )
65 .map_err(|error| CapabilityError::RemoteInvocationFailed(format!("{error:?}")))?;
66 capability_response(peer, response)
67 }
68
69 fn invoke_remote_owned(
70 &self,
71 peer: &PeerRecord,
72 request: CapabilityRequest,
73 ) -> CapabilityResult<CapabilityResponse> {
74 let endpoint_url = peer_rpc_endpoint(peer).ok_or_else(|| {
75 CapabilityError::RemoteEndpointUnavailable(request.capability.clone())
76 })?;
77 let kind = peer_rpc_call_kind(request.mode)?;
78 let CapabilityRequest {
79 request_id,
80 capability,
81 mode: _,
82 payload,
83 idempotency_key,
84 trace,
85 } = request;
86 let response = self
87 .client
88 .call_peer(
89 endpoint_url,
90 kind,
91 PeerRpcOutboundRequest::new(
92 request_id,
93 peer.identity.core_id.clone(),
94 capability,
95 payload,
96 idempotency_key,
97 trace,
98 ),
99 )
100 .map_err(|error| CapabilityError::RemoteInvocationFailed(format!("{error:?}")))?;
101 capability_response(peer, response)
102 }
103}
104
105fn peer_rpc_call_kind(mode: CapabilityMode) -> CapabilityResult<PeerRpcCallKind> {
106 match mode {
107 CapabilityMode::Query => Ok(PeerRpcCallKind::Query),
108 CapabilityMode::Command => Ok(PeerRpcCallKind::Command),
109 CapabilityMode::Stream => Err(CapabilityError::HandlerRejected(
110 "stream_remote_invocation_not_supported".to_string(),
111 )),
112 }
113}
114
115fn capability_response(
116 peer: &PeerRecord,
117 response: appcore_distributed_contracts::PeerRpcResponse,
118) -> CapabilityResult<CapabilityResponse> {
119 if response.ok {
120 return Ok(CapabilityResponse::accepted(
121 response.payload,
122 Some(peer.identity.core_id.clone()),
123 ));
124 }
125 Ok(CapabilityResponse::rejected(
126 response
127 .error
128 .unwrap_or_else(|| "remote_rejected".to_string()),
129 ))
130}
131
132fn peer_rpc_endpoint(peer: &PeerRecord) -> Option<&str> {
133 peer.endpoints
134 .iter()
135 .find(|endpoint| {
136 endpoint.name == "peer-rpc"
137 || endpoint.name == "peer_rpc"
138 || endpoint.protocol == "appcore-peer-rpc"
139 })
140 .map(|endpoint| endpoint.url.as_str())
141}