Skip to main content

appcore_gateway/
mesh.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: mesh.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/07/26 10:16:57 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/08/02 12:48:56 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Mesh relay transport for Peer RPC over the Gateway.
12
13use appcore_peer_rpc::{
14    envelope_signing_hash, payload_hash, CancellationToken, PeerRpcEnvelope, PeerRpcError,
15    PeerRpcHttpRequest, PeerRpcHttpResponse, PeerTransportProvider, PEER_COMMAND_PATH,
16    PEER_HEALTH_PATH, PEER_MANIFEST_PATH, PEER_QUERY_PATH,
17};
18use appcore_transport::{
19    send, HttpClientConfig, HttpHeader, HttpRequest, HttpTarget, TransportError,
20};
21use appcore_types::{CoreId, TenantId};
22use serde::{Deserialize, Serialize};
23use std::sync::atomic::{AtomicU64, Ordering};
24use std::time::{SystemTime, UNIX_EPOCH};
25
26/// Stable schema marker for Gateway mesh HTTP relay messages.
27pub const MESH_HTTP_SCHEMA_V1: &str = "appcore.gateway.mesh-http.v1";
28
29/// Stable Gateway mesh relay endpoint.
30pub const MESH_PEER_RELAY_PATH: &str = "/v1/gateway/mesh/peer";
31
32/// Logical Peer RPC HTTP request forwarded through a Gateway worker socket.
33#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
34pub struct MeshPeerRequest {
35    /// Schema marker. Must be [`MESH_HTTP_SCHEMA_V1`].
36    pub schema: String,
37    /// Stable request identity.
38    pub request_id: String,
39    /// Tenant boundary for the target worker.
40    pub target_tenant_id: TenantId,
41    /// Target Core connected to the Gateway as a worker.
42    pub target_core_id: CoreId,
43    /// HTTP method from the logical Peer RPC request.
44    pub method: String,
45    /// Peer RPC path from the logical request.
46    pub path: String,
47    /// Encoded logical request body.
48    pub body: Vec<u8>,
49    /// Optional bearer credential forwarded to the target peer host.
50    pub bearer_token: Option<String>,
51    /// Per-attempt timeout in milliseconds.
52    pub timeout_ms: u64,
53    /// Maximum accepted response body bytes.
54    pub max_response_bytes: usize,
55}
56
57impl std::fmt::Debug for MeshPeerRequest {
58    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
59        formatter
60            .debug_struct("MeshPeerRequest")
61            .field("schema", &self.schema)
62            .field("request_id", &self.request_id)
63            .field("target_tenant_id", &self.target_tenant_id)
64            .field("target_core_id", &self.target_core_id)
65            .field("method", &self.method)
66            .field("path", &self.path)
67            .field("body_bytes", &self.body.len())
68            .field(
69                "bearer_token",
70                &self.bearer_token.as_ref().map(|_| "REDACTED"),
71            )
72            .field("timeout_ms", &self.timeout_ms)
73            .field("max_response_bytes", &self.max_response_bytes)
74            .finish()
75    }
76}
77
78impl MeshPeerRequest {
79    /// Builds a mesh relay request from a logical Peer RPC HTTP request.
80    pub fn new(
81        request_id: impl Into<String>,
82        target_tenant_id: TenantId,
83        target_core_id: CoreId,
84        request: PeerRpcHttpRequest,
85    ) -> Self {
86        Self {
87            schema: MESH_HTTP_SCHEMA_V1.to_string(),
88            request_id: request_id.into(),
89            target_tenant_id,
90            target_core_id,
91            method: request.method,
92            path: request.path,
93            body: request.body,
94            bearer_token: request.bearer_token,
95            timeout_ms: request.timeout_ms,
96            max_response_bytes: request.max_response_bytes,
97        }
98    }
99
100    /// Converts this relay request back to a logical Peer RPC HTTP request.
101    pub fn into_peer_request(self) -> Result<PeerRpcHttpRequest, PeerRpcError> {
102        self.validate_schema()?;
103        Ok(PeerRpcHttpRequest {
104            method: self.method,
105            path: self.path,
106            body: self.body,
107            bearer_token: self.bearer_token,
108            timeout_ms: self.timeout_ms,
109            max_response_bytes: self.max_response_bytes,
110        })
111    }
112
113    /// Validates the mesh relay schema marker.
114    pub fn validate_schema(&self) -> Result<(), PeerRpcError> {
115        if self.schema != MESH_HTTP_SCHEMA_V1 {
116            return Err(PeerRpcError::ProtocolMismatch);
117        }
118        if self.request_id.is_empty()
119            || self.body.len() > crate::config::MAX_GATEWAY_MESSAGE_BYTES
120            || self.timeout_ms == 0
121            || self.timeout_ms > crate::config::MAX_GATEWAY_REQUEST_TIMEOUT.as_millis() as u64
122            || self.max_response_bytes == 0
123            || self.max_response_bytes > crate::config::MAX_GATEWAY_MESSAGE_BYTES
124        {
125            return Err(PeerRpcError::PayloadTooLarge);
126        }
127        self.expected_request_hash().map(|_| ())
128    }
129
130    pub(crate) fn expected_request_hash(&self) -> Result<Option<String>, PeerRpcError> {
131        match (self.method.as_str(), self.path.as_str()) {
132            ("GET", PEER_HEALTH_PATH | PEER_MANIFEST_PATH) if self.body.is_empty() => Ok(None),
133            ("POST", PEER_QUERY_PATH | PEER_COMMAND_PATH) => {
134                let envelope = self.peer_envelope()?;
135                Ok(Some(envelope_signing_hash(&envelope)))
136            }
137            _ => Err(PeerRpcError::InvalidEnvelope(
138                "mesh_route_not_allowed".to_string(),
139            )),
140        }
141    }
142
143    pub(crate) fn peer_envelope(&self) -> Result<PeerRpcEnvelope, PeerRpcError> {
144        let envelope = serde_json::from_slice::<PeerRpcEnvelope>(&self.body)
145            .map_err(|_| PeerRpcError::InvalidEnvelope("mesh_peer_envelope_invalid".into()))?;
146        if envelope.request_id != self.request_id
147            || envelope.tenant_id != self.target_tenant_id
148            || envelope.target_core_id != self.target_core_id
149        {
150            return Err(PeerRpcError::InvalidEnvelope(
151                "mesh_routing_metadata_mismatch".to_string(),
152            ));
153        }
154        if envelope.body_hash != payload_hash(&envelope.payload) {
155            return Err(PeerRpcError::InvalidBodyHash);
156        }
157        Ok(envelope)
158    }
159}
160
161/// Logical Peer RPC HTTP response returned through the Gateway mesh relay.
162#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
163pub struct MeshPeerResponse {
164    /// Schema marker. Must be [`MESH_HTTP_SCHEMA_V1`].
165    pub schema: String,
166    /// Stable request identity.
167    pub request_id: String,
168    /// HTTP status code returned by the target peer host.
169    pub status_code: u16,
170    /// Encoded response body returned by the target peer host.
171    pub body: Vec<u8>,
172    /// Controlled transport failure when the target worker could not complete the request.
173    pub error: Option<String>,
174}
175
176impl std::fmt::Debug for MeshPeerResponse {
177    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
178        formatter
179            .debug_struct("MeshPeerResponse")
180            .field("schema", &self.schema)
181            .field("request_id", &self.request_id)
182            .field("status_code", &self.status_code)
183            .field("body_bytes", &self.body.len())
184            .field("has_error", &self.error.is_some())
185            .finish()
186    }
187}
188
189impl MeshPeerResponse {
190    /// Creates a successful logical HTTP response.
191    pub fn ok(request_id: impl Into<String>, response: PeerRpcHttpResponse) -> Self {
192        Self {
193            schema: MESH_HTTP_SCHEMA_V1.to_string(),
194            request_id: request_id.into(),
195            status_code: response.status_code,
196            body: response.body,
197            error: None,
198        }
199    }
200
201    /// Creates a controlled relay failure.
202    pub fn rejected(request_id: impl Into<String>, error: impl Into<String>) -> Self {
203        Self {
204            schema: MESH_HTTP_SCHEMA_V1.to_string(),
205            request_id: request_id.into(),
206            status_code: 503,
207            body: Vec::new(),
208            error: Some(error.into()),
209        }
210    }
211
212    /// Converts this relay response to a logical Peer RPC HTTP response.
213    pub fn into_peer_response(self) -> Result<PeerRpcHttpResponse, PeerRpcError> {
214        if self.schema != MESH_HTTP_SCHEMA_V1 {
215            return Err(PeerRpcError::ProtocolMismatch);
216        }
217        if let Some(error) = self.error {
218            return Err(PeerRpcError::Transport(error));
219        }
220        Ok(PeerRpcHttpResponse {
221            status_code: self.status_code,
222            body: self.body,
223        })
224    }
225
226    pub(crate) fn validate_for_request(
227        &self,
228        request_id: &str,
229        max_response_bytes: usize,
230    ) -> Result<(), PeerRpcError> {
231        if self.schema != MESH_HTTP_SCHEMA_V1 || self.request_id != request_id {
232            return Err(PeerRpcError::InvalidResponse(
233                "mesh response identity mismatch".to_string(),
234            ));
235        }
236        if self.body.len() > max_response_bytes {
237            return Err(PeerRpcError::PayloadTooLarge);
238        }
239        Ok(())
240    }
241}
242
243/// Peer RPC transport that reaches Cores through a Gateway mesh relay.
244#[derive(Debug, Clone, PartialEq, Eq)]
245pub struct MeshPeerTransport {
246    relay_url: String,
247}
248
249impl MeshPeerTransport {
250    /// Creates a mesh transport using the deployment-selected Gateway relay URL.
251    pub fn new(relay_url: impl Into<String>) -> Result<Self, PeerRpcError> {
252        let relay_url = relay_url.into();
253        HttpTarget::parse(&relay_url, MESH_PEER_RELAY_PATH).map_err(map_transport_error)?;
254        Ok(Self { relay_url })
255    }
256
257    /// Returns the configured relay URL.
258    pub fn relay_url(&self) -> &str {
259        &self.relay_url
260    }
261}
262
263impl PeerTransportProvider for MeshPeerTransport {
264    fn send(
265        &self,
266        base_url: &str,
267        request: PeerRpcHttpRequest,
268    ) -> Result<PeerRpcHttpResponse, PeerRpcError> {
269        self.send_request(base_url, request, None)
270    }
271
272    fn send_cancellable(
273        &self,
274        base_url: &str,
275        request: PeerRpcHttpRequest,
276        cancellation: &CancellationToken,
277    ) -> Result<PeerRpcHttpResponse, PeerRpcError> {
278        self.send_request(base_url, request, Some(cancellation))
279    }
280}
281
282impl MeshPeerTransport {
283    fn send_request(
284        &self,
285        base_url: &str,
286        request: PeerRpcHttpRequest,
287        cancellation: Option<&CancellationToken>,
288    ) -> Result<PeerRpcHttpResponse, PeerRpcError> {
289        let target = MeshPeerTarget::parse(base_url)?;
290        validate_peer_http_request(&request)?;
291        let request_id = mesh_request_id(&request, &target);
292        let relay_request =
293            MeshPeerRequest::new(request_id, target.tenant_id, target.core_id, request);
294        let max_response_bytes = relay_request.max_response_bytes;
295        let body = serde_json::to_vec(&relay_request)
296            .map_err(|error| PeerRpcError::Transport(error.to_string()))?;
297        if body.len() > crate::config::MAX_GATEWAY_HTTP_BODY_BYTES {
298            return Err(PeerRpcError::PayloadTooLarge);
299        }
300        let encoded_request_bytes = body.len();
301        let mut http_request = HttpRequest::new("POST", body)
302            .map_err(map_transport_error)?
303            .with_header(
304                HttpHeader::new("Content-Type", "application/json").map_err(map_transport_error)?,
305            )
306            .with_header(
307                HttpHeader::new("Accept", "application/json").map_err(map_transport_error)?,
308            );
309        if let Some(token) = relay_request.bearer_token.as_ref() {
310            http_request = http_request.with_header(
311                HttpHeader::sensitive("Authorization", format!("Bearer {token}"))
312                    .map_err(map_transport_error)?,
313            );
314        }
315        let target = HttpTarget::parse(&self.relay_url, MESH_PEER_RELAY_PATH)
316            .map_err(map_transport_error)?;
317        let response = send(
318            &target,
319            &http_request,
320            HttpClientConfig {
321                timeout_ms: relay_request.timeout_ms.max(1),
322                max_request_bytes: encoded_request_bytes,
323                max_response_bytes: relay_request.max_response_bytes.saturating_add(65_536),
324                max_header_bytes: 32_768,
325            },
326            cancellation,
327        )
328        .map_err(map_transport_error)?;
329        if !(200..300).contains(&response.status_code) {
330            return Err(PeerRpcError::EndpointUnavailable);
331        }
332        let response = serde_json::from_slice::<MeshPeerResponse>(&response.body)
333            .map_err(|error| PeerRpcError::InvalidResponse(error.to_string()))?
334            .into_peer_response()?;
335        if response.body.len() > max_response_bytes {
336            return Err(PeerRpcError::PayloadTooLarge);
337        }
338        Ok(response)
339    }
340}
341
342fn validate_peer_http_request(request: &PeerRpcHttpRequest) -> Result<(), PeerRpcError> {
343    if request.timeout_ms == 0
344        || request.timeout_ms > crate::config::MAX_GATEWAY_REQUEST_TIMEOUT.as_millis() as u64
345        || request.body.len() > crate::config::MAX_GATEWAY_MESSAGE_BYTES
346        || request.max_response_bytes == 0
347        || request.max_response_bytes > crate::config::MAX_GATEWAY_MESSAGE_BYTES
348    {
349        return Err(PeerRpcError::PayloadTooLarge);
350    }
351    Ok(())
352}
353
354#[derive(Debug, Clone, PartialEq, Eq)]
355struct MeshPeerTarget {
356    tenant_id: TenantId,
357    core_id: CoreId,
358}
359
360impl MeshPeerTarget {
361    fn parse(base_url: &str) -> Result<Self, PeerRpcError> {
362        let value = base_url
363            .strip_prefix("mesh://")
364            .or_else(|| base_url.strip_prefix("appcore-mesh://"))
365            .ok_or_else(|| PeerRpcError::Transport("invalid mesh peer URL".to_string()))?;
366        let (tenant, core) = value
367            .split_once('/')
368            .ok_or_else(|| PeerRpcError::Transport("invalid mesh peer URL".to_string()))?;
369        Ok(Self {
370            tenant_id: TenantId::new(tenant)
371                .map_err(|error| PeerRpcError::Transport(format!("{error:?}")))?,
372            core_id: CoreId::new(core)
373                .map_err(|error| PeerRpcError::Transport(format!("{error:?}")))?,
374        })
375    }
376}
377
378fn mesh_request_id(request: &PeerRpcHttpRequest, target: &MeshPeerTarget) -> String {
379    if let Ok(envelope) = serde_json::from_slice::<appcore_peer_rpc::PeerRpcEnvelope>(&request.body)
380    {
381        return envelope.request_id;
382    }
383    // appcore-norm: allow(global-state) reason: atomic sequence prevents process-local mesh request collisions
384    static MESH_REQUEST_COUNTER: AtomicU64 = AtomicU64::new(0);
385    let counter = MESH_REQUEST_COUNTER.fetch_add(1, Ordering::Relaxed);
386    let now_ms = SystemTime::now()
387        .duration_since(UNIX_EPOCH)
388        .map(|duration| duration.as_millis() as u64)
389        .unwrap_or(0);
390    format!(
391        "mesh-{}-{}-{}-{}-{}",
392        target.tenant_id.as_str(),
393        target.core_id.as_str(),
394        now_ms,
395        std::process::id(),
396        counter
397    )
398}
399
400fn map_transport_error(error: TransportError) -> PeerRpcError {
401    match error {
402        TransportError::ResponseTooLarge { .. } | TransportError::RequestTooLarge { .. } => {
403            PeerRpcError::PayloadTooLarge
404        }
405        TransportError::Timeout
406        | TransportError::ConnectionRefused
407        | TransportError::Dns(_)
408        | TransportError::Cancelled => PeerRpcError::EndpointUnavailable,
409        other => PeerRpcError::Transport(other.to_string()),
410    }
411}