ive harnessed the harness
Something went wrong. Try again.
50 kB · 1477 lines
Rust
at main
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478//! 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<String> for $name { type Error = String; fn try_from(value: String) -> std::result::Result<Self, Self::Error> { 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<Self, Self::Err> { 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<TargetStamp>,}
/// 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<crate::target::EnvironmentFingerprint>, #[serde(default)] pub capabilities: Vec<Capability>,}
/// 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<String>, }, /// 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<crate::target::TargetIdentity>,}
/// 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<HelloBinding>, 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<String>, slots: Vec<HookSlot>, /// 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<WorkerIdentity>, /// 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<String>, }, 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<f64>, 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<String>, #[serde(default, skip_serializing_if = "Option::is_none")] work_item_id: Option<String>, }, #[serde(rename = "maintenance.originate")] MaintenanceOriginate { #[serde(default)] item_id: Option<String>, purpose: String, reason: String, #[serde(default)] evidence_refs: Vec<String>, #[serde(default)] impact: Option<Impact>, }, #[serde(rename = "maintenance.assess")] MaintenanceAssess { work_item_id: String, outcome: String, #[serde(default)] reason: Option<String>, #[serde(default)] missing_refs: Vec<String>, #[serde(default)] revisit_condition: Option<String>, /// Bounded revisit deadline for needs-evidence (wall-clock ms). #[serde(default)] revisit_deadline_ms: Option<i64>, /// Bounded revisit trigger for deferred (wall-clock ms). #[serde(default)] trigger_ms: Option<i64>, #[serde(default)] candidate_revision: Option<String>, #[serde(default)] check_receipts: Vec<String>, }, #[serde(rename = "evidence.resolve")] EvidenceResolve { refs: Vec<String> }, #[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<String>, cwd: PathBuf, #[serde(default)] env: std::collections::BTreeMap<String, String>, #[serde(default, skip_serializing_if = "Option::is_none")] stdin_base64: Option<String>, 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<String>, }, /// 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<String>, #[serde(default)] mode: Option<String>, #[serde(default)] profile: Option<String>, #[serde(default)] workspace: Option<String>, #[serde(default)] kernel_target: Option<serde_json::Value>, #[serde(default)] budget: Option<serde_json::Value>, #[serde(default)] branch_parent_revision: Option<String>, #[serde(default)] evidence: Vec<String>, }, #[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<String>, }, #[serde(rename = "peer.result")] PeerResult { request_id: String, #[serde(default)] peer_id: Option<String>, #[serde(default)] wait_ms: Option<u64>, }, #[serde(rename = "peer.submit_result")] PeerSubmitResult { request_id: String, #[serde(default)] peer_id: Option<String>, result: serde_json::Value, }, #[serde(rename = "peer.cancel")] PeerCancel { request_id: String, #[serde(default)] peer_id: Option<String>, }, #[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<bool>, },}
#[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<Self, Self::Error> { 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<String>, pub target_selector: TargetSelector, pub carrier_author: CarrierAuthor, #[serde(default)] pub source_refs: Vec<String>, pub carrier_placement: Option<CarrierPlacement>,}
#[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<String>, pub adoption_actor: AdoptionActor, pub operation_id: String, pub apply_at: Option<ApplicationBoundary>,}
/// 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<String>, pub adoption_actor: String, pub deadline_ms: Option<i64>, pub rationale: Option<String>, pub source_refs: Vec<String>,}
pub type HostFuture<'a> = std::pin::Pin< Box<dyn std::future::Future<Output = anyhow::Result<serde_json::Value>> + 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<String>, pub mode: Option<String>, pub profile: Option<String>, pub workspace: Option<String>, pub kernel_target: Option<serde_json::Value>, pub budget: Option<serde_json::Value>, pub branch_parent_revision: Option<String>, pub evidence: Vec<String>,}
#[derive(Clone, Debug, Serialize, Deserialize)]pub struct PeerSendReq { pub peer_id: String, pub request_id: String, pub content: String, pub evidence: Vec<String>,}
#[derive(Clone, Debug, Serialize, Deserialize)]pub struct PeerResultReq { pub request_id: String, pub peer_id: Option<String>, pub wait_ms: Option<u64>,}
#[derive(Clone, Debug, Serialize, Deserialize)]pub struct PeerSubmitResultReq { pub request_id: String, pub peer_id: Option<String>, pub result: serde_json::Value,}
#[derive(Clone, Debug, Serialize, Deserialize)]pub struct PeerCancelReq { pub request_id: String, pub peer_id: Option<String>,}
#[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<bool>,}
/// 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<String>, }, /// Work cannot proceed until the named dependencies exist or the /// bounded revisit deadline passes. NeedsEvidence { #[serde(default)] missing_refs: Vec<String>, #[serde(default)] revisit_condition: Option<String>, #[serde(default)] revisit_deadline_ms: Option<i64>, }, /// Objective retained but paused until a bounded trigger. Deferred { reason: String, #[serde(default)] trigger_ms: Option<i64>, }, /// Candidate failed validation or execution checks. CandidateFailed { error: String, #[serde(default)] candidate_revision: Option<String>, }, /// 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<String>, },}
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<Self> { let text = |name: &str| -> Option<String> { fields .get(name) .and_then(|v| v.as_str()) .map(str::to_owned) .filter(|s| !s.trim().is_empty()) }; let list = |name: &str| -> Vec<String> { 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<i64> { fields.get(name).and_then(|v| v.as_i64()) }; let checked_ms = |name: &str| -> Result<Option<i64>> { 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<String>) -> 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<WorkspaceActivity>, #[serde(default)] pub changes: Vec<FileChange>, #[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<Output>, truncated_bytes: u64, error: Option<String>, #[serde(default, skip_serializing_if = "Option::is_none")] trace: Option<CellTrace>, }, 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::<usize>() <= 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<W: AsyncWrite + Unpin>(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<R: AsyncRead + Unpin>(reader: &mut R) -> Result<Envelope> { 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)}