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}