//! A concrete vertical slice: durable cells, explicit local delivery through a //! separate policy worker, terminal waits, and hook-only revision replacement. //! This is not a model loop or a Discord adapter. use crate::{ behavior::Release, hooks::Policy, protocol::*, remote::{ callback::CallbackContext, exec::RemoteManager, lifecycle::{ BootstrapRepair, ConnectionScopedLifecycle, ConnectionState, LifecycleSnapshot, PersistenceMode, RepairReceipt, RepairRequest, }, }, store::{PreparedEffect, Store}, target::{KernelTarget, NamespaceGeneration, NamespaceId, RemoteTarget, SessionPlacement}, work::{ConsequenceMetadata, Impact, Provenance, RunCause, WorkItem}, worker::{ BehaviorLocation, BehaviorSource, HostCall, HostService, WorkCommand, Worker, WorkerBackend, WorkerConfig, }, }; use anyhow::{anyhow, bail, ensure, Context, Result}; use base64::{engine::general_purpose::STANDARD as BASE64, Engine}; use serde::Serialize; use serde_json::{json, Value}; use std::{ path::{Component, Path}, process::Stdio, sync::{ atomic::{AtomicBool, Ordering}, Arc, }, time::Duration, }; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, process::{Child, Command}, sync::{Mutex, RwLock, Semaphore}, }; #[cfg(unix)] use std::os::unix::process::CommandExt; #[derive(Debug)] struct RoutedHostError { code: String, message: String, } impl std::fmt::Display for RoutedHostError { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter.write_str(&self.message) } } impl std::error::Error for RoutedHostError {} struct ProcessGroup { child: Option, pid: i32, finished: bool, } impl ProcessGroup { fn new(child: Child) -> Result { let pid = child.id().context("workspace command has no process id")? as i32; Ok(Self { child: Some(child), pid, finished: false, }) } fn child(&mut self) -> &mut Child { self.child.as_mut().expect("live process group has a child") } async fn terminate(&mut self) -> Result { terminate_process_group(self.pid); let waited = { let child = self.child(); tokio::time::timeout(Duration::from_secs(2), child.wait()).await }; match waited { Ok(status) => status.context("waiting for terminated workspace command"), Err(_) => { kill_process_group(self.pid); self.child() .wait() .await .context("waiting for killed workspace command") } } } fn finish(&mut self) { self.finished = true; } } impl Drop for ProcessGroup { fn drop(&mut self) { if self.finished || self.child.is_none() { return; } terminate_process_group(self.pid); let pid = self.pid; let mut child = self.child.take().expect("checked above"); if let Ok(runtime) = tokio::runtime::Handle::try_current() { runtime.spawn(async move { if tokio::time::timeout(Duration::from_secs(2), child.wait()) .await .is_err() { kill_process_group(pid); let _ = child.wait().await; } }); } else { kill_process_group(pid); let _ = child.start_kill(); } } } #[cfg(unix)] fn signal_process_group(pid: i32, signal: rustix::process::Signal) { // The child is the process-group leader. ESRCH only means it won the exit race. if let Some(pid) = rustix::process::Pid::from_raw(pid) { let _ = rustix::process::kill_process_group(pid, signal); } } #[cfg(unix)] fn terminate_process_group(pid: i32) { signal_process_group(pid, rustix::process::Signal::TERM); } #[cfg(unix)] fn kill_process_group(pid: i32) { signal_process_group(pid, rustix::process::Signal::KILL); } #[cfg(not(unix))] fn terminate_process_group(_pid: i32) {} #[cfg(not(unix))] fn kill_process_group(_pid: i32) {} #[derive(Default)] struct CommandCapture { stdout: Vec, stderr: Vec, stored: usize, stdout_truncated: bool, stderr_truncated: bool, } async fn drain_command_output( mut reader: impl tokio::io::AsyncRead + Unpin, stdout: bool, limit: usize, capture: Arc>, ) { let mut chunk = [0; 65_536]; loop { let read = match reader.read(&mut chunk).await { Ok(0) | Err(_) => return, Ok(read) => read, }; let mut capture = capture.lock().await; let keep = (limit - capture.stored).min(read); let truncated = keep != read; if stdout { capture.stdout.extend_from_slice(&chunk[..keep]); capture.stdout_truncated |= truncated; } else { capture.stderr.extend_from_slice(&chunk[..keep]); capture.stderr_truncated |= truncated; } capture.stored += keep; } } fn workspace_cwd(root: &Path, cwd: &Path) -> Result { ensure!(!cwd.is_absolute(), "workspace cwd must be relative"); ensure!( !cwd.components() .any(|part| matches!(part, Component::ParentDir | Component::Prefix(_))), "workspace cwd escapes source root" ); let root = std::fs::canonicalize(root).context("cannot resolve workspace source root")?; let cwd = std::fs::canonicalize(root.join(cwd)).context("workspace cwd does not exist")?; ensure!(cwd.is_dir(), "workspace cwd is not a directory"); ensure!( cwd.starts_with(&root), "workspace cwd escapes source root through a symlink" ); Ok(cwd) } /// Session placement: which workbench executes `python(code)`, resolved once /// per workbench generation. A remote placement binds a namespace stamp that /// every envelope must echo; there is no fallback to local execution, home /// directories, or another target. #[derive(Clone)] pub enum Placement { Local, Remote { target: RemoteTarget, /// The host-owned backend for the target, supplied by the target manager. Python never /// launches SSH itself. backend: Arc, /// Namespace binding issued at launch. The worker must echo it in its hello and every /// envelope; anything else is fatal. binding: TargetStamp, /// Target-side directory of the active behavior release. behavior_bundle: String, }, } impl Placement { #[must_use] pub fn local() -> Self { Self::Local } #[must_use] pub fn remote( target: RemoteTarget, backend: Arc, namespace: NamespaceId, behavior_bundle: String, ) -> Self { Self::Remote { binding: TargetStamp { target_id: target.id.clone(), namespace_id: namespace, namespace_generation: NamespaceGeneration::fresh(), }, backend, behavior_bundle, target, } } /// The kind of target this placement names. #[must_use] pub fn kernel_target(&self) -> KernelTarget { match self { Self::Local => KernelTarget::Local, Self::Remote { target, .. } => KernelTarget::Remote(target.clone()), } } /// The durable placement descriptor persisted under this generation. fn descriptor(&self, generation: &Generation) -> Result { Ok(match self { Self::Local => SessionPlacement { kernel_target: KernelTarget::Local, namespace_id: NamespaceId::parse("main")?, namespace_generation: NamespaceGeneration::parse(generation.as_str())?, }, Self::Remote { target, binding, .. } => SessionPlacement { kernel_target: KernelTarget::Remote(target.clone()), namespace_id: binding.namespace_id.clone(), namespace_generation: binding.namespace_generation.clone(), }, }) } } /// The target/workspace/interpreter identity the agent is actually using, /// resolved ONCE per workbench generation - never flooded into prompts. #[derive(Clone, Debug, Serialize)] pub struct TargetIdentitySnapshot { pub kind: String, pub target: String, pub workspace: String, pub platform: String, pub arch: String, pub interpreter_path: String, pub interpreter_version: String, pub environment_fingerprint: Option, pub capability_revision: u16, pub generation: String, } impl TargetIdentitySnapshot { fn resolve( placement: &Placement, config: &WorkerConfig, generation: &Generation, identity: Option, ) -> Result { match (placement, identity) { (Placement::Remote { target: remote, .. }, Some(worker_identity)) => { let runtime = &worker_identity.runtime; let environment_fingerprint = runtime .environment_fingerprint .as_ref() .map(|f| f.as_str().to_string()); Ok(Self { kind: "remote".into(), target: remote.id.to_string(), workspace: remote.workspace.path.clone(), platform: runtime.platform.clone(), arch: runtime.arch.clone(), interpreter_path: runtime.interpreter_path.clone(), interpreter_version: runtime.interpreter_version.clone(), environment_fingerprint, capability_revision: worker_identity .target .as_ref() .and_then(|binding| binding.identity.as_ref()) .map(|identity| identity.capability_revision) .unwrap_or(0), generation: generation.as_str().to_string(), }) } (Placement::Remote { target, .. }, None) => bail!( "target {}: remote worker did not report its runtime identity", target.id ), (Placement::Local, _) => Ok(Self { kind: "local".into(), target: "local".into(), workspace: config.cwd.display().to_string(), platform: std::env::consts::OS.to_string(), arch: std::env::consts::ARCH.to_string(), interpreter_path: config.python.display().to_string(), interpreter_version: local_interpreter_version(config)?, environment_fingerprint: None, capability_revision: 0, generation: generation.as_str().to_string(), }), } } } /// One interpreter probe per workbench generation, not per prompt. fn local_interpreter_version(config: &WorkerConfig) -> Result { let output = std::process::Command::new(&config.python) .arg("-c") .arg("import sys;print(sys.version.split()[0])") .output() .context("local interpreter probe failed")?; ensure!(output.status.success(), "local interpreter probe failed"); let version = String::from_utf8_lossy(&output.stdout).trim().to_string(); ensure!(!version.is_empty(), "local interpreter reported no version"); Ok(version) } struct ActivePolicy { epoch: i64, policy: Arc, } struct PlacementState { placement: Placement, identity: TargetIdentitySnapshot, } pub struct SessionRuntime { pub id: SessionId, pub store: Store, config: WorkerConfig, workbench: Mutex, policy: RwLock, placement: RwLock, /// Remote namespaces are connection-scoped in M2. `None` is the local /// adapter; `Some` becomes terminally lost on transport disconnect. lifecycle: RwLock>, remote: RemoteManager, foreground: Arc, activation_gate: Arc, stopped: AtomicBool, storage_failed: AtomicBool, pub orientation_host: Arc>>>, pub peer_host: Arc>>>, } #[derive(Clone, Debug, Serialize)] pub struct ExecutionReport { pub id: OperationId, pub generation: Generation, pub hook_revision: String, pub outcome: Outcome, } /// Spawn the workbench under the placement's backend. Local is the unchanged /// adapter path; a remote launch goes through the host-owned backend and /// carries the issued namespace binding. Returns the worker and its resolved /// target identity snapshot. async fn spawn_workbench( placement: &Placement, config: &WorkerConfig, session: &SessionId, release: &Release, ) -> Result<(Worker, TargetIdentitySnapshot)> { let worker = match placement { Placement::Local => { Worker::spawn(config, session.clone(), Role::Workbench, Some(release)).await? } Placement::Remote { backend, binding, behavior_bundle, .. } => { let behavior = BehaviorSource { revision: release.id.clone(), location: BehaviorLocation::RemoteBundle(behavior_bundle.clone()), }; Worker::launch_with( backend.as_ref(), config.startup_timeout, session.clone(), Role::Workbench, Some(behavior), Some(binding.clone()), ) .await? } }; let identity = TargetIdentitySnapshot::resolve( placement, config, worker.generation(), worker.worker_identity(), )?; Ok((worker, identity)) } fn lifecycle_for_placement(placement: &Placement) -> Result> { Ok(match placement { Placement::Local => None, Placement::Remote { binding, .. } => Some(ConnectionScopedLifecycle::new( binding.namespace_id.clone(), binding.namespace_generation.clone(), )), }) } async fn run_workspace_command( root: &Path, argv: Vec, cwd: std::path::PathBuf, env: std::collections::BTreeMap, stdin_base64: Option, timeout_ms: u64, max_output_bytes: u64, ) -> Result { ensure!( !argv.is_empty() && argv.iter().all(|part| !part.is_empty()), "argv must contain non-empty strings" ); ensure!( argv.iter().all(|part| !part.contains('\0')), "argv contains NUL" ); ensure!( timeout_ms > 0 && timeout_ms <= 3_600_000, "invalid workspace command timeout" ); // Base64 expands by 4/3 and the reply carries metadata too; leave room for // its JSON envelope instead of accepting an output that cannot be framed. let output_limit = usize::try_from(max_output_bytes).context("workspace output limit is too large")?; ensure!( output_limit <= 750_000, "workspace output limit exceeds reply frame capacity" ); ensure!( env.iter().all(|(key, value)| !key.is_empty() && !key.contains('\0') && !key.contains('=') && !value.contains('\0')), "workspace environment contains an invalid key or value" ); let stdin = stdin_base64 .map(|encoded| { BASE64 .decode(encoded) .context("invalid base64 workspace stdin") }) .transpose()?; let cwd = workspace_cwd(root, &cwd)?; let mut command = Command::new(&argv[0]); command .args(&argv[1..]) .current_dir(&cwd) // This is an overlay on the daemon's already-sanitized environment; // it is not a reconstructed shell environment. .envs(&env) .stdin(if stdin.is_some() { Stdio::piped() } else { Stdio::null() }) .stdout(Stdio::piped()) .stderr(Stdio::piped()) .kill_on_drop(false); #[cfg(unix)] command.as_std_mut().process_group(0); let started = tokio::time::Instant::now(); let mut group = ProcessGroup::new(command.spawn().context("spawning workspace command")?)?; let capture = Arc::new(Mutex::new(CommandCapture::default())); let stdout = group .child() .stdout .take() .context("workspace command has no stdout")?; let stderr = group .child() .stderr .take() .context("workspace command has no stderr")?; let stdout_reader = tokio::spawn(drain_command_output( stdout, true, output_limit, capture.clone(), )); let stderr_reader = tokio::spawn(drain_command_output( stderr, false, output_limit, capture.clone(), )); let writer = group .child() .stdin .take() .zip(stdin) .map(|(mut pipe, input)| { tokio::spawn(async move { let _ = pipe.write_all(&input).await; let _ = pipe.shutdown().await; }) }); let (status, timed_out) = match tokio::time::timeout(Duration::from_millis(timeout_ms), group.child().wait()).await { Ok(status) => (status.context("waiting for workspace command")?, false), Err(_) => (group.terminate().await?, true), }; group.finish(); if let Some(writer) = writer { let _ = writer.await; } let _ = stdout_reader.await; let _ = stderr_reader.await; let capture = capture.lock().await; Ok(json!({ "argv": argv, "cwd": cwd, "exit_code": status.code(), "stdout_base64": BASE64.encode(&capture.stdout), "stderr_base64": BASE64.encode(&capture.stderr), "duration_ms": started.elapsed().as_millis() as u64, "completed": !timed_out, "timed_out": timed_out, "cancelled": false, "stdout_truncated": capture.stdout_truncated, "stderr_truncated": capture.stderr_truncated, })) } impl SessionRuntime { pub async fn start( store: Store, id: SessionId, config: WorkerConfig, initial: Release, ) -> Result> { Self::start_with_placement(store, id, config, initial, Placement::local()).await } /// Placement-aware start. A remote placement launches the workbench on the /// configured target's backend and binds its namespace stamp; refusal at /// any point is typed and precedes execution. There is no local fallback. pub async fn start_with_placement( store: Store, id: SessionId, config: WorkerConfig, initial: Release, placement: Placement, ) -> Result> { store.ensure_session(&id).await?; let (epoch, release) = match store.activation(&id).await? { Some(active) => { let release = Release::load(&active.path)?; ensure!( release.id == active.revision, "active release bytes have changed" ); (active.epoch, release) } None => (0, initial.clone()), }; let policy = Policy::start(&config, id.clone(), release).await?; let epoch = if epoch == 0 { store.activate(&id, 0, policy.release.clone()).await? } else { epoch }; let (workbench, identity) = spawn_workbench(&placement, &config, &id, &policy.release).await?; let lifecycle = lifecycle_for_placement(&placement)?; let remote = RemoteManager::new(); if let Placement::Remote { target, backend, .. } = &placement { remote.register_relay_target( target.clone(), backend.clone(), config.startup_timeout, )?; } store .set_session_placement(&id, &placement.descriptor(workbench.generation())?) .await .context("failed to persist session placement")?; if let Placement::Remote { binding, .. } = &placement { store .append( &id, "kernel.placement", json!({ "generation": workbench.generation().as_str(), "binding": binding, "reason": "workbench placed on configured target; namespace is fresh; heap not restored", }), ) .await?; } store .append( &id, "kernel.reset", json!({ "generation": workbench.generation(), "reason": "session start; heap not restored", // The instruction set this generation was spawned with, echoed by the // worker. Without it the record would only say what the roots hold now. "skills": workbench.skills_revision(), }), ) .await?; Ok(Arc::new(Self { id, store, config, workbench: Mutex::new(workbench), policy: RwLock::new(ActivePolicy { epoch, policy }), placement: RwLock::new(PlacementState { placement, identity, }), lifecycle: RwLock::new(lifecycle), remote, foreground: Arc::new(Semaphore::new(1)), activation_gate: Arc::new(Semaphore::new(1)), stopped: AtomicBool::new(false), storage_failed: AtomicBool::new(false), orientation_host: Arc::new(std::sync::RwLock::new(None)), peer_host: Arc::new(std::sync::RwLock::new(Some(Arc::new( crate::peer::LocalPeerHost, )))), })) } pub fn set_orientation_host(&self, host: Arc) { *self.orientation_host.write().unwrap() = Some(host); } pub fn set_peer_host(&self, host: Arc) { *self.peer_host.write().unwrap() = Some(host); } pub async fn hook_epoch(&self) -> i64 { self.policy.read().await.epoch } pub async fn hooks(&self) -> Value { self.policy.read().await.policy.list() } pub async fn policy(&self) -> Arc { self.policy.read().await.policy.clone() } /// Returns the advertised persistence mode for a remote workbench. /// `None` denotes the unchanged local adapter. pub async fn persistence_mode(&self) -> Option { self.lifecycle .read() .await .as_ref() .map(ConnectionScopedLifecycle::mode) } /// Returns the remote connection state. A lost state is terminal for this /// handle; only [`Self::reconnect`] may create a new namespace owner. pub async fn connection_state(&self) -> Option { self.lifecycle .read() .await .as_ref() .map(|lifecycle| lifecycle.snapshot().state) } /// The target, namespace, owner and state visible to the SDK. pub async fn lifecycle_snapshot(&self) -> Option { self.lifecycle .read() .await .as_ref() .map(ConnectionScopedLifecycle::snapshot) } pub fn workspace(&self) -> &std::path::Path { &self.config.cwd } /// Every id this session offers, read from disk on each call: a skill an agent /// just wrote is offered to the next load, not to the next restart. pub fn skill_catalog(&self) -> crate::skills::Catalog { self.config.skills.catalog() } /// Where those ids come from. Reported without the index, which has its own /// call: this is a location, not a listing. fn skill_roots_report(&self) -> Value { let catalog = self.skill_catalog(); json!({ "builtin_root": self.config.skills.builtin_root().display().to_string(), "managed_root": self.config.skills.managed_root().map(|path| path.display().to_string()), "ids": catalog.ids(), "refused_count": catalog.refused.len(), "revision": catalog.revision(), }) } /// Candidate import/handshake precedes activation. Existing cells retain the /// Arc they pinned; this operation does not replace their workbench. /// Once accepted, caller cancellation does not interrupt the DB/cache switch. pub async fn activate_candidate( self: &Arc, expected_epoch: i64, release: Release, ) -> Result { self.activate_candidate_with_work(expected_epoch, release, None) .await } pub async fn activate_candidate_with_work( self: &Arc, expected_epoch: i64, release: Release, work_item_id: Option, ) -> Result { self.activate_hooks_internal(expected_epoch, release, true, work_item_id) .await } pub async fn activate_hooks( self: &Arc, expected_epoch: i64, release: Release, ) -> Result { self.activate_hooks_internal(expected_epoch, release, false, None) .await } async fn activate_hooks_internal( self: &Arc, expected_epoch: i64, release: Release, validate_synthetic: bool, work_item_id: Option, ) -> Result { ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); // Verify release integrity from disk before anything else let rechecked = Release::load(&release.path).context("cannot reload release from disk")?; ensure!( rechecked.id == release.id, "release bytes have changed on disk: expected {}, got {}", release.id, rechecked.id ); let permit = self .activation_gate .clone() .try_acquire_owned() .context("another hook activation is in progress")?; let this = self.clone(); tokio::spawn(async move { let _permit = permit; let candidate = Policy::start(&this.config, this.id.clone(), release).await?; if validate_synthetic { if let Err(err) = candidate.validate_synthetic().await { candidate.worker.stop().await; return Err(err); } } let mut active = this.policy.write().await; ensure!( !this.stopped.load(Ordering::Acquire), "session has shut down" ); if active.epoch != expected_epoch { candidate.worker.stop().await; bail!( "activation conflict: expected epoch {expected_epoch}, current epoch {}", active.epoch ); } let epoch = match this .store .activate_with_work( &this.id, expected_epoch, candidate.release.clone(), work_item_id, ) .await { Ok(ep) => ep, Err(err) => { candidate.worker.stop().await; return Err(err); } }; *active = ActivePolicy { epoch, policy: candidate, }; Ok::<_, anyhow::Error>(epoch) }) .await .context("activation task panicked")? } /// The host task owns accepted execution, not the caller's connection. /// Dropping this future does not replay or abandon the cell. Interrupt through /// the independent stop path; every cell also has a host-owned deadline. pub async fn execute( self: &Arc, code: String, budget: Duration, ) -> Result { self.execute_with_purpose(code, budget, Purpose::Interactive) .await } pub async fn execute_with_purpose( self: &Arc, code: String, budget: Duration, purpose: Purpose, ) -> Result { ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); ensure!(code.len() <= MAX_CODE, "cell is too large"); ensure!( budget > Duration::ZERO && budget <= Duration::from_secs(3600), "invalid cell deadline" ); ensure!( !self.storage_failed.load(Ordering::Acquire), "session paused after storage failure" ); let permit = self .foreground .clone() .try_acquire_owned() .context("session already has a foreground cell")?; let this = self.clone(); tokio::spawn(async move { let _permit = permit; this.execute_owned(code, budget, purpose).await }) .await .context("execution task panicked")? } async fn execute_owned( self: Arc, code: String, budget: Duration, purpose: Purpose, ) -> Result { let policy = self.policy.read().await.policy.clone(); let placement = self.placement.read().await.placement.clone(); let remote_lifecycle = self.lifecycle.read().await.clone(); if let Some(lifecycle) = &remote_lifecycle { lifecycle.ensure_connected()?; } let worker = { let mut worker = self.workbench.lock().await; ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); if !worker.is_alive() { if remote_lifecycle.is_some() { // A remote generation is not a reconnectable heap. Refuse // this cell until the caller explicitly opens a new owner. bail!("remote connection is lost; reconnect explicitly before executing"); } // The local adapter may retain its historical restart behavior. // This branch is never used for a remote placement. let (spawned, identity) = spawn_workbench(&placement, &self.config, &self.id, &policy.release).await?; *worker = spawned; *self.placement.write().await = PlacementState { placement: placement.clone(), identity, }; self.store .append( &self.id, "kernel.reset", json!({ "generation": worker.generation(), "reason": "previous local worker ended; code was not replayed", "skills": worker.skills_revision(), }), ) .await?; } worker.clone() }; let id = OperationId::fresh(); let generation = worker.generation().clone(); let revision = policy.release.id.clone(); self.store .begin(&self.id, &id, &generation, revision.clone(), code.clone()) .await?; let service = self.service( generation.clone(), policy, placement.clone(), purpose, budget, ); let outcome = match worker .call( WorkCommand::Execute { id: id.clone(), code, purpose, }, budget, service, ) .await { Ok(outcome) => outcome, Err(error) if remote_lifecycle .as_ref() .is_some_and(|lifecycle| !lifecycle.is_connected()) || (remote_lifecycle.is_some() && !worker.is_alive()) => { let reason = format!("remote transport lost before completion: {error:#}"); if let Some(lifecycle) = &remote_lifecycle { lifecycle.mark_lost(reason.clone())?; } // The operation may already have mutated the remote process. // Persist UNKNOWN and return without issuing a second call. if let Err(storage_error) = self.store.finish_unknown(&id, reason.clone()).await { self.storage_failed.store(true, Ordering::Release); return Err(storage_error .context("cannot persist unknown remote completion; session paused")); } Outcome::Unknown { reason } } Err(error) => Outcome::Failed { error: format!("{error:#}"), }, }; if let Err(error) = self.store.finish(&id, outcome.clone()).await { self.storage_failed.store(true, Ordering::Release); worker.stop().await; return Err(error.context("cannot persist execution completion; session paused")); } Ok(ExecutionReport { id, generation, hook_revision: revision, outcome, }) } fn service( self: &Arc, generation: Generation, policy: Arc, placement: Placement, purpose: Purpose, budget: Duration, ) -> HostService { let this = self.clone(); let deadline = tokio::time::Instant::now() + budget; let behavior_path = match &placement { Placement::Local => Some(policy.release.path.to_string_lossy().into_owned()), Placement::Remote { .. } => None, }; Arc::new(move |call: HostCall| { let behavior_path = behavior_path.clone(); let (this, generation, policy, target, binding) = ( this.clone(), generation.clone(), policy.clone(), placement.kernel_target(), match &placement { Placement::Remote { binding, .. } => Some(binding.clone()), Placement::Local => None, }, ); Box::pin(async move { let result = async { match call.request { HostRequest::LocalSend { text } => { let receipt = match this .store .prepare_send( &this.id, &call.execution_id, &generation, call.id, text.clone(), ) .await? { PreparedEffect::Existing(receipt) => receipt, PreparedEffect::New(effect) => { // No database transaction or session lock is // held while awaiting the separate callback lane. let context = match (&target, &binding) { (KernelTarget::Local, None) => CallbackContext::local( this.id.clone(), call.execution_id.clone(), generation.clone(), policy.release.id.clone(), ), (KernelTarget::Remote(_), Some(stamp)) => { CallbackContext::remote( this.id.clone(), call.execution_id.clone(), stamp.clone(), generation.clone(), policy.release.id.clone(), ) } (KernelTarget::Local, Some(_)) => { return Err(anyhow!( "local callback cannot carry a remote binding" )); } (KernelTarget::Remote(remote), None) => { return Err(anyhow!( "remote target {} has no callback binding", remote.id )); } }; let review = policy.review_with_context(context, &effect, &text).await; this.store .publish_send( &this.id, &call.execution_id, &generation, effect, policy.release.id.clone(), review, ) .await? } }; Ok(serde_json::to_value(receipt)?) } HostRequest::Wait { seconds, reason } => { this.store .wait(&this.id, &call.execution_id, &generation, seconds, reason) .await } HostRequest::HooksList => Ok(policy.list()), HostRequest::HooksDescribe { slot } => Ok(policy.describe(slot)), HostRequest::BehaviorActivate { path, expected_epoch, revision, work_item_id, } => { let release = Release::load(&path).context("failed to load release from disk")?; if let Some(expected_rev) = &revision { ensure!( &release.id == expected_rev, "release revision mismatch: expected {}, got {}", expected_rev, release.id ); } let new_epoch = this .activate_candidate_with_work(expected_epoch, release.clone(), work_item_id.clone()) .await?; Ok(json!({ "epoch": new_epoch, "revision": release.id, "work_item_id": work_item_id, })) } HostRequest::MaintenanceAssess { work_item_id, outcome, reason, missing_refs, revisit_condition, revisit_deadline_ms, trigger_ms, candidate_revision, check_receipts, } => { // Closed outcome parse at the host boundary: an // unknown tag, a missing variant field or a // contradictory reference is a typed error to the // caller, never a silently accepted payload. let assessment_fields = json!({ "reason": reason, "missing_refs": missing_refs, "revisit_condition": revisit_condition, "revisit_deadline_ms": revisit_deadline_ms, "trigger_ms": trigger_ms, "candidate_revision": candidate_revision, "check_receipts": check_receipts, }); let assessment = MaintenanceAssessment::parse(&outcome, &assessment_fields)?; let request_op_str = format!("{}_{}", call.execution_id, call.id); let request_op = if request_op_str.len() <= 128 && request_op_str .bytes() .all(|b| b.is_ascii_alphanumeric() || b"_.-".contains(&b)) { OperationId::try_from(request_op_str).map_err(|e| anyhow!(e))? } else { use sha2::{Digest, Sha256}; let mut hasher = Sha256::new(); hasher.update(request_op_str.as_bytes()); let hash = format!("{:x}", hasher.finalize()); OperationId::try_from(format!("req_{}", hash)) .map_err(|e| anyhow!(e))? }; let seq = this .store .commit_maintenance_assessment( &this.id, &request_op, &work_item_id, &outcome, assessment_fields, ) .await?; Ok(json!({ "status": "recorded", "work_item_id": work_item_id, "outcome": assessment.outcome_tag(), "seq": seq, })) } HostRequest::MaintenanceOriginate { item_id, purpose, reason, evidence_refs, impact, } => { let now = crate::store::now_ms()?; let work_id = item_id .unwrap_or_else(|| format!("work-{}", call.id)); let work_item = WorkItem::new( &work_id, &purpose, RunCause::MaintenanceReady { item_id: work_id.clone(), purpose: purpose.clone(), }, ConsequenceMetadata::new(impact.unwrap_or(Impact::High)), Provenance::agent("runtime.origination", now), now, ) .with_session(this.id.to_string()) .with_evidence_refs(evidence_refs.clone()); let origination_payload = json!({ "type": "maintenance.originated", "id": work_id, "purpose": purpose, "reason": reason, "evidence_refs": evidence_refs, "impact": impact, }); let request_op_str = format!("{}_{}", call.execution_id, call.id); let request_op = if request_op_str.len() <= 128 && request_op_str .bytes() .all(|b| b.is_ascii_alphanumeric() || b"_.-".contains(&b)) { OperationId::try_from(request_op_str).map_err(|e| anyhow!(e))? } else { use sha2::{Digest, Sha256}; let mut hasher = Sha256::new(); hasher.update(request_op_str.as_bytes()); let hash = format!("{:x}", hasher.finalize()); OperationId::try_from(format!("req_{}", hash)) .map_err(|e| anyhow!(e))? }; let (runnable_id, seq) = this .store .commit_origination_with_work( &this.id, &request_op, work_item, origination_payload, ) .await?; Ok(json!({ "item_id": runnable_id.as_str(), "work_item_id": runnable_id.as_str(), "seq": seq, })) } HostRequest::EvidenceResolve { refs } => { let mut resolved = Vec::new(); let mut unavailable = Vec::new(); for r in refs { if let Some(rest) = r.strip_prefix("event:") { if let Ok(seq) = rest.parse::() { match this.store.get_event(&this.id, seq).await { Ok(Some(event)) => { let operation = event .payload .get("operation") .and_then(|v| v.as_str()) .map(ToString::to_string); let signature = event .payload .get("signature") .and_then(|v| v.as_str()) .map(ToString::to_string); let error = event .payload .get("error") .and_then(|v| v.as_str()) .or_else(|| { event .payload .get("outcome") .and_then(|o| o.get("error")) .and_then(|v| v.as_str()) }) .map(ToString::to_string); let active_epoch = this.hook_epoch().await; resolved.push(json!({ "ref": r, "kind": "event", "operation": operation, "signature": signature, "error": error, "event": { "id": event.id, "kind": event.kind, "payload": event.payload, "created_ms": event.created_ms, }, "current_lookup_context": { "workspace_path": this.workspace().display().to_string(), "session_id": this.id.to_string(), "active_release_id": policy.release.id.clone(), "active_epoch": active_epoch, }, "workspace": { "path": this.workspace().display().to_string(), "session_id": this.id.to_string(), }, "revision": { "id": policy.release.id.clone(), "epoch": active_epoch, } })); } Ok(None) => { unavailable.push(json!({"ref": r, "reason": "not_found"})); } Err(e) => { unavailable.push(json!({"ref": r, "reason": format!("{e}")})); } } } else { unavailable.push(json!({"ref": r, "reason": "invalid_seq"})); } } else if let Some(case_id) = r.strip_prefix("case:") { match this.store.get_observer_case(case_id).await { Ok(Some(case)) => { resolved.push(json!({ "ref": r, "kind": "case", "operation": case.operation, "signature": case.signature, "component_revision": case.component_revision, "workspace": { "path": case.workspace, "session_id": case.session_id, }, "count": case.count, "first_seen_ms": case.first_seen_ms, "last_seen_ms": case.last_seen_ms, "evidence_refs": case.evidence_refs, "work_item_id": case.work_item_id, })); } Ok(None) => { unavailable.push(json!({"ref": r, "reason": "not_found"})); } Err(e) => { unavailable.push(json!({"ref": r, "reason": format!("{e}")})); } } } else if let Some(work_id) = r.strip_prefix("work:") { match this.store.get_work_item(work_id).await { Ok(Some(item)) => { resolved.push(json!({ "ref": r, "kind": "work_item", "id": item.id, "session_id": item.session_id, "purpose": item.purpose, "cause": format!("{:?}", item.cause), "evidence_refs": item.provenance.source_refs, "status": format!("{:?}", item.status), "created_ms": item.created_ms, })); } Ok(None) => { unavailable.push(json!({"ref": r, "reason": "not_found"})); } Err(e) => { unavailable.push(json!({"ref": r, "reason": format!("{e}")})); } } } else if let Some(rel_path) = r.strip_prefix("source:") { let path_obj = std::path::Path::new(rel_path); if path_obj.is_absolute() || path_obj .components() .any(|c| matches!(c, std::path::Component::ParentDir | std::path::Component::Prefix(_))) { unavailable.push(json!({ "ref": r, "reason": "path escapes declared workspace root" })); continue; } let ws_root = match std::fs::canonicalize(this.workspace()) { Ok(p) => p, Err(_) => this.workspace().to_path_buf(), }; let ws_candidate = this.workspace().join(rel_path); let rel_root = match std::fs::canonicalize(&policy.release.path) { Ok(p) => p, Err(_) => policy.release.path.clone(), }; let rel_candidate = policy.release.path.join(rel_path); let candidate = if ws_candidate.is_file() { if let Ok(canon) = std::fs::canonicalize(&ws_candidate) { if canon.starts_with(&ws_root) { Some(canon) } else { None } } else { None } } else if rel_candidate.is_file() { if let Ok(canon) = std::fs::canonicalize(&rel_candidate) { if canon.starts_with(&rel_root) { Some(canon) } else { None } } else { None } } else { None }; match candidate { Some(p) => match std::fs::read_to_string(&p) { Ok(content) => { resolved.push(json!({ "ref": r, "kind": "source", "path": rel_path, "content": content, "revision": policy.release.id.clone(), })); } Err(e) => { unavailable.push(json!({"ref": r, "reason": format!("read error: {e}")})); } }, None => { unavailable.push(json!({"ref": r, "reason": "file not found in workspace or release"})); } } } else { unavailable.push(json!({"ref": r, "reason": "unknown_ref_type"})); } } Ok(json!({ "resolved": resolved, "unavailable": unavailable, })) } HostRequest::RuntimeInstructions => { // The exact instruction text this transcript is running under, so a // reader that received a patch does not have to apply it by hand // and cannot silently apply it wrongly. A patch is a change; this // is the state the change lands in. let mut events = Vec::new(); let mut after = 0i64; loop { let page = this.store.history(&this.id, after, 1024).await?; if page.is_empty() { break; } after = page.last().map(|event| event.id).unwrap_or(after); let last = page.len(); events.extend(page); if last < 1024 { break; } } let known = crate::prompt_header::known_snapshot( events .iter() .map(|event| (event.id, event.kind.as_str(), &event.payload)), ); Ok(match known { Some(snapshot) => json!({ "resolved": true, "digest": snapshot.digest, "text": snapshot.text, "sources": snapshot .sources .iter() .map(|source| json!({ "name": source.name, "digest": source.digest, })) .collect::>(), }), None => json!({ "resolved": false, "reason": "this transcript has no pinned instructions yet", }), }) } HostRequest::SkillsIndex => Ok(this.skill_catalog().report()), HostRequest::SkillsLoad { id } => { let body = this.config.skills.load(&id)?; Ok(body.report()) } HostRequest::RuntimeInspect => { let identity = this.target_identity().await; // The header records are read here rather than passed in: the pin and // the chain are durable facts about this transcript, and a value the // caller had to remember to supply would be absent exactly when it // matters most. let mut header_events = Vec::new(); let mut after = 0i64; loop { let page = this.store.history(&this.id, after, 1024).await?; if page.is_empty() { break; } after = page.last().map(|event| event.id).unwrap_or(after); let last = page.len(); header_events.extend(page); if last < 1024 { break; } } let records = || { header_events.iter().map(|event| { (event.id, event.kind.as_str(), &event.payload) }) }; let head = crate::prompt_header::newest_pin(records()); let revisions_since_head = head .as_ref() .map(|(pin_seq, _)| { header_events .iter() .filter(|event| { event.kind == crate::prompt_header::HEADER_REVISION_EVENT && event.id > *pin_seq }) .count() }) .unwrap_or(0); let known = crate::prompt_header::known_snapshot(records()); let selected = this.store.activation(&this.id).await?; // The roots below are derived from THIS session's workspace // and therefore follow the target: a peer or remote session // reports its own tree, never the daemon's. Reported paths // are real paths even when a directory has not been created // yet, so a model does not have to walk the host to find out // where it is or guess from the process cwd. let workspace = this.workspace(); let draft_root = workspace.join(".klbr").join("drafts"); // Host-sourced time, so a claim about the world can be // placed against the moment it was made rather than // reconstructed from a transcript that carries no clock. let now_ms = crate::store::now_ms().unwrap_or(0); Ok(json!({ "session_id": this.id, "execution_id": call.execution_id, "purpose": purpose, "target": identity, "now_ms": now_ms, // Same shape as the stamps the transcript carries, so the // clock and an observation can be compared by eye. "now": chrono::DateTime::from_timestamp_millis(now_ms) .map(|stamp| { stamp .to_rfc3339_opts(chrono::SecondsFormat::Millis, true) }) .unwrap_or_else(|| "unknown".into()), "source_root": { "path": workspace.display().to_string(), "exists": workspace.is_dir(), "writable": false, }, "draft_root": { "path": draft_root.display().to_string(), "target": identity.target, "exists": draft_root.is_dir(), }, // The revision the worker is actually running under is // the applied one; the store's selection is the chosen // one. They are reported separately because a selection // that has not activated is not in force. "applied_behavior_revision": policy.release.id, "pinned_behavior_revision": policy.release.id, "selected_behavior_revision": selected.as_ref().map(|s| &s.revision), "selected_behavior_epoch": selected.as_ref().map(|s| s.epoch), "behavior_source": behavior_path.map(|path| json!({"target": identity.target, "path": path, "writable": false})), "skills": this.skill_roots_report(), // Which instructions this transcript is running under, from its // own records. The head is a pinned revision and the newest text // the transcript can reconstruct, so a reader can name the head it // ran under without hashing a context it cannot reproduce byte for // byte, and can see whether a change has been observed since. "instructions": { "head_digest": head.as_ref().map(|(_, pin)| pin.digest.clone()), "head_pinned_at_ms": head.as_ref().map(|(_, pin)| pin.pinned_at_ms), "head_reason": head.as_ref().map(|(_, pin)| pin.reason.as_str()), "revisions_since_head": revisions_since_head, "known_digest": known.as_ref().map(|snapshot| snapshot.digest.clone()), }, "cell_budget": {"limit_ms": budget.as_millis() as u64, "remaining_ms": deadline.saturating_duration_since(tokio::time::Instant::now()).as_millis() as u64}, "unavailable": { "agent_id": "session-to-agent association not configured", "attempt_id": "model attempt context not supplied to workbench", "remote_behavior_source": "remote bundle has no verified mapping to the pinned host policy revision", "cumulative_model_budget": "model budget accounting not supplied to workbench" } })) } HostRequest::WorkspaceRun { argv, cwd, env, stdin_base64, timeout_ms, max_output_bytes, } => { ensure!( matches!(target, KernelTarget::Local), "workspace.run has no host executor for a remote placement" ); run_workspace_command( this.workspace(), argv, cwd, env, stdin_base64, timeout_ms, max_output_bytes, ) .await } HostRequest::WorkspaceInfo => { let epoch = this.hook_epoch().await; let identity = this.target_identity().await; Ok(json!({ "path": this.workspace().display().to_string(), "session_id": this.id.to_string(), "active_revision": policy.release.id.clone(), "epoch": epoch, "target": serde_json::to_value(&identity)?, })) } HostRequest::OrientationStage(req) => { let host = this.orientation_host.read().unwrap().clone() .context("unconfigured_orientation_host")?; host.stage_orientation(&this.store, &this.id, req).await } HostRequest::OrientationCurrent(req) => { let host = this.orientation_host.read().unwrap().clone() .context("unconfigured_orientation_host")?; host.current_orientation(&this.store, &this.id, req).await } HostRequest::OrientationAdopt(req) => { let host = this.orientation_host.read().unwrap().clone() .context("unconfigured_orientation_host")?; host.adopt_orientation(&this.store, &this.id, req).await } HostRequest::FrontierSelect(req) => { let host = this.orientation_host.read().unwrap().clone() .context("unconfigured_orientation_host")?; host.select_frontier(&this.store, &this.id, req).await } HostRequest::PeerOpen { peer_id, mode, profile, workspace, kernel_target, budget, branch_parent_revision, evidence, } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerOpenReq { peer_id, mode, profile, workspace, kernel_target, budget, branch_parent_revision, evidence, }; host.open_peer(&this.store, &this.id, req).await } HostRequest::PeerSend { peer_id, request_id, content, evidence, } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerSendReq { peer_id, request_id, content, evidence, }; host.send_peer_message(&this.store, &this.id, req).await } HostRequest::PeerResult { request_id, peer_id, wait_ms, } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerResultReq { request_id, peer_id, wait_ms, }; host.get_peer_result(&this.store, &this.id, req).await } HostRequest::PeerSubmitResult { request_id, peer_id, result, } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerSubmitResultReq { request_id, peer_id, result, }; host.submit_peer_result(&this.store, &this.id, req).await } HostRequest::PeerCancel { request_id, peer_id, } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerCancelReq { request_id, peer_id, }; host.cancel_peer_request(&this.store, &this.id, req).await } HostRequest::PeerGet { peer_id } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerGetReq { peer_id }; host.get_peer(&this.store, &this.id, req).await } HostRequest::PeerList => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; host.list_peers(&this.store, &this.id).await } HostRequest::PeerStop { peer_id, drain } => { let host_opt = this.peer_host.read().unwrap().clone(); let Some(host) = host_opt else { bail!("unconfigured_peer_host: peer host not configured for session {}", this.id); }; let req = PeerStopReq { peer_id, drain }; host.stop_peer(&this.store, &this.id, req).await } request @ (HostRequest::RemoteEnsure { .. } | HostRequest::RemoteCall { .. } | HostRequest::RemoteCancel { .. } | HostRequest::RemoteClose { .. }) => { match this .remote .handle(&this.id, &call.execution_id, request) .await { HostReply::Ok { value } => Ok(value), HostReply::Error { code, message } => { Err(RoutedHostError { code, message }.into()) } } } } } .await; match result { Ok(value) => HostReply::Ok { value }, Err(error) => { if let Some(routed) = error.downcast_ref::() { return HostReply::error(&routed.code, routed.message.clone()); } let msg = format!("{error:#}"); let code = if msg.contains("activation conflict") { "cas_conflict" } else if msg.contains("validation") || msg.contains("bytes have changed") || msg.contains("revision mismatch") { "validation_error" } else { "host_operation_failed" }; HostReply::error(code, msg) } } }) }) } /// Management bypasses policy hooks. A broken hook cannot veto its repair. pub async fn interrupt(&self) { let worker = self.workbench.lock().await.clone(); worker.stop().await; if let Some(lifecycle) = self.lifecycle.read().await.clone() { let _ = lifecycle.mark_lost("connection interrupted by host"); } } /// Wait until any active foreground cell completes execution. pub async fn wait_for_idle(&self) { if let Ok(_permit) = self.foreground.acquire().await { // Permit successfully acquired and dropped: in-flight cell finished. } } /// The resolved target identity for the CURRENT workbench generation. /// Resolved once per generation; consumers must not re-derive per prompt. pub async fn target_identity(&self) -> TargetIdentitySnapshot { self.placement.read().await.identity.clone() } /// Change the session's placement. This is a VISIBLE NEW GENERATION with /// an explicit transition: the old workbench is stopped (an in-flight cell /// is interrupted honestly), the namespace is NOT migrated, and no global /// moves. Refused while a foreground cell is running. pub async fn relocate(self: &Arc, placement: Placement) -> Result { ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); let permit = self .foreground .clone() .try_acquire_owned() .context("foreground cell is running; placement change refused")?; let policy = self.policy.read().await.policy.clone(); // Reconnect and relocation invalidate every auxiliary namespace owned by // the old connection before a fresh endpoint is admitted. self.remote.shutdown().await?; let next_lifecycle = lifecycle_for_placement(&placement)?; let old_worker = { self.workbench.lock().await.clone() }; let old_generation = old_worker.generation().clone(); let old_identity = self.target_identity().await; old_worker.stop().await; if let Some(old_lifecycle) = self.lifecycle.read().await.clone() { old_lifecycle.mark_lost("placement changed; old connection closed")?; } let (worker, identity) = spawn_workbench(&placement, &self.config, &self.id, &policy.release).await?; if let Placement::Remote { target, backend, .. } = &placement { self.remote.register_relay_target( target.clone(), backend.clone(), self.config.startup_timeout, )?; } let generation = worker.generation().clone(); self.store .append( &self.id, "kernel.placement", json!({ "generation": generation.as_str(), "previous_generation": old_generation.as_str(), "from": {"kind": old_identity.kind, "target": old_identity.target, "generation": old_identity.generation}, "to": {"kind": identity.kind, "target": identity.target, "generation": identity.generation}, "reason": "placement changed; namespace did not migrate; globals did not move", }), ) .await?; self.store .set_session_placement(&self.id, &placement.descriptor(&generation)?) .await .context("failed to persist placement change")?; *self.workbench.lock().await = worker; *self.placement.write().await = PlacementState { placement, identity, }; *self.lifecycle.write().await = next_lifecycle; drop(permit); Ok(generation) } /// Explicitly opens a fresh connection-scoped remote namespace. The old /// owner remains lost, and no operation from it is replayed. pub async fn reconnect(self: &Arc) -> Result { ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); let old_lifecycle = self .lifecycle .read() .await .clone() .context("local sessions do not have a remote connection")?; let next = old_lifecycle.reconnect()?; let current = self.placement.read().await.placement.clone(); let (target, backend, behavior_bundle) = match current { Placement::Remote { target, backend, behavior_bundle, .. } => (target, backend, behavior_bundle), Placement::Local => bail!("remote connection-scoped placement is unavailable"), }; let placement = Placement::remote( target, backend, next.namespace_id().clone(), behavior_bundle, ); self.relocate(placement).await } /// Runs the explicitly diagnosed production bootstrap repair and then /// returns ordinary execution to a fresh connection-scoped kernel. pub async fn repair_bootstrap( self: &Arc, repair: BootstrapRepair, request: RepairRequest, ) -> Result { ensure!( !self.stopped.load(Ordering::Acquire), "session has shut down" ); let lifecycle = self .lifecycle .read() .await .clone() .context("local sessions do not have a remote bootstrap repair lane")?; lifecycle.mark_lost("explicit bootstrap repair closed the old owner")?; let receipt = repair.repair(&self.store, &self.id, request).await?; // The repair action is complete; ordinary work resumes only through a // newly admitted namespace, never through the repair transport. self.reconnect().await?; Ok(receipt) } pub async fn restart_workbench(&self) -> Result<()> { let worker = self.workbench.lock().await.clone(); if worker.is_alive() { worker.stop().await; } if let Some(lifecycle) = self.lifecycle.read().await.clone() { lifecycle.mark_lost("connection closed by host")?; } Ok(()) } pub async fn shutdown(&self) { self.stopped.store(true, Ordering::Release); let _ = self.remote.shutdown().await; self.interrupt().await; let policy = self.policy.read().await.policy.clone(); policy.worker.stop().await; } }