Skip to main content

appcore_capabilities/
peer_rpc_invoker.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: peer_rpc_invoker.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/07/22 15:41:18 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/07/22 15:41:18 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Adapts capability requests to Peer RPC while preserving caller ownership.
12//!
13//! Borrowed invocation remains the compatibility path. Owned invocation moves
14//! every owned request field into the outbound DTO, translates the copyable
15//! mode into the call kind and must not duplicate its opaque payload.
16
17use crate::{
18    CapabilityError, CapabilityRequest, CapabilityResponse, CapabilityResult,
19    RemoteCapabilityInvoker,
20};
21use appcore_core::CapabilityMode;
22use appcore_distributed_contracts::{
23    PeerRecord, PeerRpcCallKind, PeerRpcClientExecutor, PeerRpcOutboundRequest,
24};
25
26/// Remote capability invoker backed by the stable peer RPC client contract.
27pub struct PeerRpcRemoteCapabilityInvoker<C> {
28    client: C,
29}
30
31impl<C> PeerRpcRemoteCapabilityInvoker<C> {
32    /// Creates an invoker using the supplied peer RPC executor.
33    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}