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