Skip to main content

Session

Struct Session 

Source
pub struct Session {
    pub station: HelloInfo,
    /* private fields */
}
Expand description

A completed, handshaked connection to a macula-station. Holds the open control stream (CONNECT/HELLO already exchanged) and the station’s identity as verified by the HELLO frame’s own signature.

Fields§

§station: HelloInfo

Implementations§

Source§

impl Session

Source

pub fn remote_address(&self) -> SocketAddr

The remote address this session’s connection is with.

Source

pub async fn open_dedicated_stream( &mut self, ) -> Result<FrameStream, ConnectionError>

Open a new dedicated QUIC stream on this same connection, separate from the control stream — the mechanism content transfer (§12) and streaming RPC (§13) both use instead of the control stream.

Source

pub async fn accept_dedicated_stream( &mut self, ) -> Result<FrameStream, ConnectionError>

Accept the next dedicated stream the peer opens toward us — e.g. the station routing an inbound STREAM_OPEN for a procedure this session has advertised (§13.2). Blocks until one arrives.

The receiving side has no advance notice of why a new stream arrived; §7 of plans/PLAN_WIRE_PROTOCOL.md says to read the stream’s own first frame to learn its purpose, which is exactly what a caller of this method does next via the returned FrameStream’s own recv_frame. The reference (quicer-backed Erlang) has a documented race here — the peer’s first bytes can arrive before the owning process is notified the stream exists at all, because its NIF stream resources start passive and only begin delivering once explicitly armed after the notification. That race doesn’t apply here: quinn/QUIC buffers inbound stream data at the transport layer regardless of whether or when the application starts reading, so nothing analogous to arm before read is needed on this side.

Source

pub fn leftover_bytes(&self) -> &[u8] ⓘ

Any bytes already read past the HELLO frame during the handshake (belonging to whatever the station sent next) that a caller building further protocol handling on top of this Session should treat as already-received.

Source

pub async fn recv_frame(&mut self) -> Result<Value, RecvFrameError>

Read the next complete application frame from the control stream.

Source

pub async fn recv_frame_timeout( &mut self, timeout: Duration, ) -> Result<Value, RecvFrameError>

As recv_frame, bounded by timeout.

Source

pub async fn call( &mut self, procedure: &str, realm: [u8; 32], payload: Value, deadline_ms: i128, identity: &KeyPair, timeout: Duration, ) -> Result<CallResponse, CallError>

Send a signed CALL on the control stream and wait for the matching RESULT or ERROR — see FrameStream::call. Announces rpc.sent_v1/rpc.completed_v1 around the call — see announce_rpc_sent for why these are always on.

Source

pub async fn call_with_ucan( &mut self, procedure: &str, realm: [u8; 32], payload: Value, deadline_ms: i128, identity: &KeyPair, timeout: Duration, ucan_token: Vec<u8>, ) -> Result<CallResponse, CallError>

As call, attaching ucan_token (e.g. from crate::ucan::create) to the outgoing CALL — for invoking a procedure gated by a crate::ucan::Policy::required policy on the provider side. See FrameStream::call_with_ucan for the full contract. Announces rpc.sent_v1/rpc.completed_v1 the same way call does.

Source

