Skip to main content

subc_protocol/
lib.rs

1//! subc wire contract.
2//!
3//! This crate is the single source of truth for the subc <-> module wire,
4//! shared by subc-core and AFT. It defines the **envelope** (the fixed
5//! 21-byte routing header subc splices on), the canonical subc-generated body
6//! schemas such as [`ErrorBody`], and the capability manifest. JSON-RPC request
7//! and response bodies remain module-owned opaque payloads to subc.
8//!
9//! ## The envelope (locked — see docs/subc-core-architecture.md §4.8)
10//!
11//! ```text
12//!  offset  size  field     type    purpose
13//!    0      4    len       u32     # of BODY bytes after this 21-byte header
14//!    4      1    ver       u8      envelope version
15//!    5      1    type      u8      frame kind (see FrameType)
16//!    6      1    flags     u8     bit0 BINARY · bits1-2 PRIORITY · bit3 LAST · bits4-5 ADMISSION · bit6 DAEMON_ORIGIN · bit7 SUBSCRIPTION
17//!    7      2    channel   u16     route = (component, session); 0 = subc itself
18//!    9      4    epoch     u32     per-slot binding epoch; 0 on channel 0
19//!   13      8    corr      u64     correlation id; CANCEL carries the target call's corr
20//!   21 -> body
21//! ```
22//!
23//! Little-endian (same-machine, native, no byte-swap on the hot path).
24//!
25//! **Frozen prefix (the versioning invariant):** `len` (u32 @ 0) and `ver`
26//! (u8 @ 4) keep fixed meaning + position in *every* future version. A reader
27//! of any version can therefore always read the first 5 bytes, learn `ver`,
28//! look up that version's header length, read the rest, and splice `len` body
29//! bytes. `decode_header` enforces this discipline.
30
31#![forbid(unsafe_code)]
32
33use std::{error::Error, fmt, path::PathBuf};
34
35use serde::{Deserialize, Serialize};
36
37pub use machine_id::{MachineId, MachineIdError};
38
39pub mod frame;
40pub mod machine_id;
41pub mod manifest;
42pub mod scope;
43pub mod session;
44pub mod tool_call;
45
46/// Canonical error codes emitted by subc.
47///
48/// Error frames remain extensible strings, but these daemon-owned route-open
49/// outcomes need identical spelling across the daemon and SDK retry policies.
50pub mod error_codes {
51    use crate::Flags;
52
53    pub const UNKNOWN_CHANNEL: &str = "unknown_channel";
54    pub const STALE_ROUTE_EPOCH: &str = "stale_route_epoch";
55    pub const UNKNOWN_MODULE: &str = "unknown_module";
56    pub const MODULE_REMOVED: &str = "module_removed";
57    /// The target module's endpoint is draining for a reload, restart or disable.
58    ///
59    /// On the data plane (an `ERROR` frame answering a `REQUEST` on a bound route)
60    /// this code is a PRE-SEND GUARANTEE: the daemon returns it only when the
61    /// request could not take a route credit because the endpoint is draining,
62    /// and it returns it BEFORE forwarding, so the module never received the
63    /// request. A caller may therefore re-dispatch the same request once the
64    /// route is reopened without risking a duplicated side effect, exactly as it
65    /// may after `unknown_channel`. The daemon keeps this guarantee:
66    /// `supervisor_reload_rejects_new_work_during_drain` asserts the module's
67    /// event journal never records the rejected request. A request the module
68    /// already received is answered by the module (or its route closes with
69    /// `route.closed`), never by this code.
70    pub const MODULE_RELOADING: &str = "module_reloading";
71    pub const MODULE_WARMING: &str = "module_warming";
72    pub const TARGET_UNAVAILABLE: &str = "target_unavailable";
73    /// A flow-scoped route targets a module that does not provide
74    /// `flow-scopes/v1`. TERMINAL: decoding `flow_id` alone does not promise
75    /// flow behaviour, and removing it would silently change the identity.
76    pub const TARGET_FLOW_UNSUPPORTED: &str = "target_flow_unsupported";
77    /// An agent-run-scoped route targets a module that does not provide
78    /// `agent-run-scopes/v1`. TERMINAL: removing `run_id` would silently change
79    /// the identity under which the route acts.
80    pub const TARGET_AGENT_RUN_UNSUPPORTED: &str = "target_agent_run_unsupported";
81    pub const MODULE_TIMEOUT: &str = "module_timeout";
82    /// The target module is declared as speaking no subc wire protocol
83    /// (`protocol: "none"` in daemon config), so it has no control lane and can
84    /// never accept a route. The daemon supervises its process and nothing else.
85    ///
86    /// TERMINAL, and deliberately neither of its two neighbours. It is not
87    /// `unknown_module`, which means "no module of this id is registered or
88    /// supervised here" and is likewise terminal; retrying here would storm the
89    /// daemon forever, because the answer is a property of the module's
90    /// declaration rather than of its current state. It is not `module_removed` either: the module is configured,
91    /// running, and supervised. Only an edit to its configuration can change
92    /// this answer, and a caller cannot wait that out.
93    pub const MODULE_NO_PROTOCOL: &str = "module_no_protocol";
94    /// A request field is malformed. The error's `detail.field` names the field
95    /// (for a `route.open`, `role_versions`). TERMINAL: the same request will be
96    /// refused the same way every time, so only a corrected request can succeed.
97    pub const INVALID_REQUEST: &str = "invalid_request";
98
99    /// A `route.open` named a scope whose owner is configured but has not
100    /// synced since this daemon incarnation started. RETRYABLE: after a daemon
101    /// restart a carrier's open can arrive before the owner re-syncs, and the
102    /// carrier waits within its own deadline. See `docs/designs/daemon-scopes.md`.
103    pub const SCOPE_NOT_SYNCED: &str = "scope_not_synced";
104    /// A scoped `route.open` was admitted, but the scope record changed before
105    /// the module's bind committed. Nothing was sent on the route, so the caller
106    /// may re-open against the current record: RETRYABLE.
107    pub const SCOPE_CHANGED: &str = "scope_changed";
108    /// A scoped `route.open` named no `scope_epoch`. Every opener names one,
109    /// the owner included, so an old call can never be carried into a newer
110    /// session that reused the ref. TERMINAL.
111    pub const SCOPE_EPOCH_REQUIRED: &str = "scope_epoch_required";
112    /// A scoped `route.open` named a ref the owner's synced set does not hold,
113    /// or an owner that is not a configured module. TERMINAL.
114    pub const SCOPE_NOT_LIVE: &str = "scope_not_live";
115    /// A scoped `route.open` named a `scope_epoch` that is not the live one, or
116    /// the scope ended between admission and commit. TERMINAL.
117    pub const SCOPE_ENDED: &str = "scope_ended";
118    /// The opener of a scoped `route.open` is neither the scope's owner nor a
119    /// listed carrier, or it is a targeted carrier and the target module is not
120    /// in its list. TERMINAL.
121    pub const SCOPE_NOT_CARRIER: &str = "scope_not_carrier";
122    /// Raised by a carrier, never by the daemon: the daemon does not advertise
123    /// `scopes/v1`, so the carrier fails the call instead of opening an
124    /// unscoped route.
125    pub const SCOPE_UNSUPPORTED: &str = "scope_unsupported";
126
127    /// `scope.sync` or `scope.apply` came from a connection that is not the
128    /// owner's sync authority: another connection of the same launch holds it,
129    /// or this connection's launch is no longer the owner's current one (a
130    /// blue/green swap candidate before cutover, or an incumbent after it).
131    pub const SCOPE_SYNC_NOT_AUTHORITY: &str = "scope_sync_not_authority";
132    /// `scope.apply` came from a connection that has not taken sync authority
133    /// through an accepted full `scope.sync`. The whole call changes nothing.
134    pub const SCOPE_SYNC_REQUIRED: &str = "scope_sync_required";
135    /// `scope.sync` or `scope.apply` carried a generation no larger than the last
136    /// one the authority accepted. The whole call is refused and nothing changes.
137    pub const SCOPE_SYNC_STALE: &str = "scope_sync_stale";
138    /// `scope.sync` or `scope.apply` exceeds the owner's live-scope limit.
139    /// The whole call is refused and nothing changes.
140    pub const SCOPE_LIVE_LIMIT_EXCEEDED: &str = "scope_live_limit_exceeded";
141    /// One record in a `scope.sync` or `scope.apply` carried more attribute bytes
142    /// than a scope may hold. The whole call is refused and nothing changes.
143    pub const SCOPE_ATTRIBUTES_TOO_LARGE: &str = "scope_attributes_too_large";
144
145    // Per-record refusals: each is reported against one record in the
146    // `scope.sync` or `scope.apply` reply. A refusal itself does not change the
147    // ref or undo expiry processing; other records in the call still apply.
148
149    /// The record lowers the `scope_epoch` the daemon holds for its ref.
150    pub const SCOPE_EPOCH_REGRESSED: &str = "scope_epoch_regressed";
151    /// The record names an `(owner, ref, scope_epoch)` that already ended in
152    /// this daemon incarnation; an ended session cannot come back.
153    pub const SCOPE_EPOCH_ENDED: &str = "scope_epoch_ended";
154    /// The record's deadline has passed, or this ref at this epoch was ended by
155    /// expiry and its tombstone is still held. The same session epoch cannot
156    /// return.
157    pub const SCOPE_EXPIRED: &str = "scope_expired";
158    /// The record adds, removes or changes the deadline of a live scope at the
159    /// same epoch. A different deadline requires a new session epoch.
160    pub const SCOPE_EXPIRY_IMMUTABLE: &str = "scope_expiry_immutable";
161    /// The record's deadline is more than `scope::MAX_SCOPE_EXPIRY_AHEAD_MS`
162    /// ahead of the daemon's current Unix wall clock.
163    pub const SCOPE_EXPIRY_TOO_FAR: &str = "scope_expiry_too_far";
164    /// The record changes `kind` at the same `scope_epoch`.
165    pub const SCOPE_KIND_CHANGED: &str = "scope_kind_changed";
166    /// The record sets `agent_id`, `delegates`, `flow_id` or `run_id` and its
167    /// owner is not listed in the daemon's `scope_authority_owners`.
168    pub const SCOPE_ATTRIBUTE_NOT_PERMITTED: &str = "scope_attribute_not_permitted";
169    /// The record's parent link is not permitted: the parent's owner has synced
170    /// and the parent is not live at the named epoch, the syncing owner is
171    /// neither the parent's owner nor in its `child_owners`, or the link would
172    /// close a cycle.
173    pub const SCOPE_PARENT_NOT_PERMITTED: &str = "scope_parent_not_permitted";
174    /// A targeted carrier entry lists no target modules, or more than
175    /// `scope::MAX_CARRIER_TARGETS`.
176    pub const SCOPE_CARRIER_TARGETS_INVALID: &str = "scope_carrier_targets_invalid";
177    /// The record sets `delegates` without an `agent_id` to delegate.
178    pub const SCOPE_DELEGATES_WITHOUT_AGENT: &str = "scope_delegates_without_agent";
179    /// The record sets `run_id` without `agent_id`, the agent this run is on
180    /// behalf of.
181    pub const SCOPE_RUN_ID_WITHOUT_AGENT: &str = "scope_run_id_without_agent";
182    /// The record sets both `run_id` and `flow_id`; a run scope is not a flow scope.
183    pub const SCOPE_RUN_ID_WITH_FLOW_ID: &str = "scope_run_id_with_flow_id";
184    /// The record sets `run_id` with `delegates`; an agent-run scope cannot
185    /// delegate the agent's authority.
186    pub const SCOPE_RUN_ID_DELEGATES: &str = "scope_run_id_delegates";
187
188    /// An `operator.confirm` was declined or withdrawn. `detail.reason` is
189    /// `person`, `backoff`, `route_closed`, `module_closed` or `caller_cancelled`.
190    pub const OPERATOR_DECLINED: &str = "operator_declined";
191    /// An `operator.confirm` could not prompt or complete. `detail.reason` is
192    /// `no_presence`, `timeout`, `module_limit`, `queue_full`, `queue_wait`,
193    /// `provider_stuck`, `unsupported_platform` or `provider_error`.
194    pub const OPERATOR_PRESENCE_UNAVAILABLE: &str = "operator_presence_unavailable";
195    /// The summary of an `operator.confirm` breaks the daemon's summary rules.
196    pub const OPERATOR_SUMMARY_INVALID: &str = "operator_summary_invalid";
197    /// An `operator.confirm` is not from a nonce-proven module connection for
198    /// an open route with a verified opener on that connection.
199    pub const OPERATOR_REQUEST_NOT_PERMITTED: &str = "operator_request_not_permitted";
200
201    /// Whether a `route.open` refusal carrying `code` may be retried in place
202    /// within the caller's deadline, or is terminal for the target as named.
203    ///
204    /// This lives beside the codes because every consumer with its own
205    /// connection layer needs the same answer: a copied list breaks loudly on
206    /// a renamed code and silently on an added one. The SDKs call this; the
207    /// golden `decision_tables.json` (`route_open_retryable`) is the record
208    /// the daemon and every SDK are tested against, and the test in
209    /// `golden_json.rs` holds this function to it.
210    ///
211    /// Unknown codes are terminal: a refusal this crate has never heard of
212    /// must not be retried on the strength of a match-all arm.
213    ///
214    /// `unknown_module` is terminal: it now means only "no module of this id
215    /// is registered or supervised here" — a typo, or a peer not deployed on
216    /// this host. The daemon reports a configured-but-late target with its own
217    /// retryable codes (`module_warming`, `target_unavailable`), so a caller
218    /// that races an unsupervised module's HELLO owns its own retry; retrying
219    /// in place only papers over that race for one narrow window.
220    pub fn is_retryable_route_open(code: &str) -> bool {
221        if matches!(code, TARGET_FLOW_UNSUPPORTED | TARGET_AGENT_RUN_UNSUPPORTED) {
222            return false;
223        }
224        matches!(
225            code,
226            MODULE_RELOADING
227                | MODULE_WARMING
228                | TARGET_UNAVAILABLE
229                | MODULE_TIMEOUT
230                | SCOPE_NOT_SYNCED
231                | SCOPE_CHANGED
232        )
233    }
234
235    /// Whether the caller should evict this established route, reopen it, and
236    /// resend the request once. Consumers mapping route death to their own
237    /// custody, provider, or suspect verdict keep their own named code lists;
238    /// sharing this predicate for those meanings can change a verdict on a
239    /// protocol bump (prefrontal#59).
240    ///
241    /// Takes envelope flags so the daemon-origin check can be enabled here in
242    /// one place. Phase 2 requires first that every daemon a consumer can meet
243    /// sets DAEMON_ORIGIN on these error frames, both landed and deployed;
244    /// only then may this predicate require the bit. Requiring it sooner makes
245    /// a new SDK against an older daemon stop evicting on a genuine
246    /// `stale_route_epoch` — the failure this predicate is meant to prevent.
247    ///
248    /// The daemon emits these codes in `RouterError::to_error_frame` at
249    /// `crates/subc-daemon/src/router.rs` (UnknownChannel and StaleRouteEpoch).
250    pub fn is_established_route_dead(_flags: Flags, code: &str) -> bool {
251        matches!(code, UNKNOWN_CHANNEL | STALE_ROUTE_EPOCH)
252    }
253}
254
255pub use frame::{Frame, FrameBuildError};
256
257/// Why subc is closing a module's client routes.
258#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
259#[serde(rename_all = "snake_case")]
260pub enum RouteCloseReason {
261    Reload,
262    Restart,
263    Disable,
264    Crash,
265    /// A live route became forbidden because newly attested capability metadata
266    /// matched its supervised opening module's deny edge.
267    CapabilityDenied,
268    /// The route's scope ended: it expired, its owner ended it, or replaced it
269    /// with a higher `scope_epoch` (a new session under the same ref).
270    ScopeEnded,
271    /// The route's opener is no longer a listed carrier of its scope, or the
272    /// route's target was removed from that carrier's target list.
273    ScopeCarrierRemoved,
274    /// The scope's `delegates` went from true to false, or its `agent_id`,
275    /// `flow_id` or `run_id` changed.
276    ScopeDelegationChanged,
277    /// The scope's parent ended, so its stamp no longer names a live parent.
278    ScopeParentEnded,
279}
280
281/// Per-route bind identity shared by client-facing and module-facing control.
282///
283/// EVERY FIELD HERE IS CLIENT-SUPPLIED AND UNATTESTED. The daemon canonicalizes
284/// `project_root` as a path but does not verify that the caller has any relation
285/// to it, and `harness`, `session`, and `project_id` are strings the caller chose.
286/// A client holding the connection key can present any values it likes.
287///
288/// This sits directly above `Principal`, which is the opposite: stamped BY the
289/// daemon from a launch nonce it minted. The two travel together on every
290/// `route.bind`, so a module reading them side by side is reading one fact it can
291/// trust and four it cannot. THE DISTINCTION IS INVISIBLE FROM THE TYPES, which
292/// is why it is written here.
293///
294/// So these fields are for SCOPING AND ATTRIBUTION -- which project's state to
295/// open, which session to thread, what to log -- and never for authorization. A
296/// module that grants capability on `harness` or trusts `project_root` to bound
297/// what a caller may reach has built an authorization check on a value the caller
298/// controls. Gate on `Principal` instead, and where a module needs a caller fact
299/// subc does not stamp, it must establish that fact itself rather than believe
300/// this struct.
301#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
302#[non_exhaustive]
303pub struct BindIdentity {
304    pub project_root: PathBuf,
305    pub harness: String,
306    pub session: String,
307    /// The entorhinal-registered project id (`pj-…`) for `project_root`, when
308    /// the root is a registered project. Aliases count as registered projects;
309    /// implicit roots do not. Absent means "no stable id, key on the triple",
310    /// not "unknown".
311    ///
312    /// A producer sends the id on every bind of a session or on none. A producer
313    /// that alternates between `Some(id)` and `None` across binds silently forks
314    /// the consumer's lineage into separate stores, with no error at either end.
315    /// Therefore, a producer that cannot answer consistently must answer `None`
316    /// consistently.
317    ///
318    /// Resolve this at most once per session, before its first bind, and persist
319    /// the outcome with the session. ALF's resolver has real `Resolved`, `Unavailable`,
320    /// and `Disabled` outcomes: if unavailable at cold start is re-resolved on a
321    /// later bind, the session can alternate from `None` to `Some(id)`. Send only
322    /// registered or alias resolutions, never implicit, unavailable, or disabled
323    /// fallback ids.
324    #[serde(default, skip_serializing_if = "Option::is_none")]
325    pub project_id: Option<String>,
326}
327
328impl BindIdentity {
329    /// Constructs an identity with no registered project id.
330    ///
331    /// Use this instead of a struct literal so future additive identity fields do
332    /// not force construction-site migrations across the fleet.
333    pub fn new(
334        project_root: impl Into<PathBuf>,
335        harness: impl Into<String>,
336        session: impl Into<String>,
337    ) -> Self {
338        Self {
339            project_root: project_root.into(),
340            harness: harness.into(),
341            session: session.into(),
342            project_id: None,
343        }
344    }
345}
346
347/// Caller fact stamped by subc on each route.bind relayed to a module.
348#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
349#[serde(tag = "kind", rename_all = "snake_case")]
350pub enum Principal {
351    /// A daemon-spawned module proved possession of its launch nonce.
352    Reserved { module_id: String },
353    /// No consumer identity was presented; the caller is a direct key-holder.
354    Direct,
355    /// Reserved vocabulary for a future degraded/no-key-auth mode.
356    Unverified,
357}
358
359/// Explicit target for a route open/bind operation.
360///
361/// RouteTarget.kind ↔ ProviderRole mapping:
362///
363/// | RouteTarget.kind | required ProviderRole | disambiguator |
364/// |---|---|---|
365/// | `tool_provider` | `ToolProvider` | v1: ≤1 per module |
366/// | `management_surface` | `ManagementSurface` | v1: ≤1 per module |
367/// | `internal_service` | `InternalService` | `service_id` (multiple allowed) |
368///
369/// `ProviderRole::PipelineStage` is intentionally unroutable; pipeline modules
370/// are wired by an orchestrator rather than opened directly by clients.
371#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
372#[serde(tag = "kind", rename_all = "snake_case")]
373pub enum RouteTarget {
374    ToolProvider {
375        module_id: String,
376    },
377    ManagementSurface {
378        module_id: String,
379    },
380    InternalService {
381        module_id: String,
382        service_id: String,
383    },
384}
385
386/// Envelope protocol version this build speaks.
387pub const PROTOCOL_VERSION: u8 = 2;
388
389/// The version of THIS crate (`subc-protocol`) as compiled into the linking
390/// binary — the fleet's shared wire-vocabulary version, and the value
391/// `ManifestProvenance.wire_crate_version` declares. Not the version of a
392/// module's own envelope or payload crates: those are different numbering
393/// spaces, and declaring one here produces a confident wrong answer at any
394/// version gate (insula shipped exactly that before the referent was written
395/// down). `env!` makes it a property of the compiled binary, not of whatever
396/// source tree sits beside it at run time.
397pub const SUBC_PROTOCOL_CRATE_VERSION: &str = env!("CARGO_PKG_VERSION");
398
399/// Oldest envelope protocol version this build accepts.
400pub const MIN_SUPPORTED_VERSION: u8 = 2;
401
402/// Env var subc sets on each supervised child telling it the module_id it is
403/// supervised under, so it can register under that id.
404pub const SUBC_MODULE_ID_ENV: &str = "SUBC_MODULE_ID";
405
406/// Env var subc sets, on each spawn of a `reserved` module only, to a fresh
407/// one-time launch nonce. The child echoes it in `ModuleHelloBody::launch_nonce`;
408/// subc accepts a reserved module_id's HELLO only when the nonce matches the one it
409/// last injected for that id. Non-reserved modules never receive it.
410pub const SUBC_LAUNCH_NONCE_ENV: &str = "SUBC_LAUNCH_NONCE";
411
412/// Fixed header length for `PROTOCOL_VERSION` 2.
413pub const HEADER_LEN: usize = 21;
414
415/// Bytes of the frozen prefix (`len` u32 + `ver` u8) that are stable across
416/// every envelope version. A reader needs only these to learn the version and
417/// thus the full header length.
418pub const FROZEN_PREFIX_LEN: usize = 5;
419
420/// Maximum frame body accepted before allocation.
421///
422/// This 64 MiB starting cap prevents a malformed header from forcing an
423/// unbounded allocation. Future protocol versions can negotiate or encode a
424/// different cap while preserving the frozen prefix.
425pub const MAX_FRAME_BODY_LEN: u32 = 64 * 1024 * 1024;
426
427/// Canonical JSON body for all subc-generated `ERROR` frames.
428///
429/// `detail` is an optional machine-parsable surface for refusals whose remedy
430/// needs more than a code (e.g. a producer-published backoff number, an
431/// observed-vs-configured size pair). Absent detail serializes to nothing, so
432/// bodies without it are byte-identical to the pre-detail wire and older
433/// readers simply never see the field. Producers document each code's detail
434/// fields where the code is defined; `detail` must never carry secrets.
435#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
436pub struct ErrorBody {
437    pub code: String,
438    pub message: String,
439    #[serde(default, skip_serializing_if = "Option::is_none")]
440    pub detail: Option<serde_json::Value>,
441}
442
443impl ErrorBody {
444    /// A detail-less error body; the common case.
445    pub fn new(code: impl Into<String>, message: impl Into<String>) -> Self {
446        Self {
447            code: code.into(),
448            message: message.into(),
449            detail: None,
450        }
451    }
452
453    /// Attach a machine-parsable detail object to this error.
454    pub fn with_detail(mut self, detail: serde_json::Value) -> Self {
455        self.detail = Some(detail);
456        self
457    }
458}
459
460/// Module-to-subc `HELLO` body used during module registration.
461#[derive(Clone, Serialize, Deserialize, PartialEq)]
462pub struct ModuleHelloBody {
463    pub manifest: manifest::ModuleManifest,
464    pub protocol_ver: u8,
465    #[serde(default)]
466    pub control_ops: Option<Vec<String>>,
467    /// One-time launch nonce, echoed back from the `SUBC_LAUNCH_NONCE` environment
468    /// variable the daemon injected when it spawned this process. Only a daemon-spawned
469    /// process for a `reserved` module receives a nonce; subc accepts a reserved
470    /// `module_id`'s HELLO only when this matches the nonce it last injected for that
471    /// id, so a different process cannot register as a reserved module while the real
472    /// one is down/restarting. Absent (`serde(default)`) for non-reserved modules and
473    /// self-connecting providers, which are never nonce-checked.
474    #[serde(default, skip_serializing_if = "Option::is_none")]
475    pub launch_nonce: Option<String>,
476}
477
478// Hand-written so the launch nonce is never printed. The nonce is the credential
479// that attributes a connection to a supervised module, and a derived Debug would
480// write it into any log line or panic message that formats this value. Same
481// reasoning as ConnectionInfo's Debug in subc-transport.
482impl fmt::Debug for ModuleHelloBody {
483    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
484        f.debug_struct("ModuleHelloBody")
485            .field("manifest", &self.manifest)
486            .field("protocol_ver", &self.protocol_ver)
487            .field("control_ops", &self.control_ops)
488            .field(
489                "launch_nonce",
490                &self
491                    .launch_nonce
492                    .as_ref()
493                    .map(|nonce| format!("<{} bytes redacted>", nonce.len())),
494            )
495            .finish()
496    }
497}
498
499/// subc-to-module `HELLO_ACK` body used during module registration.
500#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
501pub struct ModuleHelloAckBody {
502    pub negotiated_ver: u8,
503    pub subc_ops: Vec<String>,
504    pub subc_capabilities: Vec<String>,
505    /// The module's resolved storage descriptor, when the daemon's central config
506    /// configures managed storage. Carried opaquely here (subc-protocol stays a
507    /// thin wire crate with no storage/database dependency); a module that uses
508    /// managed storage deserializes it into `cortexkit_store_types::StorageDescriptor`
509    /// and hands it to `cortexkit-store`. Absent when no storage is configured, and
510    /// `serde(default)` so an older module simply ignores it.
511    #[serde(default, skip_serializing_if = "Option::is_none")]
512    pub storage: Option<serde_json::Value>,
513    /// The daemon's machine id (see [`MachineId`]): a name for this machine,
514    /// never an authority. Nothing may admit a peer, grant trust or skip a check
515    /// because two messages carry the same value.
516    ///
517    /// Carried as a plain string so one malformed value cannot fail the whole
518    /// registration reply; a module validates it with [`MachineId::parse`].
519    /// Absent from a daemon that predates the machine id, which a module must
520    /// read as exactly that, never as "no machine" and never as a reason to mint
521    /// its own.
522    #[serde(default, skip_serializing_if = "Option::is_none")]
523    pub machine_id: Option<String>,
524}
525
526/// Frame kind (`type` byte at offset 5).
527///
528/// `CANCEL`, `PING`, `PONG`, and `GOODBYE` are pure-header frames (`len == 0`);
529/// only `HELLO`/`HELLO_ACK` and the RPC payloads carry bodies.
530#[derive(Debug, Clone, Copy, PartialEq, Eq)]
531#[repr(u8)]
532pub enum FrameType {
533    Request = 0,
534    Response = 1,
535    Push = 2,
536    StreamData = 3,
537    StreamEnd = 4,
538    Error = 5,
539    Cancel = 6,
540    Ping = 7,
541    Pong = 8,
542    Hello = 9,
543    HelloAck = 10,
544    Goodbye = 11,
545}
546
547impl FrameType {
548    /// Map the raw `type` byte to a `FrameType`, or `None` if unknown.
549    pub fn from_u8(b: u8) -> Option<Self> {
550        Some(match b {
551            0 => Self::Request,
552            1 => Self::Response,
553            2 => Self::Push,
554            3 => Self::StreamData,
555            4 => Self::StreamEnd,
556            5 => Self::Error,
557            6 => Self::Cancel,
558            7 => Self::Ping,
559            8 => Self::Pong,
560            9 => Self::Hello,
561            10 => Self::HelloAck,
562            11 => Self::Goodbye,
563            _ => return None,
564        })
565    }
566
567    pub fn is_pure_header(self) -> bool {
568        matches!(self, Self::Cancel | Self::Ping | Self::Pong | Self::Goodbye)
569    }
570}
571
572/// Scheduling priority carried in `flags` bits 1-2. subc schedules on this
573/// without parsing the body.
574#[derive(Debug, Clone, Copy, PartialEq, Eq)]
575#[repr(u8)]
576pub enum Priority {
577    Passive = 0,
578    Interactive = 1,
579    Background = 2,
580}
581
582impl Priority {
583    fn from_bits(bits: u8) -> Option<Self> {
584        Some(match bits {
585            0 => Self::Passive,
586            1 => Self::Interactive,
587            2 => Self::Background,
588            _ => return None,
589        })
590    }
591}
592
593/// Admission behavior carried in `flags` bits 4-5.
594#[derive(Debug, Clone, Copy, PartialEq, Eq)]
595#[repr(u8)]
596pub enum AdmissionClass {
597    Normal = 0,
598    Expedite = 1,
599    Sheddable = 2,
600}
601
602impl AdmissionClass {
603    fn from_bits(bits: u8) -> Option<Self> {
604        Some(match bits {
605            0 => Self::Normal,
606            1 => Self::Expedite,
607            2 => Self::Sheddable,
608            _ => return None,
609        })
610    }
611}
612
613const FLAG_BINARY: u8 = 0b0000_0001; // bit 0
614const FLAG_PRIORITY_MASK: u8 = 0b0000_0110; // bits 1-2
615const FLAG_PRIORITY_SHIFT: u8 = 1;
616const FLAG_LAST: u8 = 0b0000_1000; // bit 3
617const FLAG_ADMISSION_MASK: u8 = 0b0011_0000; // bits 4-5
618const FLAG_ADMISSION_SHIFT: u8 = 4;
619pub const FLAG_DAEMON_ORIGIN: u8 = 0b0100_0000;
620/// A request credit the client explicitly declares as a held-open subscription.
621pub const FLAG_SUBSCRIPTION: u8 = 0b1000_0000;
622
623/// The `flags` byte (offset 6): binary, priority, last, admission, daemon origin, subscription.
624#[derive(Debug, Clone, Copy, PartialEq, Eq)]
625pub struct Flags(pub u8);
626
627impl Flags {
628    /// Build flags with the default [`AdmissionClass::Normal`] class.
629    pub fn new(binary: bool, priority: Priority, last: bool) -> Self {
630        let mut b = 0u8;
631        if binary {
632            b |= FLAG_BINARY;
633        }
634        b |= (priority as u8) << FLAG_PRIORITY_SHIFT;
635        if last {
636            b |= FLAG_LAST;
637        }
638        Flags(b)
639    }
640
641    /// Return these flags with a typed admission class.
642    pub fn with_admission_class(mut self, admission_class: AdmissionClass) -> Self {
643        self.0 =
644            (self.0 & !FLAG_ADMISSION_MASK) | ((admission_class as u8) << FLAG_ADMISSION_SHIFT);
645        self
646    }
647
648    /// Body is raw bytes (bulk lane) rather than JSON-RPC.
649    pub fn is_binary(self) -> bool {
650        self.0 & FLAG_BINARY != 0
651    }
652
653    /// Final frame of a streamed message.
654    pub fn is_last(self) -> bool {
655        self.0 & FLAG_LAST != 0
656    }
657
658    /// Decode the priority bits, or `None` if they hold a reserved value.
659    pub fn priority(self) -> Option<Priority> {
660        Priority::from_bits((self.0 & FLAG_PRIORITY_MASK) >> FLAG_PRIORITY_SHIFT)
661    }
662
663    /// Decode the admission-class bits, or `None` if they hold `0b11`.
664    pub fn admission_class(self) -> Option<AdmissionClass> {
665        AdmissionClass::from_bits((self.0 & FLAG_ADMISSION_MASK) >> FLAG_ADMISSION_SHIFT)
666    }
667
668    /// True when a request was explicitly opened as a held-open subscription.
669    pub fn is_subscription(self) -> bool {
670        self.0 & FLAG_SUBSCRIPTION != 0
671    }
672
673    /// True when the frame was authored by the daemon.
674    pub fn is_daemon_origin(self) -> bool {
675        self.0 & FLAG_DAEMON_ORIGIN != 0
676    }
677
678    /// Return these flags with daemon origin asserted.
679    pub fn with_daemon_origin(mut self) -> Self {
680        self.0 |= FLAG_DAEMON_ORIGIN;
681        self
682    }
683
684    /// Return these flags with daemon origin cleared.
685    pub fn without_daemon_origin(self) -> Self {
686        Self(self.0 & !FLAG_DAEMON_ORIGIN)
687    }
688}
689
690/// A decoded envelope header. The body is the `len` bytes that follow it.
691#[derive(Debug, Clone, Copy, PartialEq, Eq)]
692pub struct EnvelopeHeader {
693    /// Number of body bytes after the header.
694    pub len: u32,
695    /// Envelope version.
696    pub ver: u8,
697    /// Frame kind.
698    pub ty: FrameType,
699    /// Flag bits.
700    pub flags: Flags,
701    /// Sender-local route slot; 0 is the control channel.
702    pub channel: u16,
703    /// Sender-local binding epoch; 0 is reserved for the control channel.
704    pub epoch: u32,
705    /// Correlation id.
706    pub corr: u64,
707}
708
709impl EnvelopeHeader {
710    /// Serialize the header to its fixed 21-byte little-endian form.
711    pub fn encode(&self) -> [u8; HEADER_LEN] {
712        let mut buf = [0u8; HEADER_LEN];
713        buf[0..4].copy_from_slice(&self.len.to_le_bytes());
714        buf[4] = self.ver;
715        buf[5] = self.ty as u8;
716        buf[6] = self.flags.0;
717        buf[7..9].copy_from_slice(&self.channel.to_le_bytes());
718        buf[9..13].copy_from_slice(&self.epoch.to_le_bytes());
719        buf[13..21].copy_from_slice(&self.corr.to_le_bytes());
720        buf
721    }
722}
723
724/// Why a header could not be decoded.
725#[derive(Debug, Clone, Copy, PartialEq, Eq)]
726pub enum DecodeError {
727    /// Fewer than `FROZEN_PREFIX_LEN` bytes — cannot even read `len`/`ver`.
728    TooShortForPrefix { have: usize },
729    /// `ver` is not a version this build understands.
730    UnsupportedVersion { ver: u8 },
731    /// Version known but fewer than its header length is present.
732    TooShortForHeader { have: usize, need: usize },
733    /// `type` byte is not a known `FrameType`.
734    UnknownFrameType { byte: u8 },
735    /// A reserved flag bit is set (retained for older decoder error compatibility).
736    ReservedFlagBits { flags: u8 },
737    /// Priority bits 1-2 hold the reserved value `0b11`.
738    ReservedPriorityBits { flags: u8 },
739    /// Admission bits 4-5 hold the reserved value `0b11`.
740    ReservedAdmissionClass { flags: u8 },
741    /// SHEDDABLE is set on a frame type that must be delivered.
742    SheddableIllegalFrameType { ty: FrameType, flags: u8 },
743    /// Channel 0 carried an epoch other than its reserved epoch 0.
744    NonzeroEpochOnControlChannel { epoch: u32 },
745    /// A pure-header frame declared body bytes.
746    PureHeaderFrameWithBody { ty: FrameType, len: u32 },
747}
748
749impl fmt::Display for DecodeError {
750    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
751        match self {
752            Self::TooShortForPrefix { have } => {
753                write!(f, "header shorter than frozen prefix: have {have} bytes")
754            }
755            Self::UnsupportedVersion { ver } => write!(f, "unsupported envelope version {ver}"),
756            Self::TooShortForHeader { have, need } => {
757                write!(
758                    f,
759                    "header too short for version: have {have} bytes, need {need}"
760                )
761            }
762            Self::UnknownFrameType { byte } => write!(f, "unknown frame type byte {byte}"),
763            Self::ReservedFlagBits { flags } => {
764                write!(f, "reserved flag bits set in flags 0b{flags:08b}")
765            }
766            Self::ReservedPriorityBits { flags } => {
767                write!(f, "reserved priority bits set in flags 0b{flags:08b}")
768            }
769            Self::ReservedAdmissionClass { flags } => {
770                write!(f, "reserved admission class set in flags 0b{flags:08b}")
771            }
772            Self::SheddableIllegalFrameType { ty, flags } => write!(
773                f,
774                "SHEDDABLE admission class is illegal on {ty:?} in flags 0b{flags:08b}"
775            ),
776            Self::NonzeroEpochOnControlChannel { epoch } => {
777                write!(f, "control channel carried nonzero epoch {epoch}")
778            }
779            Self::PureHeaderFrameWithBody { ty, len } => {
780                write!(
781                    f,
782                    "pure-header frame {ty:?} declared non-zero body length {len}"
783                )
784            }
785        }
786    }
787}
788
789impl Error for DecodeError {}
790
791/// How many header bytes a given envelope version occupies. Driven by the
792/// frozen prefix: read `ver`, then learn the full header length here.
793fn header_len_for_version(ver: u8) -> Option<usize> {
794    match ver {
795        PROTOCOL_VERSION => Some(HEADER_LEN),
796        _ => None,
797    }
798}
799
800/// Decode an envelope header from the front of `bytes`, following the
801/// frozen-prefix discipline:
802/// 1. need at least the 5-byte prefix to read `len` + `ver`;
803/// 2. dispatch the full header length on `ver`;
804/// 3. need the full header present; then parse the rest.
805///
806/// Never panics on malformed input — returns a typed [`DecodeError`].
807pub fn decode_header(bytes: &[u8]) -> Result<EnvelopeHeader, DecodeError> {
808    if bytes.len() < FROZEN_PREFIX_LEN {
809        return Err(DecodeError::TooShortForPrefix { have: bytes.len() });
810    }
811    let ver = bytes[4];
812    let need = header_len_for_version(ver).ok_or(DecodeError::UnsupportedVersion { ver })?;
813    if bytes.len() < need {
814        return Err(DecodeError::TooShortForHeader {
815            have: bytes.len(),
816            need,
817        });
818    }
819
820    let len = u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
821    let ty =
822        FrameType::from_u8(bytes[5]).ok_or(DecodeError::UnknownFrameType { byte: bytes[5] })?;
823    let flags = Flags(bytes[6]);
824    if flags.priority().is_none() {
825        return Err(DecodeError::ReservedPriorityBits { flags: bytes[6] });
826    }
827    let admission_class = flags
828        .admission_class()
829        .ok_or(DecodeError::ReservedAdmissionClass { flags: bytes[6] })?;
830    if admission_class == AdmissionClass::Sheddable
831        && !matches!(ty, FrameType::Push | FrameType::StreamData)
832    {
833        return Err(DecodeError::SheddableIllegalFrameType {
834            ty,
835            flags: bytes[6],
836        });
837    }
838    if ty.is_pure_header() && len != 0 {
839        return Err(DecodeError::PureHeaderFrameWithBody { ty, len });
840    }
841    let channel = u16::from_le_bytes([bytes[7], bytes[8]]);
842    let epoch = u32::from_le_bytes([bytes[9], bytes[10], bytes[11], bytes[12]]);
843    if channel == 0 && epoch != 0 {
844        return Err(DecodeError::NonzeroEpochOnControlChannel { epoch });
845    }
846    let corr = u64::from_le_bytes([
847        bytes[13], bytes[14], bytes[15], bytes[16], bytes[17], bytes[18], bytes[19], bytes[20],
848    ]);
849
850    Ok(EnvelopeHeader {
851        len,
852        ver,
853        ty,
854        flags,
855        channel,
856        epoch,
857        corr,
858    })
859}
860
861#[cfg(test)]
862mod tests {
863    use super::*;
864
865    fn hdr(len: u32, ty: FrameType, flags: Flags, channel: u16, corr: u64) -> EnvelopeHeader {
866        hdr_with_epoch(len, ty, flags, channel, u32::from(channel != 0), corr)
867    }
868
869    fn hdr_with_epoch(
870        len: u32,
871        ty: FrameType,
872        flags: Flags,
873        channel: u16,
874        epoch: u32,
875        corr: u64,
876    ) -> EnvelopeHeader {
877        EnvelopeHeader {
878            len,
879            ver: PROTOCOL_VERSION,
880            ty,
881            flags,
882            channel,
883            epoch,
884            corr,
885        }
886    }
887
888    #[test]
889    fn bind_identity_with_project_id_round_trips_json() {
890        let mut identity = BindIdentity::new("/tmp/project", "opencode", "session-1");
891        identity.project_id = Some("pj-a1b2c3d4".to_string());
892
893        let encoded = serde_json::to_vec(&identity).unwrap();
894        let decoded: BindIdentity = serde_json::from_slice(&encoded).unwrap();
895
896        assert_eq!(decoded, identity);
897    }
898
899    #[test]
900    fn bind_identity_without_project_id_round_trips_json() {
901        let identity = BindIdentity::new("/tmp/project", "opencode", "session-1");
902
903        let encoded = serde_json::to_vec(&identity).unwrap();
904        let decoded: BindIdentity = serde_json::from_slice(&encoded).unwrap();
905
906        assert_eq!(decoded, identity);
907    }
908
909    #[test]
910    fn legacy_bind_identity_without_project_id_decodes() {
911        let decoded: BindIdentity = serde_json::from_value(serde_json::json!({
912            "project_root": "/tmp/project",
913            "harness": "opencode",
914            "session": "session-1"
915        }))
916        .unwrap();
917
918        assert_eq!(decoded.project_id, None);
919    }
920
921    #[test]
922    fn bind_identity_none_omits_project_id_instead_of_serializing_null() {
923        let encoded =
924            serde_json::to_value(BindIdentity::new("/tmp/project", "opencode", "session-1"))
925                .unwrap();
926
927        assert!(encoded.get("project_id").is_none());
928    }
929
930    #[test]
931    fn wire_crate_version_is_a_numeric_three_component_version() {
932        let components = SUBC_PROTOCOL_CRATE_VERSION.split('.').collect::<Vec<_>>();
933
934        assert!(!SUBC_PROTOCOL_CRATE_VERSION.is_empty());
935        assert_eq!(components.len(), 3);
936        assert!(components
937            .iter()
938            .all(|component| !component.is_empty() && component.parse::<u64>().is_ok()));
939    }
940
941    #[test]
942    fn route_target_variants_round_trip_json() {
943        let targets = [
944            RouteTarget::ToolProvider {
945                module_id: "aft".to_string(),
946            },
947            RouteTarget::ManagementSurface {
948                module_id: "memory".to_string(),
949            },
950            RouteTarget::InternalService {
951                module_id: "bus".to_string(),
952                service_id: "dm".to_string(),
953            },
954        ];
955
956        for target in targets {
957            let encoded = serde_json::to_vec(&target).unwrap();
958            let decoded: RouteTarget = serde_json::from_slice(&encoded).unwrap();
959            assert_eq!(decoded, target);
960        }
961    }
962
963    #[test]
964    fn error_body_round_trips_json() {
965        let body = ErrorBody {
966            code: "config_divergence".to_string(),
967            message: "active config differs".to_string(),
968            detail: None,
969        };
970
971        let encoded = serde_json::to_vec(&body).unwrap();
972        let decoded: ErrorBody = serde_json::from_slice(&encoded).unwrap();
973
974        assert_eq!(decoded, body);
975    }
976
977    #[test]
978    fn round_trip_request() {
979        let h = hdr(
980            1234,
981            FrameType::Request,
982            Flags::new(false, Priority::Interactive, false),
983            42,
984            0xDEAD_BEEF_0000_0001,
985        );
986        let decoded = decode_header(&h.encode()).unwrap();
987        assert_eq!(h, decoded);
988    }
989
990    #[test]
991    fn round_trip_all_frame_types() {
992        for b in 0u8..=11 {
993            let ty = FrameType::from_u8(b).unwrap();
994            let h = hdr(0, ty, Flags::new(false, Priority::Passive, false), 0, 0);
995            assert_eq!(decode_header(&h.encode()).unwrap().ty, ty);
996        }
997    }
998
999    #[test]
1000    fn pure_header_frame_has_zero_len() {
1001        // CANCEL carries only header (len = 0) + the target corr.
1002        let h = hdr(
1003            0,
1004            FrameType::Cancel,
1005            Flags::new(false, Priority::Passive, false),
1006            7,
1007            99,
1008        );
1009        let d = decode_header(&h.encode()).unwrap();
1010        assert_eq!(d.len, 0);
1011        assert_eq!(d.corr, 99);
1012    }
1013
1014    #[test]
1015    fn flags_round_trip() {
1016        let f = Flags::new(true, Priority::Background, true)
1017            .with_admission_class(AdmissionClass::Expedite);
1018        assert!(f.is_binary());
1019        assert!(f.is_last());
1020        assert_eq!(f.priority(), Some(Priority::Background));
1021        assert_eq!(f.admission_class(), Some(AdmissionClass::Expedite));
1022        let h = hdr(8, FrameType::StreamData, f, 1, 1);
1023        assert_eq!(decode_header(&h.encode()).unwrap().flags, f);
1024    }
1025
1026    #[test]
1027    fn daemon_origin_flags_decode_and_round_trip() {
1028        let old = hdr(0, FrameType::Error, Flags(0), 7, 1);
1029        let old_decoded = decode_header(&old.encode()).unwrap();
1030        assert!(!old_decoded.flags.is_daemon_origin());
1031
1032        let daemon = hdr(0, FrameType::Error, Flags(0).with_daemon_origin(), 7, 1);
1033        let daemon_decoded = decode_header(&daemon.encode()).unwrap();
1034        assert!(daemon_decoded.flags.is_daemon_origin());
1035        assert_eq!(daemon_decoded.flags.without_daemon_origin(), Flags(0));
1036        assert!(Flags(0).with_daemon_origin().is_daemon_origin());
1037    }
1038
1039    #[test]
1040    fn little_endian_and_frozen_prefix_layout() {
1041        let h = hdr_with_epoch(
1042            0x0403_0201,
1043            FrameType::Request,
1044            Flags(0),
1045            0x0605,
1046            0x0a09_0807,
1047            0x1211_100f_0e0d_0c0b,
1048        );
1049        let buf = h.encode();
1050        assert_eq!(&buf[0..4], &[1, 2, 3, 4]);
1051        assert_eq!(buf[4], PROTOCOL_VERSION);
1052        assert_eq!(&buf[7..9], &[5, 6]);
1053        assert_eq!(&buf[9..13], &[7, 8, 9, 10]);
1054        assert_eq!(&buf[13..21], &[11, 12, 13, 14, 15, 16, 17, 18]);
1055        assert_eq!(buf.len(), HEADER_LEN);
1056    }
1057
1058    #[test]
1059    fn reject_too_short_for_prefix() {
1060        assert_eq!(
1061            decode_header(&[0, 0, 0, 0]),
1062            Err(DecodeError::TooShortForPrefix { have: 4 })
1063        );
1064    }
1065
1066    #[test]
1067    fn reject_too_short_for_header() {
1068        // Valid 5-byte prefix but the v2 header is truncated.
1069        let mut b = [0u8; 10];
1070        b[4] = PROTOCOL_VERSION;
1071        assert_eq!(
1072            decode_header(&b),
1073            Err(DecodeError::TooShortForHeader {
1074                have: 10,
1075                need: HEADER_LEN
1076            })
1077        );
1078    }
1079
1080    #[test]
1081    fn reject_unsupported_version() {
1082        let mut b = [0u8; HEADER_LEN];
1083        b[4] = 1;
1084        assert_eq!(
1085            decode_header(&b),
1086            Err(DecodeError::UnsupportedVersion { ver: 1 })
1087        );
1088    }
1089
1090    #[test]
1091    fn reject_unknown_frame_type() {
1092        let mut b = [0u8; HEADER_LEN];
1093        b[4] = PROTOCOL_VERSION;
1094        b[5] = 99;
1095        assert_eq!(
1096            decode_header(&b),
1097            Err(DecodeError::UnknownFrameType { byte: 99 })
1098        );
1099    }
1100
1101    #[test]
1102    fn subscription_flag_decodes_and_tags_the_request() {
1103        let mut b = [0u8; HEADER_LEN];
1104        b[4] = PROTOCOL_VERSION;
1105        b[5] = FrameType::Request as u8;
1106        b[6] = FLAG_SUBSCRIPTION;
1107        let decoded = decode_header(&b).unwrap();
1108        assert!(decoded.flags.is_subscription());
1109    }
1110
1111    #[test]
1112    fn reject_reserved_priority_bits() {
1113        let mut b = [0u8; HEADER_LEN];
1114        b[4] = PROTOCOL_VERSION;
1115        b[5] = FrameType::Request as u8;
1116        b[6] = 0b0000_0110; // priority bits 1-2 are reserved value 0b11
1117        assert_eq!(
1118            decode_header(&b),
1119            Err(DecodeError::ReservedPriorityBits { flags: 0b0000_0110 })
1120        );
1121    }
1122
1123    #[test]
1124    fn reject_pure_header_frame_with_body_len() {
1125        let h = hdr(
1126            1,
1127            FrameType::Ping,
1128            Flags::new(false, Priority::Passive, false),
1129            0,
1130            1,
1131        );
1132        assert_eq!(
1133            decode_header(&h.encode()),
1134            Err(DecodeError::PureHeaderFrameWithBody {
1135                ty: FrameType::Ping,
1136                len: 1
1137            })
1138        );
1139    }
1140
1141    #[test]
1142    fn epoch_boundaries_round_trip() {
1143        for (channel, epoch) in [(0, 0), (1, 1), (u16::MAX, u32::MAX)] {
1144            let h = hdr_with_epoch(
1145                0,
1146                FrameType::Request,
1147                Flags::new(false, Priority::Passive, false),
1148                channel,
1149                epoch,
1150                9,
1151            );
1152            assert_eq!(decode_header(&h.encode()).unwrap(), h);
1153        }
1154    }
1155
1156    #[test]
1157    fn admission_classes_accept_three_values_and_reject_reserved_value() {
1158        for (ty, admission_class) in [
1159            (FrameType::Request, AdmissionClass::Normal),
1160            (FrameType::Request, AdmissionClass::Expedite),
1161            (FrameType::Push, AdmissionClass::Sheddable),
1162            (FrameType::StreamData, AdmissionClass::Sheddable),
1163        ] {
1164            let flags = Flags::new(false, Priority::Interactive, false)
1165                .with_admission_class(admission_class);
1166            let h = hdr(0, ty, flags, 1, 2);
1167            assert_eq!(decode_header(&h.encode()).unwrap().flags, flags);
1168        }
1169
1170        let mut h = hdr(
1171            0,
1172            FrameType::Push,
1173            Flags::new(false, Priority::Passive, false),
1174            1,
1175            2,
1176        )
1177        .encode();
1178        h[6] |= 0b0011_0000;
1179        assert_eq!(
1180            decode_header(&h),
1181            Err(DecodeError::ReservedAdmissionClass { flags: h[6] })
1182        );
1183    }
1184
1185    #[test]
1186    fn sheddable_rejected_on_every_illegal_frame_type() {
1187        let flags = Flags::new(false, Priority::Passive, false)
1188            .with_admission_class(AdmissionClass::Sheddable);
1189        for ty in [
1190            FrameType::Request,
1191            FrameType::Response,
1192            FrameType::StreamEnd,
1193            FrameType::Error,
1194            FrameType::Cancel,
1195            FrameType::Ping,
1196            FrameType::Pong,
1197            FrameType::Hello,
1198            FrameType::HelloAck,
1199            FrameType::Goodbye,
1200        ] {
1201            let h = hdr(0, ty, flags, 1, 2);
1202            assert_eq!(
1203                decode_header(&h.encode()),
1204                Err(DecodeError::SheddableIllegalFrameType { ty, flags: flags.0 })
1205            );
1206        }
1207    }
1208
1209    #[test]
1210    fn nonzero_epoch_on_control_channel_is_rejected() {
1211        let h = hdr_with_epoch(
1212            0,
1213            FrameType::Request,
1214            Flags::new(false, Priority::Passive, false),
1215            0,
1216            u32::MAX,
1217            2,
1218        );
1219        assert_eq!(
1220            decode_header(&h.encode()),
1221            Err(DecodeError::NonzeroEpochOnControlChannel { epoch: u32::MAX })
1222        );
1223    }
1224}
1225
1226#[cfg(test)]
1227mod launch_nonce_redaction_tests {
1228    use super::*;
1229
1230    #[test]
1231    fn hello_body_debug_never_prints_the_nonce() {
1232        let body = ModuleHelloBody {
1233            manifest: manifest::ModuleManifest::builder("broca", "0.1.0").build(),
1234            protocol_ver: 2,
1235            control_ops: None,
1236            launch_nonce: Some("nonce-f00dfeed1234abcd".to_string()),
1237        };
1238        let printed = format!("{body:?}");
1239        assert!(printed.contains("broca"), "{printed}");
1240        assert!(
1241            !printed.contains("nonce-f00dfeed1234abcd"),
1242            "launch nonce printed: {printed}"
1243        );
1244    }
1245}