Skip to main content

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}