pub async fn publish( &mut self, spec: &PublishSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>

Send a signed PUBLISH, carrying the end-to-end publisher_sig (over topic/realm/publisher/seq/payload, independent of frame type) so the resulting EVENT survives being relayed beyond one hop — a station verifies an EVENT’s per-hop signature against whichever station forwarded it, which only matches on hop 1; every hop after that needs publisher_sig instead. Matches the Erlang reference SDK’s own default (pubsub_emit_publisher_sig, true since macula 4.6.0). Fire-and-forget — no reply is expected on the wire; a subscriber (this session included, if subscribed to the same topic/realm) receives an EVENT asynchronously, read via recv_frame / recv_event.

Source

pub async fn subscribe( &mut self, spec: &SubscribeSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>

Send a signed SUBSCRIBE. Fire-and-forget.

Source

pub async fn unsubscribe( &mut self, spec: &UnsubscribeSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>

Send a signed UNSUBSCRIBE. Fire-and-forget.

Source

pub async fn advertise( &mut self, spec: &AdvertiseSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>

Send a signed ADVERTISE (§6.9) — registers this connection as the handler for spec’s (realm, procedure). Fire-and-forget on the wire; the station then routes inbound CALLs (control stream) and STREAM_OPENs (a fresh dedicated stream — see accept_dedicated_stream) for that procedure back to this connection.

Source

pub async fn unadvertise( &mut self, spec: &UnadvertiseSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>

Send a signed UNADVERTISE. Fire-and-forget.

Source

pub async fn keep_advertised<F>( &mut self, spec: &AdvertiseSpec, identity: &KeyPair, interval: Duration, stop: F, on_error: impl Fn(SendFrameError), )
where F: Future<Output = ()>,

Sends an ADVERTISE for spec immediately, then again every interval, until stop resolves. advertise’s own doc notes the station’s registration is tied to the connection that sent it — a long-lived server needs to keep re-asserting it. advertise is a stateless, side-effect-free-on- repeat wire send (unlike the Erlang reference’s advertise/5, which spawns a real per-call OTP supervisor and so needs a reuse_sup option to avoid leaking one per tick), so there is nothing equivalent to worry about leaking here — same reasoning macula-go’s KeepAdvertised already applied and verified live.

A failed tick is reported via on_error but does not stop the loop — it tries again at the next interval regardless. This cannot detect or repair a dead session on its own; if the underlying connection has actually gone down, every tick will keep failing until stop resolves. See crate::direct_dial::keep_advertised_direct for the direct-dial equivalent (same shape, same reasoning).

Source

pub async fn recv_event( &mut self, timeout: Duration, ) -> Result<EventInfo, RecvEventError>

Read the next frame and parse it as an EVENT, bounded by timeout. Any non-EVENT frame received first is an error, not silently skipped — unlike call’s response wait, a caller waiting specifically for a pubsub delivery has no reason to expect anything else to legitimately arrive first.

Source

pub async fn serve_one_call<L>( &mut self, lookup: L, identity: &KeyPair, timeout: Duration, ) -> Result<(), ServeCallError>
where L: Fn(&[u8; 32], &str) -> Option<CallHandler>,

The provider role’s counterpart to call: block for the next inbound CALL frame on the control stream, bounded by timeout, look it up via lookup, invoke the matching handler, and send the resulting RESULT or ERROR back over this same connection — see plans/PLAN_WIRE_PROTOCOL.md §6.9’s routing description and macula_station_link.erl’s handle_inbound_call/2, which this mirrors field for field, including its BOLT#4 error-code mapping.

Any non-CALL frame that arrives first (e.g. a stray EVENT from an active subscribe, or a RESULT/ERROR for some other in-flight call) is discarded, not queued — the same “control stream, one thing at a time” limitation call’s own doc already carries. A session that needs to serve CALLs and also act as a caller/subscriber concurrently should use a second Session, exactly like this crate’s own streaming-provider live test does.

A caller wanting a long-lived server loops on this:

loop {
    if let Err(e) = session.serve_one_call(&lookup, identity, Duration::from_secs(30)).await {
        // ServeCallError::Timeout just means nothing arrived -- keep looping.
        eprintln!("{e}");
    }
}
Source

pub async fn serve_one_call_gated<L, P>( &mut self, lookup: L, policy: P, identity: &KeyPair, timeout: Duration, ) -> Result<(), ServeCallError>
where L: Fn(&[u8; 32], &str) -> Option<CallHandler>, P: Fn(&[u8; 32], &str) -> Policy,

serve_one_call, additionally gating each inbound CALL through policy BEFORE lookup runs — mirrors macula_station_link.erl’s handle_inbound_call/2 exactly: an open policy (the default serve_one_call uses) behaves identically; a crate::ucan::Policy::required policy demands a CALL’s ucan_token verify against the required issuer, and refuses with BOLT#4 unauthorized WITHOUT ever invoking lookup or a handler if it doesn’t — a CallHandler never sees the raw token either way, matching the reference’s own handler contract (payload only).

Source

pub async fn close(self, reason: &str, detail: Option<&str>, identity: &KeyPair)

Close the control stream and connection gracefully with a GOODBYE frame, matching macula_peering_conn.erl’s connected -> draining transition (minus the full drain-timeout bookkeeping, since this crate isn’t holding a supervisor to clean up).

write_all(...).await and finish() both only guarantee the data was handed to quinn’s own send-scheduling machinery, not that it reached the peer – Connection::close is abrupt and does not wait for outstanding stream data to be delivered. Found live 2026-08-29 in the Go port of this exact pattern (macula-go’s connection.Session.Close): a PUBLISH sent immediately before Close intermittently never reached the peer, root-caused to this race. Fixed proactively here before it was independently rediscovered against this crate – same doc comment (“minus the drain-timeout bookkeeping”), same write-then-immediately-abort-connection shape, so the same race applies. Closing the stream via finish() first, then giving the background sender a bounded window before hard-closing the connection, mirrors the Erlang reference’s own bounded-drain approach.

Source

pub async fn run_publisher( &mut self, spec: &PublishSpec, identity: &KeyPair, announce: bool, ) -> Result<(), SendFrameError>

The supervised counterpart to the bare publish primitive, matching macula_publisher.erl in spirit: publishes pubsub.publish_started_v1 before the publish and pubsub.publish_completed_v1 after, both under spec’s own realm. Fact-publish failures are silently discarded — matching macula_publisher.erl’s own publish/5 helper, which throws away its result unconditionally (_ = macula:publish(...), ok).

Unlike Erlang’s version — a supervised worker process a caller can kill mid-flight — this crate’s bare publish is already a synchronous, near-instant frame send (no ack on this wire, no network round-trip to await), so there is no meaningful “cancel before it starts” window worth a dedicated mechanism. Await this directly, or wrap it in tokio::select!/tokio::time::timeout yourself if you need to abandon it early — dropping a Future IS real cancellation in Rust; Erlang has to simulate that by killing a worker process.

Source

pub async fn run_subscriber<F>( &mut self, spec: &SubscribeSpec, identity: &KeyPair, stop: F, handler: impl FnMut(EventInfo), ) -> Result<(), RunSubscriberError>
where F: Future<Output = ()>,

The supervised counterpart to the bare subscribe/recv_event primitives, matching macula_subscriber.erl in spirit: subscribes once, then dispatches every inbound EVENT to handler until stop resolves. Unsubscribes on return, including on cancellation.

Mirrors serve_one_call’s own frame loop, not recv_event: a shared control stream can carry other frame types between one EVENT and the next, so a wrong-frame-type parse failure is skipped and polling continues, exactly like serve_one_call skips a non-“call” frame — it is NOT treated as fatal the way recv_event’s own contract treats any parse failure. Confirmed live in the Go port of this exact pattern (macula-go’s Session.RunSubscriber): without this, a single non-EVENT frame arriving on the control stream aborted the whole subscriber loop.

No OTP pid to address a running subscriber by; stop plays that role — matches keep_advertised’s own cancellation shape exactly, not a new one. handler cannot itself stop the loop (no return value) — by the same design keep_advertised already established, where on_error can only report, not halt; stopping is always external, via stop.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more