macula_rust/connection.rs
1//! The CONNECT/HELLO handshake and the application-frame stream
2//! abstraction, ported from `src/peering/macula_peering_conn.erl`
3//! (`macula-io/macula`) — see `plans/PLAN_WIRE_PROTOCOL.md` §3.
4//!
5//! Only the client role's `connecting -> handshaking -> connected` path
6//! is implemented. [`FrameStream`] is the reusable "send/receive signed
7//! application frames on one QUIC stream" primitive — [`Session`] wraps
8//! one for the control stream, and [`Session::open_dedicated_stream`]
9//! hands out fresh ones for content transfer (§12) and streaming RPC
10//! (§13), which both run on dedicated streams rather than the control
11//! stream.
12
13use std::sync::Arc;
14use std::time::Duration;
15
16use crate::bolt4;
17use crate::cbor::Value;
18use crate::frame::{self, Decoded, HelloInfo};
19use crate::identity::KeyPair;
20use crate::transport::{self, ConnectError, Trust};
21
22/// A boxed, `'static` future — hand-rolled rather than pulling in the
23/// `futures` crate for one type alias.
24pub type BoxFuture<'a, T> = std::pin::Pin<Box<dyn std::future::Future<Output = T> + Send + 'a>>;
25
26/// Answers one inbound CALL. `Ok(payload)` sends a RESULT; `Err(reason)`
27/// sends an ERROR (BOLT#4 `unknown_error`, `detail = reason`); a panic
28/// inside the handler (caught via [`tokio::spawn`], the same "one
29/// transient task per call" shape `macula_station_link.erl` uses one
30/// process per call for) is sent as ERROR `temporary_relay_failure` —
31/// matching that module's own `safe_invoke_handler/4` mapping exactly
32/// (including sending no `detail` on a crash, since the reference
33/// doesn't either — it only logs locally).
34pub type CallHandler =
35 Arc<dyn Fn(Value) -> BoxFuture<'static, Result<Value, String>> + Send + Sync>;
36
37/// Matches `HANDSHAKE_TIMEOUT_MS` in `macula_peering_conn.erl`: CONNECT
38/// -> HELLO is sub-second on a healthy peer; this is generous. The most
39/// common real-world trigger for hitting it is a protocol version
40/// mismatch — bytes accumulate but never form a valid frame, so the
41/// station-side symptom and this crate's symptom are the same shape.
42pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30);
43
44/// Default timeout for a single CALL awaiting its RESULT/ERROR. Not from
45/// the reference source (macula's own CALL timeout is caller-supplied
46/// per-call via `deadline_ms` inside the frame itself, not a transport-
47/// level default) — a reasonable local default for this crate's API.
48pub const DEFAULT_CALL_TIMEOUT: Duration = Duration::from_secs(30);
49
50/// Bound on a single read from a QUIC stream while accumulating a frame.
51/// Not a protocol limit — just how much to ask the stream for at once;
52/// `frame::decode`'s own `MAX_FRAME_BYTES` is the real cap.
53const READ_CHUNK: usize = 64 * 1024;
54
55// ---------------------------------------------------------------------
56// FrameStream — send/receive signed application frames on one QUIC
57// stream. The control stream (inside Session) and every dedicated
58// stream (content transfer, streaming RPC) are each one of these.
59// ---------------------------------------------------------------------
60
61pub struct FrameStream {
62 send: quinn::SendStream,
63 recv: quinn::RecvStream,
64 /// Bytes read but not yet consumed by a decoded frame — carried
65 /// over between reads so nothing is ever dropped.
66 buf: Vec<u8>,
67}
68
69#[derive(Debug)]
70pub enum SendFrameError {
71 Encode(frame::EncodeFrameError),
72 Write(quinn::WriteError),
73}
74
75impl std::fmt::Display for SendFrameError {
76 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77 match self {
78 SendFrameError::Encode(e) => write!(f, "encoding frame: {e}"),
79 SendFrameError::Write(e) => write!(f, "writing to stream: {e}"),
80 }
81 }
82}
83
84impl std::error::Error for SendFrameError {}
85
86#[derive(Debug)]
87pub enum RecvFrameError {
88 Read(quinn::ReadError),
89 StreamClosed,
90 Decode(frame::DecodeFrameError),
91 Timeout,
92}
93
94impl std::fmt::Display for RecvFrameError {
95 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
96 match self {
97 RecvFrameError::Read(e) => write!(f, "reading from stream: {e}"),
98 RecvFrameError::StreamClosed => write!(f, "peer closed the stream"),
99 RecvFrameError::Decode(e) => write!(f, "decoding a frame: {e}"),
100 RecvFrameError::Timeout => write!(f, "timed out waiting for a frame"),
101 }
102 }
103}
104
105impl std::error::Error for RecvFrameError {}
106
107#[derive(Debug)]
108pub enum CallError {
109 Send(SendFrameError),
110 Recv(RecvFrameError),
111}
112
113impl std::fmt::Display for CallError {
114 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
115 match self {
116 CallError::Send(e) => write!(f, "sending CALL: {e}"),
117 CallError::Recv(e) => write!(f, "awaiting RESULT/ERROR: {e}"),
118 }
119 }
120}
121
122impl std::error::Error for CallError {}
123
124impl FrameStream {
125 fn new(send: quinn::SendStream, recv: quinn::RecvStream) -> Self {
126 Self::with_buf(send, recv, Vec::new())
127 }
128
129 fn with_buf(send: quinn::SendStream, recv: quinn::RecvStream, buf: Vec<u8>) -> Self {
130 Self { send, recv, buf }
131 }
132
133 /// Any bytes already read past the last decoded frame — for the
134 /// control stream specifically, this starts with whatever was left
135 /// over from the handshake itself.
136 pub fn leftover_bytes(&self) -> &[u8] {
137 &self.buf
138 }
139
140 pub async fn send_frame(&mut self, frame: Value) -> Result<(), SendFrameError> {
141 let encoded = frame::encode(&frame).map_err(SendFrameError::Encode)?;
142 self.send
143 .write_all(&encoded)
144 .await
145 .map_err(SendFrameError::Write)
146 }
147
148 /// Read the next complete application frame, using (and updating)
149 /// any bytes already buffered.
150 pub async fn recv_frame(&mut self) -> Result<Value, RecvFrameError> {
151 let mut chunk = vec![0u8; READ_CHUNK];
152 loop {
153 match frame::decode(&self.buf) {
154 Ok(Decoded::Frame(value, consumed)) => {
155 self.buf.drain(..consumed);
156 return Ok(value);
157 }
158 Ok(Decoded::More(_)) => {}
159 Err(e) => return Err(RecvFrameError::Decode(e)),
160 }
161 let n = self
162 .recv
163 .read(&mut chunk)
164 .await
165 .map_err(RecvFrameError::Read)?
166 .ok_or(RecvFrameError::StreamClosed)?;
167 self.buf.extend_from_slice(&chunk[..n]);
168 }
169 }
170
171 /// As [`recv_frame`](Self::recv_frame), bounded by `timeout`.
172 pub async fn recv_frame_timeout(&mut self, timeout: Duration) -> Result<Value, RecvFrameError> {
173 tokio::time::timeout(timeout, self.recv_frame())
174 .await
175 .unwrap_or(Err(RecvFrameError::Timeout))
176 }
177
178 /// Send a signed CALL for `procedure` and wait for the matching
179 /// RESULT or ERROR, correlated by `call_id`.
180 ///
181 /// **Known v1 limitation (control stream only):** any frame that
182 /// arrives before the match (e.g. an EVENT from an active
183 /// SUBSCRIBE) is discarded, not queued or dispatched elsewhere —
184 /// correct for a client doing one thing at a time on the control
185 /// stream, not yet correct for CALL and PUBLISH/SUBSCRIBE used
186 /// concurrently on it. Harmless on a **dedicated** stream (content
187 /// transfer, streaming RPC), since nothing else ever arrives there
188 /// to discard.
189 pub async fn call(
190 &mut self,
191 procedure: &str,
192 realm: [u8; 32],
193 payload: Value,
194 deadline_ms: i128,
195 identity: &KeyPair,
196 timeout: Duration,
197 ) -> Result<frame::CallResponse, CallError> {
198 let call_id: [u8; 16] = rand::random();
199 let spec = frame::CallSpec::new(
200 call_id,
201 procedure,
202 realm,
203 payload,
204 deadline_ms,
205 identity.node_id(),
206 );
207 let signed = frame::sign(frame::call(&spec), identity);
208 self.send_frame(signed).await.map_err(CallError::Send)?;
209
210 tokio::time::timeout(timeout, self.await_call_response(call_id))
211 .await
212 .unwrap_or(Err(CallError::Recv(RecvFrameError::Timeout)))
213 }
214
215 /// As [`call`](Self::call), additionally attaching `ucan_token` to the
216 /// outgoing CALL frame — for invoking a procedure gated by a
217 /// [`crate::ucan::Policy::required`] policy. A procedure that isn't
218 /// gated ignores the token; one that is checks it (see
219 /// [`Session::serve_one_call_gated`]) before ever running its
220 /// handler, so an invalid/missing token comes back as a BOLT#4
221 /// `unauthorized` error frame, not a Rust error from this call.
222 ///
223 /// One parameter over [`call`](Self::call)'s own count, for the one
224 /// new thing this adds — same reasoning
225 /// [`crate::direct_dial::keep_advertised_direct`] already gives for
226 /// its own allow.
227 #[allow(clippy::too_many_arguments)]
228 pub async fn call_with_ucan(
229 &mut self,
230 procedure: &str,
231 realm: [u8; 32],
232 payload: Value,
233 deadline_ms: i128,
234 identity: &KeyPair,
235 timeout: Duration,
236 ucan_token: Vec<u8>,
237 ) -> Result<frame::CallResponse, CallError> {
238 let call_id: [u8; 16] = rand::random();
239 let mut spec = frame::CallSpec::new(
240 call_id,
241 procedure,
242 realm,
243 payload,
244 deadline_ms,
245 identity.node_id(),
246 );
247 spec.ucan_token = ucan_token;
248 let signed = frame::sign(frame::call(&spec), identity);
249 self.send_frame(signed).await.map_err(CallError::Send)?;
250
251 tokio::time::timeout(timeout, self.await_call_response(call_id))
252 .await
253 .unwrap_or(Err(CallError::Recv(RecvFrameError::Timeout)))
254 }
255
256 async fn await_call_response(
257 &mut self,
258 call_id: [u8; 16],
259 ) -> Result<frame::CallResponse, CallError> {
260 loop {
261 let value = self.recv_frame().await.map_err(CallError::Recv)?;
262 if frame::frame_call_id(&value) != Some(call_id) {
263 continue; // not ours — see call()'s doc on this limitation
264 }
265 if let Ok(response) = frame::parse_call_response(&value) {
266 return Ok(response);
267 }
268 // Matching call_id but not a result/error shape: keep
269 // waiting rather than erroring, since nothing else in the
270 // protocol is expected to carry this call's id.
271 }
272 }
273}
274
275// ---------------------------------------------------------------------
276// Session — the handshaked connection, wrapping the control stream.
277// ---------------------------------------------------------------------
278
279/// A completed, handshaked connection to a macula-station. Holds the
280/// open control stream (CONNECT/HELLO already exchanged) and the
281/// station's identity as verified by the HELLO frame's own signature.
282pub struct Session {
283 connection: quinn::Connection,
284 control: FrameStream,
285 pub station: HelloInfo,
286}
287
288#[derive(Debug)]
289pub enum HandshakeError {
290 Transport(ConnectError),
291 OpenStream(quinn::ConnectionError),
292 Write(quinn::WriteError),
293 Read(quinn::ReadError),
294 /// The peer closed the stream before a complete frame arrived.
295 StreamClosed,
296 Timeout,
297 Encode(frame::EncodeFrameError),
298 Decode(frame::DecodeFrameError),
299 /// Received a frame, but it wasn't a HELLO — a station is never
300 /// expected to send anything else at this point in the handshake.
301 UnexpectedFrameType(frame::ParseHelloError),
302 /// The HELLO frame's own signature didn't verify against the
303 /// node_id it claims — proves nothing about who actually sent it.
304 SignatureInvalid(frame::VerifyError),
305 /// The station completed the handshake but refused the connection
306 /// (`accepted = false`), e.g. a puzzle-invalid or unrecognized
307 /// identity.
308 Refused {
309 refusal_code: Option<i128>,
310 },
311}
312
313impl std::fmt::Display for HandshakeError {
314 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
315 match self {
316 HandshakeError::Transport(e) => write!(f, "transport: {e}"),
317 HandshakeError::OpenStream(e) => write!(f, "opening control stream: {e}"),
318 HandshakeError::Write(e) => write!(f, "sending CONNECT: {e}"),
319 HandshakeError::Read(e) => write!(f, "reading from control stream: {e}"),
320 HandshakeError::StreamClosed => {
321 write!(f, "station closed the stream before HELLO arrived")
322 }
323 HandshakeError::Timeout => write!(
324 f,
325 "no HELLO within {HANDSHAKE_TIMEOUT:?} (likely a protocol mismatch)"
326 ),
327 HandshakeError::Encode(e) => write!(f, "encoding CONNECT: {e}"),
328 HandshakeError::Decode(e) => write!(f, "decoding the station's response: {e}"),
329 HandshakeError::UnexpectedFrameType(e) => write!(f, "expected a HELLO frame: {e}"),
330 HandshakeError::SignatureInvalid(e) => write!(f, "HELLO signature check failed: {e}"),
331 HandshakeError::Refused { refusal_code } => {
332 write!(
333 f,
334 "station refused the connection (refusal_code = {refusal_code:?})"
335 )
336 }
337 }
338 }
339}
340
341impl std::error::Error for HandshakeError {}
342
343/// Dial `host:port` and complete the full CONNECT/HELLO handshake:
344/// open a QUIC connection, open the control stream, send a signed
345/// CONNECT built from `identity`, and wait for a HELLO whose own
346/// signature verifies against the node_id it claims.
347///
348/// `identity` **must** be puzzle-hardened
349/// ([`KeyPair::generate_with_puzzle`](crate::identity::KeyPair::generate_with_puzzle))
350/// — see that function's own doc and `plans/PLAN_WIRE_PROTOCOL.md` §5's
351/// callout: an unhardened identity fails this handshake silently (the
352/// QUIC/TLS layer looks healthy right up until the HELLO never accepts).
353pub async fn connect(
354 host: &str,
355 port: u16,
356 trust: Trust,
357 identity: &KeyPair,
358) -> Result<Session, HandshakeError> {
359 tokio::time::timeout(
360 HANDSHAKE_TIMEOUT,
361 connect_inner(host, port, trust, identity),
362 )
363 .await
364 .unwrap_or(Err(HandshakeError::Timeout))
365}
366
367async fn connect_inner(
368 host: &str,
369 port: u16,
370 trust: Trust,
371 identity: &KeyPair,
372) -> Result<Session, HandshakeError> {
373 let connection = transport::connect(host, port, trust)
374 .await
375 .map_err(HandshakeError::Transport)?;
376
377 let (mut send, mut recv) = connection
378 .open_bi()
379 .await
380 .map_err(HandshakeError::OpenStream)?;
381
382 let connect_spec =
383 crate::frame::ConnectSpec::new(identity.node_id(), identity.puzzle_evidence());
384 let connect_frame = frame::sign(frame::connect(&connect_spec), identity);
385 let encoded = frame::encode(&connect_frame).map_err(HandshakeError::Encode)?;
386 send.write_all(&encoded)
387 .await
388 .map_err(HandshakeError::Write)?;
389
390 let (hello_value, buf) = read_one_frame(&mut recv).await?;
391
392 let station = frame::parse_hello(&hello_value).map_err(HandshakeError::UnexpectedFrameType)?;
393 frame::verify(&hello_value, &station.node_id).map_err(HandshakeError::SignatureInvalid)?;
394
395 if !station.accepted {
396 return Err(HandshakeError::Refused {
397 refusal_code: station.refusal_code,
398 });
399 }
400
401 Ok(Session {
402 connection,
403 control: FrameStream::with_buf(send, recv, buf),
404 station,
405 })
406}
407
408/// Read from `recv` until one complete frame has been decoded, returning
409/// it along with any leftover bytes already read that belong to the
410/// *next* frame (so a caller can carry them forward instead of losing
411/// them). Handshake-only — [`FrameStream::recv_frame`] is the
412/// post-handshake equivalent.
413async fn read_one_frame(recv: &mut quinn::RecvStream) -> Result<(Value, Vec<u8>), HandshakeError> {
414 let mut buf = Vec::new();
415 let mut chunk = vec![0u8; READ_CHUNK];
416 loop {
417 match frame::decode(&buf) {
418 Ok(Decoded::Frame(value, consumed)) => {
419 buf.drain(..consumed);
420 return Ok((value, buf));
421 }
422 Ok(Decoded::More(_)) => {}
423 Err(e) => return Err(HandshakeError::Decode(e)),
424 }
425 let n = recv
426 .read(&mut chunk)
427 .await
428 .map_err(HandshakeError::Read)?
429 .ok_or(HandshakeError::StreamClosed)?;
430 buf.extend_from_slice(&chunk[..n]);
431 }
432}
433
434impl Session {
435 /// The remote address this session's connection is with.
436 pub fn remote_address(&self) -> std::net::SocketAddr {
437 self.connection.remote_address()
438 }
439
440 /// Open a new dedicated QUIC stream on this same connection, separate
441 /// from the control stream — the mechanism content transfer (§12)
442 /// and streaming RPC (§13) both use instead of the control stream.
443 pub async fn open_dedicated_stream(&mut self) -> Result<FrameStream, quinn::ConnectionError> {
444 let (send, recv) = self.connection.open_bi().await?;
445 Ok(FrameStream::new(send, recv))
446 }
447
448 /// Accept the next dedicated stream the *peer* opens toward us —
449 /// e.g. the station routing an inbound STREAM_OPEN for a procedure
450 /// this session has [`advertise`](Self::advertise)d (§13.2). Blocks
451 /// until one arrives.
452 ///
453 /// The receiving side has no advance notice of why a new stream
454 /// arrived; §7 of `plans/PLAN_WIRE_PROTOCOL.md` says to read the
455 /// stream's own first frame to learn its purpose, which is exactly
456 /// what a caller of this method does next via the returned
457 /// `FrameStream`'s own `recv_frame`. The reference (`quicer`-backed
458 /// Erlang) has a documented race here — the peer's first bytes can
459 /// arrive before the owning process is notified the stream exists at
460 /// all, because its NIF stream resources start passive and only
461 /// begin delivering once explicitly armed *after* the notification.
462 /// That race doesn't apply here: `quinn`/QUIC buffers inbound stream
463 /// data at the transport layer regardless of whether or when the
464 /// application starts reading, so nothing analogous to arm before
465 /// read is needed on this side.
466 pub async fn accept_dedicated_stream(&mut self) -> Result<FrameStream, quinn::ConnectionError> {
467 let (send, recv) = self.connection.accept_bi().await?;
468 Ok(FrameStream::new(send, recv))
469 }
470
471 /// Any bytes already read past the HELLO frame during the handshake
472 /// (belonging to whatever the station sent next) that a caller
473 /// building further protocol handling on top of this `Session`
474 /// should treat as already-received.
475 pub fn leftover_bytes(&self) -> &[u8] {
476 self.control.leftover_bytes()
477 }
478
479 /// Read the next complete application frame from the control stream.
480 pub async fn recv_frame(&mut self) -> Result<Value, RecvFrameError> {
481 self.control.recv_frame().await
482 }
483
484 /// As [`recv_frame`](Self::recv_frame), bounded by `timeout`.
485 pub async fn recv_frame_timeout(&mut self, timeout: Duration) -> Result<Value, RecvFrameError> {
486 self.control.recv_frame_timeout(timeout).await
487 }
488
489 /// Send a signed CALL on the control stream and wait for the
490 /// matching RESULT or ERROR — see [`FrameStream::call`]. Announces
491 /// `rpc.sent_v1`/`rpc.completed_v1` around the call — see
492 /// `announce_rpc_sent` for why these are always on.
493 pub async fn call(
494 &mut self,
495 procedure: &str,
496 realm: [u8; 32],
497 payload: Value,
498 deadline_ms: i128,
499 identity: &KeyPair,
500 timeout: Duration,
501 ) -> Result<frame::CallResponse, CallError> {
502 let request_id: [u8; 16] = rand::random();
503 announce_rpc_sent(&mut *self, realm, identity, request_id).await;
504 let result = self
505 .control
506 .call(procedure, realm, payload, deadline_ms, identity, timeout)
507 .await;
508 announce_rpc_completed(&mut *self, realm, identity, request_id, &result).await;
509 result
510 }
511
512 /// As [`call`](Self::call), attaching `ucan_token` (e.g. from
513 /// [`crate::ucan::create`]) to the outgoing CALL — for invoking a
514 /// procedure gated by a [`crate::ucan::Policy::required`] policy on
515 /// the provider side. See [`FrameStream::call_with_ucan`] for the
516 /// full contract. Announces `rpc.sent_v1`/`rpc.completed_v1` the same
517 /// way [`call`](Self::call) does.
518 #[allow(clippy::too_many_arguments)]
519 pub async fn call_with_ucan(
520 &mut self,
521 procedure: &str,
522 realm: [u8; 32],
523 payload: Value,
524 deadline_ms: i128,
525 identity: &KeyPair,
526 timeout: Duration,
527 ucan_token: Vec<u8>,
528 ) -> Result<frame::CallResponse, CallError> {
529 let request_id: [u8; 16] = rand::random();
530 announce_rpc_sent(&mut *self, realm, identity, request_id).await;
531 let result = self
532 .control
533 .call_with_ucan(
534 procedure,
535 realm,
536 payload,
537 deadline_ms,
538 identity,
539 timeout,
540 ucan_token,
541 )
542 .await;
543 announce_rpc_completed(&mut *self, realm, identity, request_id, &result).await;
544 result
545 }
546
547 /// Send a signed PUBLISH, carrying the end-to-end `publisher_sig`
548 /// (over topic/realm/publisher/seq/payload, independent of frame
549 /// type) so the resulting EVENT survives being relayed beyond one
550 /// hop — a station verifies an EVENT's per-hop `signature` against
551 /// whichever station forwarded it, which only matches on hop 1;
552 /// every hop after that needs `publisher_sig` instead. Matches the
553 /// Erlang reference SDK's own default (`pubsub_emit_publisher_sig`,
554 /// true since macula 4.6.0). Fire-and-forget — no reply is expected
555 /// on the wire; a subscriber (this session included, if subscribed
556 /// to the same topic/realm) receives an EVENT asynchronously, read
557 /// via [`recv_frame`](Self::recv_frame) /
558 /// [`recv_event`](Self::recv_event).
559 pub async fn publish(
560 &mut self,
561 spec: &frame::PublishSpec,
562 identity: &KeyPair,
563 ) -> Result<(), SendFrameError> {
564 let unsigned = frame::publish(spec);
565 let with_publisher_sig = frame::sign_publisher(unsigned, identity);
566 let signed = frame::sign(with_publisher_sig, identity);
567 self.control.send_frame(signed).await
568 }
569
570 /// Send a signed SUBSCRIBE. Fire-and-forget.
571 pub async fn subscribe(
572 &mut self,
573 spec: &frame::SubscribeSpec,
574 identity: &KeyPair,
575 ) -> Result<(), SendFrameError> {
576 let signed = frame::sign(frame::subscribe(spec), identity);
577 self.control.send_frame(signed).await
578 }
579
580 /// Send a signed UNSUBSCRIBE. Fire-and-forget.
581 pub async fn unsubscribe(
582 &mut self,
583 spec: &frame::UnsubscribeSpec,
584 identity: &KeyPair,
585 ) -> Result<(), SendFrameError> {
586 let signed = frame::sign(frame::unsubscribe(spec), identity);
587 self.control.send_frame(signed).await
588 }
589
590 /// Send a signed ADVERTISE (§6.9) — registers this connection as the
591 /// handler for `spec`'s `(realm, procedure)`. Fire-and-forget on the
592 /// wire; the station then routes inbound CALLs (control stream) and
593 /// STREAM_OPENs (a fresh dedicated stream — see
594 /// [`accept_dedicated_stream`](Self::accept_dedicated_stream)) for
595 /// that procedure back to this connection.
596 pub async fn advertise(
597 &mut self,
598 spec: &frame::AdvertiseSpec,
599 identity: &KeyPair,
600 ) -> Result<(), SendFrameError> {
601 let signed = frame::sign(frame::advertise(spec), identity);
602 self.control.send_frame(signed).await
603 }
604
605 /// Send a signed UNADVERTISE. Fire-and-forget.
606 pub async fn unadvertise(
607 &mut self,
608 spec: &frame::UnadvertiseSpec,
609 identity: &KeyPair,
610 ) -> Result<(), SendFrameError> {
611 let signed = frame::sign(frame::unadvertise(spec), identity);
612 self.control.send_frame(signed).await
613 }
614
615 /// Sends an ADVERTISE for `spec` immediately, then again every
616 /// `interval`, until `stop` resolves. [`advertise`](Self::advertise)'s
617 /// own doc notes the station's registration is tied to the connection
618 /// that sent it — a long-lived server needs to keep re-asserting it.
619 /// [`advertise`](Self::advertise) is a stateless, side-effect-free-on-
620 /// repeat wire send (unlike the Erlang reference's `advertise/5`, which
621 /// spawns a real per-call OTP supervisor and so needs a `reuse_sup`
622 /// option to avoid leaking one per tick), so there is nothing
623 /// equivalent to worry about leaking here — same reasoning
624 /// `macula-go`'s `KeepAdvertised` already applied and verified
625 /// live.
626 ///
627 /// A failed tick is reported via `on_error` but does not stop the
628 /// loop — it tries again at the next interval regardless. This cannot
629 /// detect or repair a dead session on its own; if the underlying
630 /// connection has actually gone down, every tick will keep failing
631 /// until `stop` resolves. See
632 /// [`crate::direct_dial::keep_advertised_direct`] for the direct-dial
633 /// equivalent (same shape, same reasoning).
634 pub async fn keep_advertised<F>(
635 &mut self,
636 spec: &frame::AdvertiseSpec,
637 identity: &KeyPair,
638 interval: Duration,
639 stop: F,
640 on_error: impl Fn(SendFrameError),
641 ) where
642 F: std::future::Future<Output = ()>,
643 {
644 tokio::pin!(stop);
645 let mut ticker = tokio::time::interval(interval);
646 loop {
647 tokio::select! {
648 _ = &mut stop => return,
649 _ = ticker.tick() => {
650 if let Err(e) = self.advertise(spec, identity).await {
651 on_error(e);
652 }
653 }
654 }
655 }
656 }
657
658 /// Read the next frame and parse it as an EVENT, bounded by
659 /// `timeout`. Any non-EVENT frame received first is an error, not
660 /// silently skipped — unlike [`call`](Self::call)'s response wait,
661 /// a caller waiting specifically for a pubsub delivery has no reason
662 /// to expect anything else to legitimately arrive first.
663 pub async fn recv_event(
664 &mut self,
665 timeout: Duration,
666 ) -> Result<frame::EventInfo, RecvEventError> {
667 let value = self
668 .control
669 .recv_frame_timeout(timeout)
670 .await
671 .map_err(RecvEventError::Recv)?;
672 frame::parse_event(&value).map_err(RecvEventError::Parse)
673 }
674
675 /// The provider role's counterpart to [`call`](Self::call): block
676 /// for the next inbound CALL frame on the control stream, bounded
677 /// by `timeout`, look it up via `lookup`, invoke the matching
678 /// handler, and send the resulting RESULT or ERROR back over this
679 /// same connection — see `plans/PLAN_WIRE_PROTOCOL.md` §6.9's
680 /// routing description and `macula_station_link.erl`'s
681 /// `handle_inbound_call/2`, which this mirrors field for field,
682 /// including its BOLT#4 error-code mapping.
683 ///
684 /// Any non-CALL frame that arrives first (e.g. a stray EVENT from
685 /// an active [`subscribe`](Self::subscribe), or a RESULT/ERROR for
686 /// some other in-flight [`call`](Self::call)) is discarded, not
687 /// queued — the same "control stream, one thing at a time"
688 /// limitation [`call`](Self::call)'s own doc already carries. A
689 /// session that needs to serve CALLs and also act as a
690 /// caller/subscriber concurrently should use a second `Session`,
691 /// exactly like this crate's own streaming-provider live test does.
692 ///
693 /// A caller wanting a long-lived server loops on this:
694 ///
695 /// ```no_run
696 /// # use std::time::Duration;
697 /// # async fn example(session: &mut macula_rust::connection::Session, identity: &macula_rust::identity::KeyPair, lookup: impl Fn(&[u8; 32], &str) -> Option<macula_rust::connection::CallHandler>) {
698 /// loop {
699 /// if let Err(e) = session.serve_one_call(&lookup, identity, Duration::from_secs(30)).await {
700 /// // ServeCallError::Timeout just means nothing arrived -- keep looping.
701 /// eprintln!("{e}");
702 /// }
703 /// }
704 /// # }
705 /// ```
706 pub async fn serve_one_call<L>(
707 &mut self,
708 lookup: L,
709 identity: &KeyPair,
710 timeout: Duration,
711 ) -> Result<(), ServeCallError>
712 where
713 L: Fn(&[u8; 32], &str) -> Option<CallHandler>,
714 {
715 self.serve_one_call_gated(
716 lookup,
717 |_, _| crate::ucan::Policy::open(),
718 identity,
719 timeout,
720 )
721 .await
722 }
723
724 /// [`serve_one_call`](Self::serve_one_call), additionally gating each
725 /// inbound CALL through `policy` BEFORE `lookup` runs — mirrors
726 /// `macula_station_link.erl`'s `handle_inbound_call/2` exactly: an
727 /// open policy (the default [`serve_one_call`](Self::serve_one_call)
728 /// uses) behaves identically; a [`crate::ucan::Policy::required`]
729 /// policy demands a CALL's `ucan_token` verify against the required
730 /// issuer, and refuses with BOLT#4 `unauthorized` WITHOUT ever
731 /// invoking `lookup` or a handler if it doesn't — a [`CallHandler`]
732 /// never sees the raw token either way, matching the reference's own
733 /// handler contract (payload only).
734 pub async fn serve_one_call_gated<L, P>(
735 &mut self,
736 lookup: L,
737 policy: P,
738 identity: &KeyPair,
739 timeout: Duration,
740 ) -> Result<(), ServeCallError>
741 where
742 L: Fn(&[u8; 32], &str) -> Option<CallHandler>,
743 P: Fn(&[u8; 32], &str) -> crate::ucan::Policy,
744 {
745 tokio::time::timeout(
746 timeout,
747 self.serve_one_call_gated_inner(lookup, policy, identity),
748 )
749 .await
750 .unwrap_or(Err(ServeCallError::Timeout))
751 }
752
753 async fn serve_one_call_gated_inner<L, P>(
754 &mut self,
755 lookup: L,
756 policy: P,
757 identity: &KeyPair,
758 ) -> Result<(), ServeCallError>
759 where
760 L: Fn(&[u8; 32], &str) -> Option<CallHandler>,
761 P: Fn(&[u8; 32], &str) -> crate::ucan::Policy,
762 {
763 loop {
764 let value = self
765 .control
766 .recv_frame()
767 .await
768 .map_err(ServeCallError::Recv)?;
769 let Ok(call_info) = frame::parse_call(&value) else {
770 continue; // not ours -- see this method's doc on the limitation
771 };
772 let reply =
773 build_call_reply(call_info, &lookup, &policy, identity, Some(&mut *self)).await;
774 let signed = frame::sign(reply, identity);
775 self.control
776 .send_frame(signed)
777 .await
778 .map_err(ServeCallError::Send)?;
779 return Ok(());
780 }
781 }
782
783 /// Bounds how long [`close`](Self::close) waits after its last write
784 /// before hard-closing the connection -- see that method's own doc
785 /// for why this exists at all. Short relative to the Erlang
786 /// reference's own 5s draining-state upper bound
787 /// (`macula_peering.erl`, `?DRAIN_TIMEOUT_MS`): this side only needs
788 /// to cover quinn's own internal send-scheduling latency, not a full
789 /// round trip's worth of protocol drain.
790 const CLOSE_DRAIN: Duration = Duration::from_millis(250);
791
792 /// Close the control stream and connection gracefully with a GOODBYE
793 /// frame, matching `macula_peering_conn.erl`'s `connected -> draining`
794 /// transition (minus the full drain-timeout bookkeeping, since this
795 /// crate isn't holding a supervisor to clean up).
796 ///
797 /// `write_all(...).await` and `finish()` both only guarantee the data
798 /// was handed to quinn's own send-scheduling machinery, not that it
799 /// reached the peer -- `Connection::close` is abrupt and does not
800 /// wait for outstanding stream data to be delivered. Found live
801 /// 2026-08-29 in the Go port of this exact pattern
802 /// (macula-go's connection.Session.Close): a PUBLISH sent
803 /// immediately before Close intermittently never reached the peer,
804 /// root-caused to this race. Fixed proactively here before it was
805 /// independently rediscovered against this crate -- same doc
806 /// comment ("minus the drain-timeout bookkeeping"), same
807 /// write-then-immediately-abort-connection shape, so the same race
808 /// applies. Closing the stream via `finish()` first, then giving the
809 /// background sender a bounded window before hard-closing the
810 /// connection, mirrors the Erlang reference's own bounded-drain
811 /// approach.
812 pub async fn close(mut self, reason: &str, detail: Option<&str>, identity: &KeyPair) {
813 let goodbye = frame::sign(frame::goodbye(reason, detail), identity);
814 if let Ok(encoded) = frame::encode(&goodbye) {
815 let _ = self.control.send.write_all(&encoded).await;
816 }
817 let _ = self.control.send.finish();
818 tokio::time::sleep(Self::CLOSE_DRAIN).await;
819 self.connection.close(0u32.into(), reason.as_bytes());
820 }
821
822 /// The supervised counterpart to the bare [`publish`](Self::publish)
823 /// primitive, matching `macula_publisher.erl` in spirit: publishes
824 /// `pubsub.publish_started_v1` before the publish and
825 /// `pubsub.publish_completed_v1` after, both under `spec`'s own realm.
826 /// Fact-publish failures are silently discarded — matching
827 /// `macula_publisher.erl`'s own `publish/5` helper, which throws away
828 /// its result unconditionally (`_ = macula:publish(...), ok`).
829 ///
830 /// Unlike Erlang's version — a supervised worker process a caller can
831 /// kill mid-flight — this crate's bare `publish` is already a
832 /// synchronous, near-instant frame send (no ack on this wire, no
833 /// network round-trip to await), so there is no meaningful "cancel
834 /// before it starts" window worth a dedicated mechanism. Await this
835 /// directly, or wrap it in `tokio::select!`/`tokio::time::timeout`
836 /// yourself if you need to abandon it early — dropping a `Future` IS
837 /// real cancellation in Rust; Erlang has to simulate that by killing a
838 /// worker process.
839 pub async fn run_publisher(
840 &mut self,
841 spec: &frame::PublishSpec,
842 identity: &KeyPair,
843 announce: bool,
844 ) -> Result<(), SendFrameError> {
845 let publish_id: [u8; 16] = rand::random();
846 if announce {
847 let payload = Value::Map(vec![])
848 .with_field("publish_id", Value::Bytes(publish_id.to_vec()))
849 .with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()));
850 let fact = frame::PublishSpec::new(
851 "pubsub.publish_started_v1",
852 spec.realm,
853 identity.node_id(),
854 rand::random(),
855 payload,
856 now_ms(),
857 );
858 let _ = self.publish(&fact, identity).await;
859 }
860
861 let result = self.publish(spec, identity).await;
862
863 if announce {
864 let payload =
865 Value::Map(vec![]).with_field("publish_id", Value::Bytes(publish_id.to_vec()));
866 let payload = match &result {
867 Ok(()) => payload.with_field("outcome", Value::text("completed")),
868 Err(e) => payload
869 .with_field("outcome", Value::text("failed"))
870 .with_field("reason", Value::text(e.to_string())),
871 };
872 let fact = frame::PublishSpec::new(
873 "pubsub.publish_completed_v1",
874 spec.realm,
875 identity.node_id(),
876 rand::random(),
877 payload,
878 now_ms(),
879 );
880 let _ = self.publish(&fact, identity).await;
881 }
882
883 result
884 }
885
886 /// The supervised counterpart to the bare
887 /// [`subscribe`](Self::subscribe)/[`recv_event`](Self::recv_event)
888 /// primitives, matching `macula_subscriber.erl` in spirit: subscribes
889 /// once, then dispatches every inbound EVENT to `handler` until `stop`
890 /// resolves. Unsubscribes on return, including on cancellation.
891 ///
892 /// Mirrors [`serve_one_call`](Self::serve_one_call)'s own frame loop,
893 /// not [`recv_event`](Self::recv_event): a shared control stream can
894 /// carry other frame types between one EVENT and the next, so a
895 /// wrong-frame-type parse failure is skipped and polling continues,
896 /// exactly like `serve_one_call` skips a non-"call" frame — it is NOT
897 /// treated as fatal the way `recv_event`'s own contract treats any
898 /// parse failure. Confirmed live in the Go port of this exact pattern
899 /// (`macula-go`'s `Session.RunSubscriber`): without this, a single
900 /// non-EVENT frame arriving on the control stream aborted the whole
901 /// subscriber loop.
902 ///
903 /// No OTP pid to address a running subscriber by; `stop` plays that
904 /// role — matches [`keep_advertised`](Self::keep_advertised)'s own
905 /// cancellation shape exactly, not a new one. `handler` cannot itself
906 /// stop the loop (no return value) — by the same design `keep_advertised`
907 /// already established, where `on_error` can only report, not halt;
908 /// stopping is always external, via `stop`.
909 pub async fn run_subscriber<F>(
910 &mut self,
911 spec: &frame::SubscribeSpec,
912 identity: &KeyPair,
913 stop: F,
914 mut handler: impl FnMut(frame::EventInfo),
915 ) -> Result<(), RunSubscriberError>
916 where
917 F: std::future::Future<Output = ()>,
918 {
919 self.subscribe(spec, identity)
920 .await
921 .map_err(RunSubscriberError::Subscribe)?;
922
923 tokio::pin!(stop);
924 let result = loop {
925 tokio::select! {
926 _ = &mut stop => break Ok(()),
927 frame_result = self.control.recv_frame() => {
928 let value = match frame_result {
929 Ok(v) => v,
930 Err(e) => break Err(RunSubscriberError::Recv(e)),
931 };
932 let Ok(evt) = frame::parse_event(&value) else {
933 continue; // not ours -- see this method's doc on the limitation
934 };
935 handler(evt);
936 }
937 }
938 };
939
940 let _ = self
941 .unsubscribe(
942 &frame::UnsubscribeSpec::new(spec.topic.clone(), spec.realm, spec.subscriber),
943 identity,
944 )
945 .await;
946
947 result
948 }
949}
950
951fn now_ms() -> u64 {
952 std::time::SystemTime::now()
953 .duration_since(std::time::UNIX_EPOCH)
954 .expect("system clock before 1970")
955 .as_millis() as u64
956}
957
958#[derive(Debug)]
959pub enum RecvEventError {
960 Recv(RecvFrameError),
961 Parse(frame::ParseEventError),
962}
963
964impl std::fmt::Display for RecvEventError {
965 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
966 match self {
967 RecvEventError::Recv(e) => write!(f, "{e}"),
968 RecvEventError::Parse(e) => write!(f, "expected an EVENT frame: {e}"),
969 }
970 }
971}
972
973impl std::error::Error for RecvEventError {}
974
975#[derive(Debug)]
976pub enum ServeCallError {
977 Recv(RecvFrameError),
978 Send(SendFrameError),
979 /// No inbound CALL arrived within the requested timeout.
980 Timeout,
981}
982
983impl std::fmt::Display for ServeCallError {
984 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
985 match self {
986 ServeCallError::Recv(e) => write!(f, "{e}"),
987 ServeCallError::Send(e) => write!(f, "sending the reply: {e}"),
988 ServeCallError::Timeout => write!(f, "timed out waiting for an inbound CALL"),
989 }
990 }
991}
992
993impl std::error::Error for ServeCallError {}
994
995/// Errors from [`Session::run_subscriber`].
996#[derive(Debug)]
997pub enum RunSubscriberError {
998 Subscribe(SendFrameError),
999 Recv(RecvFrameError),
1000}
1001
1002impl std::fmt::Display for RunSubscriberError {
1003 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1004 match self {
1005 RunSubscriberError::Subscribe(e) => write!(f, "subscribing: {e}"),
1006 RunSubscriberError::Recv(e) => write!(f, "{e}"),
1007 }
1008 }
1009}
1010
1011impl std::error::Error for RunSubscriberError {}
1012
1013/// Build the RESULT/ERROR reply for one inbound CALL — mirrors
1014/// `macula_station_link.erl`'s `handle_inbound_call/2` +
1015/// `safe_invoke_handler/4` exactly: `policy` is checked FIRST (a
1016/// rejection is BOLT#4 `unauthorized`, and `lookup`/a handler never run
1017/// at all); then a lookup miss is `unknown_next_peer`; the handler
1018/// running to completion produces a RESULT (`Ok`) or `unknown_error` with
1019/// `detail` (`Err`); a handler panic — caught via `tokio::spawn`, the
1020/// same "one transient task per call" shape the reference's own "one
1021/// process per call" uses — is `temporary_relay_failure`, with no
1022/// `detail`, matching the reference not sending one on a crash either.
1023// RPC telemetry auto-facts, matching `macula_request.erl` (caller side:
1024// rpc.sent_v1/rpc.completed_v1) and `macula_response.erl` (provider side:
1025// rpc.received_v1/rpc.replied_v1) exactly -- same topic names, same
1026// `request_id` field (16 fresh random bytes per call, independent of the
1027// wire CALL frame's own `call_id` -- the reference tracks its own request
1028// lifecycle separately from the wire frame, and this does too), same realm
1029// as the call itself, fire-and-forget (a fact-publish failure here never
1030// fails the underlying `call`/`serve_one_call_gated`, matching
1031// `macula_response.erl`'s own `_ = macula:publish(...), ok` and
1032// `macula_request.erl`'s identical `publish/5` helper -- same pattern this
1033// crate's own `run_publisher` already uses for its pubsub facts).
1034//
1035// Always on, matching the reference's ACTUAL behavior on each side, not
1036// just a blanket claim -- checked directly rather than assumed:
1037// `macula_request.erl`'s `start_link/7` and `start_link_direct/8` both
1038// hardcode `true` literally at the tuple-construction call site; there is
1039// no `Opts` key or parameter that reaches it at all on the caller side.
1040// `macula_response.erl`'s `advertise/6` DOES read `announce` from its
1041// `Opts` map with a `true` default (`maps:get(announce, Opts, true)`) --
1042// technically overridable -- but the one real caller in this workspace
1043// (`hecate_om_capabilities.erl`'s `advertise_opts/1`) never sets it to
1044// `false`. Matching Go's `macula-go` decision here: no toggle exposed
1045// on either side, since exposing one on `call`/`serve_one_call_gated` --
1046// this crate's two most heavily used functions -- for an option nothing
1047// in the reference ecosystem actually flips would be a real-blast-radius
1048// signature change for no practical benefit.
1049const RPC_SENT_TOPIC: &str = "rpc.sent_v1";
1050const RPC_COMPLETED_TOPIC: &str = "rpc.completed_v1";
1051const RPC_RECEIVED_TOPIC: &str = "rpc.received_v1";
1052const RPC_REPLIED_TOPIC: &str = "rpc.replied_v1";
1053
1054fn request_id_payload(request_id: [u8; 16]) -> Value {
1055 Value::Map(vec![]).with_field("request_id", Value::Bytes(request_id.to_vec()))
1056}
1057
1058async fn announce_fact(
1059 session: &mut Session,
1060 realm: [u8; 32],
1061 identity: &KeyPair,
1062 topic: &str,
1063 payload: Value,
1064) {
1065 let fact = frame::PublishSpec::new(
1066 topic,
1067 realm,
1068 identity.node_id(),
1069 rand::random(),
1070 payload,
1071 now_ms(),
1072 );
1073 let _ = session.publish(&fact, identity).await;
1074}
1075
1076async fn announce_rpc_sent(
1077 session: &mut Session,
1078 realm: [u8; 32],
1079 identity: &KeyPair,
1080 request_id: [u8; 16],
1081) {
1082 announce_fact(
1083 session,
1084 realm,
1085 identity,
1086 RPC_SENT_TOPIC,
1087 request_id_payload(request_id),
1088 )
1089 .await;
1090}
1091
1092/// Matches `macula_request.erl`'s `outcome_fields/2`: `completed` (no Rust
1093/// error, not a bolt4 ERROR frame) or `failed` (either). Erlang
1094/// additionally has a `cancelled` outcome from its own
1095/// gen_server-cancellable `macula_request:cancel/1` -- this crate's plain
1096/// `call` has no cancellation concept independent of an ordinary
1097/// error/timeout at this layer, so that outcome is not reachable here and
1098/// is not fabricated (same reasoning Go's port already documented).
1099async fn announce_rpc_completed(
1100 session: &mut Session,
1101 realm: [u8; 32],
1102 identity: &KeyPair,
1103 request_id: [u8; 16],
1104 result: &Result<frame::CallResponse, CallError>,
1105) {
1106 let payload = request_id_payload(request_id);
1107 let payload = match result {
1108 Err(e) => payload
1109 .with_field("outcome", Value::text("failed"))
1110 .with_field("reason", Value::text(e.to_string())),
1111 Ok(frame::CallResponse::Error { name, .. }) => payload
1112 .with_field("outcome", Value::text("failed"))
1113 .with_field("reason", Value::text(name.clone())),
1114 Ok(frame::CallResponse::Result { .. }) => {
1115 payload.with_field("outcome", Value::text("completed"))
1116 }
1117 };
1118 announce_fact(session, realm, identity, RPC_COMPLETED_TOPIC, payload).await;
1119}
1120
1121async fn announce_rpc_received(
1122 session: &mut Session,
1123 realm: [u8; 32],
1124 identity: &KeyPair,
1125 request_id: [u8; 16],
1126) {
1127 announce_fact(
1128 session,
1129 realm,
1130 identity,
1131 RPC_RECEIVED_TOPIC,
1132 request_id_payload(request_id),
1133 )
1134 .await;
1135}
1136
1137/// Matches `macula_response.erl`'s `outcome_fields/2`: `replied` (`{ok,
1138/// _}`) or `failed` (`{error, Reason}`). A handler panic is deliberately
1139/// NOT announced here at all -- matching the reference exactly, where a
1140/// crashing `Module:handle_request/2` crashes the whole per-request child
1141/// before its own `publish_replied/2` call is ever reached, so
1142/// `REQUEST_REPLIED` is never published for a crash there either.
1143async fn announce_rpc_replied(
1144 session: &mut Session,
1145 realm: [u8; 32],
1146 identity: &KeyPair,
1147 request_id: [u8; 16],
1148 handler_err: Option<&str>,
1149) {
1150 let payload = request_id_payload(request_id);
1151 let payload = match handler_err {
1152 Some(reason) => payload
1153 .with_field("outcome", Value::text("failed"))
1154 .with_field("reason", Value::text(reason)),
1155 None => payload.with_field("outcome", Value::text("replied")),
1156 };
1157 announce_fact(session, realm, identity, RPC_REPLIED_TOPIC, payload).await;
1158}
1159
1160/// Fires `rpc.received_v1`/`rpc.replied_v1` around dispatch when `session`
1161/// is `Some` -- `None` for the pure dispatch-logic unit tests in
1162/// `ucan_gating_tests` below, which deliberately exercise this function
1163/// with no network at all (mirrors `macula-go`'s identical
1164/// nil-session-safe `announceFact`). `rpc.received_v1` fires only after
1165/// `policy` and `lookup` both pass, matching `macula_response.erl`'s own
1166/// per-request child only starting once the raw advertise mechanism
1167/// already decided to dispatch to a real handler -- a UCAN-rejected or
1168/// unadvertised-procedure CALL announces neither fact.
1169///
1170/// `#[allow(needless_option_as_deref)]`: the lint is right that
1171/// `Option<&mut Session>::as_deref_mut()`'s RETURN type is identical to
1172/// its input type, but wrong that the call is needless here -- three
1173/// separate call sites below each need their own short-lived reborrow
1174/// from the same `session` local; using `session` directly at any one of
1175/// them would move it out for the rest of the function, breaking the
1176/// other two.
1177#[allow(clippy::needless_option_as_deref)]
1178async fn build_call_reply<L, P>(
1179 call_info: frame::CallInfo,
1180 lookup: &L,
1181 policy: &P,
1182 identity: &KeyPair,
1183 mut session: Option<&mut Session>,
1184) -> Value
1185where
1186 L: Fn(&[u8; 32], &str) -> Option<CallHandler>,
1187 P: Fn(&[u8; 32], &str) -> crate::ucan::Policy,
1188{
1189 let self_pub = identity.node_id();
1190
1191 if policy(&call_info.realm, &call_info.procedure)
1192 .check(&call_info.ucan_token)
1193 .is_err()
1194 {
1195 return frame::call_error(&frame::CallErrorSpec::new(
1196 call_info.call_id,
1197 bolt4::Code::Unauthorized,
1198 self_pub,
1199 ));
1200 }
1201
1202 let Some(handler) = lookup(&call_info.realm, &call_info.procedure) else {
1203 return frame::call_error(&frame::CallErrorSpec::new(
1204 call_info.call_id,
1205 bolt4::Code::UnknownNextPeer,
1206 self_pub,
1207 ));
1208 };
1209
1210 let request_id: [u8; 16] = rand::random();
1211 if let Some(s) = session.as_deref_mut() {
1212 announce_rpc_received(s, call_info.realm, identity, request_id).await;
1213 }
1214
1215 let payload = call_info.payload;
1216 let outcome = tokio::spawn(async move { handler(payload).await }).await;
1217 match outcome {
1218 Ok(Ok(value)) => {
1219 if let Some(s) = session.as_deref_mut() {
1220 announce_rpc_replied(s, call_info.realm, identity, request_id, None).await;
1221 }
1222 frame::result(&frame::ResultSpec::new(call_info.call_id, value, self_pub))
1223 }
1224 Ok(Err(reason)) => {
1225 if let Some(s) = session.as_deref_mut() {
1226 announce_rpc_replied(s, call_info.realm, identity, request_id, Some(&reason)).await;
1227 }
1228 let mut spec =
1229 frame::CallErrorSpec::new(call_info.call_id, bolt4::Code::UnknownError, self_pub);
1230 spec.detail = Some(reason);
1231 frame::call_error(&spec)
1232 }
1233 Err(_join_error) => frame::call_error(&frame::CallErrorSpec::new(
1234 call_info.call_id,
1235 bolt4::Code::TemporaryRelayFailure,
1236 self_pub,
1237 )),
1238 }
1239}
1240
1241#[cfg(test)]
1242mod ucan_gating_tests {
1243 //! Proves `serve_one_call_gated`'s policy wiring end-to-end WITHOUT a
1244 //! network — `build_call_reply` is a plain async function of
1245 //! `(CallInfo, lookup, policy, self_pub)`, so its dispatch/reply logic
1246 //! is fully testable in isolation. Mirrors `macula-go`'s own 4
1247 //! connection-level UCAN-gating unit tests (`serve_ucan_test.go`).
1248 use super::*;
1249 use crate::identity::KeyPair;
1250 use crate::ucan::{self, Policy};
1251
1252 fn call_info(ucan_token: Vec<u8>) -> frame::CallInfo {
1253 frame::CallInfo {
1254 call_id: [1; 16],
1255 procedure: "test.proc".into(),
1256 realm: [0; 32],
1257 payload: Value::Null,
1258 deadline_ms: 0,
1259 caller: [2; 32],
1260 ucan_token,
1261 }
1262 }
1263
1264 fn never_called_lookup() -> impl Fn(&[u8; 32], &str) -> Option<CallHandler> {
1265 |_, _| panic!("handler lookup must not run when policy rejects the call")
1266 }
1267
1268 fn echo_lookup() -> impl Fn(&[u8; 32], &str) -> Option<CallHandler> {
1269 |_, _| {
1270 Some(Arc::new(|payload: Value| {
1271 Box::pin(async move { Ok(payload) })
1272 }))
1273 }
1274 }
1275
1276 #[tokio::test]
1277 async fn open_policy_never_gates_dispatch() {
1278 let identity = KeyPair::generate();
1279 let reply = build_call_reply(
1280 call_info(Vec::new()),
1281 &echo_lookup(),
1282 &|_, _| Policy::open(),
1283 &identity,
1284 None,
1285 )
1286 .await;
1287 assert!(matches!(
1288 frame::parse_call_response(&reply),
1289 Ok(frame::CallResponse::Result { .. })
1290 ));
1291 }
1292
1293 #[tokio::test]
1294 async fn required_policy_refuses_a_call_with_no_token_before_lookup_runs() {
1295 let id = KeyPair::generate();
1296 let identity = KeyPair::generate();
1297 let reply = build_call_reply(
1298 call_info(Vec::new()),
1299 &never_called_lookup(),
1300 &move |_, _| Policy::required(id.node_id()),
1301 &identity,
1302 None,
1303 )
1304 .await;
1305 match frame::parse_call_response(&reply) {
1306 Ok(frame::CallResponse::Error { code, .. }) => {
1307 assert_eq!(code, bolt4::Code::Unauthorized as u8)
1308 }
1309 other => panic!("expected an Unauthorized ERROR frame, got {other:?}"),
1310 }
1311 }
1312
1313 #[tokio::test]
1314 async fn required_policy_refuses_a_token_from_the_wrong_issuer_before_lookup_runs() {
1315 let required_issuer = KeyPair::generate();
1316 let impostor = KeyPair::generate();
1317 let bad_token = ucan::create(
1318 "did:iss",
1319 "did:aud",
1320 vec![],
1321 &impostor,
1322 ucan::CreateOpts::default(),
1323 )
1324 .unwrap();
1325 let identity = KeyPair::generate();
1326 let reply = build_call_reply(
1327 call_info(bad_token),
1328 &never_called_lookup(),
1329 &move |_, _| Policy::required(required_issuer.node_id()),
1330 &identity,
1331 None,
1332 )
1333 .await;
1334 match frame::parse_call_response(&reply) {
1335 Ok(frame::CallResponse::Error { code, .. }) => {
1336 assert_eq!(code, bolt4::Code::Unauthorized as u8)
1337 }
1338 other => panic!("expected an Unauthorized ERROR frame, got {other:?}"),
1339 }
1340 }
1341
1342 #[tokio::test]
1343 async fn required_policy_lets_a_valid_token_reach_the_handler() {
1344 let id = KeyPair::generate();
1345 let good_token = ucan::create(
1346 "did:iss",
1347 "did:aud",
1348 vec![],
1349 &id,
1350 ucan::CreateOpts::default(),
1351 )
1352 .unwrap();
1353 let identity = KeyPair::generate();
1354 let reply = build_call_reply(
1355 call_info(good_token),
1356 &echo_lookup(),
1357 &move |_, _| Policy::required(id.node_id()),
1358 &identity,
1359 None,
1360 )
1361 .await;
1362 assert!(matches!(
1363 frame::parse_call_response(&reply),
1364 Ok(frame::CallResponse::Result { .. })
1365 ));
1366 }
1367}