//! Host-owned wire vocabulary. Python values never deserialize into host code. use crate::work::Impact; use anyhow::{anyhow, bail, ensure, Result}; use serde::{Deserialize, Serialize}; use serde_json::Value; use std::{fmt, path::PathBuf, str::FromStr}; use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; pub const VERSION: u16 = 1; /// The relay protocol is a DISTINCT version plane from the worker protocol: /// a remote launcher multiplexes framed control over a different transport, /// and the two must never be conflated (REMOTE-KERNELS.md section 6). pub const RELAY_PROTOCOL_VERSION: u16 = 1; pub const MAX_FRAME: usize = 1_048_576; pub const MAX_CODE: usize = 262_144; pub const MAX_OUTPUT: usize = 65_536; /// Remote dispatcher values travel in the worker frame, not cell output. pub const MAX_REMOTE_VALUE: usize = 1_040_000; pub const MAX_REMOTE_ERROR: usize = 16_384; macro_rules! identity { ($name:ident) => { #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(try_from = "String", into = "String")] pub struct $name(String); impl $name { pub fn fresh() -> Self { Self(uuid::Uuid::new_v4().to_string()) } pub fn as_str(&self) -> &str { &self.0 } } impl TryFrom for $name { type Error = String; fn try_from(value: String) -> std::result::Result { if value.is_empty() || value.len() > 128 || !value .bytes() .all(|b| b.is_ascii_alphanumeric() || b"_.-".contains(&b)) { return Err(concat!("invalid ", stringify!($name)).into()); } Ok(Self(value)) } } impl From<$name> for String { fn from(id: $name) -> Self { id.0 } } impl FromStr for $name { type Err = String; fn from_str(value: &str) -> std::result::Result { value.to_owned().try_into() } } impl fmt::Display for $name { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { write!(f, "{}", self.0) } } }; } identity!(SessionId); identity!(Generation); identity!(OperationId); identity!(RequestId); identity!(EffectId); #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum Role { Workbench, Hooks, Remote, } impl Role { pub fn argument(self) -> &'static str { match self { Self::Workbench => "workbench", Self::Hooks => "hooks", Self::Remote => "remote", } } } #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] pub enum HookSlot { #[serde(rename = "attention.plan")] Attention, #[serde(rename = "delivery.review")] Delivery, } impl HookSlot { pub fn name(self) -> &'static str { match self { Self::Attention => "attention.plan", Self::Delivery => "delivery.review", } } } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct Envelope { pub version: u16, pub session_id: SessionId, pub generation: Generation, pub message: Message, /// Target/namespace binding. Local workers omit it; a remote worker /// envelope must carry it, and a LOCAL envelope claiming one is a /// forged identity. Checked with validate_envelope_target. #[serde(default, skip_serializing_if = "Option::is_none")] pub target: Option, } /// The target/namespace identity bound into every envelope of one worker. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct TargetStamp { pub target_id: crate::target::TargetId, pub namespace_id: crate::target::NamespaceId, pub namespace_generation: crate::target::NamespaceGeneration, } impl fmt::Display for TargetStamp { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { write!( f, "{}[{}@{}]", self.target_id, self.namespace_id, self.namespace_generation ) } } /// Typed envelope-level target/namespace mismatch kinds. These join /// version/generation/session mismatches in the fatal class. #[derive(Clone, Debug, PartialEq, Eq)] pub enum BindingMismatch { /// A worker whose launch was local claimed a remote target identity. ForgedTargetStamp, /// A remote envelope omitted, or altered, its target/namespace binding. MissingTargetStamp { expected: TargetStamp }, /// The stamp does not match the placement the host issued. ForeignTargetStamp { expected: TargetStamp, actual: TargetStamp, }, } impl fmt::Display for BindingMismatch { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::ForgedTargetStamp => { f.write_str("local worker envelope must not claim a remote target stamp") } Self::MissingTargetStamp { expected } => write!( f, "remote envelope is missing its target binding: expected {expected}" ), Self::ForeignTargetStamp { expected, actual } => { write!(f, "foreign target stamp: expected {expected}, got {actual}") } } } } impl std::error::Error for BindingMismatch {} /// Host-side check: every envelope must agree with the placement it was /// launched under, exactly as session/generation already must. A mismatch /// is fatal; nothing is silently re-routed. pub fn validate_envelope_target(envelope: &Envelope, expected: Option<&TargetStamp>) -> Result<()> { match (envelope.target.clone(), expected) { (Some(_), None) => bail!(BindingMismatch::ForgedTargetStamp), (None, Some(expected)) => bail!(BindingMismatch::MissingTargetStamp { expected: expected.clone(), }), (Some(actual), Some(expected)) if &actual != expected => { bail!(BindingMismatch::ForeignTargetStamp { expected: expected.clone(), actual, }) } _ => Ok(()), } } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum Purpose { #[default] Interactive, Maintenance, Eval, } impl Purpose { #[must_use] pub const fn as_str(self) -> &'static str { match self { Self::Interactive => "interactive", Self::Maintenance => "maintenance", Self::Eval => "eval", } } } impl fmt::Display for Purpose { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str(self.as_str()) } } /// The runtime identity a worker reports in its Hello, distinct from the /// target identity the HOST verified at launch (klbr_runtime::target). /// The host cross-checks the two before admitting work. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct RuntimeIdentity { pub platform: String, pub arch: String, pub interpreter_path: String, pub interpreter_version: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub environment_fingerprint: Option, #[serde(default)] pub capabilities: Vec, } /// One capability the endpoint advertises. Advertised-but-unusable is /// expressible as such (e.g. Closures { dill_available: false }); "healthy /// because the name exists" is not a valid reading of this set. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "capability", rename_all = "snake_case", deny_unknown_fields)] pub enum Capability { /// A structured command process can be started in the target. Shell { /// Resolved shell, when the endpoint knows it (e.g. "/bin/zsh"). shell: Option, }, /// Bounded remote filesystem I/O under the checked workspace. FileSystem, /// Third-party imports are available in the worker namespace. Imports, /// Trusted function shipment. Unusable without dill; the host must /// treat it as Unsupported rather than routing around it silently. Closures { dill_available: bool }, /// A skill manifest with a source revision is served. SkillManifest { revision: String }, /// The relay speaks this framing version. RelayFraming { version: u16 }, } impl Capability { #[must_use] pub fn is_usable(&self) -> bool { match self { Self::Shell { .. } | Self::FileSystem | Self::Imports => true, Self::Closures { dill_available } => *dill_available, Self::SkillManifest { revision } => !revision.trim().is_empty(), Self::RelayFraming { version } => *version <= RELAY_PROTOCOL_VERSION, } } #[must_use] pub fn unusable_reason(&self) -> Option<&'static str> { match self { Self::Closures { dill_available: false, } => Some("dill is not importable in the worker namespace"), Self::SkillManifest { revision } if revision.trim().is_empty() => { Some("skill manifest advertised with an empty revision") } Self::RelayFraming { version } if *version > RELAY_PROTOCOL_VERSION => { Some("relay framing version is newer than the host supports") } _ => None, } } } impl RuntimeIdentity { /// Capabilities that were advertised but cannot actually be used. #[must_use] pub fn unusable_capabilities(&self) -> Vec<&Capability> { self.capabilities .iter() .filter(|c| !c.is_usable()) .collect() } /// Typed support check used BEFORE admitting work that needs a /// capability. Exact-typed match plus usability: a name match with an /// unusable payload is not support. #[must_use] pub fn supports(&self, required: &Capability) -> bool { self.capabilities.iter().any(|c| c == required) && required.is_usable() } fn empty_field(kind: &'static str, value: &str) -> Result<()> { if value.trim().is_empty() { bail!(NegotiationError::RuntimeIncomplete { field: kind }); } Ok(()) } } /// The target/namespace binding a REMOTE worker reports in its Hello. /// /// `stamp` is the host-issued launch binding echoed by the worker (a /// binding check, the same evidence class as --token). `identity` is NOT /// the echo: it is the identity MEASURED on the target by the launcher's /// own probe (m1-relay's StartupReport) and attached at the relay boundary. /// A worker frame can carry the stamp alone; only a frame with a measured /// identity may be admitted to the host, and a frame that CLAIMS an /// identity without the launcher measuring it is refused, not merged. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct HelloBinding { pub stamp: TargetStamp, #[serde(default, skip_serializing_if = "Option::is_none")] pub identity: Option, } /// What a worker reports about itself and its placement at Hello time. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct WorkerIdentity { pub protocol_version: u16, pub relay_protocol_version: u16, /// Remote workers report this; local workers omit it. #[serde(default, skip_serializing_if = "Option::is_none")] pub target: Option, pub runtime: RuntimeIdentity, } /// Typed negotiation mismatch kinds, each its own check failure class. #[derive(Clone, Debug, PartialEq, Eq)] pub enum NegotiationError { WorkerProtocolVersionMismatch { reported: u16, supported: u16, }, RelayProtocolVersionMismatch { reported: u16, supported: u16, }, /// A remote launch requires a reported target binding. TargetBindingMissing, /// A remote worker frame claimed a placement stamp but carried no /// MEASURED target identity. The echo is a binding check, never /// verified identity; only the launcher's own probe may supply it. TargetIdentityUnmeasured, /// The reported binding contradicts the placement the host issued. BindingMismatch { field: &'static str, reported: String, resolved: String, }, /// The runtime identity contradicts the resolved target identity. RuntimeMismatch { field: &'static str, reported: String, resolved: String, }, /// The runtime identity is internally incomplete. RuntimeIncomplete { field: &'static str, }, } impl fmt::Display for NegotiationError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::WorkerProtocolVersionMismatch { reported, supported, } => write!( f, "worker protocol version mismatch: reported {reported}, host supports {supported}" ), Self::RelayProtocolVersionMismatch { reported, supported, } => write!( f, "relay protocol version mismatch: reported {reported}, host supports {supported}" ), Self::TargetBindingMissing => f.write_str( "remote worker must report its target/namespace binding in hello identity", ), Self::TargetIdentityUnmeasured => f.write_str( "remote hello carries a placement stamp but no MEASURED target identity; \ the echoed binding is not evidence about the target", ), Self::BindingMismatch { field, reported, resolved, } => write!( f, "worker target binding mismatch on {field}: reported {reported:?}, host resolved {resolved:?}" ), Self::RuntimeMismatch { field, reported, resolved, } => write!( f, "worker runtime mismatch on {field}: reported {reported:?}, host resolved {resolved:?}" ), Self::RuntimeIncomplete { field } => { write!(f, "worker runtime identity {field} is missing") } } } } impl std::error::Error for NegotiationError {} impl WorkerIdentity { /// Host-side check of a reported worker identity. `expected` is the /// binding the host ISSUED at launch: None for a local launch (any /// reported stamp is a forged identity), Some for a remote launch /// (a missing or differing binding is fatal, field by field). A local /// worker that reports no identity at all stays fully compatible. pub fn validate(&self, expected: Option<&HelloBinding>) -> Result<()> { ensure!( self.protocol_version == VERSION, NegotiationError::WorkerProtocolVersionMismatch { reported: self.protocol_version, supported: VERSION, } ); ensure!( self.relay_protocol_version == RELAY_PROTOCOL_VERSION, NegotiationError::RelayProtocolVersionMismatch { reported: self.relay_protocol_version, supported: RELAY_PROTOCOL_VERSION, } ); RuntimeIdentity::empty_field("platform", &self.runtime.platform)?; RuntimeIdentity::empty_field("arch", &self.runtime.arch)?; RuntimeIdentity::empty_field("interpreter_path", &self.runtime.interpreter_path)?; RuntimeIdentity::empty_field("interpreter_version", &self.runtime.interpreter_version)?; let mismatch = |field, reported: String, resolved: String| { if reported == resolved { Ok(()) } else { bail!(NegotiationError::BindingMismatch { field, reported, resolved, }) } }; match (&self.target, expected) { (Some(_), None) => bail!(BindingMismatch::ForgedTargetStamp), (None, Some(_)) => bail!(NegotiationError::TargetBindingMissing), (Some(reported), Some(expected)) => { mismatch( "target_id", reported.stamp.target_id.to_string(), expected.stamp.target_id.to_string(), )?; mismatch( "namespace_id", reported.stamp.namespace_id.as_str().to_owned(), expected.stamp.namespace_id.as_str().to_owned(), )?; mismatch( "namespace_generation", reported.stamp.namespace_generation.as_str().to_owned(), expected.stamp.namespace_generation.as_str().to_owned(), )?; let measured = reported .identity .as_ref() .ok_or(NegotiationError::TargetIdentityUnmeasured)?; // The host's expectation is its own resolved identity; a // placement expectation without one is a host-side bug. let b = expected.identity.as_ref().ok_or_else(|| { anyhow!("host issued a remote placement without a resolved target identity") })?; let a = measured; mismatch("account", a.account.clone(), b.account.clone())?; mismatch("machine", a.machine.clone(), b.machine.clone())?; mismatch("platform", a.platform.clone(), b.platform.clone())?; mismatch("arch", a.arch.clone(), b.arch.clone())?; mismatch( "interpreter_path", a.interpreter_path.clone(), b.interpreter_path.clone(), )?; mismatch( "interpreter_version", a.interpreter_version.clone(), b.interpreter_version.clone(), )?; mismatch("cwd", a.cwd.clone(), b.cwd.clone())?; mismatch( "capability_revision", a.capability_revision.to_string(), b.capability_revision.to_string(), )?; mismatch( "environment_fingerprint", a.environment_fingerprint .as_ref() .map(ToString::to_string) .unwrap_or_default(), b.environment_fingerprint .as_ref() .map(ToString::to_string) .unwrap_or_default(), )?; // The reported runtime must also agree with the expected identity. let runtime_mismatch = |field, reported: String, resolved: String| { if reported == resolved { Ok(()) } else { bail!(NegotiationError::RuntimeMismatch { field, reported, resolved, }) } }; runtime_mismatch( "platform", self.runtime.platform.clone(), b.platform.clone(), )?; runtime_mismatch("arch", self.runtime.arch.clone(), b.arch.clone())?; runtime_mismatch( "interpreter_path", self.runtime.interpreter_path.clone(), b.interpreter_path.clone(), )?; } (None, None) => { // Local worker without identity: the compatible local path. } } Ok(()) } /// Stamp-level check for launch sites that hold the placement but not /// (yet) a verified identity: the placement the host issued must be /// echoed exactly, and a local launch must not receive any stamp. The /// full identity cross-check is validate(expected: Option<&HelloBinding>). pub fn validate_stamp(&self, expected: Option<&TargetStamp>) -> Result<()> { ensure!( self.protocol_version == VERSION, NegotiationError::WorkerProtocolVersionMismatch { reported: self.protocol_version, supported: VERSION, } ); ensure!( self.relay_protocol_version == RELAY_PROTOCOL_VERSION, NegotiationError::RelayProtocolVersionMismatch { reported: self.relay_protocol_version, supported: RELAY_PROTOCOL_VERSION, } ); match (&self.target, expected) { (Some(_), None) => bail!(BindingMismatch::ForgedTargetStamp), (None, Some(_)) => bail!(NegotiationError::TargetBindingMissing), (Some(reported), Some(expected)) => { ensure!( reported.stamp == *expected, BindingMismatch::ForeignTargetStamp { expected: expected.clone(), actual: reported.stamp.clone(), } ); } (None, None) => {} } Ok(()) } } #[allow(clippy::large_enum_variant)] #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum Message { Hello { token: String, role: Role, revision: Option, slots: Vec, /// Optional runtime/target identity. Local workers may omit it /// (the field defaults, so the local handshake is byte-compatible); /// a remote worker MUST report it and the host validates it before /// admitting work. #[serde(default, skip_serializing_if = "Option::is_none")] identity: Option, /// Content identity of the skill bodies this worker was spawned with. /// Echoed rather than assumed, so the durable reset can name the /// instruction set this generation actually loaded, not what the roots /// happen to hold when the record is written. #[serde(default, skip_serializing_if = "Option::is_none")] skills: Option, }, Execute { id: OperationId, code: String, #[serde(default)] purpose: Purpose, }, Invoke { id: OperationId, slot: HookSlot, input: Value, }, RemoteDispatch { id: OperationId, operation: String, payload: Value, }, Interrupt { id: OperationId, }, Shutdown, HostRequest { id: RequestId, execution_id: OperationId, request: HostRequest, }, HostReply { id: RequestId, execution_id: OperationId, reply: HostReply, }, Completed { id: OperationId, outcome: Outcome, }, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "method", deny_unknown_fields)] pub enum HostRequest { #[serde(rename = "local.send")] LocalSend { text: String }, #[serde(rename = "runtime.wait")] Wait { seconds: Option, reason: String, }, #[serde(rename = "hooks.list")] HooksList, #[serde(rename = "hooks.describe")] HooksDescribe { slot: HookSlot }, #[serde(rename = "behavior.activate")] BehaviorActivate { path: PathBuf, expected_epoch: i64, #[serde(default, skip_serializing_if = "Option::is_none")] revision: Option, #[serde(default, skip_serializing_if = "Option::is_none")] work_item_id: Option, }, #[serde(rename = "maintenance.originate")] MaintenanceOriginate { #[serde(default)] item_id: Option, purpose: String, reason: String, #[serde(default)] evidence_refs: Vec, #[serde(default)] impact: Option, }, #[serde(rename = "maintenance.assess")] MaintenanceAssess { work_item_id: String, outcome: String, #[serde(default)] reason: Option, #[serde(default)] missing_refs: Vec, #[serde(default)] revisit_condition: Option, /// Bounded revisit deadline for needs-evidence (wall-clock ms). #[serde(default)] revisit_deadline_ms: Option, /// Bounded revisit trigger for deferred (wall-clock ms). #[serde(default)] trigger_ms: Option, #[serde(default)] candidate_revision: Option, #[serde(default)] check_receipts: Vec, }, #[serde(rename = "evidence.resolve")] EvidenceResolve { refs: Vec }, #[serde(rename = "workspace.info")] WorkspaceInfo, /// Run one argv directly in the session's immutable source tree. The host /// validates and owns the process; workers never choose a shell. #[serde(rename = "workspace.run")] WorkspaceRun { argv: Vec, cwd: PathBuf, #[serde(default)] env: std::collections::BTreeMap, #[serde(default, skip_serializing_if = "Option::is_none")] stdin_base64: Option, timeout_ms: u64, max_output_bytes: u64, }, #[serde(rename = "runtime.instructions")] RuntimeInstructions, #[serde(rename = "runtime.inspect")] RuntimeInspect, /// Every skill id this host offers, with the digest each body had when it was /// indexed, plus every id it refused and why. #[serde(rename = "skills.index")] SkillsIndex, /// One skill body. The host reads it: the id a model learns is a handle, and a /// body that changed since it was indexed is refused rather than half-read. #[serde(rename = "skills.load")] SkillsLoad { id: String }, /// Ensure an auxiliary remote Python namespace owned by the host. #[serde(rename = "remote.ensure")] RemoteEnsure { target: String, #[serde(default, skip_serializing_if = "Option::is_none")] workspace: Option, }, /// Invoke one reflected operation in an already ensured namespace. The /// payload is data-only JSON; operation identity is supplied by the /// caller and checked by the host before dispatch. #[serde(rename = "remote.call")] RemoteCall { target: String, namespace_id: String, namespace_generation: String, operation_id: String, operation: String, payload: Value, }, #[serde(rename = "remote.cancel")] RemoteCancel { namespace_id: String, namespace_generation: String, operation_id: String, }, #[serde(rename = "remote.close")] RemoteClose { namespace_id: String, namespace_generation: String, }, #[serde(rename = "orientation.stage")] OrientationStage(OrientationStageReq), #[serde(rename = "orientation.current")] OrientationCurrent(OrientationCurrentReq), #[serde(rename = "orientation.adopt")] OrientationAdopt(OrientationAdoptReq), #[serde(rename = "frontier.select")] FrontierSelect(FrontierSelectReq), #[serde(rename = "peer.open", alias = "agents.open")] PeerOpen { #[serde(default)] peer_id: Option, #[serde(default)] mode: Option, #[serde(default)] profile: Option, #[serde(default)] workspace: Option, #[serde(default)] kernel_target: Option, #[serde(default)] budget: Option, #[serde(default)] branch_parent_revision: Option, #[serde(default)] evidence: Vec, }, #[serde(rename = "peer.send", alias = "peer.message")] PeerSend { peer_id: String, request_id: String, #[serde(alias = "text")] content: String, #[serde(default, alias = "evidence_refs")] evidence: Vec, }, #[serde(rename = "peer.result")] PeerResult { request_id: String, #[serde(default)] peer_id: Option, #[serde(default)] wait_ms: Option, }, #[serde(rename = "peer.submit_result")] PeerSubmitResult { request_id: String, #[serde(default)] peer_id: Option, result: serde_json::Value, }, #[serde(rename = "peer.cancel")] PeerCancel { request_id: String, #[serde(default)] peer_id: Option, }, #[serde(rename = "peer.get")] PeerGet { peer_id: String }, #[serde(rename = "peer.list")] PeerList, #[serde(rename = "peer.stop")] PeerStop { peer_id: String, #[serde(default)] drain: Option, }, } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum CarrierMode { #[default] Form, Situate, Elaborate, RepairReturn, } impl CarrierMode { #[must_use] pub const fn as_str(self) -> &'static str { match self { Self::Form => "form", Self::Situate => "situate", Self::Elaborate => "elaborate", Self::RepairReturn => "repair_return", } } } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum CarrierPlacement { Header, #[default] #[serde(alias = "prefix")] LeadingSystem, Tail, ScopedDirective, } impl CarrierPlacement { #[must_use] pub const fn as_str(self) -> &'static str { match self { Self::Header => "header", Self::LeadingSystem => "leading_system", Self::Tail => "tail", Self::ScopedDirective => "scoped_directive", } } } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum CarrierRole { #[default] System, } impl CarrierRole { #[must_use] pub const fn as_str(self) -> &'static str { match self { Self::System => "system", } } } #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum TargetSelector { Agent, User { user_id: String }, Peer { peer_id: String }, Inherited { source_ref: String }, } #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum CarrierAuthor { Model { model_id: String }, Actor { actor_id: String }, Human { name: String }, } #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum AdoptionActor { Agent, Operator { user_id: String }, Peer { peer_id: String }, } #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum ApplicationBoundary { #[default] NextProviderBoundary, StoredOnly, } impl TryFrom<&str> for ApplicationBoundary { type Error = String; fn try_from(value: &str) -> Result { match value { "next_provider_boundary" => Ok(Self::NextProviderBoundary), "stored_only" => Ok(Self::StoredOnly), _ => Err(value.into()), } } } #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum OrientationScope { Global, Project { id: String }, Session { id: String }, Work { id: String }, Peer { peer_id: String }, } impl OrientationScope { #[must_use] pub fn is_global(&self) -> bool { matches!(self, Self::Global) } #[must_use] pub fn is_peer(&self) -> bool { matches!(self, Self::Peer { .. }) } #[must_use] pub fn peer_id(&self) -> Option<&str> { match self { Self::Peer { peer_id } => Some(peer_id.as_str()), _ => None, } } #[must_use] pub fn label(&self) -> String { match self { Self::Global => "global".to_string(), Self::Project { id } => format!("project:{id}"), Self::Session { id } => format!("session:{id}"), Self::Work { id } => format!("work:{id}"), Self::Peer { peer_id } => format!("peer:{peer_id}"), } } } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct OrientationStageReq { pub candidate_id: String, pub text: String, pub scope: OrientationScope, pub mode: CarrierMode, pub target: String, pub parent_revision: Option, pub target_selector: TargetSelector, pub carrier_author: CarrierAuthor, #[serde(default)] pub source_refs: Vec, pub carrier_placement: Option, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct OrientationCurrentReq { pub scope: OrientationScope, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct OrientationAdoptReq { pub candidate_id: String, pub expected_parent: Option, pub adoption_actor: AdoptionActor, pub operation_id: String, pub apply_at: Option, } /// Request payload for actor frontier choice selection. #[derive(Clone, Debug, Serialize, Deserialize)] pub struct FrontierSelectReq { pub frontier_id: String, pub action: String, pub candidate_carrier: Option, pub adoption_actor: String, pub deadline_ms: Option, pub rationale: Option, pub source_refs: Vec, } pub type HostFuture<'a> = std::pin::Pin< Box> + Send + 'a>, >; /// Injected host capability allowing klbr-runtime to dispatch orientation and frontier /// operations to klbr-core/klbr-daemon implementations without a cyclic dependency. pub trait OrientationHost: Send + Sync { fn stage_orientation<'a>( &'a self, store: &'a crate::store::Store, session_id: &'a SessionId, req: OrientationStageReq, ) -> HostFuture<'a>; fn current_orientation<'a>( &'a self, store: &'a crate::store::Store, session_id: &'a SessionId, req: OrientationCurrentReq, ) -> HostFuture<'a>; fn adopt_orientation<'a>( &'a self, store: &'a crate::store::Store, session_id: &'a SessionId, req: OrientationAdoptReq, ) -> HostFuture<'a>; fn select_frontier<'a>( &'a self, store: &'a crate::store::Store, session_id: &'a SessionId, req: FrontierSelectReq, ) -> HostFuture<'a>; } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerOpenReq { pub peer_id: Option, pub mode: Option, pub profile: Option, pub workspace: Option, pub kernel_target: Option, pub budget: Option, pub branch_parent_revision: Option, pub evidence: Vec, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerSendReq { pub peer_id: String, pub request_id: String, pub content: String, pub evidence: Vec, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerResultReq { pub request_id: String, pub peer_id: Option, pub wait_ms: Option, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerSubmitResultReq { pub request_id: String, pub peer_id: Option, pub result: serde_json::Value, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerCancelReq { pub request_id: String, pub peer_id: Option, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerGetReq { pub peer_id: String, } #[derive(Clone, Debug, Serialize, Deserialize)] pub struct PeerStopReq { pub peer_id: String, pub drain: Option, } /// Injected host capability allowing klbr-runtime to dispatch peer /// operations to klbr-daemon / session registry implementations. pub trait PeerHost: Send + Sync { fn open_peer<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerOpenReq, ) -> HostFuture<'a>; fn get_peer<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerGetReq, ) -> HostFuture<'a>; fn list_peers<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, ) -> HostFuture<'a>; fn send_peer_message<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerSendReq, ) -> HostFuture<'a>; fn get_peer_result<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerResultReq, ) -> HostFuture<'a>; fn submit_peer_result<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerSubmitResultReq, ) -> HostFuture<'a>; fn cancel_peer_request<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerCancelReq, ) -> HostFuture<'a>; fn stop_peer<'a>( &'a self, store: &'a crate::store::Store, caller_session_id: &'a SessionId, req: PeerStopReq, ) -> HostFuture<'a>; } /// Closed, host-validated outcome of a maintenance assessment. /// /// Parsed at the host boundary before anything is persisted. An unknown /// outcome tag, a missing required variant field, or a contradictory /// reference is a typed error to the caller, never a silent default. The /// variant fields here are the durable contract; the daemon's /// MaintenanceDisposition consumes them. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "outcome", rename_all = "kebab-case", deny_unknown_fields)] pub enum MaintenanceAssessment { /// Explicit, evidence-backed judgment that no change is warranted. /// This is the only autonomous close path for an objective. NoChange { reason: String, #[serde(default)] evidence_refs: Vec, }, /// Work cannot proceed until the named dependencies exist or the /// bounded revisit deadline passes. NeedsEvidence { #[serde(default)] missing_refs: Vec, #[serde(default)] revisit_condition: Option, #[serde(default)] revisit_deadline_ms: Option, }, /// Objective retained but paused until a bounded trigger. Deferred { reason: String, #[serde(default)] trigger_ms: Option, }, /// Candidate failed validation or execution checks. CandidateFailed { error: String, #[serde(default)] candidate_revision: Option, }, /// Candidate activated for this work item. The host validates the /// candidate revision, executed check receipts and the activation /// record for this work/attempt before accepting it. Activated { candidate_revision: String, #[serde(default)] check_receipts: Vec, }, } impl MaintenanceAssessment { /// The closed outcome tag as serialized on the wire and in the store. pub fn outcome_tag(&self) -> &'static str { match self { Self::NoChange { .. } => "no-change", Self::NeedsEvidence { .. } => "needs-evidence", Self::Deferred { .. } => "deferred", Self::CandidateFailed { .. } => "candidate-failed", Self::Activated { .. } => "activated", } } /// Parse the closed outcome from an outcome tag plus flat variant fields. /// One parser for both the host boundary and the durable store commit. pub fn parse(outcome_tag: &str, fields: &Value) -> Result { let text = |name: &str| -> Option { fields .get(name) .and_then(|v| v.as_str()) .map(str::to_owned) .filter(|s| !s.trim().is_empty()) }; let list = |name: &str| -> Vec { fields .get(name) .and_then(|v| v.as_array()) .map(|a| { a.iter() .filter_map(|x| x.as_str().map(str::to_owned)) .collect() }) .unwrap_or_default() }; let ms = |name: &str| -> Option { fields.get(name).and_then(|v| v.as_i64()) }; let checked_ms = |name: &str| -> Result> { match ms(name) { Some(v) if v > 0 => Ok(Some(v)), Some(_) => bail!("{name} must be a positive wall-clock ms value"), None => Ok(None), } }; match outcome_tag { "no-change" => { let reason = text("reason") .ok_or_else(|| anyhow!("no-change assessment requires a non-empty reason"))?; Ok(Self::NoChange { reason, evidence_refs: list("evidence_refs"), }) } "needs-evidence" => { let missing_refs = list("missing_refs"); let revisit_deadline_ms = checked_ms("revisit_deadline_ms")?; ensure!( !missing_refs.is_empty() || revisit_deadline_ms.is_some(), "needs-evidence assessment requires the actual missing refs or a bounded revisit deadline" ); Ok(Self::NeedsEvidence { missing_refs, revisit_condition: text("revisit_condition"), revisit_deadline_ms, }) } "deferred" => { let reason = text("reason") .ok_or_else(|| anyhow!("deferred assessment requires a non-empty reason"))?; Ok(Self::Deferred { reason, trigger_ms: checked_ms("trigger_ms")?, }) } "candidate-failed" => { let error = text("error").or_else(|| text("reason")).ok_or_else(|| { anyhow!("candidate-failed assessment requires a non-empty error") })?; Ok(Self::CandidateFailed { error, candidate_revision: text("candidate_revision"), }) } "activated" => { let candidate_revision = text("candidate_revision").ok_or_else(|| { anyhow!("activated assessment requires the candidate revision it activated") })?; Ok(Self::Activated { candidate_revision, check_receipts: list("check_receipts"), }) } other => bail!( "unknown maintenance assessment outcome {other:?}; expected one of no-change, needs-evidence, deferred, candidate-failed, activated" ), } } } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum HostReply { Ok { value: Value }, Error { code: String, message: String }, } impl HostReply { pub fn error(code: &str, message: impl Into) -> Self { Self::Error { code: code.into(), message: message.into(), } } } #[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)] pub struct CellTrace { #[serde(default)] pub activities: Vec, #[serde(default)] pub changes: Vec, #[serde(default)] pub truncated: bool, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct WorkspaceActivity { pub kind: String, pub target: String, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct FileChange { pub path: String, pub diff: String, pub added: usize, pub removed: usize, #[serde(default)] pub truncated: bool, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum Outcome { Cell { status: CellStatus, output: Vec, truncated_bytes: u64, error: Option, #[serde(default, skip_serializing_if = "Option::is_none")] trace: Option, }, Hook { value: Value, }, Failed { error: String, }, /// A completed opaque remote dispatcher value. It is bounded separately /// from cell output because it is returned as a control-frame value. Remote { value: Value, }, /// A typed remote dispatcher refusal or execution error. RemoteError { code: String, message: String, }, /// The host observed an uncertain transport/process boundary. This is /// distinct from a clean process failure: the operation may have run, so /// callers must not replay it automatically. Unknown { reason: String, }, } impl Outcome { pub fn validate(&self, role: Role) -> Result<()> { match self { Self::Cell { output, error, .. } => { ensure!(role == Role::Workbench, "cell result from a hook worker"); ensure!(output.len() <= 256, "too many output records"); ensure!( output.iter().map(|o| o.text.len()).sum::() <= MAX_OUTPUT, "output limit exceeded" ); ensure!( error.as_ref().is_none_or(|e| e.len() <= 65_536), "oversized error" ); } Self::Hook { .. } => ensure!(role == Role::Hooks, "hook result from workbench"), Self::Remote { value } => { ensure!(role == Role::Remote, "remote result from non-remote worker"); ensure!( serde_json::to_vec(value)?.len() <= MAX_REMOTE_VALUE, "oversized remote value" ); } Self::RemoteError { code, message } => { ensure!(role == Role::Remote, "remote error from non-remote worker"); ensure!( !code.is_empty() && code.len() <= 128, "invalid remote error code" ); ensure!(message.len() <= MAX_REMOTE_ERROR, "oversized remote error"); } Self::Failed { error } => ensure!(error.len() <= 65_536, "oversized failure"), Self::Unknown { reason } => ensure!( !reason.trim().is_empty() && reason.len() <= 65_536, "unknown outcome requires a bounded reason" ), } Ok(()) } } #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum CellStatus { Ok, Error, Yielded, Interrupted, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct Output { pub stream: OutputStream, pub text: String, } #[derive(Clone, Copy, Debug, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum OutputStream { Stdout, Stderr, Result, } pub async fn write_frame(writer: &mut W, value: &Envelope) -> Result<()> { let bytes = serde_json::to_vec(value)?; ensure!( !bytes.is_empty() && bytes.len() <= MAX_FRAME, "outbound frame exceeds limit" ); writer.write_u32(bytes.len() as u32).await?; writer.write_all(&bytes).await?; writer.flush().await?; Ok(()) } /// A timeout/drop during this operation invalidates the connection. The worker /// actor destroys it rather than reading a new header from a partial payload. pub async fn read_frame(reader: &mut R) -> Result { let length = reader.read_u32().await? as usize; ensure!( (1..=MAX_FRAME).contains(&length), "invalid inbound frame length" ); let mut bytes = vec![0; length]; reader.read_exact(&mut bytes).await?; let envelope: Envelope = serde_json::from_slice(&bytes)?; if envelope.version != VERSION { bail!("unsupported protocol version"); } Ok(envelope) }