1use 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
26pub const MESH_HTTP_SCHEMA_V1: &str = "appcore.gateway.mesh-http.v1";
28
29pub const MESH_PEER_RELAY_PATH: &str = "/v1/gateway/mesh/peer";
31
32#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
34pub struct MeshPeerRequest {
35 pub schema: String,
37 pub request_id: String,
39 pub target_tenant_id: TenantId,
41 pub target_core_id: CoreId,
43 pub method: String,
45 pub path: String,
47 pub body: Vec<u8>,
49 pub bearer_token: Option<String>,
51 pub timeout_ms: u64,
53 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 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 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 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#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
163pub struct MeshPeerResponse {
164 pub schema: String,
166 pub request_id: String,
168 pub status_code: u16,
170 pub body: Vec<u8>,
172 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 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 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 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#[derive(Debug, Clone, PartialEq, Eq)]
245pub struct MeshPeerTransport {
246 relay_url: String,
247}
248
249impl MeshPeerTransport {
250 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 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 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}