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: HelloInfoImplementations§
Source§impl Session
impl Session
Sourcepub fn remote_address(&self) -> SocketAddr
pub fn remote_address(&self) -> SocketAddr
The remote address this session’s connection is with.
Sourcepub async fn open_dedicated_stream(
&mut self,
) -> Result<FrameStream, ConnectionError>
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.
Sourcepub async fn accept_dedicated_stream(
&mut self,
) -> Result<FrameStream, ConnectionError>
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.
Sourcepub fn leftover_bytes(&self) -> &[u8] ⓘ
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.
Sourcepub async fn recv_frame(&mut self) -> Result<Value, RecvFrameError>
pub async fn recv_frame(&mut self) -> Result<Value, RecvFrameError>
Read the next complete application frame from the control stream.
Sourcepub async fn recv_frame_timeout(
&mut self,
timeout: Duration,
) -> Result<Value, RecvFrameError>
pub async fn recv_frame_timeout( &mut self, timeout: Duration, ) -> Result<Value, RecvFrameError>
As recv_frame, bounded by timeout.
Sourcepub async fn call(
&mut self,
procedure: &str,
realm: [u8; 32],
payload: Value,
deadline_ms: i128,
identity: &KeyPair,
timeout: Duration,
) -> Result<CallResponse, CallError>
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.
Sourcepub 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>
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.
Sourcepub async fn publish(
&mut self,
spec: &PublishSpec,
identity: &KeyPair,
) -> Result<(), SendFrameError>
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.
Sourcepub async fn subscribe(
&mut self,
spec: &SubscribeSpec,
identity: &KeyPair,
) -> Result<(), SendFrameError>
pub async fn subscribe( &mut self, spec: &SubscribeSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>
Send a signed SUBSCRIBE. Fire-and-forget.
Sourcepub async fn unsubscribe(
&mut self,
spec: &UnsubscribeSpec,
identity: &KeyPair,
) -> Result<(), SendFrameError>
pub async fn unsubscribe( &mut self, spec: &UnsubscribeSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>
Send a signed UNSUBSCRIBE. Fire-and-forget.
Sourcepub async fn advertise(
&mut self,
spec: &AdvertiseSpec,
identity: &KeyPair,
) -> Result<(), SendFrameError>
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.
Sourcepub async fn unadvertise(
&mut self,
spec: &UnadvertiseSpec,
identity: &KeyPair,
) -> Result<(), SendFrameError>
pub async fn unadvertise( &mut self, spec: &UnadvertiseSpec, identity: &KeyPair, ) -> Result<(), SendFrameError>
Send a signed UNADVERTISE. Fire-and-forget.
Sourcepub async fn keep_advertised<F>(
&mut self,
spec: &AdvertiseSpec,
identity: &KeyPair,
interval: Duration,
stop: F,
on_error: impl Fn(SendFrameError),
)
pub async fn keep_advertised<F>( &mut self, spec: &AdvertiseSpec, identity: &KeyPair, interval: Duration, stop: F, on_error: impl Fn(SendFrameError), )
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).
Sourcepub async fn recv_event(
&mut self,
timeout: Duration,
) -> Result<EventInfo, RecvEventError>
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.
Sourcepub async fn serve_one_call<L>(
&mut self,
lookup: L,
identity: &KeyPair,
timeout: Duration,
) -> Result<(), ServeCallError>
pub async fn serve_one_call<L>( &mut self, lookup: L, identity: &KeyPair, timeout: Duration, ) -> Result<(), ServeCallError>
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}");
}
}Sourcepub async fn serve_one_call_gated<L, P>(
&mut self,
lookup: L,
policy: P,
identity: &KeyPair,
timeout: Duration,
) -> Result<(), ServeCallError>
pub async fn serve_one_call_gated<L, P>( &mut self, lookup: L, policy: P, identity: &KeyPair, timeout: Duration, ) -> Result<(), ServeCallError>
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).
Sourcepub async fn close(self, reason: &str, detail: Option<&str>, identity: &KeyPair)
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.
Sourcepub async fn run_publisher(
&mut self,
spec: &PublishSpec,
identity: &KeyPair,
announce: bool,
) -> Result<(), SendFrameError>
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.
Sourcepub async fn run_subscriber<F>(
&mut self,
spec: &SubscribeSpec,
identity: &KeyPair,
stop: F,
handler: impl FnMut(EventInfo),
) -> Result<(), RunSubscriberError>
pub async fn run_subscriber<F>( &mut self, spec: &SubscribeSpec, identity: &KeyPair, stop: F, handler: impl FnMut(EventInfo), ) -> Result<(), RunSubscriberError>
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.