Skip to main content

subc_protocol/
session.rs

1//! Session route control wire contract.
2//!
3//! subc has two distinct channel-0 handshakes. Module registration is the
4//! module-to-subc `HELLO`/`HELLO_ACK` handshake that registers the manifest and
5//! liveness. Route bind is the client-to-subc-to-module request/response
6//! handshake that binds one client route to a module route channel.
7
8use std::collections::BTreeMap;
9
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12
13use crate::{
14    manifest::{CapabilityDeclarations, ProviderRole},
15    scope::{ScopeEnded, ScopeRecord, ScopeRecordResult, ScopeStamp, ScopeStatus},
16    BindIdentity, Principal, RouteCloseReason, RouteTarget,
17};
18
19pub const MODULE_CONTROL_OP_HEALTH_CHECK: &str = "health.check";
20pub const MODULE_TO_SUBC_OP_CATALOG_UPDATE: &str = "catalog.update";
21
22#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
23#[serde(rename_all = "snake_case")]
24pub enum HealthStatus {
25    Ok,
26    Degraded,
27    Failing,
28}
29
30#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
31pub struct HealthReport {
32    pub status: HealthStatus,
33    #[serde(default, skip_serializing_if = "Option::is_none")]
34    pub detail: Option<String>,
35    #[serde(default, skip_serializing_if = "Option::is_none")]
36    pub metrics: Option<Value>,
37}
38
39impl HealthReport {
40    pub fn ok() -> Self {
41        Self {
42            status: HealthStatus::Ok,
43            detail: None,
44            metrics: None,
45        }
46    }
47}
48
49/// subc-to-module channel-0 control RPC body.
50#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
51#[serde(tag = "op")]
52// RouteBind carries the complete bind metadata, while HealthCheck is a marker;
53// preserving the direct wire shape is more useful than boxing every bind field.
54#[allow(clippy::large_enum_variant)]
55pub enum ModuleControlRequest {
56    #[serde(rename = "route.bind")]
57    RouteBind {
58        route_channel: u16,
59        epoch: u32,
60        target: RouteTarget,
61        identity: BindIdentity,
62        /// The daemon's attestation of the consumer, and the only field here a
63        /// provider may grant privilege on.
64        ///
65        /// `Reserved` is minted at exactly one place in the daemon, on the branch
66        /// where the consumer's launch nonce matched a supervised spawn — the
67        /// function that checks is the function that mints, so the value cannot
68        /// exist without the check having run. That property is what a provider is
69        /// relying on, and it is the reason to key authority on this rather than on
70        /// `identity`, which is client-supplied and unattested (see BindIdentity).
71        ///
72        /// Absent means the daemon made no attestation, which is not the same as a
73        /// denial: it is the shape a pre-attestation peer sends. Treat it as
74        /// unattested rather than as trusted-by-default.
75        #[serde(default, skip_serializing_if = "Option::is_none")]
76        principal: Option<Principal>,
77        /// Consumer-declared reverse-request capabilities for the route. This is
78        /// an unverified declaration, not a privilege grant; if a consumer
79        /// over-declares, providers may still send reverse requests that later
80        /// time out or deny. Providers must treat an absent field as no
81        /// reverse-request capability. The vocabulary is open strings; known MCP
82        /// method-family values today are "elicitation", "sampling", and
83        /// "roots".
84        #[serde(default, skip_serializing_if = "Option::is_none")]
85        consumer_capabilities: Option<Vec<String>>,
86        /// The versions of provider roles the consumer speaks on this route,
87        /// role name to version (`{"tool-provider": "v1"}`), copied from the
88        /// consumer's `route.open` unchanged. Like `consumer_capabilities` it
89        /// is the consumer's unverified declaration and grants nothing; a
90        /// provider uses it to choose which version of a role's wire shape to
91        /// speak. The daemon has checked it with [`validate_role_versions`] and
92        /// never sends an empty map. Absent means the consumer declared none:
93        /// a legacy consumer, or a daemon that predates the field.
94        #[serde(default, skip_serializing_if = "Option::is_none")]
95        role_versions: Option<BTreeMap<String, String>>,
96        /// Opaque admission facts supplied by the configured carrier module.
97        #[serde(default, skip_serializing_if = "Option::is_none")]
98        admission_facts: Option<Value>,
99        /// The daemon's stamp of the scope the route was admitted under, taken
100        /// from the owner's synced record at admission. Like `principal`, it is
101        /// the daemon's, never the opener's: a provider may act on it (on
102        /// `owner_authorized`, `delegates` and `agent_id` together), and must
103        /// treat it as fixed for the route's life, because a change that
104        /// revokes authority closes the route.
105        ///
106        /// Absent means the route was opened without a scope, or by a daemon
107        /// that predates scopes. A provider that needs a scope refuses the bind.
108        #[serde(default, skip_serializing_if = "Option::is_none")]
109        scope: Option<ScopeStamp>,
110    },
111    #[serde(rename = "health.check")]
112    HealthCheck {},
113}
114
115/// One-way subc-to-module channel-0 control command.
116#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
117#[serde(tag = "op")]
118pub enum ModuleControlCommand {
119    #[serde(rename = "module.draining")]
120    Draining {
121        reason: RouteCloseReason,
122        /// Absolute Unix-millisecond deadline for this drain.
123        ///
124        /// WALL CLOCK, WHILE THE DAEMON ENFORCES THE CEILING ON A
125        /// SUSPEND-EXCLUDING MONOTONIC CLOCK (`Instant`, supervise.rs). Both
126        /// processes share one host so `CLOCK_REALTIME` agrees exactly, and the
127        /// two clocks diverge only across host sleep: `Instant` stops, wall does
128        /// not. So a module that sleeps mid-drain wakes to a deadline further in
129        /// the past than the daemon's own ceiling, computes LESS remaining time
130        /// than it has, and seals early.
131        ///
132        /// That direction is deliberate and is the safe one — a module stopping
133        /// early loses nothing, since the daemon kills at its own ceiling
134        /// regardless. The reverse (a module believing it has time the daemon
135        /// has already spent) is the failure this ordering avoids. A module must
136        /// therefore treat this as "no later than", never as a grant.
137        deadline_ms: u64,
138    },
139}
140
141/// Module-to-subc channel-0 response body.
142#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
143#[serde(tag = "op")]
144pub enum ModuleControlResponse {
145    /// ACK-only success. Rejections use the `FrameType::Error` lane.
146    #[serde(rename = "route.bind")]
147    RouteBindAck {},
148    #[serde(rename = "health.check")]
149    HealthCheck {
150        status: HealthStatus,
151        #[serde(default, skip_serializing_if = "Option::is_none")]
152        detail: Option<String>,
153        #[serde(default, skip_serializing_if = "Option::is_none")]
154        metrics: Option<Value>,
155    },
156}
157
158/// Module-originated channel-0 control RPC body.
159///
160/// This is intentionally separate from [`ModuleControlRequest`]: that enum is the
161/// daemon-to-module direction (`route.bind`, `health.check`), while these bodies
162/// are sent by an already-registered module to subc on a `REQUEST` frame.
163#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
164#[serde(tag = "op")]
165pub enum ModuleControlRequestFromModule {
166    #[serde(rename = "catalog.update")]
167    CatalogUpdate {
168        provides: Vec<ProviderRole>,
169        /// An attested replacement for the static capability declaration emitted
170        /// by the module's current manifest. `None` preserves the prior
171        /// declaration so existing role-only catalog updates remain byte-identical.
172        #[serde(default, skip_serializing_if = "Option::is_none")]
173        capabilities: Option<CapabilityDeclarations>,
174        /// Updates readiness without re-registering. `None` leaves it unchanged.
175        ///
176        /// Both directions are allowed, but repeatedly flapping readiness looks
177        /// like a restart storm to callers and is a defect in the module.
178        #[serde(default, skip_serializing_if = "Option::is_none")]
179        ready: Option<bool>,
180    },
181    #[serde(rename = "supervisor.live_roots")]
182    LiveRoots {},
183    /// Register this module's full scope set. The owner is the module whose
184    /// registered connection sends it; nothing in the body names the owner.
185    /// Per-record refusals come back in the reply; a refusal of the whole sync
186    /// (not the owner's sync authority, a stale generation, a bound exceeded)
187    /// is an `Error` frame and changes nothing.
188    #[serde(rename = "scope.sync")]
189    ScopeSync {
190        generation: u64,
191        scopes: Vec<ScopeRecord>,
192    },
193    /// Read one scope's current state.
194    #[serde(rename = "scope.describe")]
195    ScopeDescribe {
196        owner: Principal,
197        #[serde(rename = "ref")]
198        scope_ref: String,
199    },
200}
201
202/// Counts of routes for one canonical project root.
203#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
204pub struct LiveRoot {
205    pub project_root: std::path::PathBuf,
206    pub bound: u64,
207    pub pending: u64,
208}
209
210/// subc's channel-0 response body for module-originated control RPCs.
211#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
212#[serde(tag = "op")]
213pub enum ModuleControlResponseToModule {
214    #[serde(rename = "catalog.update")]
215    CatalogUpdate {},
216    #[serde(rename = "supervisor.live_roots")]
217    LiveRoots {
218        roots: Vec<LiveRoot>,
219        unknown_root_bindings: u64,
220        total_bindings: u64,
221    },
222    #[serde(rename = "scope.sync")]
223    ScopeSync {
224        generation: u64,
225        results: Vec<ScopeRecordResult>,
226        #[serde(default, skip_serializing_if = "Vec::is_empty")]
227        ended: Vec<ScopeEnded>,
228    },
229    #[serde(rename = "scope.describe")]
230    ScopeDescribe {
231        status: ScopeStatus,
232        /// The live epoch, or for `ended` the most recent epoch that ended.
233        #[serde(default, skip_serializing_if = "Option::is_none")]
234        scope_epoch: Option<u64>,
235        daemon_incarnation: String,
236        /// Whether the owner has synced since this daemon incarnation started.
237        owner_synced: bool,
238        /// Whether the owner is a module in the daemon's supervised roster.
239        owner_configured: bool,
240        /// The stamp fields, present only when `status` is `live`.
241        #[serde(default, skip_serializing_if = "Option::is_none")]
242        scope: Option<ScopeStamp>,
243    },
244}
245
246/// A module's channel-0 request to confirm one in-flight write with the person.
247///
248/// This standalone body keeps the existing exhaustive module-control enums
249/// unchanged. The daemon validates the summary and the asking connection's
250/// permission before showing a prompt; constructing a request grants nothing.
251#[non_exhaustive]
252#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
253pub struct OperatorConfirmRequest {
254    op: OperatorConfirmOp,
255    pub summary: String,
256    /// The caller's route channel as seen on the asking module's connection.
257    pub route_channel: u16,
258    /// The epoch of that route on the asking module's connection.
259    pub route_epoch: u32,
260}
261
262impl OperatorConfirmRequest {
263    pub fn new(summary: impl Into<String>, route_channel: u16, route_epoch: u32) -> Self {
264        Self {
265            op: OperatorConfirmOp::Confirm,
266            summary: summary.into(),
267            route_channel,
268            route_epoch,
269        }
270    }
271}
272
273/// A successful channel-0 reply to an [`OperatorConfirmRequest`].
274///
275/// Only `outcome: "confirmed"` confirms the write. Refusals use [`crate::ErrorBody`]
276/// instead. Readers must not treat an unknown outcome as confirmation.
277#[non_exhaustive]
278#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
279pub struct OperatorConfirmReply {
280    op: OperatorConfirmOp,
281    pub outcome: String,
282}
283
284impl OperatorConfirmReply {
285    pub fn confirmed() -> Self {
286        Self {
287            op: OperatorConfirmOp::Confirm,
288            outcome: "confirmed".to_string(),
289        }
290    }
291}
292
293// A required, single-valued field both emits the op and refuses a missing or
294// different op when decoding, without exposing a caller-editable tag.
295#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
296enum OperatorConfirmOp {
297    #[serde(rename = "operator.confirm")]
298    Confirm,
299}
300
301impl From<HealthReport> for ModuleControlResponse {
302    fn from(report: HealthReport) -> Self {
303        Self::HealthCheck {
304            status: report.status,
305            detail: report.detail,
306            metrics: report.metrics,
307        }
308    }
309}
310
311impl ModuleControlResponse {
312    pub fn health_report(&self) -> Option<HealthReport> {
313        match self {
314            Self::HealthCheck {
315                status,
316                detail,
317                metrics,
318            } => Some(HealthReport {
319                status: *status,
320                detail: detail.clone(),
321                metrics: metrics.clone(),
322            }),
323            Self::RouteBindAck {} => None,
324        }
325    }
326}
327
328/// Module-to-subc channel-0 push body.
329#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
330#[serde(tag = "op")]
331pub enum ModuleControlPush {
332    #[serde(rename = "route.status")]
333    RouteStatus {
334        route_channel: u16,
335        route_epoch: u32,
336        status: String,
337    },
338}
339
340/// The wire name of the `role_versions` field on `route.open` and
341/// `route.bind`, for the `detail.field` of the `invalid_request` error that
342/// refuses a malformed one.
343pub const ROLE_VERSIONS_FIELD: &str = "role_versions";
344
345/// Most entries a `role_versions` map may hold.
346pub const MAX_ROLE_VERSIONS: usize = 8;
347
348/// The longest role name accepted in `role_versions`, in bytes.
349pub const MAX_ROLE_NAME_LEN: usize = 64;
350
351/// Why a `role_versions` map was refused.
352#[derive(Clone, Debug, PartialEq, Eq)]
353#[non_exhaustive]
354pub enum RoleVersionsError {
355    /// The map has more than [`MAX_ROLE_VERSIONS`] entries.
356    TooMany { count: usize },
357    /// A role name is not lowercase ASCII letters and digits in words joined
358    /// by single hyphens (`tool-provider`), or is longer than
359    /// [`MAX_ROLE_NAME_LEN`] bytes.
360    InvalidRole { role: String },
361    /// A version is not `v` followed by a positive integer without leading
362    /// zeros (`v1`, `v12`).
363    InvalidVersion { role: String, version: String },
364}
365
366impl RoleVersionsError {
367    /// The request field the error is about: always [`ROLE_VERSIONS_FIELD`].
368    pub fn field(&self) -> &'static str {
369        ROLE_VERSIONS_FIELD
370    }
371}
372
373impl std::fmt::Display for RoleVersionsError {
374    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
375        match self {
376            Self::TooMany { count } => write!(
377                f,
378                "{ROLE_VERSIONS_FIELD} has {count} entries; at most {MAX_ROLE_VERSIONS} are allowed"
379            ),
380            Self::InvalidRole { role } => write!(
381                f,
382                "{ROLE_VERSIONS_FIELD} names role {role:?}, which is not lowercase letters and \
383                 digits in words joined by '-', at most {MAX_ROLE_NAME_LEN} bytes"
384            ),
385            Self::InvalidVersion { role, version } => write!(
386                f,
387                "{ROLE_VERSIONS_FIELD} gives role {role:?} version {version:?}, which is not 'v' \
388                 followed by a positive integer without leading zeros"
389            ),
390        }
391    }
392}
393
394impl std::error::Error for RoleVersionsError {}
395
396/// Check a `role_versions` map: at most [`MAX_ROLE_VERSIONS`] entries, each
397/// role name matching `^[a-z0-9]+(-[a-z0-9]+)*$` in at most
398/// [`MAX_ROLE_NAME_LEN`] bytes, and each version matching `^v[1-9][0-9]*$`.
399///
400/// The daemon refuses a `route.open` that fails this, so a consumer can run
401/// the same check before sending. An empty map passes; the daemon treats it
402/// as no declaration at all.
403pub fn validate_role_versions(
404    role_versions: &BTreeMap<String, String>,
405) -> Result<(), RoleVersionsError> {
406    if role_versions.len() > MAX_ROLE_VERSIONS {
407        return Err(RoleVersionsError::TooMany {
408            count: role_versions.len(),
409        });
410    }
411    for (role, version) in role_versions {
412        if !is_role_name(role) {
413            return Err(RoleVersionsError::InvalidRole { role: role.clone() });
414        }
415        if !is_role_version(version) {
416            return Err(RoleVersionsError::InvalidVersion {
417                role: role.clone(),
418                version: version.clone(),
419            });
420        }
421    }
422    Ok(())
423}
424
425fn is_role_name(role: &str) -> bool {
426    !role.is_empty()
427        && role.len() <= MAX_ROLE_NAME_LEN
428        && role.split('-').all(|word| {
429            !word.is_empty()
430                && word
431                    .bytes()
432                    .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
433        })
434}
435
436fn is_role_version(version: &str) -> bool {
437    let bytes = version.as_bytes();
438    bytes.len() >= 2
439        && bytes[0] == b'v'
440        && (b'1'..=b'9').contains(&bytes[1])
441        && bytes[2..].iter().all(u8::is_ascii_digit)
442}
443
444#[cfg(test)]
445mod tests {
446    use super::*;
447
448    fn map(entries: &[(&str, &str)]) -> BTreeMap<String, String> {
449        entries
450            .iter()
451            .map(|(role, version)| (role.to_string(), version.to_string()))
452            .collect()
453    }
454
455    #[test]
456    fn role_versions_accept_role_names_and_positive_versions() {
457        for entries in [
458            vec![],
459            vec![("tool-provider", "v1")],
460            vec![("a", "v9"), ("b2", "v10"), ("x-1-y", "v1203")],
461        ] {
462            assert_eq!(
463                validate_role_versions(&map(&entries)),
464                Ok(()),
465                "{entries:?}"
466            );
467        }
468        let longest = "a".repeat(MAX_ROLE_NAME_LEN);
469        assert_eq!(validate_role_versions(&map(&[(&longest, "v1")])), Ok(()));
470        let full: BTreeMap<String, String> = (0..MAX_ROLE_VERSIONS)
471            .map(|index| (format!("role-{index}"), "v1".to_string()))
472            .collect();
473        assert_eq!(validate_role_versions(&full), Ok(()));
474    }
475
476    #[test]
477    fn role_versions_refuse_malformed_role_names() {
478        let too_long = "a".repeat(MAX_ROLE_NAME_LEN + 1);
479        for role in [
480            "",
481            "Tool-provider",
482            "tool_provider",
483            "tool provider",
484            "-tool",
485            "tool-",
486            "tool--provider",
487            "tool.provider",
488            "outil-é",
489            too_long.as_str(),
490        ] {
491            let error = validate_role_versions(&map(&[(role, "v1")])).unwrap_err();
492            assert_eq!(
493                error,
494                RoleVersionsError::InvalidRole {
495                    role: role.to_string()
496                },
497                "{role:?}"
498            );
499            assert_eq!(error.field(), "role_versions");
500        }
501    }
502
503    #[test]
504    fn role_versions_refuse_malformed_versions() {
505        for version in [
506            "", "v", "v0", "v01", "1", "V1", "v1.0", "v-1", "v1 ", " v1", "vx",
507        ] {
508            let error = validate_role_versions(&map(&[("tool-provider", version)])).unwrap_err();
509            assert_eq!(
510                error,
511                RoleVersionsError::InvalidVersion {
512                    role: "tool-provider".to_string(),
513                    version: version.to_string(),
514                },
515                "{version:?}"
516            );
517            assert_eq!(error.field(), "role_versions");
518        }
519    }
520
521    #[test]
522    fn role_versions_refuse_more_than_eight_entries() {
523        let nine: BTreeMap<String, String> = (0..=MAX_ROLE_VERSIONS)
524            .map(|index| (format!("role-{index}"), "v1".to_string()))
525            .collect();
526        let error = validate_role_versions(&nine).unwrap_err();
527        assert_eq!(error, RoleVersionsError::TooMany { count: 9 });
528        assert_eq!(error.field(), "role_versions");
529        assert!(error.to_string().starts_with("role_versions"), "{error}");
530    }
531
532    #[test]
533    fn route_bind_omits_absent_role_versions_and_carries_present_ones_verbatim() {
534        let bind = |role_versions| ModuleControlRequest::RouteBind {
535            route_channel: 1,
536            epoch: 1,
537            target: crate::RouteTarget::ToolProvider {
538                module_id: "aft".to_string(),
539            },
540            identity: crate::BindIdentity::new("/tmp/p", "h", "s"),
541            principal: None,
542            consumer_capabilities: None,
543            role_versions,
544            admission_facts: None,
545            scope: None,
546        };
547        let absent = serde_json::to_value(bind(None)).unwrap();
548        assert!(absent.get("role_versions").is_none(), "{absent}");
549        let present = bind(Some(map(&[("tool-provider", "v1")])));
550        let encoded = serde_json::to_value(&present).unwrap();
551        assert_eq!(
552            encoded["role_versions"],
553            serde_json::json!({ "tool-provider": "v1" })
554        );
555        let decoded: ModuleControlRequest = serde_json::from_value(encoded).unwrap();
556        assert_eq!(decoded, present);
557    }
558}