use crate::catalog::{ ConfigBundleRef, CreateWorkerRequest, ProfileSelector, ProfileSourceArchiveSource, RepositoryRefObservation, RepositoryRefObservationRequest, WorkerDetail, WorkerLifecycleAck, WorkerRestoreIntent, WorkerStatus, WorkerSummary, WorkingDirectoryRepositoryAccessRequest, WorkingDirectoryRequest, WorkingDirectoryStatus as CatalogWorkingDirectoryStatus, WorkspaceApiRef, }; use crate::config_bundle::{ ConfigBundle, ConfigBundleAvailability, ConfigBundleSummary, validate_config_bundle, validate_config_bundle_ref, }; use crate::diagnostics::{DiagnosticSeverity, RuntimeDiagnostic}; use crate::error::RuntimeError; use crate::execution::WorkerExecutionRestoreRequest; use crate::execution::{ WorkerExecutionBackend, WorkerExecutionBackendRef, WorkerExecutionHandle, WorkerExecutionOperation, WorkerExecutionResult, WorkerExecutionSpawnRequest, WorkerExecutionSpawnResult, WorkspaceConfigFetchRequest, WorkspaceConfigFetchResult, }; #[cfg(feature = "fs-store")] use crate::fs_store::{ FsRuntimeStore, FsRuntimeStoreOptions, PersistedRuntimeState, PersistedWorkerExecution, PersistedWorkerExecutionBinding, PersistedWorkerRecord, }; use crate::identity::{WorkerId, WorkerRef}; use crate::interaction::{WorkerInput, WorkerInputKind, WorkerInteractionAck}; use crate::management::{ RuntimeBackendKind, RuntimeOptions, RuntimeStatus, RuntimeSummary, WorkerDeleteResult, }; #[cfg(feature = "ws-server")] use crate::observation::{WorkerObservationCursor, WorkerObservationEvent}; use crate::resource::{ BackendResourceClient, BackendResourceError, BackendResourceFetchRequest, BackendResourceKind, REPOSITORY_SSH_ACCESS_CONTENT_TYPE, RepositorySshAccessSecret, }; #[cfg(feature = "fs-store")] use crate::retention::{ FsWorkerRetentionProvider, WorkerRetentionExecutionRequest, WorkerRetentionExecutionResult, WorkerRetentionInventory, WorkerRetentionInventorySnapshot, WorkerRetentionProvider, }; use protocol::subscription::{ EventSubscriptionSelector, SubscriptionEventPayload, SubscriptionSnapshot, SubscriptionValidationError, SubscriptionWorkdirId, SubscriptionWorker, SubscriptionWorkerId, SubscriptionWorkerState, }; use protocol::{Event, Method}; use std::collections::BTreeMap; #[cfg(feature = "ws-server")] use std::collections::VecDeque; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, Weak}; #[cfg(feature = "ws-server")] use tokio::sync::broadcast; use tokio::sync::mpsc; use uuid::Uuid; /// Workspace-scoped Runtime authorization context supplied by a trusted backend. #[derive(Clone, Debug, Eq, PartialEq)] pub struct RuntimeWorkspaceScope { pub workspace_id: String, pub server_id: String, } impl RuntimeWorkspaceScope { pub fn new(workspace_id: impl Into, server_id: impl Into) -> Self { Self { workspace_id: workspace_id.into(), server_id: server_id.into(), } } } const SUBSCRIPTION_QUEUE_CAPACITY: usize = 256; #[derive(Clone, Debug)] pub struct RuntimeSubscriptionUpdate { pub subject_revision: u64, pub payload: SubscriptionEventPayload, } #[derive(Debug, thiserror::Error)] pub enum RuntimeSubscriptionRecvError { #[error("Runtime event subscription lagged and requires a fresh snapshot")] Lagged, #[error("Runtime event subscription closed")] Closed, } /// A gap-free snapshot/live subscription owned by one Runtime connection. /// Dropping the subscription removes its bounded producer queue. pub struct RuntimeEventSelectorSubscription { subscription_id: u64, selector: EventSubscriptionSelector, snapshot_revision: u64, snapshot: SubscriptionSnapshot, receiver: mpsc::Receiver, lagged: Arc, runtime: Weak>, } impl RuntimeEventSelectorSubscription { pub fn selector(&self) -> &EventSubscriptionSelector { &self.selector } pub fn snapshot_revision(&self) -> u64 { self.snapshot_revision } pub fn snapshot(&self) -> &SubscriptionSnapshot { &self.snapshot } pub async fn recv( &mut self, ) -> Result { if self.lagged.load(Ordering::Acquire) { self.receiver.close(); return Err(RuntimeSubscriptionRecvError::Lagged); } match self.receiver.recv().await { Some(_) if self.lagged.load(Ordering::Acquire) => { self.receiver.close(); Err(RuntimeSubscriptionRecvError::Lagged) } Some(update) => Ok(update), None if self.lagged.load(Ordering::Acquire) => { Err(RuntimeSubscriptionRecvError::Lagged) } None => Err(RuntimeSubscriptionRecvError::Closed), } } } impl Drop for RuntimeEventSelectorSubscription { fn drop(&mut self) { let Some(runtime) = self.runtime.upgrade() else { return; }; if let Ok(mut state) = runtime.lock() { state.subscriptions.remove(&self.subscription_id); } } } /// Concrete embedded Runtime domain entity. /// /// The default implementation is memory-backed and tools/provider-less by /// design. An optional `fs-store` feature adds filesystem persistence while /// preserving the same typed authority boundary. It can later be adapted by /// backend registries or web servers without making sockets, sessions, or paths /// public authority. #[derive(Clone, Debug)] pub struct Runtime { inner: Arc>, worker_operations: Arc>>>>, } impl Runtime { /// Create a memory-backed Runtime with generated identity. pub fn new_memory() -> Self { Self::with_options(RuntimeOptions::default()) } /// Create a memory-backed Runtime with explicit options. pub fn with_options(options: RuntimeOptions) -> Self { let state = RuntimeState::new(options.display_name); Self { inner: Arc::new(Mutex::new(state)), worker_operations: Arc::new(Mutex::new(BTreeMap::new())), } } /// Create a memory-backed Runtime with an attached execution backend. pub fn with_execution_backend( options: RuntimeOptions, backend: Arc, ) -> Result { let runtime = Self::with_options(options); runtime.install_execution_backend(backend)?; Ok(runtime) } pub fn install_backend_resource_client( &self, client: Arc, ) -> Result<(), RuntimeError> { self.lock()?.backend_resource_client = Some(BackendResourceClientRef(client)); Ok(()) } pub fn install_workspace_backend_resource_client( &self, workspace_id: impl Into, client: Arc, ) -> Result<(), RuntimeError> { let workspace_id = workspace_id.into(); if workspace_id.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "Backend resource client Workspace id is empty".to_string(), )); } self.lock()? .workspace_backend_resource_clients .insert(workspace_id, BackendResourceClientRef(client)); Ok(()) } /// Create or restore a filesystem-backed Runtime. /// /// The store is scoped by `options.root`; if the directory already exists, /// persisted state is loaded and validated. If it does not exist, a fresh /// Runtime is initialized and durable files are created before return. #[cfg(feature = "fs-store")] pub fn with_fs_store(options: FsRuntimeStoreOptions) -> Result { Self::with_fs_store_inner(options, None) } /// Create or restore a filesystem-backed Runtime with an execution backend. #[cfg(feature = "fs-store")] pub fn with_fs_store_and_execution_backend( options: FsRuntimeStoreOptions, backend: Arc, ) -> Result { Self::with_fs_store_inner(options, Some(WorkerExecutionBackendRef::new(backend)?)) } #[cfg(feature = "fs-store")] fn with_fs_store_inner( options: FsRuntimeStoreOptions, execution_backend: Option, ) -> Result { let opened = FsRuntimeStore::open_or_create(options.root, &options.runtime_id)?; let mut state = if let Some(persisted) = opened.state { RuntimeState::from_persisted(persisted, opened.store)? } else { let state = RuntimeState::new_fs_backed(options.display_name, opened.store); state.persist_runtime_snapshot()?; state }; state.execution_backend = execution_backend; let runtime = Self { inner: Arc::new(Mutex::new(state)), worker_operations: Arc::new(Mutex::new(BTreeMap::new())), }; runtime.restore_persisted_worker_executions()?; Ok(runtime) } /// Management-plane summary. pub fn summary(&self) -> Result { let state = self.lock()?; let mut active_worker_count = 0; let mut stopped_worker_count = 0; for worker in state.workers.values() { match worker.status { WorkerStatus::Idle | WorkerStatus::Running | WorkerStatus::Paused => { active_worker_count += 1; } WorkerStatus::Stopped => stopped_worker_count += 1, } } Ok(RuntimeSummary { display_name: state.display_name.clone(), backend: state.backend, status: state.status, worker_count: state.workers.len(), active_worker_count, stopped_worker_count, diagnostic_count: state.diagnostics.len(), os: std::env::consts::OS.to_string(), arch: std::env::consts::ARCH.to_string(), worker_creation_available: state.execution_backend.is_some(), }) } /// Current Runtime lifecycle state. pub fn status(&self) -> Result { Ok(self.lock()?.status) } /// Store a backend-synced Profile/config bundle for later Worker creation. pub fn store_config_bundle( &self, bundle: ConfigBundle, ) -> Result { validate_config_bundle(&bundle)?; let mut state = self.lock()?; state.ensure_running()?; let reference = ConfigBundleRef { id: bundle.metadata.id.clone(), digest: bundle.metadata.digest.clone(), }; let summary = bundle.summary(); if let Some(existing) = state.config_bundles.get(&bundle.metadata.id) { if existing.metadata.digest != bundle.metadata.digest { return Err(RuntimeError::ConfigBundleDigestMismatch { bundle_id: bundle.metadata.id.clone(), expected_digest: existing.metadata.digest.clone(), actual_digest: bundle.metadata.digest.clone(), }); } return Ok(ConfigBundleAvailability { reference, summary: existing.summary(), }); } state .config_bundles .insert(bundle.metadata.id.clone(), bundle); state.persist_runtime_snapshot()?; Ok(ConfigBundleAvailability { reference, summary }) } /// List synced config bundles known to this Runtime. pub fn list_config_bundles(&self) -> Result, RuntimeError> { Ok(self .lock()? .config_bundles .values() .map(ConfigBundle::summary) .collect()) } /// Validate that a config bundle reference is present and digest-matched. pub fn check_config_bundle( &self, reference: &ConfigBundleRef, ) -> Result { let state = self.lock()?; state.check_config_bundle_ref(reference) } /// Notify the execution backend of the Workspace's current immutable /// Prompt projection. The Runtime keeps this cache outside persisted Worker /// restore authority. pub fn observe_workspace_prompt_projection( &self, projection: worker::WorkspacePromptProjection, ) -> Result<(), RuntimeError> { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "Workspace Prompt projection notification requires an execution backend" .to_string(), } })? }; backend .observe_workspace_prompt_projection(projection) .map_err(|message| RuntimeError::ExecutionBackendUnavailable { message }) } /// Stop the Runtime. v0 keeps data readable after stop, but rejects new /// create/send/worker lifecycle mutations. pub fn stop_runtime(&self) -> Result<(), RuntimeError> { let mut state = self.lock()?; if state.status == RuntimeStatus::Stopped { return Ok(()); } state.status = RuntimeStatus::Stopped; state.persist_runtime_snapshot()?; state.persist_workers()?; Ok(()) } /// Create a Runtime-owned working directory through the attached execution backend. pub fn create_working_directory( &self, request: WorkingDirectoryRequest, ) -> Result { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "working directory creation requires an execution backend".to_string(), } })? }; backend .create_working_directory(&request) .map_err(RuntimeError::from) } pub async fn create_working_directory_from_resource( &self, mut request: WorkingDirectoryRequest, ) -> Result { if let Some(ssh) = request .materialization .as_mut() .and_then(|materialization| materialization.ssh.as_mut()) { self.resolve_repository_access_resource(ssh).await?; } self.create_working_directory(request) } pub async fn observe_repository_ref_from_resource( &self, mut request: RepositoryRefObservationRequest, ) -> Result { if let Some(ssh) = request .materialization .as_mut() .and_then(|materialization| materialization.ssh.as_mut()) { self.resolve_repository_access_resource(ssh).await?; } self.observe_repository_ref(request) } pub fn observe_repository_ref( &self, request: RepositoryRefObservationRequest, ) -> Result { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "Repository ref observation requires an execution backend".to_string(), } })? }; backend .observe_repository_ref(&request) .map_err(RuntimeError::from) } pub fn authorize_working_directory_repository_access( &self, request: WorkingDirectoryRepositoryAccessRequest, ) -> Result<(), RuntimeError> { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "working directory Repository access requires an execution backend" .to_string(), } })? }; backend .authorize_working_directory_repository_access(&request) .map_err(RuntimeError::from) } async fn resolve_repository_access_resource( &self, ssh: &mut crate::catalog::RepositorySshMaterializationAccess, ) -> Result<(), RuntimeError> { if ssh.credential_candidates.is_empty() { return Err(RuntimeError::InvalidRequest( "Repository SSH access requires at least one credential candidate".to_string(), )); } if ssh .credential_candidates .iter() .all(|candidate| !candidate.private_key.expose().is_empty()) && !ssh.known_hosts_entry.expose().is_empty() { return Ok(()); } let (client, runtime_id) = { let state = self.lock()?; let client = state .workspace_backend_resource_clients .get(&ssh.secret_resource.workspace_id) .cloned() .or_else(|| state.backend_resource_client.clone()) .ok_or_else(|| { RuntimeError::InvalidRequest(format!( "Backend Repository access resource client is unavailable for Workspace `{}`", ssh.secret_resource.workspace_id )) })?; let runtime_id = state.runtime_identity.clone().ok_or_else(|| { RuntimeError::InvalidRequest("Runtime identity is unavailable".to_string()) })?; (client, runtime_id) }; tracing::info!( target: "yoi::repository_access", event = "repository_access_resource_fetch_started", workspace_id = %ssh.secret_resource.workspace_id, resource_id = %ssh.secret_resource.resource_id, runtime_id = %runtime_id, credential_candidate_count = ssh.credential_candidates.len(), "fetching Repository SSH access resource from Workspace Backend" ); let mut response = client .0 .fetch_resource(BackendResourceFetchRequest { handle: ssh.secret_resource.clone(), runtime_id: runtime_id.clone(), worker_id: None, audit_correlation_id: ssh.secret_resource.audit_correlation_id.clone(), }) .await .map_err(repository_resource_error)?; if response.kind != BackendResourceKind::RepositorySshAccess || response.content_type != REPOSITORY_SSH_ACCESS_CONTENT_TYPE || response.resource_id != ssh.secret_resource.resource_id || response.digest != ssh.secret_resource.digest || response.bytes.len() as u64 > ssh.secret_resource.max_bytes { return Err(RuntimeError::InvalidRequest( "Backend Repository SSH access resource response was invalid".to_string(), )); } let secret = serde_json::from_slice::(&response.bytes); response.bytes.fill(0); let mut secret = secret.map_err(|_| { RuntimeError::InvalidRequest( "Backend Repository SSH access resource payload was invalid".to_string(), ) })?; if secret.credential_candidates.len() != ssh.credential_candidates.len() || secret .credential_candidates .iter() .zip(&ssh.credential_candidates) .any(|(secret, metadata)| { secret.credential_id != metadata.credential_id || secret.credential_revision != metadata.credential_revision }) { return Err(RuntimeError::InvalidRequest( "Backend Repository SSH access resource credential metadata was invalid" .to_string(), )); } for (candidate, secret) in ssh .credential_candidates .iter_mut() .zip(&mut secret.credential_candidates) { candidate.private_key = crate::catalog::SensitiveString::new(std::mem::take(&mut secret.private_key)); } ssh.known_hosts_entry = crate::catalog::SensitiveString::new(std::mem::take(&mut secret.known_hosts_entry)); tracing::info!( target: "yoi::repository_access", event = "repository_access_resource_fetch_succeeded", workspace_id = %ssh.secret_resource.workspace_id, resource_id = %ssh.secret_resource.resource_id, runtime_id = %runtime_id, credential_candidate_count = ssh.credential_candidates.len(), "fetched Repository SSH access resource from Workspace Backend" ); Ok(()) } pub async fn authorize_working_directory_repository_access_from_resource( &self, mut request: WorkingDirectoryRepositoryAccessRequest, ) -> Result<(), RuntimeError> { let materialization_runtime_id = request.materialization.runtime_id.clone(); let ssh = request.materialization.ssh.as_mut().ok_or_else(|| { RuntimeError::InvalidRequest("Repository SSH access metadata is missing".to_string()) })?; if let Err(error) = self.resolve_repository_access_resource(ssh).await { tracing::warn!( target: "yoi::repository_access", event = "repository_access_resource_fetch_failed", workspace_id = %ssh.secret_resource.workspace_id, resource_id = %ssh.secret_resource.resource_id, runtime_id = %materialization_runtime_id, error = %error, "failed to fetch Repository SSH access resource from Workspace Backend" ); return Err(error); } self.authorize_working_directory_repository_access(request) } /// List Runtime-owned working directories through the attached execution backend. pub fn list_working_directories( &self, ) -> Result, RuntimeError> { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "working directory listing requires an execution backend".to_string(), } })? }; let statuses = backend.list_working_directories(); self.annotate_working_directory_statuses(statuses) } /// Get a Runtime-owned working directory status. pub fn working_directory( &self, working_directory_id: &str, ) -> Result { let backend = { let state = self.lock()?; state.ensure_running()?; state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "working directory lookup requires an execution backend".to_string(), } })? }; let status = backend .working_directory(working_directory_id) .map_err(RuntimeError::from)?; self.annotate_working_directory_status(status) } /// Open a fresh Workdir operation session in this Runtime's authorized Workspace. /// /// A same-Runtime owner Worker can be supplied as an additional persisted-binding check. /// Cross-Runtime callers rely on the Runtime's Workspace capability scope and the existence /// of the Runtime-owned materialization; Backend attachment remains occupancy authority. pub fn open_workdir_session_scoped( &self, scope: &RuntimeWorkspaceScope, working_directory_id: &str, owner_worker_ref: Option<&WorkerRef>, ) -> Result { let backend = { let mut state = self.lock()?; state.ensure_running()?; state.ensure_workspace_owner(scope, false)?; if let Some(worker_ref) = owner_worker_ref { let owns_workdir = state .workers .get(&worker_ref.worker_id) .is_some_and(|worker| { worker.belongs_to_workspace(&scope.workspace_id) && worker.working_directory.as_ref().is_some_and(|status| { status.summary.working_directory_id == working_directory_id }) }); if !owns_workdir { return Err(RuntimeError::WorkingDirectory( crate::working_directory::WorkingDirectoryDiagnostic::rejected( "working_directory_not_found", "working directory was not assigned to the authorized owner Worker", ), )); } } state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "opening a Workdir session requires an execution backend".to_string(), } })? }; backend .working_directory(working_directory_id) .map_err(RuntimeError::WorkingDirectory)?; backend .open_workdir_session(working_directory_id) .map_err(RuntimeError::WorkingDirectory) } /// Cleanup a Runtime-owned working directory. pub fn cleanup_working_directory( &self, working_directory_id: &str, ) -> Result { let backend = { let state = self.lock()?; state.ensure_running()?; if let Some(worker_id) = state.primary_worker_id_for_workdir(working_directory_id) { return Err(RuntimeError::InvalidRequest(format!( "working directory {working_directory_id} is assigned to worker {worker_id}" ))); } state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "working directory cleanup requires an execution backend".to_string(), } })? }; backend .cleanup_working_directory(working_directory_id) .map_err(RuntimeError::from) } fn annotate_working_directory_statuses( &self, statuses: Vec, ) -> Result, RuntimeError> { statuses .into_iter() .map(|status| self.annotate_working_directory_status(status)) .collect() } fn annotate_working_directory_status( &self, mut status: CatalogWorkingDirectoryStatus, ) -> Result { let state = self.lock()?; status.summary.primary_worker_id = state .primary_worker_id_for_workdir(status.summary.working_directory_id.as_str()) .map(|worker_id| worker_id.to_string()); Ok(status) } /// Create a Worker through the canonical profile-source + execution backend path. pub fn create_worker( &self, request: CreateWorkerRequest, ) -> Result { self.create_worker_with_workspace(request, None) } /// Create a Worker scoped to a workspace authorization context. pub fn create_worker_scoped( &self, scope: &RuntimeWorkspaceScope, request: CreateWorkerRequest, ) -> Result { self.create_worker_with_workspace(request, Some(scope)) } fn existing_worker_for_create( &self, request: &CreateWorkerRequest, scope: Option<&RuntimeWorkspaceScope>, ) -> Result, RuntimeError> { let mut state = self.lock()?; state.ensure_running()?; validate_create_worker_request(request)?; validate_create_workspace_scope(request, scope.map(|scope| scope.workspace_id.as_str()))?; if let Some(scope) = scope { state.ensure_workspace_owner(scope, false)?; } let workspace_id = scope.map(|scope| scope.workspace_id.as_str()); if let Some(existing) = state.workers.get(&request.worker_id) { if existing.workspace_id.as_deref() != workspace_id { return Err(RuntimeError::InvalidRequest(format!( "worker {} already belongs to another Workspace scope", request.worker_id ))); } if existing.request.create_fingerprint != request.create_fingerprint { return Err(RuntimeError::InvalidRequest(format!( "worker {} was already created with a different fingerprint", request.worker_id ))); } return Ok(Some(existing.detail())); } Ok(None) } fn refresh_workspace_config(&self, request: &CreateWorkerRequest) -> Result<(), RuntimeError> { let ProfileSourceArchiveSource::WorkspaceConfig { archive: expected_archive, } = &request.profile_source else { return Ok(()); }; let expected = request.config_bundle.as_ref().ok_or_else(|| { RuntimeError::InvalidRequest( "Workspace Config profile source requires a config_bundle reference".to_string(), ) })?; let workspace_api = request.workspace_api.clone().ok_or_else(|| { RuntimeError::InvalidRequest( "Workspace Config profile source requires a Workspace API reference".to_string(), ) })?; let profile_key = match &request.profile { ProfileSelector::Builtin(value) | ProfileSelector::Named(value) => value.clone(), }; let cache_key = format!( "{}\u{0}{}\u{0}{}", workspace_api.base_url, workspace_api.workspace_id, profile_key ); let fetch_gate = { let mut state = self.lock()?; state .workspace_config_fetch_gates .entry(cache_key.clone()) .or_insert_with(|| Arc::new(Mutex::new(()))) .clone() }; let _fetch_guard = fetch_gate.lock().map_err(|_| RuntimeError::StatePoisoned)?; let (backend, cached) = { let state = self.lock()?; let cached_reference = state .workspace_config_latest .get(&cache_key) .unwrap_or(expected); let cached = state .config_bundles .get(&cached_reference.id) .filter(|bundle| bundle.metadata.digest == cached_reference.digest) .map(|bundle| ConfigBundleRef { id: bundle.metadata.id.clone(), digest: bundle.metadata.digest.clone(), }); (state.execution_backend.clone(), cached) }; let backend = backend.ok_or_else(|| RuntimeError::ExecutionBackendUnavailable { message: "Workspace Config refresh requires an execution backend".to_string(), })?; let fetched = backend .fetch_workspace_config(WorkspaceConfigFetchRequest { workspace_api, profile: request.profile.clone(), expected: expected.clone(), cached: cached.clone(), }) .map_err(|message| { RuntimeError::InvalidRequest(format!( "failed to refresh latest Workspace Config: {message}" )) })?; let resolved = match fetched { WorkspaceConfigFetchResult::NotModified => cached.ok_or_else(|| { RuntimeError::InvalidRequest( "latest Workspace Config returned not-modified without a cached bundle" .to_string(), ) })?, WorkspaceConfigFetchResult::Modified(bundle) => { self.store_config_bundle(bundle)?.reference } }; if &resolved != expected { return Err(RuntimeError::InvalidRequest(format!( "latest Workspace Config changed while Worker creation was being prepared: expected '{}@{}', got '{}@{}'", expected.id, expected.digest, resolved.id, resolved.digest ))); } let mut state = self.lock()?; let bundle = state.config_bundles.get(&resolved.id).ok_or_else(|| { RuntimeError::InvalidRequest( "latest Workspace Config was not available after refresh".to_string(), ) })?; let archive = bundle.profile_source_archive.as_ref().ok_or_else(|| { RuntimeError::InvalidRequest( "latest Workspace Config is missing its profile source archive".to_string(), ) })?; if &archive.reference != expected_archive { return Err(RuntimeError::InvalidRequest( "latest Workspace Config profile source archive reference mismatch".to_string(), )); } state.workspace_config_latest.insert(cache_key, resolved); Ok(()) } fn create_worker_with_workspace( &self, request: CreateWorkerRequest, scope: Option<&RuntimeWorkspaceScope>, ) -> Result { let worker_id = request.worker_id; let workspace_id = scope.map(|scope| scope.workspace_id.as_str()); let started_at = std::time::Instant::now(); let result = self.create_worker_with_workspace_inner(request, scope); if let Err(error) = &result { write_runtime_worker_create_failure( worker_id, workspace_id, error, started_at.elapsed(), ); } result } fn create_worker_with_workspace_inner( &self, request: CreateWorkerRequest, scope: Option<&RuntimeWorkspaceScope>, ) -> Result { let operation_lock = self.worker_operation_lock(request.worker_id)?; let _operation_guard = operation_lock .lock() .map_err(|_| RuntimeError::StatePoisoned)?; if let Some(existing) = self.existing_worker_for_create(&request, scope)? { return Ok(existing); } self.refresh_workspace_config(&request)?; let (backend, worker_ref, spawn_request) = { let mut state = self.lock()?; state.ensure_running()?; validate_create_worker_request(&request)?; validate_create_workspace_scope( &request, scope.map(|scope| scope.workspace_id.as_str()), )?; if let Some(scope) = scope { state.ensure_workspace_owner(scope, true)?; }; let workspace_id = scope.map(|scope| scope.workspace_id.as_str()); if let Some(existing) = state.workers.get(&request.worker_id) { if existing.workspace_id.as_deref() != workspace_id { return Err(RuntimeError::InvalidRequest(format!( "worker {} already belongs to another Workspace scope", request.worker_id ))); } if existing.request.create_fingerprint != request.create_fingerprint { return Err(RuntimeError::InvalidRequest(format!( "worker {} was already created with a different fingerprint", request.worker_id ))); } return Ok(existing.detail()); } state.validate_worker_config_boundary(&request)?; if let Some(working_directory_id) = requested_primary_workdir_id(&request) { if let Some(owner_worker_id) = state.primary_worker_id_for_workdir(working_directory_id) { return Err(RuntimeError::InvalidRequest(format!( "working directory {working_directory_id} is already assigned to worker {owner_worker_id}" ))); } } let backend = state.execution_backend.clone().ok_or_else(|| { RuntimeError::ExecutionBackendUnavailable { message: "worker creation requires an execution backend".to_string(), } })?; let config_bundle = state.resolve_config_bundle_ref(request.config_bundle.as_ref())?; let worker_id = request.worker_id; let worker_ref = WorkerRef::new(worker_id); let durable_request = durable_create_worker_request(&request); let record = WorkerRecord { worker_ref: worker_ref.clone(), worker_id: worker_id.clone(), status: WorkerStatus::Stopped, worker_state: None, workspace_id: scope.map(|scope| scope.workspace_id.clone()), request: durable_request, run_generation: 1, execution_bound: true, restore_intent: WorkerRestoreIntent::Explicit, working_directory: None, execution_handle: None, internal_workers: BTreeMap::new(), }; state.workers.insert(worker_id, record); state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; let spawn_request = WorkerExecutionSpawnRequest { worker_ref: worker_ref.clone(), run_generation: 1, request, workspace_scope: scope.cloned(), context: self.execution_context(worker_ref.clone()), working_directory: None, config_bundle, }; (backend, worker_ref, spawn_request) }; let spawn_result = backend.spawn_worker(spawn_request); let (handle, initial_worker_state, working_directory) = match spawn_result { WorkerExecutionSpawnResult::Connected { handle, worker_state, working_directory, } => (handle, worker_state, working_directory), WorkerExecutionSpawnResult::Rejected(result) | WorkerExecutionSpawnResult::Errored(result) => { self.rollback_failed_create(&worker_ref)?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }); } }; if let Some(mut initial_input) = { let state = self.lock()?; state.worker(&worker_ref)?.request.initial_input.clone() } { let expected_submission_id = initial_input .submission_request_id .clone() .filter(|request_id| !request_id.trim().is_empty()) .unwrap_or_else(|| Uuid::now_v7().to_string()); initial_input.submission_request_id = Some(expected_submission_id.clone()); let dispatch_result = backend.dispatch_input(&handle, initial_input.clone()); if !dispatch_result.is_accepted() { self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: dispatch_result.operation, outcome: dispatch_result.outcome, message: dispatch_result.message_or_default(), result: dispatch_result, }); } let has_durable_acceptance = dispatch_result .submission .as_ref() .is_some_and(|ack| ack.submission_request_id == expected_submission_id); if !has_durable_acceptance { self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; let result = WorkerExecutionResult::rejected( WorkerExecutionOperation::Input, "execution backend accepted initial input without a durable submission acknowledgement", ); return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }); } let detail = match self.commit_created_worker( &worker_ref, handle.clone(), initial_worker_state.clone(), working_directory, dispatch_result, ) { Ok(detail) => detail, Err(error) => { self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; return Err(error); } }; if let Err(error) = self.record_input_observation(&worker_ref, initial_input) { self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; return Err(error); } Ok(detail) } else { match self.commit_created_worker( &worker_ref, handle.clone(), initial_worker_state, working_directory, WorkerExecutionResult::accepted(WorkerExecutionOperation::Spawn), ) { Ok(detail) => Ok(detail), Err(error) => { self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; Err(error) } } } } /// List Workers known to this Runtime. pub fn list_workers(&self) -> Result, RuntimeError> { let state = self.lock()?; Ok(state.workers.values().map(WorkerRecord::summary).collect()) } pub fn subscribe_event_selector( &self, selector: EventSubscriptionSelector, ) -> Result { self.subscribe_event_selector_for_workspace(None, selector) } pub fn subscribe_event_selector_scoped( &self, scope: &RuntimeWorkspaceScope, selector: EventSubscriptionSelector, ) -> Result { self.subscribe_event_selector_for_workspace(Some(scope), selector) } fn subscribe_event_selector_for_workspace( &self, scope: Option<&RuntimeWorkspaceScope>, selector: EventSubscriptionSelector, ) -> Result { selector.validate().map_err(subscription_validation_error)?; if matches!( selector, EventSubscriptionSelector::WorkerProtocol { .. } | EventSubscriptionSelector::WorkspaceWorkers | EventSubscriptionSelector::WorkspaceWorkdirs ) { return Err(RuntimeError::InvalidRequest( "Runtime event subscriptions support only runtime_workers and worker_lifecycle selectors" .to_string(), )); } let mut state = self.lock()?; if let Some(scope) = scope { state.ensure_workspace_owner_for_existing_workers(scope)?; state.persist_runtime_snapshot()?; } let snapshot = state.subscription_snapshot(scope, &selector)?; let snapshot_revision = state.subscription_revision; let subscription_id = state.next_event_subscription_id; state.next_event_subscription_id = state.next_event_subscription_id.saturating_add(1); let (sender, receiver) = mpsc::channel(SUBSCRIPTION_QUEUE_CAPACITY); let lagged = Arc::new(AtomicBool::new(false)); state.subscriptions.insert( subscription_id, SubscriptionSink { selector: selector.clone(), workspace_id: scope.map(|scope| scope.workspace_id.clone()), sender, lagged: lagged.clone(), }, ); Ok(RuntimeEventSelectorSubscription { subscription_id, selector, snapshot_revision, snapshot, receiver, lagged, runtime: Arc::downgrade(&self.inner), }) } /// List stopped Workers known to this Runtime. pub fn list_stopped_workers(&self) -> Result, RuntimeError> { let state = self.lock()?; Ok(state .workers .values() .filter(|worker| worker.status == WorkerStatus::Stopped) .map(WorkerRecord::summary) .collect()) } /// List Workers visible to a workspace-scoped Runtime authorization context. pub fn list_workers_scoped( &self, scope: &RuntimeWorkspaceScope, ) -> Result, RuntimeError> { let mut state = self.lock()?; let visible_workspace = state.ensure_workspace_owner_for_existing_workers(scope)?; if !visible_workspace { return Ok(Vec::new()); } state.persist_runtime_snapshot()?; Ok(state .workers .values() .filter(|worker| worker.belongs_to_workspace(&scope.workspace_id)) .map(WorkerRecord::summary) .collect()) } /// List stopped Workers visible to a workspace-scoped Runtime authorization context. pub fn list_stopped_workers_scoped( &self, scope: &RuntimeWorkspaceScope, ) -> Result, RuntimeError> { let mut state = self.lock()?; let visible_workspace = state.ensure_workspace_owner_for_existing_workers(scope)?; if !visible_workspace { return Ok(Vec::new()); } state.persist_runtime_snapshot()?; Ok(state .workers .values() .filter(|worker| { worker.status == WorkerStatus::Stopped && worker.belongs_to_workspace(&scope.workspace_id) }) .map(WorkerRecord::summary) .collect()) } /// Fetch Worker detail through a workspace-scoped Runtime authorization context. pub fn worker_detail_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, ) -> Result { let mut state = self.lock()?; let worker = state.worker(worker_ref)?; if !worker.belongs_to_workspace(&scope.workspace_id) { return Err(RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }); } state.ensure_workspace_owner(scope, true)?; state.persist_runtime_snapshot()?; Ok(state.worker(worker_ref)?.detail()) } /// Fetch Worker detail. The supplied [`WorkerRef`] must match this Runtime. pub fn worker_detail(&self, worker_ref: &WorkerRef) -> Result { let state = self.lock()?; let worker = state.worker(worker_ref)?; Ok(worker.detail()) } fn ensure_worker_in_workspace( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, ) -> Result<(), RuntimeError> { let mut state = self.lock()?; let worker = state.worker(worker_ref)?; if !worker.belongs_to_workspace(&scope.workspace_id) { return Err(RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }); } state.ensure_workspace_owner(scope, true)?; state.persist_runtime_snapshot()?; Ok(()) } /// Replace the Workspace API identity binding persisted for a Worker. pub fn replace_worker_workspace_api_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, workspace_api: WorkspaceApiRef, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; if workspace_api.workspace_id != scope.workspace_id { return Err(RuntimeError::InvalidRequest(format!( "Workspace API scope `{}` does not match authorized workspace `{}`", workspace_api.workspace_id, scope.workspace_id ))); } self.replace_worker_workspace_api(worker_ref, workspace_api) } pub fn replace_worker_workspace_api( &self, worker_ref: &WorkerRef, workspace_api: WorkspaceApiRef, ) -> Result { let previous_workspace_api = { let state = self.lock()?; let worker = state.worker(worker_ref)?; if let Some(existing) = worker.request.workspace_api.as_ref() && (existing.workspace_id != workspace_api.workspace_id || existing.base_url.trim_end_matches('/') != workspace_api.base_url.trim_end_matches('/')) { return Err(RuntimeError::InvalidRequest( "Workspace API replacement cannot change Worker Workspace identity or base URL" .to_string(), )); } worker.request.workspace_api.clone() }; { let mut state = self.lock()?; state.worker_mut(worker_ref)?.request.workspace_api = Some(workspace_api); if let Err(error) = state.persist_runtime_snapshot() { state.worker_mut(worker_ref)?.request.workspace_api = previous_workspace_api; return Err(error); } } let state = self.lock()?; Ok(state.worker(worker_ref)?.detail()) } /// Attach a live execution through a workspace-scoped Runtime authorization context. pub fn restore_worker_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.restore_worker(worker_ref) } /// Attach a live execution to a persisted Worker definition. /// /// Current liveness is never read from disk. If a handle is already /// present this is idempotent; otherwise the configured backend is tried. pub fn restore_worker(&self, worker_ref: &WorkerRef) -> Result { let (backend, request) = { let mut state = self.lock()?; state.ensure_running()?; let (worker_request, previous_working_directory, run_generation) = { let worker = state.worker(worker_ref)?; if worker.execution_handle.is_some() { return Ok(worker.detail()); } ( worker.request.clone(), worker.working_directory.clone(), worker.run_generation.saturating_add(1).max(1), ) }; let backend = state.execution_backend.clone().ok_or_else(|| { RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id.clone(), message: "runtime has no execution backend".to_string(), } })?; { let worker = state.worker_mut(worker_ref)?; worker.run_generation = run_generation; worker.execution_bound = true; worker.restore_intent = WorkerRestoreIntent::Automatic; } state.persist_worker(&worker_ref.worker_id)?; let workspace_scope = worker_request.workspace_api.as_ref().and_then(|api| { state .workspace_owners .get(&api.workspace_id) .map(|server_id| RuntimeWorkspaceScope::new(&api.workspace_id, server_id)) }); let request = WorkerExecutionRestoreRequest { worker_ref: worker_ref.clone(), run_generation, request: worker_request, workspace_scope, context: self.execution_context(worker_ref.clone()), previous_working_directory, working_directory: None, config_bundle: None, }; (backend, request) }; match backend.restore_worker(request) { WorkerExecutionSpawnResult::Connected { handle, worker_state, working_directory, } => { self.commit_restored_worker_execution( worker_ref, handle, worker_state, WorkerStatus::Idle, working_directory, )?; self.worker_detail(worker_ref) } WorkerExecutionSpawnResult::Rejected(result) | WorkerExecutionSpawnResult::Errored(result) => { #[cfg(feature = "fs-store")] { self.lock()? .record_restore_failure(worker_ref, result.clone())?; } Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }) } } } fn ensure_worker_execution(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> { let state = self.lock()?; let worker = state.worker(worker_ref)?; if worker.execution_handle.is_some() { return Ok(()); } let message = if worker.status == WorkerStatus::Stopped { "stopped worker requires an explicit restore" } else { "worker has no live execution handle" }; Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id, message: message.to_string(), }) } /// Accept input into a Worker through a workspace-scoped Runtime authorization context. pub fn send_input_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, input: WorkerInput, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.send_input(worker_ref, input) } /// Accept input into a Worker. pub fn send_input( &self, worker_ref: &WorkerRef, mut input: WorkerInput, ) -> Result { validate_worker_input(&input)?; let expected_submission_id = if matches!(input.kind, WorkerInputKind::User | WorkerInputKind::Notify) { let submission_id = input .submission_request_id .clone() .filter(|request_id| !request_id.trim().is_empty()) .unwrap_or_else(|| Uuid::now_v7().to_string()); input.submission_request_id = Some(submission_id.clone()); Some(submission_id) } else { None }; self.ensure_worker_execution(worker_ref)?; let (backend, handle) = { let state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; if !worker.status.is_active() { return Err(RuntimeError::InvalidRequest(format!( "worker {} is not running", worker_ref.worker_id ))); } let backend = state.execution_backend.clone(); let handle = worker.execution_handle.clone(); match (backend, handle) { (Some(backend), Some(handle)) => (backend, handle), _ => { return Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id.clone(), message: "worker has no live execution handle".to_string(), }); } } }; let dispatch_result = backend.dispatch_input(&handle, input.clone()); if !dispatch_result.is_accepted() { self.record_execution_result(worker_ref, dispatch_result.clone())?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: dispatch_result.operation, outcome: dispatch_result.outcome, message: dispatch_result.message_or_default(), result: dispatch_result, }); } if let Some(expected_submission_id) = expected_submission_id && dispatch_result .submission .as_ref() .is_none_or(|ack| ack.submission_request_id != expected_submission_id) { let result = WorkerExecutionResult::rejected( WorkerExecutionOperation::Input, "execution backend did not acknowledge the committed Runtime submission request id", ); self.record_execution_result(worker_ref, result.clone())?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }); } let submission = dispatch_result.submission.clone(); let mut state = self.lock()?; state.ensure_running()?; let worker = state.worker_mut(worker_ref)?; if let Some(snapshot) = dispatch_result.worker_state.as_ref() { let _ = worker.apply_worker_state(snapshot); } let status = worker.status; #[cfg(feature = "ws-server")] if let Some(payload) = input_protocol_event(&input) { state.push_worker_observation_event(worker_ref.clone(), payload); } state.publish_worker_upsert(worker_ref.worker_id)?; state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(WorkerInteractionAck { worker_ref: worker_ref.clone(), status, submission, }) } /// Store a client-local file in the owning Worker session before input submit. pub fn upload_worker_file_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, file_name: &str, media_type: &str, content: &[u8], ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.upload_worker_file(worker_ref, file_name, media_type, content) } pub fn upload_worker_file( &self, worker_ref: &WorkerRef, file_name: &str, media_type: &str, content: &[u8], ) -> Result { self.upload_worker_file_inner(worker_ref, file_name, media_type, content, None) } pub fn upload_worker_file_with_context_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, file_name: &str, media_type: &str, content: &[u8], context: &session_store::UploadedFileUploadContext, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.upload_worker_file_inner(worker_ref, file_name, media_type, content, Some(context)) } pub fn upload_worker_file_with_context( &self, worker_ref: &WorkerRef, file_name: &str, media_type: &str, content: &[u8], context: &session_store::UploadedFileUploadContext, ) -> Result { self.upload_worker_file_inner(worker_ref, file_name, media_type, content, Some(context)) } fn upload_worker_file_inner( &self, worker_ref: &WorkerRef, file_name: &str, media_type: &str, content: &[u8], context: Option<&session_store::UploadedFileUploadContext>, ) -> Result { let (backend, handle) = { let state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; match ( state.execution_backend.clone(), worker.execution_handle.clone(), ) { (Some(backend), Some(handle)) => (backend, handle), _ => { return Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id.clone(), message: "worker has no live execution handle".to_string(), }); } } }; backend .upload_file(&handle, file_name, media_type, content, context) .map_err(|result| RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }) } /// Delete an unsubmitted uploaded file from the owning Worker session. pub fn delete_worker_uploaded_file_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, artifact_id: &str, ) -> Result<(), RuntimeError> { self.ensure_worker_in_workspace(scope, worker_ref)?; self.delete_worker_uploaded_file(worker_ref, artifact_id) } pub fn delete_worker_uploaded_file( &self, worker_ref: &WorkerRef, artifact_id: &str, ) -> Result<(), RuntimeError> { let (backend, handle) = { let state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; match ( state.execution_backend.clone(), worker.execution_handle.clone(), ) { (Some(backend), Some(handle)) => (backend, handle), _ => { return Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id.clone(), message: "worker has no live execution handle".to_string(), }); } } }; let result = backend.delete_uploaded_file(&handle, artifact_id); if result.is_accepted() { Ok(()) } else { Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }) } } /// Return live completion entries through a workspace-scoped Runtime authorization context. pub fn worker_completions_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, kind: protocol::CompletionKind, prefix: &str, ) -> Result, RuntimeError> { self.ensure_worker_in_workspace(scope, worker_ref)?; self.worker_completions(worker_ref, kind, prefix) } /// Return live completion entries for the Worker composer. pub fn worker_completions( &self, worker_ref: &WorkerRef, kind: protocol::CompletionKind, prefix: &str, ) -> Result, RuntimeError> { let (backend, handle) = { let state = self.lock()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; ( state.execution_backend.clone(), worker.execution_handle.clone(), ) }; let Some((backend, handle)) = backend.zip(handle) else { return Ok(Vec::new()); }; Ok(backend.worker_completions(&handle, kind, prefix)) } /// Accept a protocol method through a workspace-scoped Runtime authorization context. pub fn send_protocol_method_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, method: Method, ) -> Result, RuntimeError> { self.ensure_worker_in_workspace(scope, worker_ref)?; self.send_protocol_method(worker_ref, method) } /// Accept a protocol method for a Worker through a Backend/runtime transport. /// /// Most methods are delivered to the execution backend unchanged. Methods with /// direct same-connection replies in the local socket protocol return those /// events from this function so WebSocket transports can write them back to the /// requesting client without rebroadcasting them. pub fn send_protocol_method( &self, worker_ref: &WorkerRef, method: Method, ) -> Result, RuntimeError> { if let Method::ListCompletions { kind, prefix } = method { let entries = self.worker_completions(worker_ref, kind, &prefix)?; return Ok(vec![Event::Completions { kind, entries }]); } if matches!(&method, Method::Shutdown { .. }) { self.stop_worker(worker_ref, Some("worker protocol shutdown".to_string()))?; return Ok(Vec::new()); } self.ensure_worker_execution(worker_ref)?; let (backend, handle) = { let state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; if !worker.status.is_active() { return Err(RuntimeError::InvalidRequest(format!( "worker {} is not running", worker_ref.worker_id ))); } let backend = state.execution_backend.clone(); let handle = worker.execution_handle.clone(); match (backend, handle) { (Some(backend), Some(handle)) => (backend, handle), _ => { return Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id.clone(), message: "worker has no live execution handle".to_string(), }); } } }; let dispatch_result = backend.dispatch_method(&handle, method); if !dispatch_result.is_accepted() { self.record_execution_result(worker_ref, dispatch_result.clone())?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: dispatch_result.operation, outcome: dispatch_result.outcome, message: dispatch_result.message_or_default(), result: dispatch_result, }); } self.record_execution_result(worker_ref, dispatch_result)?; Ok(Vec::new()) } fn commit_created_worker( &self, worker_ref: &WorkerRef, handle: WorkerExecutionHandle, initial_worker_state: protocol::WorkerStateSnapshot, working_directory: Option, result: WorkerExecutionResult, ) -> Result { let mut state = self.lock()?; let detail = { let worker = state.worker_mut(worker_ref)?; worker.execution_handle = Some(handle); worker.execution_bound = true; worker.status = WorkerStatus::Idle; let _ = worker.apply_worker_state(&initial_worker_state); if let Some(snapshot) = result.worker_state.as_ref() { let _ = worker.apply_worker_state(snapshot); } worker.restore_intent = restore_intent_for_status(worker.status); worker.working_directory = working_directory; worker.detail() }; state.publish_worker_upsert(worker_ref.worker_id)?; state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(detail) } fn cleanup_connected_failed_create( &self, backend: &WorkerExecutionBackendRef, worker_ref: &WorkerRef, handle: &WorkerExecutionHandle, ) -> Result<(), RuntimeError> { let stop_result = backend.stop_worker(handle); if stop_result.is_accepted() { return self.rollback_failed_create(worker_ref); } let mut state = self.lock()?; let record = state.worker_mut(worker_ref)?; record.execution_handle = Some(handle.clone()); record.worker_state = stop_result.worker_state.clone(); state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(()) } fn rollback_failed_create(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> { let mut state = self.lock()?; if state.workers.contains_key(&worker_ref.worker_id) { state.delete_worker_snapshot(&worker_ref.worker_id)?; let record = state .workers .remove(&worker_ref.worker_id) .expect("Worker existence checked before failed-create rollback"); let workspace_id = record.workspace_id.clone(); if let Some(workspace_id) = workspace_id.as_deref() { state.forget_workspace_owner_if_unused(workspace_id); } #[cfg(feature = "ws-server")] state .observation_events .retain(|event| event.worker_ref != *worker_ref); state.publish_worker_removed(worker_ref.worker_id, workspace_id.as_deref())?; state.persist_runtime_snapshot()?; } Ok(()) } fn record_execution_result( &self, worker_ref: &WorkerRef, result: WorkerExecutionResult, ) -> Result<(), RuntimeError> { // Accepted dispatch without a state snapshot is transport evidence only; // the revisioned protocol stream remains live authority. Test/detached // backends may return an exact full snapshot as their acknowledgement. if !result.is_accepted() { return Ok(()); } let Some(snapshot) = result.worker_state else { return Ok(()); }; let mut state = self.lock()?; let worker = state.worker_mut(worker_ref)?; let applied = worker .apply_worker_state(&snapshot) .is_ok_and(|result| matches!(result, protocol::WorkerStateSnapshotApply::Applied)); if !applied { return Ok(()); } state.publish_worker_upsert(worker_ref.worker_id)?; state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(()) } fn dispatch_lifecycle_to_backend( &self, worker_ref: &WorkerRef, operation: WorkerExecutionOperation, ) -> Result<(), RuntimeError> { let Some((backend, handle)) = ({ let state = self.lock()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; if !worker.status.is_active() { return Ok(()); } match ( state.execution_backend.clone(), worker.execution_handle.clone(), ) { (Some(backend), Some(handle)) => Some((backend, handle)), _ => None, } }) else { return Ok(()); }; let result = match operation { WorkerExecutionOperation::Stop => backend.stop_worker(&handle), WorkerExecutionOperation::Cancel => backend.cancel_worker(&handle), WorkerExecutionOperation::Spawn | WorkerExecutionOperation::Restore | WorkerExecutionOperation::Input | WorkerExecutionOperation::UploadFile | WorkerExecutionOperation::DeleteUploadedFile | WorkerExecutionOperation::ProtocolMethod => return Ok(()), }; if result.is_accepted() { return Ok(()); } Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }) } /// Stop a Worker through a workspace-scoped Runtime authorization context. pub fn stop_worker_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, reason: Option, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.stop_worker(worker_ref, reason) } /// Stop a Worker. Repeated stops are idempotent. pub fn stop_worker( &self, worker_ref: &WorkerRef, reason: Option, ) -> Result { self.dispatch_lifecycle_to_backend(worker_ref, WorkerExecutionOperation::Stop)?; let _ = reason; self.transition_worker(worker_ref, WorkerStatus::Stopped) } /// Cancel a Worker through a workspace-scoped Runtime authorization context. pub fn cancel_worker_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, reason: Option, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.cancel_worker(worker_ref, reason) } /// Cancel the current run while keeping the Worker session available. pub fn cancel_worker( &self, worker_ref: &WorkerRef, reason: Option, ) -> Result { { let state = self.lock()?; state.ensure_running()?; if state.worker(worker_ref)?.status == WorkerStatus::Stopped { return Ok(WorkerLifecycleAck { worker_ref: worker_ref.clone(), status: WorkerStatus::Stopped, worker_state: None, }); } } self.dispatch_lifecycle_to_backend(worker_ref, WorkerExecutionOperation::Cancel)?; let _ = reason; let state = self.lock()?; let worker = state.worker(worker_ref)?; Ok(WorkerLifecycleAck { worker_ref: worker_ref.clone(), status: worker.status, worker_state: worker.worker_state.clone(), }) } /// Delete a non-running Worker through a workspace-scoped Runtime authorization context. pub fn delete_worker_scoped( &self, scope: &RuntimeWorkspaceScope, worker_ref: &WorkerRef, ) -> Result { self.ensure_worker_in_workspace(scope, worker_ref)?; self.delete_worker(worker_ref) } /// Delete a non-running Worker from Runtime state and persisted Worker storage. pub fn delete_worker( &self, worker_ref: &WorkerRef, ) -> Result { let operation_lock = self.worker_operation_lock(worker_ref.worker_id)?; let _operation_guard = operation_lock .lock() .map_err(|_| RuntimeError::StatePoisoned)?; let (backend, execution_handle) = { let state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; if worker.status.is_active() { return Err(RuntimeError::InvalidRequest(format!( "worker {} is running and must be stopped before deletion", worker_ref.worker_id ))); } ( state.execution_backend.clone(), worker.execution_handle.clone(), ) }; if let Some(handle) = execution_handle { let backend = backend.ok_or_else(|| RuntimeError::ExecutionBackendUnavailable { message: "Worker deletion requires its execution backend to confirm shutdown" .to_string(), })?; let result = backend.stop_worker(&handle); if !result.is_accepted() { return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id, operation: result.operation, outcome: result.outcome, message: result.message_or_default(), result, }); } } let mut state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker(worker_ref)?; if worker.status.is_active() { return Err(RuntimeError::InvalidRequest(format!( "worker {} became active before deletion", worker_ref.worker_id ))); } state.delete_worker_snapshot(&worker_ref.worker_id)?; let removed = state.workers.remove(&worker_ref.worker_id).ok_or_else(|| { RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, } })?; let removed_workspace_id = removed.workspace_id.clone(); if let Some(workspace_id) = removed_workspace_id.as_deref() { state.forget_workspace_owner_if_unused(workspace_id); } #[cfg(feature = "ws-server")] state .observation_events .retain(|event| event.worker_ref != *worker_ref); state.publish_worker_removed(worker_ref.worker_id, removed_workspace_id.as_deref())?; state.persist_runtime_snapshot()?; Ok(WorkerDeleteResult { worker_id: removed.worker_id, deleted: true, }) } /// Cursor pointing after the current worker-scoped protocol observation event. #[cfg(feature = "ws-server")] pub fn worker_observation_cursor_now( &self, worker_ref: &WorkerRef, ) -> Result { let state = self.lock()?; state.ensure_worker_ref(worker_ref)?; let sequence = state .observation_events .iter() .rev() .find(|event| &event.worker_ref == worker_ref) .map(|event| event.sequence) .unwrap_or(0); Ok(WorkerObservationCursor::new(sequence)) } /// Build the current Worker Snapshot event used as the first observation frame. #[cfg(feature = "ws-server")] pub fn worker_observation_snapshot( &self, worker_ref: &WorkerRef, ) -> Result { let (backend, handle) = { let state = self.lock()?; let worker = state.worker(worker_ref)?; ( state.execution_backend.clone(), worker.execution_handle.clone(), ) }; if let (Some(backend), Some(handle)) = (backend, handle) { if let Some(snapshot) = backend.worker_snapshot(&handle) { return Ok(snapshot); } } Err(RuntimeError::WorkerExecutionUnavailable { worker_id: worker_ref.worker_id, message: "authoritative Worker snapshot is unavailable".to_string(), }) } /// Replay retained worker-scoped protocol observation events after a cursor. #[cfg(feature = "ws-server")] pub fn read_worker_observation_events( &self, worker_ref: &WorkerRef, cursor: WorkerObservationCursor, ) -> Result, RuntimeError> { let state = self.lock()?; state.ensure_worker_ref(worker_ref)?; state.validate_worker_observation_cursor(worker_ref, cursor)?; Ok(state .observation_events .iter() .filter(|event| &event.worker_ref == worker_ref && event.sequence > cursor.sequence) .cloned() .collect()) } /// Subscribe to live protocol observation events. #[cfg(feature = "ws-server")] pub fn subscribe_worker_observation( &self, ) -> Result, RuntimeError> { Ok(self.lock()?.observation_tx.subscribe()) } /// Append a Worker protocol event to the observation bus. #[cfg(feature = "ws-server")] pub fn observe_worker_event( &self, worker_ref: &WorkerRef, payload: protocol::Event, ) -> Result { let mut state = self.lock()?; state.ensure_worker_ref(worker_ref)?; let worker_state_changed = state.project_protocol_event_to_worker_state(worker_ref, &payload); let activity_changed = state.project_internal_worker_activity(worker_ref, &payload); if worker_state_changed || activity_changed { state.publish_worker_upsert(worker_ref.worker_id)?; } if worker_state_changed { state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; } let event = state.push_worker_observation_event(worker_ref.clone(), payload); Ok(event) } /// Snapshot current diagnostics. pub fn diagnostics(&self) -> Result, RuntimeError> { Ok(self.lock()?.diagnostics.clone()) } #[cfg(feature = "ws-server")] fn record_input_observation( &self, worker_ref: &WorkerRef, input: WorkerInput, ) -> Result<(), RuntimeError> { let mut state = self.lock()?; state.ensure_worker_ref(worker_ref)?; if let Some(payload) = input_protocol_event(&input) { state.push_worker_observation_event(worker_ref.clone(), payload); } Ok(()) } #[cfg(not(feature = "ws-server"))] fn record_input_observation( &self, _worker_ref: &WorkerRef, _input: WorkerInput, ) -> Result<(), RuntimeError> { Ok(()) } fn transition_worker( &self, worker_ref: &WorkerRef, status: WorkerStatus, ) -> Result { let mut state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; let worker = state.worker_mut(worker_ref)?; worker.status = status; worker.worker_state = None; worker.restore_intent = restore_intent_for_status(status); worker.execution_handle = None; worker.internal_workers.clear(); let status = worker.status; state.publish_worker_upsert(worker_ref.worker_id)?; state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(WorkerLifecycleAck { worker_ref: worker_ref.clone(), status, worker_state: None, }) } fn install_execution_backend( &self, backend: Arc, ) -> Result<(), RuntimeError> { let backend = WorkerExecutionBackendRef::new(backend)?; let mut state = self.lock()?; state.execution_backend = Some(backend); Ok(()) } #[cfg(feature = "ws-server")] fn execution_context(&self, worker_ref: WorkerRef) -> crate::execution::WorkerExecutionContext { let runtime = self.clone(); crate::execution::WorkerExecutionContext::new( worker_ref, Arc::new(move |worker_ref, payload| runtime.observe_worker_event(&worker_ref, payload)), ) } #[cfg(not(feature = "ws-server"))] fn execution_context(&self, worker_ref: WorkerRef) -> crate::execution::WorkerExecutionContext { crate::execution::WorkerExecutionContext::new(worker_ref) } #[cfg(feature = "fs-store")] fn restore_persisted_worker_executions(&self) -> Result<(), RuntimeError> { #[derive(Clone)] struct RestoreCandidate { worker_ref: WorkerRef, request: CreateWorkerRequest, run_generation: u64, previous_working_directory: Option, config_bundle: Option, } let candidates = { let mut state = self.lock()?; if state.execution_backend.is_none() { return Ok(()); } let worker_ids = state .workers .values() .filter(|worker| { worker.execution_handle.is_none() && worker.execution_bound && worker.status.is_active() && worker.restore_intent == WorkerRestoreIntent::Automatic }) .map(|worker| worker.worker_id) .collect::>(); let mut candidates = Vec::with_capacity(worker_ids.len()); for worker_id in worker_ids { let (worker_ref, request, previous_working_directory, run_generation) = { let worker = state .workers .get(&worker_id) .expect("collected Worker exists"); ( worker.worker_ref.clone(), worker.request.clone(), worker.working_directory.clone(), worker.run_generation.saturating_add(1).max(1), ) }; state .workers .get_mut(&worker_id) .expect("collected Worker exists") .run_generation = run_generation; state.persist_worker(&worker_id)?; candidates.push(RestoreCandidate { worker_ref, request, run_generation, previous_working_directory, config_bundle: None, }); } candidates }; for candidate in candidates { let (backend, workspace_scope) = { let state = self.lock()?; let workspace_scope = candidate.request.workspace_api.as_ref().and_then(|api| { state .workspace_owners .get(&api.workspace_id) .map(|server_id| RuntimeWorkspaceScope::new(&api.workspace_id, server_id)) }); (state.execution_backend.clone(), workspace_scope) }; let Some(backend) = backend else { return Ok(()); }; let request = WorkerExecutionRestoreRequest { worker_ref: candidate.worker_ref.clone(), run_generation: candidate.run_generation, request: candidate.request, workspace_scope, context: self.execution_context(candidate.worker_ref.clone()), previous_working_directory: candidate.previous_working_directory, working_directory: None, config_bundle: candidate.config_bundle, }; match backend.restore_worker(request) { WorkerExecutionSpawnResult::Connected { handle, worker_state, working_directory, } => self.commit_restored_worker_execution( &candidate.worker_ref, handle, worker_state, WorkerStatus::Idle, working_directory, )?, WorkerExecutionSpawnResult::Rejected(result) | WorkerExecutionSpawnResult::Errored(result) => { let mut state = self.lock()?; state.record_restore_failure(&candidate.worker_ref, result)?; } } } Ok(()) } fn commit_restored_worker_execution( &self, worker_ref: &WorkerRef, handle: WorkerExecutionHandle, worker_state: protocol::WorkerStateSnapshot, status: WorkerStatus, working_directory: Option, ) -> Result<(), RuntimeError> { let mut state = self.lock()?; state.ensure_worker_ref(worker_ref)?; { let worker = state.worker_mut(worker_ref)?; worker.execution_handle = Some(handle); worker.execution_bound = true; worker.status = status; let _ = worker.apply_worker_state(&worker_state); worker.restore_intent = restore_intent_for_status(worker.status); worker.working_directory = working_directory; } state.publish_worker_upsert(worker_ref.worker_id)?; state.persist_runtime_snapshot()?; state.persist_worker(&worker_ref.worker_id)?; Ok(()) } /// Bind the Backend registry identity once. Retention evidence fails closed /// until the Runtime host supplies this trusted configuration. pub fn bind_runtime_identity(&self, runtime_id: &str) -> Result<(), RuntimeError> { if runtime_id.trim().is_empty() || runtime_id.len() > 160 { return Err(RuntimeError::InvalidRequest( "Runtime identity must be non-empty and bounded".to_string(), )); } let mut state = self.lock()?; match state.runtime_identity.as_deref() { Some(current) if current == runtime_id => Ok(()), Some(_) => Err(RuntimeError::InvalidRequest( "Runtime identity is already bound".to_string(), )), None => { state.runtime_identity = Some(runtime_id.to_string()); Ok(()) } } } /// Read canonical aggregate facts needed by a Backend removal plan. #[cfg(feature = "fs-store")] pub fn worker_retention_inventory( &self, workspace_id: &str, worker_ref: &WorkerRef, ) -> Result { let state = self.lock()?; let runtime_id = state.runtime_identity.as_deref().ok_or_else(|| { RuntimeError::InvalidRequest( "Runtime identity is not bound for Worker retention".to_string(), ) })?; let worker = state.worker(worker_ref)?; if worker.workspace_id.as_deref() != Some(workspace_id) { return Err(RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }); } let store = state.fs_store().ok_or_else(|| { RuntimeError::InvalidRequest( "Worker retention archive authority requires an fs-backed Runtime".to_string(), ) })?; FsWorkerRetentionProvider::new(store.runtime_dir()).inventory( workspace_id, runtime_id, worker.worker_id, worker.run_generation, ) } /// Enumerate host-authoritative Runtime inventory for Backend orphan /// reconciliation. Runtime identity and Workspace scope are derived here, /// not accepted in a diagnostic payload. #[cfg(feature = "fs-store")] pub fn list_worker_retention_inventory( &self, workspace_id: &str, ) -> Result { let state = self.lock()?; let runtime_id = state.runtime_identity.as_deref().ok_or_else(|| { RuntimeError::InvalidRequest( "Runtime identity is not bound for Worker retention".to_string(), ) })?; let store = state.fs_store().ok_or_else(|| { RuntimeError::InvalidRequest( "Worker retention inventory requires an fs-backed Runtime".to_string(), ) })?; let provider = FsWorkerRetentionProvider::new(store.runtime_dir()); provider.snapshot(workspace_id, runtime_id) } /// Execute a Backend-resolved retention plan. Only stopped Workers are /// eligible. Provider receipt lookup happens before live lookup so exact /// retries converge after aggregate removal. #[cfg(feature = "fs-store")] pub fn execute_worker_retention( &self, request: &WorkerRetentionExecutionRequest, ) -> Result { let mut state = self.lock()?; let runtime_id = state.runtime_identity.clone().ok_or_else(|| { RuntimeError::InvalidRequest( "Runtime identity is not bound for Worker retention".to_string(), ) })?; if request.source_runtime_id != runtime_id { return Err(RuntimeError::InvalidRequest( "Worker retention Runtime identity mismatch".to_string(), )); } let store = state.fs_store().ok_or_else(|| { RuntimeError::InvalidRequest( "Worker retention execution requires an fs-backed Runtime".to_string(), ) })?; let provider = FsWorkerRetentionProvider::new(store.runtime_dir()); if let Some(completed) = provider.completed_for(request)? { state.workers.remove(&request.worker_id); state.persist_runtime_snapshot()?; return Ok(completed); } let Some(worker) = state.workers.get(&request.worker_id) else { // Recover a pending receipt after a crash between aggregate removal // and final receipt/Runtime catalog commit. return provider.recover_after_source_removal(request); }; if worker.workspace_id.as_deref() != Some(request.workspace_id.as_str()) { return Err(RuntimeError::WorkerNotFound { worker_id: request.worker_id, }); } if worker.status != WorkerStatus::Stopped { return Err(RuntimeError::InvalidRequest( "Worker retention requires a stopped Worker".to_string(), )); } if worker.run_generation != request.expected_run_generation { return Err(RuntimeError::InvalidRequest(format!( "Worker retention plan expected generation {}, current generation is {}", request.expected_run_generation, worker.run_generation ))); } let result = provider.execute(request)?; state.workers.remove(&request.worker_id); state.persist_runtime_snapshot()?; Ok(result) } fn worker_operation_lock(&self, worker_id: WorkerId) -> Result>, RuntimeError> { let mut operations = self .worker_operations .lock() .map_err(|_| RuntimeError::StatePoisoned)?; Ok(operations .entry(worker_id) .or_insert_with(|| Arc::new(Mutex::new(()))) .clone()) } fn lock(&self) -> Result, RuntimeError> { self.inner.lock().map_err(|_| RuntimeError::StatePoisoned) } } #[cfg_attr(not(feature = "fs-store"), allow(dead_code))] #[derive(Clone, Debug)] enum RuntimePersistence { Memory, #[cfg(feature = "fs-store")] Fs(FsRuntimeStore), } #[derive(Debug)] struct SubscriptionSink { selector: EventSubscriptionSelector, workspace_id: Option, sender: mpsc::Sender, lagged: Arc, } #[derive(Clone)] struct BackendResourceClientRef(Arc); impl std::fmt::Debug for BackendResourceClientRef { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { formatter.write_str("BackendResourceClientRef(..)") } } #[derive(Debug)] struct RuntimeState { display_name: Option, backend: RuntimeBackendKind, /// Backend-bound stable identity used for cross-boundary retention evidence. /// It is configured once by the Runtime host and never model input. runtime_identity: Option, #[cfg_attr(not(feature = "fs-store"), allow(dead_code))] persistence: RuntimePersistence, status: RuntimeStatus, execution_backend: Option, backend_resource_client: Option, workspace_backend_resource_clients: BTreeMap, #[cfg(feature = "fs-store")] next_diagnostic_id: u64, workers: BTreeMap, workspace_owners: BTreeMap, config_bundles: BTreeMap, workspace_config_latest: BTreeMap, workspace_config_fetch_gates: BTreeMap>>, diagnostics: Vec, subscription_revision: u64, worker_subject_revisions: BTreeMap, next_event_subscription_id: u64, subscriptions: BTreeMap, #[cfg(feature = "ws-server")] next_observation_sequence: u64, #[cfg(feature = "ws-server")] observation_events: VecDeque, #[cfg(feature = "ws-server")] observation_tx: broadcast::Sender, } impl RuntimeState { fn new(display_name: Option) -> Self { Self { display_name, backend: RuntimeBackendKind::Memory, runtime_identity: None, persistence: RuntimePersistence::Memory, status: RuntimeStatus::Running, execution_backend: None, backend_resource_client: None, workspace_backend_resource_clients: BTreeMap::new(), #[cfg(feature = "fs-store")] next_diagnostic_id: 1, workers: BTreeMap::new(), workspace_owners: BTreeMap::new(), config_bundles: BTreeMap::new(), workspace_config_latest: BTreeMap::new(), workspace_config_fetch_gates: BTreeMap::new(), diagnostics: Vec::new(), subscription_revision: 0, worker_subject_revisions: BTreeMap::new(), next_event_subscription_id: 1, subscriptions: BTreeMap::new(), #[cfg(feature = "ws-server")] next_observation_sequence: 1, #[cfg(feature = "ws-server")] observation_events: VecDeque::new(), #[cfg(feature = "ws-server")] observation_tx: broadcast::channel(256).0, } } #[cfg(feature = "fs-store")] fn new_fs_backed(display_name: Option, store: FsRuntimeStore) -> Self { Self { display_name, backend: RuntimeBackendKind::FsStore, runtime_identity: None, persistence: RuntimePersistence::Fs(store), status: RuntimeStatus::Running, execution_backend: None, backend_resource_client: None, workspace_backend_resource_clients: BTreeMap::new(), #[cfg(feature = "fs-store")] next_diagnostic_id: 1, workers: BTreeMap::new(), workspace_owners: BTreeMap::new(), config_bundles: BTreeMap::new(), workspace_config_latest: BTreeMap::new(), workspace_config_fetch_gates: BTreeMap::new(), diagnostics: Vec::new(), subscription_revision: 0, worker_subject_revisions: BTreeMap::new(), next_event_subscription_id: 1, subscriptions: BTreeMap::new(), #[cfg(feature = "ws-server")] next_observation_sequence: 1, #[cfg(feature = "ws-server")] observation_events: VecDeque::new(), #[cfg(feature = "ws-server")] observation_tx: broadcast::channel(256).0, } } #[cfg(feature = "fs-store")] fn from_persisted( persisted: PersistedRuntimeState, store: FsRuntimeStore, ) -> Result { let mut workers = BTreeMap::new(); let diagnostics = persisted.diagnostics; let next_diagnostic_id = persisted.next_diagnostic_id; for (worker_id, worker) in persisted.workers { let run_generation = worker.execution.last_run_generation; workers.insert( worker_id, WorkerRecord { worker_ref: worker.worker_ref, worker_id: worker.worker_id, status: worker.status, worker_state: None, workspace_id: worker.workspace_id, request: worker.request, run_generation, execution_bound: worker.execution.binding.is_some(), restore_intent: worker.execution.restore_intent, working_directory: worker.working_directory, execution_handle: None, internal_workers: BTreeMap::new(), }, ); } Ok(Self { display_name: persisted.display_name, backend: RuntimeBackendKind::FsStore, runtime_identity: None, persistence: RuntimePersistence::Fs(store), status: persisted.status, execution_backend: None, backend_resource_client: None, workspace_backend_resource_clients: BTreeMap::new(), next_diagnostic_id, workers, config_bundles: BTreeMap::new(), workspace_config_latest: BTreeMap::new(), workspace_config_fetch_gates: BTreeMap::new(), workspace_owners: persisted.workspace_owners, diagnostics, subscription_revision: 0, worker_subject_revisions: BTreeMap::new(), next_event_subscription_id: 1, subscriptions: BTreeMap::new(), #[cfg(feature = "ws-server")] next_observation_sequence: 1, #[cfg(feature = "ws-server")] observation_events: VecDeque::new(), #[cfg(feature = "ws-server")] observation_tx: broadcast::channel(256).0, }) } #[cfg(feature = "fs-store")] fn persisted_state(&self) -> PersistedRuntimeState { PersistedRuntimeState { display_name: self.display_name.clone(), status: self.status, next_diagnostic_id: self.next_diagnostic_id, workers: self .workers .iter() .map(|(worker_id, worker)| (worker_id.clone(), worker.persisted_record())) .collect(), workspace_owners: self.workspace_owners.clone(), diagnostics: self.diagnostics.clone(), } } #[cfg(feature = "fs-store")] fn fs_store(&self) -> Option<&FsRuntimeStore> { match &self.persistence { RuntimePersistence::Memory => None, RuntimePersistence::Fs(store) => Some(store), } } #[cfg(feature = "fs-store")] fn persist_runtime_snapshot(&self) -> Result<(), RuntimeError> { if let Some(store) = self.fs_store() { store.write_runtime_snapshot(&self.persisted_state())?; } Ok(()) } #[cfg(feature = "fs-store")] fn persist_worker(&self, worker_id: &WorkerId) -> Result<(), RuntimeError> { if let Some(store) = self.fs_store() { let worker = self.workers .get(worker_id) .ok_or_else(|| RuntimeError::WorkerNotFound { worker_id: *worker_id, })?; store.write_worker_snapshot(&worker.persisted_record())?; } Ok(()) } #[cfg(feature = "fs-store")] fn delete_worker_snapshot(&self, worker_id: &WorkerId) -> Result<(), RuntimeError> { if let Some(store) = self.fs_store() { store.delete_worker_snapshot(worker_id)?; } Ok(()) } #[cfg(feature = "fs-store")] fn persist_workers(&self) -> Result<(), RuntimeError> { if self.fs_store().is_some() { for worker_id in self.workers.keys() { self.persist_worker(worker_id)?; } } Ok(()) } #[cfg(not(feature = "fs-store"))] fn persist_runtime_snapshot(&self) -> Result<(), RuntimeError> { Ok(()) } #[cfg(not(feature = "fs-store"))] fn persist_worker(&self, _worker_id: &WorkerId) -> Result<(), RuntimeError> { Ok(()) } #[cfg(not(feature = "fs-store"))] fn delete_worker_snapshot(&self, _worker_id: &WorkerId) -> Result<(), RuntimeError> { Ok(()) } #[cfg(not(feature = "fs-store"))] fn persist_workers(&self) -> Result<(), RuntimeError> { Ok(()) } fn ensure_running(&self) -> Result<(), RuntimeError> { if self.status == RuntimeStatus::Stopped { Err(RuntimeError::RuntimeStopped) } else { Ok(()) } } fn check_config_bundle_ref( &self, reference: &ConfigBundleRef, ) -> Result { validate_config_bundle_ref(reference)?; let bundle = self.config_bundles.get(&reference.id).ok_or_else(|| { RuntimeError::ConfigBundleMissing { bundle_id: reference.id.clone(), } })?; if bundle.metadata.digest != reference.digest { return Err(RuntimeError::ConfigBundleDigestMismatch { bundle_id: reference.id.clone(), expected_digest: reference.digest.clone(), actual_digest: bundle.metadata.digest.clone(), }); } Ok(ConfigBundleAvailability { reference: reference.clone(), summary: bundle.summary(), }) } fn resolve_config_bundle_ref( &self, reference: Option<&ConfigBundleRef>, ) -> Result, RuntimeError> { let Some(reference) = reference else { return Ok(None); }; self.check_config_bundle_ref(reference)?; Ok(self.config_bundles.get(&reference.id).cloned()) } fn validate_worker_config_boundary( &self, _request: &CreateWorkerRequest, ) -> Result<(), RuntimeError> { Ok(()) } fn ensure_worker_ref(&self, worker_ref: &WorkerRef) -> Result<(), RuntimeError> { if !self.workers.contains_key(&worker_ref.worker_id) { return Err(RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }); } Ok(()) } fn ensure_workspace_owner( &mut self, scope: &RuntimeWorkspaceScope, claim_if_missing: bool, ) -> Result { if scope.workspace_id.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "Runtime auth workspace_id must not be empty".to_string(), )); } if scope.server_id.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "Runtime auth server_id must not be empty".to_string(), )); } match self.workspace_owners.get(&scope.workspace_id) { Some(owner_server_id) if owner_server_id == &scope.server_id => Ok(true), Some(owner_server_id) => Err(RuntimeError::WorkspaceOwnerMismatch { workspace_id: scope.workspace_id.clone(), owner_server_id: owner_server_id.clone(), requester_server_id: scope.server_id.clone(), }), None if claim_if_missing => { self.workspace_owners .insert(scope.workspace_id.clone(), scope.server_id.clone()); Ok(true) } None => Ok(false), } } fn ensure_workspace_owner_for_existing_workers( &mut self, scope: &RuntimeWorkspaceScope, ) -> Result { let has_workspace_worker = self .workers .values() .any(|worker| worker.belongs_to_workspace(&scope.workspace_id)); self.ensure_workspace_owner(scope, has_workspace_worker) } fn forget_workspace_owner_if_unused(&mut self, workspace_id: &str) -> bool { if self .workers .values() .any(|worker| worker.belongs_to_workspace(workspace_id)) { return false; } self.workspace_owners.remove(workspace_id).is_some() } fn worker(&self, worker_ref: &WorkerRef) -> Result<&WorkerRecord, RuntimeError> { self.ensure_worker_ref(worker_ref)?; self.workers .get(&worker_ref.worker_id) .ok_or_else(|| RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }) } fn worker_mut(&mut self, worker_ref: &WorkerRef) -> Result<&mut WorkerRecord, RuntimeError> { self.ensure_worker_ref(worker_ref)?; self.workers .get_mut(&worker_ref.worker_id) .ok_or_else(|| RuntimeError::WorkerNotFound { worker_id: worker_ref.worker_id, }) } fn subscription_snapshot( &self, scope: Option<&RuntimeWorkspaceScope>, selector: &EventSubscriptionSelector, ) -> Result { let workers = match selector { EventSubscriptionSelector::RuntimeWorkers => self .workers .values() .filter(|worker| { scope.is_none_or(|scope| worker.belongs_to_workspace(&scope.workspace_id)) }) .map(|worker| self.subscription_worker(worker)) .collect::, _>>()?, EventSubscriptionSelector::WorkerLifecycle { worker_ids } => { let mut selected = Vec::with_capacity(worker_ids.as_slice().len()); for worker_id in worker_ids.as_slice() { let runtime_worker_id = WorkerId::parse(worker_id.as_str()).ok_or_else(|| { RuntimeError::InvalidRequest(format!( "worker_lifecycle selector contains invalid Runtime Worker id {worker_id}" )) })?; let worker = self.workers.get(&runtime_worker_id).ok_or( RuntimeError::WorkerNotFound { worker_id: runtime_worker_id, }, )?; if scope.is_some_and(|scope| !worker.belongs_to_workspace(&scope.workspace_id)) { return Err(RuntimeError::WorkerNotFound { worker_id: runtime_worker_id, }); } selected.push(self.subscription_worker(worker)?); } selected } EventSubscriptionSelector::WorkerProtocol { .. } | EventSubscriptionSelector::WorkspaceWorkers | EventSubscriptionSelector::WorkspaceWorkdirs => { return Err(RuntimeError::InvalidRequest( "selector is not produced by the Runtime lifecycle subscription".to_string(), )); } }; Ok(SubscriptionSnapshot::Workers { workers }) } fn subscription_worker( &self, worker: &WorkerRecord, ) -> Result { let worker_id = SubscriptionWorkerId::new(worker.worker_id.to_string()) .map_err(subscription_validation_error)?; let repository_id = worker .working_directory .as_ref() .map(|working_directory| working_directory.summary.repository_id.clone()); let working_directory_id = worker .working_directory .as_ref() .map(|working_directory| { SubscriptionWorkdirId::new(working_directory.summary.working_directory_id.clone()) .map_err(subscription_validation_error) }) .transpose()?; let profile = match &worker.request.profile { ProfileSelector::Builtin(name) | ProfileSelector::Named(name) => Some(name.clone()), }; Ok(SubscriptionWorker { worker_id, runtime_id: None, resource_key: None, subject_revision: self .worker_subject_revisions .get(&worker.worker_id) .copied() .unwrap_or(0), worker_state: worker.worker_state.clone(), state: subscription_worker_state(worker.status), has_running_internal_workers: worker .internal_workers .values() .any(|worker| worker.status == protocol::WorkerStatus::Running), workspace_id: worker.workspace_id.clone(), display_name: worker.request.display_name.clone(), profile, repository_id, repository_key: None, working_directory_id, }) } fn publish_worker_upsert(&mut self, worker_id: WorkerId) -> Result<(), RuntimeError> { self.subscription_revision = self.subscription_revision.saturating_add(1); let subject_revision = { let revision = self.worker_subject_revisions.entry(worker_id).or_insert(0); *revision = revision.saturating_add(1); *revision }; let worker = self .workers .get(&worker_id) .ok_or(RuntimeError::WorkerNotFound { worker_id })?; let workspace_id = worker.workspace_id.clone(); let mut projected = self.subscription_worker(worker)?; projected.subject_revision = subject_revision; let projected_worker_id = projected.worker_id.clone(); self.deliver_worker_subscription_update( &projected_worker_id, workspace_id.as_deref(), RuntimeSubscriptionUpdate { subject_revision, payload: SubscriptionEventPayload::WorkerUpserted { worker: projected }, }, ); Ok(()) } fn publish_worker_removed( &mut self, worker_id: WorkerId, workspace_id: Option<&str>, ) -> Result<(), RuntimeError> { self.subscription_revision = self.subscription_revision.saturating_add(1); let subject_revision = { let revision = self.worker_subject_revisions.entry(worker_id).or_insert(0); *revision = revision.saturating_add(1); *revision }; let worker_id = SubscriptionWorkerId::new(worker_id.to_string()) .map_err(subscription_validation_error)?; self.deliver_worker_subscription_update( &worker_id, workspace_id, RuntimeSubscriptionUpdate { subject_revision, payload: SubscriptionEventPayload::WorkerRemoved { worker_id: worker_id.clone(), runtime_id: None, }, }, ); Ok(()) } fn deliver_worker_subscription_update( &mut self, worker_id: &SubscriptionWorkerId, workspace_id: Option<&str>, update: RuntimeSubscriptionUpdate, ) { let mut closed = Vec::new(); for (subscription_id, sink) in &self.subscriptions { if sink .workspace_id .as_deref() .is_some_and(|expected| workspace_id != Some(expected)) { continue; } let selected = match &sink.selector { EventSubscriptionSelector::RuntimeWorkers => true, EventSubscriptionSelector::WorkerLifecycle { worker_ids } => { worker_ids.contains(worker_id) } EventSubscriptionSelector::WorkerProtocol { .. } | EventSubscriptionSelector::WorkspaceWorkers | EventSubscriptionSelector::WorkspaceWorkdirs => false, }; if !selected { continue; } match sink.sender.try_send(update.clone()) { Ok(()) => {} Err(mpsc::error::TrySendError::Full(_)) => { sink.lagged.store(true, Ordering::Release); closed.push(*subscription_id); } Err(mpsc::error::TrySendError::Closed(_)) => closed.push(*subscription_id), } } for subscription_id in closed { self.subscriptions.remove(&subscription_id); } } fn primary_worker_id_for_workdir(&self, working_directory_id: &str) -> Option { self.workers.values().find_map(|worker| { if worker .working_directory .as_ref() .is_some_and(|binding| binding.summary.working_directory_id == working_directory_id) || requested_primary_workdir_id(&worker.request) == Some(working_directory_id) { Some(worker.worker_id) } else { None } }) } #[cfg(feature = "fs-store")] fn record_restore_failure( &mut self, worker_ref: &WorkerRef, result: WorkerExecutionResult, ) -> Result<(), RuntimeError> { let message = result .message .clone() .unwrap_or_else(|| "worker execution restore failed".to_string()); let diagnostic_id = self.next_diagnostic_id; self.next_diagnostic_id += 1; let workspace_id = self .workers .get(&worker_ref.worker_id) .and_then(|worker| worker.workspace_id.as_deref()) .unwrap_or(""); eprintln!( "yoi-runtime: Worker execution restore failed: diagnostic_id={diagnostic_id} runtime_id={} workspace_id={workspace_id} worker_id={} operation={:?} outcome={:?}: {message}", self.runtime_identity.as_deref().unwrap_or(""), worker_ref.worker_id, result.operation, result.outcome, ); self.diagnostics.push(RuntimeDiagnostic { id: diagnostic_id, severity: DiagnosticSeverity::Warning, code: "worker_execution_restore_failed".to_string(), message: format!( "worker {} execution restore failed: {message}", worker_ref.worker_id ), worker_ref: Some(worker_ref.clone()), }); let worker = self.worker_mut(worker_ref)?; worker.execution_handle = None; worker.status = WorkerStatus::Stopped; worker.restore_intent = WorkerRestoreIntent::Explicit; worker.internal_workers.clear(); self.publish_worker_upsert(worker_ref.worker_id)?; self.persist_runtime_snapshot()?; self.persist_worker(&worker_ref.worker_id)?; Ok(()) } #[cfg(feature = "ws-server")] fn validate_worker_observation_cursor( &self, worker_ref: &WorkerRef, cursor: WorkerObservationCursor, ) -> Result<(), RuntimeError> { if let Some(first) = self .observation_events .iter() .find(|event| &event.worker_ref == worker_ref) { if cursor.sequence != 0 && cursor.sequence < first.sequence { return Err(RuntimeError::InvalidRequest(format!( "worker observation cursor {} is expired for worker {}", cursor.encode(), worker_ref.worker_id ))); } } if cursor.sequence >= self.next_observation_sequence { return Err(RuntimeError::InvalidRequest(format!( "worker observation cursor {} is unknown for worker {}", cursor.encode(), worker_ref.worker_id ))); } Ok(()) } #[cfg(feature = "ws-server")] fn push_worker_observation_event( &mut self, worker_ref: WorkerRef, payload: protocol::Event, ) -> WorkerObservationEvent { const MAX_OBSERVATION_BACKLOG: usize = 1024; let sequence = self.next_observation_sequence; self.next_observation_sequence += 1; let event = WorkerObservationEvent::new(sequence, worker_ref, payload); self.observation_events.push_back(event.clone()); while self.observation_events.len() > MAX_OBSERVATION_BACKLOG { self.observation_events.pop_front(); } let _ = self.observation_tx.send(event.clone()); event } fn internal_worker_snapshot_statuses( statuses: &mut BTreeMap, snapshot: &protocol::InternalWorkerSnapshot, ) { statuses.insert( snapshot.worker.session_id.clone(), InternalWorkerActivity { status: snapshot.status, parent_session_id: snapshot.worker.parent_session_id.clone(), }, ); for child in &snapshot.internal_workers { Self::internal_worker_snapshot_statuses(statuses, child); } } fn remove_internal_worker_subtree( statuses: &mut BTreeMap, root_session_id: &str, ) { let mut removed = vec![root_session_id.to_string()]; while let Some(parent_session_id) = removed.pop() { let children = statuses .iter() .filter_map(|(session_id, worker)| { (worker.parent_session_id.as_deref() == Some(parent_session_id.as_str())) .then(|| session_id.clone()) }) .collect::>(); statuses.remove(&parent_session_id); removed.extend(children); } } fn project_internal_worker_event( statuses: &mut BTreeMap, worker: &protocol::InternalWorkerRef, event: &protocol::Event, ) { match event { protocol::Event::Snapshot { state, internal_workers, .. } => { Self::remove_internal_worker_subtree(statuses, &worker.session_id); statuses.insert( worker.session_id.clone(), InternalWorkerActivity { status: state.catalog_status(), parent_session_id: worker.parent_session_id.clone(), }, ); for child in internal_workers { Self::internal_worker_snapshot_statuses(statuses, child); } } protocol::Event::InternalWorker { worker: nested_worker, event, .. } => Self::project_internal_worker_event(statuses, nested_worker, event), protocol::Event::WorkerState { snapshot } | protocol::Event::CommandAcknowledged { acknowledgement: protocol::WorkerCommandAcknowledgement { state: snapshot, .. }, } => { statuses.insert( worker.session_id.clone(), InternalWorkerActivity { status: snapshot.catalog_status(), parent_session_id: worker.parent_session_id.clone(), }, ); } _ => {} } } fn update_internal_worker_activity( statuses: &mut BTreeMap, event: &protocol::Event, ) -> bool { let was_running = statuses .values() .any(|worker| worker.status == protocol::WorkerStatus::Running); match event { protocol::Event::Snapshot { internal_workers, .. } => { statuses.clear(); for child in internal_workers { Self::internal_worker_snapshot_statuses(statuses, child); } } protocol::Event::InternalWorker { worker, event, .. } => { Self::project_internal_worker_event(statuses, worker, event); } _ => {} } let is_running = statuses .values() .any(|worker| worker.status == protocol::WorkerStatus::Running); was_running != is_running } fn project_internal_worker_activity( &mut self, worker_ref: &WorkerRef, event: &protocol::Event, ) -> bool { let Some(worker) = self.workers.get_mut(&worker_ref.worker_id) else { return false; }; Self::update_internal_worker_activity(&mut worker.internal_workers, event) } fn project_protocol_event_to_worker_state( &mut self, worker_ref: &WorkerRef, event: &protocol::Event, ) -> bool { let Some(worker) = self.workers.get_mut(&worker_ref.worker_id) else { return false; }; let incoming = match event { protocol::Event::WorkerState { snapshot } | protocol::Event::Snapshot { state: snapshot, .. } | protocol::Event::CommandAcknowledged { acknowledgement: protocol::WorkerCommandAcknowledgement { state: snapshot, .. }, } => snapshot, _ => return false, }; match worker.apply_worker_state(incoming) { Ok(protocol::WorkerStateSnapshotApply::Applied) => true, Ok( protocol::WorkerStateSnapshotApply::Duplicate | protocol::WorkerStateSnapshotApply::Stale, ) | Err(_) => false, } } } #[derive(Debug, Clone)] struct InternalWorkerActivity { status: protocol::WorkerStatus, parent_session_id: Option, } #[derive(Debug)] struct WorkerRecord { worker_ref: WorkerRef, worker_id: WorkerId, status: WorkerStatus, worker_state: Option, workspace_id: Option, request: CreateWorkerRequest, run_generation: u64, execution_bound: bool, restore_intent: WorkerRestoreIntent, working_directory: Option, execution_handle: Option, internal_workers: BTreeMap, } impl WorkerRecord { fn apply_worker_state( &mut self, incoming: &protocol::WorkerStateSnapshot, ) -> Result { match self.worker_state.as_mut() { Some(current) => protocol::apply_worker_state_snapshot(current, incoming), None => { self.worker_state = Some(incoming.clone()); Ok(protocol::WorkerStateSnapshotApply::Applied) } } } fn belongs_to_workspace(&self, workspace_id: &str) -> bool { self.workspace_id.as_deref() == Some(workspace_id) } fn summary(&self) -> WorkerSummary { WorkerSummary { worker_ref: self.worker_ref.clone(), worker_id: self.worker_id, status: self.status, worker_state: self.worker_state.clone(), workspace_id: self.workspace_id.clone(), working_directory: self.working_directory.clone(), profile: self.request.profile.clone(), display_name: self.request.display_name.clone(), profile_source: self.request.profile_source.reference(), config_bundle: self.request.config_bundle.clone(), } } fn detail(&self) -> WorkerDetail { WorkerDetail { worker_ref: self.worker_ref.clone(), worker_id: self.worker_id, status: self.status, worker_state: self.worker_state.clone(), workspace_id: self.workspace_id.clone(), working_directory: self.working_directory.clone(), profile: self.request.profile.clone(), display_name: self.request.display_name.clone(), profile_source: self.request.profile_source.reference(), config_bundle: self.request.config_bundle.clone(), } } #[cfg(feature = "fs-store")] fn persisted_record(&self) -> PersistedWorkerRecord { PersistedWorkerRecord { worker_ref: self.worker_ref.clone(), worker_id: self.worker_id.clone(), request: self.request.clone(), status: self.status, execution: PersistedWorkerExecution { last_run_generation: self.run_generation, binding: self .execution_bound .then_some(PersistedWorkerExecutionBinding { run_generation: self.run_generation, }), restore_intent: self.restore_intent, }, workspace_id: self.workspace_id.clone(), working_directory: self.working_directory.clone(), } } } fn restore_intent_for_status(status: WorkerStatus) -> WorkerRestoreIntent { if status.is_active() { WorkerRestoreIntent::Automatic } else { WorkerRestoreIntent::Explicit } } fn repository_resource_error(error: BackendResourceError) -> RuntimeError { let (code, message) = match error { BackendResourceError::Expired => ( "repository_access_credential_expired", "Repository access credential lease expired", ), BackendResourceError::Unauthorized { .. } => ( "repository_access_credential_unauthorized", "Repository access credential lease was rejected", ), BackendResourceError::MissingResource => ( "repository_access_credential_unavailable", "Repository access credential lease is unavailable or already consumed", ), BackendResourceError::Timeout => ( "repository_access_resource_fetch_timeout", "Timed out while fetching Repository SSH access from Workspace Backend", ), BackendResourceError::Transport { .. } => ( "repository_access_provider_unavailable", "Repository access credential provider is unavailable", ), BackendResourceError::UnsupportedKind | BackendResourceError::Oversized { .. } | BackendResourceError::DigestMismatch { .. } | BackendResourceError::ContentTypeMismatch { .. } | BackendResourceError::InvalidResponse { .. } => ( "repository_access_credential_invalid", "Repository access credential response is invalid", ), }; RuntimeError::WorkingDirectory( crate::working_directory::WorkingDirectoryDiagnostic::rejected(code, message), ) } fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest { let mut durable = request.clone(); if let Some(working_directory) = durable.working_directory_request.as_mut() && let Some(materialization) = working_directory.materialization.as_mut() { materialization.ssh = None; } durable } fn requested_primary_workdir_id(request: &CreateWorkerRequest) -> Option<&str> { request .working_directory .as_ref() .map(|claim| claim.working_directory_id.as_str()) .or_else(|| { request .working_directory_request .as_ref() .and_then(|request| request.backend_workdir_id.as_deref()) }) } fn validate_create_worker_request(request: &CreateWorkerRequest) -> Result<(), RuntimeError> { if request.create_fingerprint.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "create_fingerprint must not be empty".to_string(), )); } match &request.profile_source { crate::catalog::ProfileSourceArchiveSource::Embedded { archive } => { archive.verify().map_err(|err| { RuntimeError::InvalidRequest(format!("profile_source archive is invalid: {err}")) })?; } crate::catalog::ProfileSourceArchiveSource::WorkspaceConfig { archive } => { if archive.digest.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "profile_source.archive.digest must not be empty".to_string(), )); } } } if let Some(input) = &request.initial_input { if input.kind != WorkerInputKind::User { return Err(RuntimeError::InvalidInitialInputKind { kind: format!("{:?}", input.kind), }); } let has_segments = input .segments .as_ref() .is_some_and(|segments| !segments.is_empty()); if input.content.trim().is_empty() && !has_segments { return Err(RuntimeError::InvalidRequest( "initial_input.content must not be empty".to_string(), )); } } Ok(()) } fn worker_execution_failure_code(result: &WorkerExecutionResult) -> &'static str { let message = result.message.as_deref().unwrap_or_default(); if message.contains("timed out waiting for durable Worker Submit acceptance") { "input_commit_timeout" } else if message.contains("worker adapter task did not complete within") { "adapter_task_timeout" } else if message.contains("failed to send Worker method") { "worker_method_send_failed" } else if message.contains("worker rejected Submit") { "worker_submit_rejected" } else if message.contains("before durable acceptance") { "worker_failed_before_input_commit" } else if message.contains("event stream closed") { "worker_event_stream_closed" } else { "worker_execution_failed" } } fn runtime_worker_create_failure_fields( error: &RuntimeError, ) -> (&'static str, Option, Option) { match error { RuntimeError::RuntimeStopped => ("runtime_stopped", None, None), RuntimeError::InvalidInitialInputKind { .. } => ("invalid_initial_input_kind", None, None), RuntimeError::WorkerNotFound { .. } => ("worker_not_found", None, None), RuntimeError::WorkerExecutionUnavailable { .. } => { ("worker_execution_unavailable", None, None) } RuntimeError::ExecutionBackendUnavailable { .. } => { ("execution_backend_unavailable", None, None) } RuntimeError::WorkerExecutionRejected { operation, outcome, .. } => ( "worker_execution_rejected", Some(format!("{operation:?}")), Some(format!("{outcome:?}")), ), RuntimeError::LimitTooLarge { .. } => ("limit_too_large", None, None), RuntimeError::InvalidRequest(_) => ("invalid_request", None, None), RuntimeError::WorkspaceOwnerMismatch { .. } => ("workspace_owner_mismatch", None, None), RuntimeError::WorkingDirectory(_) => ("working_directory", None, None), RuntimeError::ConfigBundleMissing { .. } => ("config_bundle_missing", None, None), RuntimeError::ConfigBundleDigestMismatch { .. } => { ("config_bundle_digest_mismatch", None, None) } RuntimeError::InvalidProfileSelector { .. } => ("invalid_profile_selector", None, None), RuntimeError::UnsupportedConfigDeclaration { .. } => { ("unsupported_config_declaration", None, None) } RuntimeError::StoreIo { .. } => ("store_io", None, None), RuntimeError::StoreMissing { .. } => ("store_missing", None, None), RuntimeError::StoreCorrupt { .. } => ("store_corrupt", None, None), RuntimeError::StatePoisoned => ("state_poisoned", None, None), } } fn write_runtime_worker_create_failure( worker_id: WorkerId, workspace_id: Option<&str>, error: &RuntimeError, elapsed: std::time::Duration, ) { let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(error); let execution_failure_code = match error { RuntimeError::WorkerExecutionRejected { result, .. } => { worker_execution_failure_code(result) } _ => "", }; tracing::error!( target: "yoi::worker_create", event = "worker_create_failed", component = "runtime", workspace_id = workspace_id.unwrap_or(""), worker_id = %worker_id, error_kind, operation = operation.as_deref().unwrap_or(""), outcome = outcome.as_deref().unwrap_or(""), execution_failure_code, duration_ms = u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX), "Worker creation failed" ); } fn validate_create_workspace_scope( request: &CreateWorkerRequest, workspace_id: Option<&str>, ) -> Result<(), RuntimeError> { let Some(workspace_id) = workspace_id else { return Ok(()); }; if workspace_id.trim().is_empty() { return Err(RuntimeError::InvalidRequest( "Runtime auth workspace_id must not be empty".to_string(), )); } if let Some(request_workspace_id) = request .workspace_api .as_ref() .map(|workspace_api| workspace_api.workspace_id.as_str()) { if request_workspace_id != workspace_id { return Err(RuntimeError::InvalidRequest(format!( "request workspace_id {request_workspace_id} does not match Runtime auth workspace_id {workspace_id}" ))); } } let snapshot = request.memory_settings.as_ref().ok_or_else(|| { RuntimeError::InvalidRequest( "Workspace-scoped Worker create requires a bound Memory settings snapshot".to_string(), ) })?; if snapshot.workspace_id != workspace_id { return Err(RuntimeError::InvalidRequest(format!( "Memory settings workspace_id {} does not match Runtime auth workspace_id {workspace_id}", snapshot.workspace_id ))); } if snapshot.settings_revision == 0 { return Err(RuntimeError::InvalidRequest( "Memory settings revision must be at least 1".to_string(), )); } if !manifest::is_normalized_workspace_memory_language(&snapshot.language) { return Err(RuntimeError::InvalidRequest( "Memory settings language must be a normalized bounded UTF-8 value".to_string(), )); } Ok(()) } fn validate_worker_input(input: &WorkerInput) -> Result<(), RuntimeError> { let has_segments = input .segments .as_ref() .is_some_and(|segments| !segments.is_empty()); if !input.kind.is_empty_content_allowed() && input.content.trim().is_empty() && !has_segments { return Err(RuntimeError::InvalidRequest( "worker input content must not be empty".to_string(), )); } if input .segments .as_ref() .is_some_and(|segments| segments.is_empty()) { return Err(RuntimeError::InvalidRequest( "worker input segments must not be empty".to_string(), )); } Ok(()) } #[cfg(feature = "ws-server")] fn input_protocol_event(input: &WorkerInput) -> Option { match input.kind { // Submit is projected only after the Worker commits UserInput. Queued // payloads must never become model- or client-visible history early. WorkerInputKind::User | WorkerInputKind::Notify => None, WorkerInputKind::Compact | WorkerInputKind::ListRewindTargets | WorkerInputKind::RegisterPeer => Some(protocol::Event::SystemItem { item: serde_json::json!({ "kind": "embedded_worker_command_input", "command": input.kind, "content": input.content.clone(), }), }), } } fn subscription_validation_error(error: SubscriptionValidationError) -> RuntimeError { RuntimeError::InvalidRequest(format!("invalid event subscription: {error}")) } fn subscription_worker_state(status: WorkerStatus) -> SubscriptionWorkerState { match status { WorkerStatus::Idle => SubscriptionWorkerState::Idle, WorkerStatus::Running => SubscriptionWorkerState::Running, WorkerStatus::Paused => SubscriptionWorkerState::Paused, WorkerStatus::Stopped => SubscriptionWorkerState::Stopped, } } #[cfg(test)] mod tests { use super::*; use crate::catalog::{ ConfigBundleRef, MaterializerKind, ProfileSelector, RepositoryMaterializationContext, RepositorySshCredentialCandidate, RepositorySshMaterializationAccess, SensitiveString, WorkingDirectoryClaim, WorkingDirectoryRepository, WorkingDirectoryRequest, WorkspaceApiRef, }; use crate::config_bundle::{ ConfigBundle, ConfigBundleMetadata, ConfigBundleProvenance, ConfigDeclaration, ConfigDeclarationKind, ConfigProfileDescriptor, }; use crate::execution::{ WorkerExecutionBackend, WorkerExecutionContext, WorkerExecutionHandle, WorkerExecutionRestoreRequest, }; use crate::working_directory::WorkingDirectoryDiagnostic; use async_trait::async_trait; use std::collections::BTreeMap; #[cfg(feature = "fs-store")] use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; #[test] fn worker_create_failure_fields_exclude_raw_error_messages() { let (error_kind, operation, outcome) = runtime_worker_create_failure_fields( &RuntimeError::InvalidRequest("private-token".to_string()), ); assert_eq!(error_kind, "invalid_request"); assert_eq!(operation, None); assert_eq!(outcome, None); let timeout = WorkerExecutionResult::errored( WorkerExecutionOperation::Input, "timed out waiting for durable Worker Submit acceptance; private-token", ); assert_eq!( worker_execution_failure_code(&timeout), "input_commit_timeout" ); } fn test_command() -> protocol::WorkerCommandEnvelope { protocol::WorkerCommandEnvelope { command_id: 1, expected_execution_generation: 1, expected_worker_state_revision: 0, } } #[test] fn repository_resource_failures_keep_typed_credential_diagnostics() { let cases = [ ( BackendResourceError::Expired, "repository_access_credential_expired", ), ( BackendResourceError::MissingResource, "repository_access_credential_unavailable", ), ( BackendResourceError::Timeout, "repository_access_resource_fetch_timeout", ), ( BackendResourceError::Unauthorized { message: "denied".to_string(), }, "repository_access_credential_unauthorized", ), ]; for (error, expected_code) in cases { let RuntimeError::WorkingDirectory(diagnostic) = repository_resource_error(error) else { panic!("Repository resource failure lost its typed diagnostic") }; assert_eq!(diagnostic.code, expected_code); } } fn internal_worker_ref( session_id: &str, parent_session_id: Option<&str>, ) -> protocol::InternalWorkerRef { protocol::InternalWorkerRef { session_id: session_id.to_string(), parent_session_id: parent_session_id.map(str::to_string), name: session_id.to_string(), kind: protocol::InternalWorkerKind::SubWorker, } } fn internal_worker_status_event( worker: protocol::InternalWorkerRef, status: protocol::WorkerStatus, ) -> protocol::Event { protocol::Event::InternalWorker { worker, revision: 1, event: Box::new(protocol::Event::WorkerState { snapshot: status.into(), }), } } #[test] fn internal_worker_activity_tracks_running_children_independently() { let mut activity = BTreeMap::new(); assert!(RuntimeState::update_internal_worker_activity( &mut activity, &internal_worker_status_event( internal_worker_ref("child-a", None), protocol::WorkerStatus::Running, ), )); assert!(!RuntimeState::update_internal_worker_activity( &mut activity, &internal_worker_status_event( internal_worker_ref("child-b", None), protocol::WorkerStatus::Running, ), )); assert!(!RuntimeState::update_internal_worker_activity( &mut activity, &internal_worker_status_event( internal_worker_ref("child-a", None), protocol::WorkerStatus::Idle, ), )); assert!(RuntimeState::update_internal_worker_activity( &mut activity, &internal_worker_status_event( internal_worker_ref("child-b", None), protocol::WorkerStatus::Stopped, ), )); } #[test] fn nested_internal_worker_activity_reaches_the_parent_projection() { let mut activity = BTreeMap::new(); let direct_child = internal_worker_ref("child", None); let nested_running = protocol::Event::InternalWorker { worker: direct_child.clone(), revision: 1, event: Box::new(internal_worker_status_event( internal_worker_ref("grandchild", Some("child")), protocol::WorkerStatus::Running, )), }; assert!(RuntimeState::update_internal_worker_activity( &mut activity, &nested_running, )); let nested_idle = protocol::Event::InternalWorker { worker: direct_child, revision: 2, event: Box::new(internal_worker_status_event( internal_worker_ref("grandchild", Some("child")), protocol::WorkerStatus::Idle, )), }; assert!(RuntimeState::update_internal_worker_activity( &mut activity, &nested_idle, )); } #[test] fn parent_snapshot_replaces_stale_internal_worker_activity() { let mut activity = BTreeMap::new(); RuntimeState::update_internal_worker_activity( &mut activity, &internal_worker_status_event( internal_worker_ref("child-a", None), protocol::WorkerStatus::Running, ), ); let snapshot = protocol::Event::Snapshot { session: protocol::SessionSnapshot { pending_submissions: protocol::PendingSubmissionsSnapshot::default(), entries: Vec::new(), }, greeting: protocol::Greeting { worker_name: "parent".to_string(), cwd: "/tmp".to_string(), provider: "test".to_string(), model: "test".to_string(), scope_summary: String::new(), tools: Vec::new(), context_window: 0, context_tokens: 0, }, state: protocol::WorkerStatus::Idle.into(), in_flight: protocol::InFlightSnapshot::default(), internal_workers: Vec::new(), }; assert!(RuntimeState::update_internal_worker_activity( &mut activity, &snapshot, )); assert!(activity.is_empty()); } #[test] fn runtime_identity_binding_is_immutable_and_host_owned() { let runtime = Runtime::new_memory(); runtime.bind_runtime_identity("runtime-a").unwrap(); runtime.bind_runtime_identity("runtime-a").unwrap(); assert!(runtime.bind_runtime_identity("runtime-b").is_err()); assert_eq!( runtime.lock().unwrap().runtime_identity.as_deref(), Some("runtime-a") ); } #[test] fn typed_segments_allow_empty_flat_content() { let input = WorkerInput { kind: WorkerInputKind::User, content: String::new(), submission_request_id: None, segments: Some(vec![protocol::Segment::Flow { selector: "builtin:coder-review".to_string(), }]), }; assert!(validate_worker_input(&input).is_ok()); } #[test] fn empty_user_input_without_segments_is_rejected() { let input = WorkerInput { kind: WorkerInputKind::User, content: String::new(), submission_request_id: None, segments: Some(Vec::new()), }; assert!(matches!( validate_worker_input(&input), Err(RuntimeError::InvalidRequest(_)) )); } #[test] fn typed_flow_segments_allow_empty_initial_flat_content() { let mut request = task_request("flow"); request.initial_input = Some(WorkerInput { kind: WorkerInputKind::User, content: String::new(), submission_request_id: None, segments: Some(vec![protocol::Segment::Flow { selector: "builtin:coder-review".to_string(), }]), }); assert!(validate_create_worker_request(&request).is_ok()); } fn test_profile_source_archive() -> crate::profile_archive::ProfileSourceArchive { crate::profile_archive::ProfileSourceArchive::build( crate::profile_archive::ProfileSourceArchiveInput { id: "test-profile-source".to_string(), entrypoints: BTreeMap::from([( "builtin:coder".to_string(), "profiles/coder.dcdl".to_string(), )]), imports: BTreeMap::new(), sources: BTreeMap::from([("profiles/coder.dcdl".to_string(), "{}".to_string())]), }, ) .unwrap() } fn task_request(_objective: &str) -> CreateWorkerRequest { let profile = ProfileSelector::Builtin("builtin:coder".to_string()); let bundle = test_bundle_for_profile(profile.clone()); CreateWorkerRequest { worker_id: WorkerId::now_v7(), create_fingerprint: "test-create".to_string(), profile, display_name: None, profile_source: crate::catalog::ProfileSourceArchiveSource::Embedded { archive: test_profile_source_archive(), }, config_bundle: Some(ConfigBundleRef { id: bundle.metadata.id, digest: bundle.metadata.digest, }), initial_input: None, working_directory_request: None, working_directory: None, worker_observation_enabled: false, worker_observation_grants: Vec::new(), workspace_api: None, memory_settings: Some(manifest::WorkspaceMemorySettingsSnapshot { workspace_id: "local".to_string(), settings_revision: 1, language: "English".to_string(), }), } } #[test] fn durable_worker_request_omits_repository_credentials() { let mut request = task_request("worker-secret-redaction"); request.working_directory_request = Some(WorkingDirectoryRequest { repository: WorkingDirectoryRepository { id: "repository-1".to_string(), provider: "git".to_string(), source: workspace_api::RepositorySource { kind: workspace_api::RepositorySourceKind::Ssh, uri: "ssh://git@example.test/repo.git".to_string(), }, source_revision: 1, source_fingerprint: "sha256:source".to_string(), selector: None, }, materializer: MaterializerKind::RuntimeGitClone, backend_workdir_id: Some("working-directory-1".to_string()), materialization: Some(RepositoryMaterializationContext { workspace_id: "workspace-1".to_string(), runtime_id: "runtime-1".to_string(), operation_id: "operation-1".to_string(), config_revision: 1, config_projection_digest: "sha256:projection".to_string(), ssh: Some(RepositorySshMaterializationAccess { credential_candidates: vec![RepositorySshCredentialCandidate { credential_id: "credential-1".to_string(), credential_revision: 1, private_key: SensitiveString::new("private-key-bytes"), }], host_trust_id: "host-trust-1".to_string(), host_trust_revision: 1, access: workspace_api::RepositoryAccessMode::ReadOnly, expires_at_epoch_seconds: u64::MAX, repository_id: "repository-1".to_string(), repository_source_fingerprint: "sha256:source".to_string(), repository_uri: "ssh://git@example.test/repo.git".to_string(), secret_resource: repository_resource_handle(), known_hosts_entry: SensitiveString::new("known-hosts-entry"), }), }), }); let durable = durable_create_worker_request(&request); assert!( request .working_directory_request .as_ref() .and_then(|working_directory| working_directory.materialization.as_ref()) .and_then(|materialization| materialization.ssh.as_ref()) .is_some() ); assert!( durable .working_directory_request .as_ref() .and_then(|working_directory| working_directory.materialization.as_ref()) .and_then(|materialization| materialization.ssh.as_ref()) .is_none() ); let serialized = serde_json::to_string(&durable).unwrap(); assert!(!serialized.contains("private-key-bytes")); assert!(!serialized.contains("known-hosts-entry")); } fn repository_resource_handle() -> crate::resource::BackendResourceHandle { crate::resource::BackendResourceHandle { kind: crate::resource::BackendResourceKind::RepositorySshAccess, workspace_id: "workspace-1".to_string(), scope_id: Some("repository-ssh-access".to_string()), runtime_id: Some("runtime-1".to_string()), worker_id: None, resource_id: "repository-access-1".to_string(), digest: "opaque:repository-access-1".to_string(), operation: crate::resource::BackendResourceOperation::FetchOnce, expires_at_unix_seconds: i64::MAX, nonce: "repository-access-1".to_string(), revision: "1".to_string(), generation: None, max_bytes: crate::resource::DEFAULT_REPOSITORY_SSH_ACCESS_MAX_BYTES, content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE.to_string(), redaction: crate::resource::ResourceRedactionPolicy::RuntimeInternalOnly, audit_correlation_id: "repository-access-1".to_string(), profile_source_graph: None, } } #[tokio::test] async fn repository_access_resource_is_fetched_before_provider_authorization() { let (runtime, backend) = runtime_and_backend(); backend .repository_access_available .store(true, Ordering::SeqCst); runtime.bind_runtime_identity("runtime-1").unwrap(); let handle = repository_resource_handle(); runtime .install_workspace_backend_resource_client( "workspace-1", Arc::new(TestRepositoryResourceClient { response: Mutex::new(Some(crate::resource::BackendResourceFetchResponse { kind: crate::resource::BackendResourceKind::RepositorySshAccess, resource_id: handle.resource_id.clone(), digest: handle.digest.clone(), content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE .to_string(), bytes: serde_json::to_vec(&RepositorySshAccessSecret { credential_candidates: vec![ crate::resource::RepositorySshAccessSecretCandidate { credential_id: "credential-1".to_string(), credential_revision: 1, private_key: "private-key-bytes-1".to_string(), }, crate::resource::RepositorySshAccessSecretCandidate { credential_id: "credential-2".to_string(), credential_revision: 3, private_key: "private-key-bytes-2".to_string(), }, ], known_hosts_entry: "known-hosts-entry".to_string(), }) .unwrap(), audit_correlation_id: handle.audit_correlation_id.clone(), })), }), ) .unwrap(); let request = WorkingDirectoryRepositoryAccessRequest { working_directory_id: "working-directory-1".to_string(), materialization: RepositoryMaterializationContext { workspace_id: "workspace-1".to_string(), runtime_id: "runtime-1".to_string(), operation_id: "operation-1".to_string(), config_revision: 1, config_projection_digest: "sha256:projection".to_string(), ssh: Some(RepositorySshMaterializationAccess { credential_candidates: vec![ RepositorySshCredentialCandidate { credential_id: "credential-1".to_string(), credential_revision: 1, private_key: SensitiveString::default(), }, RepositorySshCredentialCandidate { credential_id: "credential-2".to_string(), credential_revision: 3, private_key: SensitiveString::default(), }, ], host_trust_id: "host-trust-1".to_string(), host_trust_revision: 1, access: workspace_api::RepositoryAccessMode::ReadOnly, expires_at_epoch_seconds: u64::MAX, repository_id: "repository-1".to_string(), repository_source_fingerprint: "sha256:source".to_string(), repository_uri: "ssh://git@example.test/repo.git".to_string(), secret_resource: handle, known_hosts_entry: SensitiveString::default(), }), }, }; let replay = request.clone(); runtime .authorize_working_directory_repository_access_from_resource(request) .await .unwrap(); assert!( runtime .authorize_working_directory_repository_access_from_resource(replay) .await .is_err() ); let accesses = backend.repository_accesses.lock().unwrap(); assert_eq!(accesses.len(), 1); let access = accesses[0].materialization.ssh.as_ref().unwrap(); assert_eq!(access.credential_candidates.len(), 2); assert_eq!( access.credential_candidates[0].credential_id, "credential-1" ); assert_eq!( access.credential_candidates[1].credential_id, "credential-2" ); assert_eq!( access.credential_candidates[0].private_key.expose(), "private-key-bytes-1" ); assert_eq!( access.credential_candidates[1].private_key.expose(), "private-key-bytes-2" ); assert_eq!(access.known_hosts_entry.expose(), "known-hosts-entry"); } #[tokio::test] async fn working_directory_create_fetches_repository_access_before_provider_call() { let (runtime, backend) = runtime_and_backend(); backend .repository_access_available .store(true, Ordering::SeqCst); runtime.bind_runtime_identity("runtime-1").unwrap(); let handle = repository_resource_handle(); runtime .install_workspace_backend_resource_client( "workspace-1", Arc::new(TestRepositoryResourceClient { response: Mutex::new(Some(crate::resource::BackendResourceFetchResponse { kind: crate::resource::BackendResourceKind::RepositorySshAccess, resource_id: handle.resource_id.clone(), digest: handle.digest.clone(), content_type: crate::resource::REPOSITORY_SSH_ACCESS_CONTENT_TYPE .to_string(), bytes: serde_json::to_vec(&RepositorySshAccessSecret { credential_candidates: vec![ crate::resource::RepositorySshAccessSecretCandidate { credential_id: "credential-1".to_string(), credential_revision: 1, private_key: "create-private-key-bytes-1".to_string(), }, crate::resource::RepositorySshAccessSecretCandidate { credential_id: "credential-2".to_string(), credential_revision: 3, private_key: "create-private-key-bytes-2".to_string(), }, ], known_hosts_entry: "create-known-hosts-entry".to_string(), }) .unwrap(), audit_correlation_id: handle.audit_correlation_id.clone(), })), }), ) .unwrap(); let request = WorkingDirectoryRequest { repository: WorkingDirectoryRepository { id: "repository-1".to_string(), provider: "git".to_string(), source: workspace_api::RepositorySource { kind: workspace_api::RepositorySourceKind::Ssh, uri: "ssh://git@example.test/repo.git".to_string(), }, source_revision: 1, source_fingerprint: "sha256:source".to_string(), selector: None, }, materializer: MaterializerKind::RuntimeGitClone, backend_workdir_id: Some("working-directory-1".to_string()), materialization: Some(RepositoryMaterializationContext { workspace_id: "workspace-1".to_string(), runtime_id: "runtime-1".to_string(), operation_id: "operation-create".to_string(), config_revision: 1, config_projection_digest: "sha256:projection".to_string(), ssh: Some(RepositorySshMaterializationAccess { credential_candidates: vec![ RepositorySshCredentialCandidate { credential_id: "credential-1".to_string(), credential_revision: 1, private_key: SensitiveString::default(), }, RepositorySshCredentialCandidate { credential_id: "credential-2".to_string(), credential_revision: 3, private_key: SensitiveString::default(), }, ], host_trust_id: "host-trust-1".to_string(), host_trust_revision: 1, access: workspace_api::RepositoryAccessMode::ReadOnly, expires_at_epoch_seconds: u64::MAX, repository_id: "repository-1".to_string(), repository_source_fingerprint: "sha256:source".to_string(), repository_uri: "ssh://git@example.test/repo.git".to_string(), secret_resource: handle, known_hosts_entry: SensitiveString::default(), }), }), }; assert!( runtime .create_working_directory_from_resource(request) .await .is_err() ); let requests = backend.working_directory_requests.lock().unwrap(); let access = requests[0] .materialization .as_ref() .and_then(|materialization| materialization.ssh.as_ref()) .unwrap(); assert_eq!( access.credential_candidates[0].private_key.expose(), "create-private-key-bytes-1" ); assert_eq!( access.credential_candidates[1].private_key.expose(), "create-private-key-bytes-2" ); assert_eq!( access.known_hosts_entry.expose(), "create-known-hosts-entry" ); } fn scoped_task_request(objective: &str, workspace_id: &str) -> CreateWorkerRequest { let mut request = task_request(objective); request.workspace_api = Some(WorkspaceApiRef { workspace_id: workspace_id.to_string(), base_url: format!("https://workspace.example/{workspace_id}"), }); request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot { workspace_id: workspace_id.to_string(), settings_revision: 1, language: "English".to_string(), }); request } #[test] fn workspace_create_requires_matching_normalized_memory_settings_snapshot() { let mut request = scoped_task_request("memory-snapshot", "workspace-a"); assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok()); request.memory_settings.as_mut().unwrap().language = "Français".to_string(); assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_ok()); request.memory_settings = None; assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err()); request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot { workspace_id: "workspace-b".to_string(), settings_revision: 1, language: "English".to_string(), }); assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err()); request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot { workspace_id: "workspace-a".to_string(), settings_revision: 2, language: " english ".to_string(), }); assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err()); } fn scope(workspace_id: &str, server_id: &str) -> RuntimeWorkspaceScope { RuntimeWorkspaceScope::new(workspace_id, server_id) } fn test_bundle_for_profile(profile: ProfileSelector) -> ConfigBundle { ConfigBundle { metadata: ConfigBundleMetadata { id: "bundle-1".to_string(), digest: String::new(), revision: "rev-1".to_string(), workspace_id: "workspace-1".to_string(), created_at: "2026-06-26T00:00:00Z".to_string(), provenance: ConfigBundleProvenance { source: "workspace-backend".to_string(), detail: Some("profile-sync".to_string()), }, }, profiles: vec![ConfigProfileDescriptor { selector: profile, label: Some("Coder".to_string()), }], declarations: vec![ConfigDeclaration { kind: ConfigDeclarationKind::CapabilityGrant, name: "read".to_string(), reference: "capability:read".to_string(), }], prompt_catalog: None, profile_source_archive: None, profile_source_archive_handle: None, } .with_computed_digest() } #[derive(Default)] struct TestExecutionBackend { dispatch_result: Mutex>, stop_result: Mutex>, restore_result: Mutex>, restore_count: Mutex, run_generations: Mutex>, config_bundles: Mutex>>, workspace_config_fetches: Mutex>, workspace_config_results: Mutex>, contexts: Mutex>, dispatched_inputs: Mutex>, repository_accesses: Mutex>, repository_access_available: AtomicBool, working_directory_requests: Mutex>, preserve_submission_acknowledgement_id: AtomicBool, #[cfg(feature = "ws-server")] snapshots: Mutex>, } impl TestExecutionBackend { fn set_dispatch_result(&self, result: WorkerExecutionResult) { *self.dispatch_result.lock().unwrap() = Some(result); } fn set_stop_result(&self, result: WorkerExecutionResult) { *self.stop_result.lock().unwrap() = Some(result); } fn preserve_submission_acknowledgement_id(&self) { self.preserve_submission_acknowledgement_id .store(true, Ordering::SeqCst); } #[cfg(feature = "ws-server")] fn set_worker_snapshot(&self, worker_ref: &WorkerRef, snapshot: protocol::Event) { self.snapshots .lock() .unwrap() .insert(worker_ref.worker_id.clone(), snapshot); } #[cfg(feature = "ws-server")] fn publish_text_delta( &self, worker_ref: &WorkerRef, text: &str, ) -> Result { let contexts = self.contexts.lock().unwrap(); let context = contexts.get(&worker_ref.worker_id).expect("context stored"); context.publish_protocol_event(protocol::Event::TextDelta { text: text.into() }) } } impl WorkerExecutionBackend for TestExecutionBackend { fn backend_id(&self) -> &str { "test-execution-backend" } fn fetch_workspace_config( &self, request: WorkspaceConfigFetchRequest, ) -> Result { self.workspace_config_fetches.lock().unwrap().push(request); let mut results = self.workspace_config_results.lock().unwrap(); if results.is_empty() { return Err("no Workspace Config fetch result configured".to_string()); } Ok(results.remove(0)) } fn create_working_directory( &self, request: &WorkingDirectoryRequest, ) -> Result { self.working_directory_requests .lock() .unwrap() .push(request.clone()); Err(WorkingDirectoryDiagnostic::rejected( "working_directory_unsupported", "Worker execution backend does not support working directory materialization", )) } fn authorize_working_directory_repository_access( &self, request: &WorkingDirectoryRepositoryAccessRequest, ) -> Result<(), WorkingDirectoryDiagnostic> { self.repository_accesses .lock() .unwrap() .push(request.clone()); if self.repository_access_available.load(Ordering::SeqCst) { Ok(()) } else { Err(WorkingDirectoryDiagnostic::rejected( "working_directory_repository_access_unsupported", "Worker execution backend does not support Repository access authorization", )) } } fn spawn_worker(&self, request: WorkerExecutionSpawnRequest) -> WorkerExecutionSpawnResult { self.run_generations .lock() .unwrap() .push(request.run_generation); self.config_bundles .lock() .unwrap() .push(request.config_bundle.clone()); self.contexts .lock() .unwrap() .insert(request.worker_ref.worker_id.clone(), request.context); WorkerExecutionSpawnResult::Connected { handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()), worker_state: protocol::WorkerStateSnapshot { execution_generation: request.run_generation, ..protocol::WorkerStatus::Idle.into() }, working_directory: request .working_directory .as_ref() .map(|binding| binding.status()), } } fn restore_worker( &self, request: WorkerExecutionRestoreRequest, ) -> WorkerExecutionSpawnResult { *self.restore_count.lock().unwrap() += 1; self.run_generations .lock() .unwrap() .push(request.run_generation); self.config_bundles .lock() .unwrap() .push(request.config_bundle.clone()); if let Some(result) = self.restore_result.lock().unwrap().clone() { return result; } self.contexts .lock() .unwrap() .insert(request.worker_ref.worker_id.clone(), request.context); WorkerExecutionSpawnResult::Connected { handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()), worker_state: protocol::WorkerStateSnapshot { execution_generation: request.run_generation, ..protocol::WorkerStatus::Idle.into() }, working_directory: request .working_directory .as_ref() .map(|binding| binding.status()), } } fn dispatch_input( &self, _handle: &WorkerExecutionHandle, input: WorkerInput, ) -> WorkerExecutionResult { let submission_id = input.submission_request_id.clone(); self.dispatched_inputs.lock().unwrap().push(input); let mut result = self .dispatch_result .lock() .unwrap() .clone() .unwrap_or_else(|| { WorkerExecutionResult::accepted_submission( WorkerExecutionOperation::Input, "request-test", "test-submission", protocol::SubmissionDisposition::Started, ) }); if !self .preserve_submission_acknowledgement_id .load(Ordering::SeqCst) && let (Some(ack), Some(submission_id)) = (result.submission.as_mut(), submission_id) { ack.submission_request_id = submission_id; } result } fn stop_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult { self.stop_result .lock() .unwrap() .take() .unwrap_or_else(|| WorkerExecutionResult::accepted(WorkerExecutionOperation::Stop)) } fn cancel_worker(&self, _handle: &WorkerExecutionHandle) -> WorkerExecutionResult { WorkerExecutionResult::accepted(WorkerExecutionOperation::Cancel) } #[cfg(feature = "ws-server")] fn worker_snapshot(&self, handle: &WorkerExecutionHandle) -> Option { self.snapshots .lock() .unwrap() .get(&handle.worker_ref().worker_id) .cloned() } } struct TestRepositoryResourceClient { response: Mutex>, } #[async_trait] impl BackendResourceClient for TestRepositoryResourceClient { async fn fetch_resource( &self, _request: BackendResourceFetchRequest, ) -> Result { self.response .lock() .unwrap() .take() .ok_or(BackendResourceError::MissingResource) } } fn runtime_with_backend() -> Runtime { let runtime = Runtime::with_execution_backend( RuntimeOptions::default(), Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); runtime } fn runtime_and_backend() -> (Runtime, Arc) { let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); (runtime, backend) } fn test_bundle() -> ConfigBundle { test_bundle_for_profile(ProfileSelector::Builtin("builtin:coder".to_string())) } fn bundled_task_request(objective: &str, bundle: &ConfigBundle) -> CreateWorkerRequest { let mut request = task_request(objective); request.config_bundle = Some(ConfigBundleRef { id: bundle.metadata.id.clone(), digest: bundle.metadata.digest.clone(), }); request } fn receive_subscription_update( subscription: &mut RuntimeEventSelectorSubscription, ) -> Result { tokio::runtime::Builder::new_current_thread() .build() .unwrap() .block_on(subscription.recv()) } #[test] fn runtime_worker_subscription_has_gap_free_snapshot_and_live_updates() { let runtime = runtime_with_backend(); let mut subscription = runtime .subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers) .unwrap(); assert_eq!(subscription.snapshot_revision(), 0); let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else { panic!("runtime_workers must return a Worker snapshot"); }; assert!(workers.is_empty()); let created = runtime.create_worker(task_request("live")).unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert_eq!(update.subject_revision, 1); match update.payload { SubscriptionEventPayload::WorkerUpserted { worker } => { assert_eq!(worker.worker_id.as_str(), created.worker_id.to_string()); assert_eq!(worker.subject_revision, 1); assert_eq!(worker.state, SubscriptionWorkerState::Idle); } payload => panic!("unexpected subscription payload: {payload:?}"), } runtime .observe_worker_event( &created.worker_ref, internal_worker_status_event( internal_worker_ref("child-live", None), protocol::WorkerStatus::Running, ), ) .unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert_eq!(update.subject_revision, 2); match update.payload { SubscriptionEventPayload::WorkerUpserted { worker } => { assert_eq!(worker.state, SubscriptionWorkerState::Idle); assert!(worker.has_running_internal_workers); } payload => panic!("unexpected subscription payload: {payload:?}"), } runtime .observe_worker_event( &created.worker_ref, internal_worker_status_event( internal_worker_ref("child-live", None), protocol::WorkerStatus::Idle, ), ) .unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert_eq!(update.subject_revision, 3); match update.payload { SubscriptionEventPayload::WorkerUpserted { worker } => { assert_eq!(worker.state, SubscriptionWorkerState::Idle); assert!(!worker.has_running_internal_workers); } payload => panic!("unexpected subscription payload: {payload:?}"), } runtime .observe_worker_event( &created.worker_ref, internal_worker_status_event( internal_worker_ref("child-live", None), protocol::WorkerStatus::Running, ), ) .unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert_eq!(update.subject_revision, 4); match update.payload { SubscriptionEventPayload::WorkerUpserted { worker } => { assert_eq!(worker.state, SubscriptionWorkerState::Idle); assert!(worker.has_running_internal_workers); } payload => panic!("unexpected subscription payload: {payload:?}"), } runtime.stop_worker(&created.worker_ref, None).unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert_eq!(update.subject_revision, 5); match update.payload { SubscriptionEventPayload::WorkerUpserted { worker } => { assert_eq!(worker.worker_id.as_str(), created.worker_id.to_string()); assert_eq!(worker.state, SubscriptionWorkerState::Stopped); assert!(!worker.has_running_internal_workers); } payload => panic!("unexpected subscription payload: {payload:?}"), } runtime.stop_runtime().unwrap(); } #[test] fn worker_lifecycle_subscription_delivers_only_selected_workers() { let runtime = runtime_with_backend(); let first = runtime.create_worker(task_request("first")).unwrap(); let second = runtime.create_worker(task_request("second")).unwrap(); let selected_id = SubscriptionWorkerId::new(first.worker_id.to_string()).unwrap(); let mut subscription = runtime .subscribe_event_selector(EventSubscriptionSelector::WorkerLifecycle { worker_ids: protocol::subscription::SubscriptionWorkerIds::new([ selected_id.clone() ]) .unwrap(), }) .unwrap(); let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else { panic!("worker_lifecycle must return a Worker snapshot"); }; assert_eq!(workers.len(), 1); assert_eq!(workers[0].worker_id, selected_id); runtime.stop_worker(&second.worker_ref, None).unwrap(); assert!(matches!( subscription.receiver.try_recv(), Err(mpsc::error::TryRecvError::Empty) )); runtime.stop_worker(&first.worker_ref, None).unwrap(); let update = receive_subscription_update(&mut subscription).unwrap(); assert!(matches!( update.payload, SubscriptionEventPayload::WorkerUpserted { worker } if worker.worker_id == selected_id )); } #[test] fn scoped_runtime_worker_subscription_hides_other_workspaces() { let runtime = runtime_with_backend(); let workspace_a = runtime .create_worker_scoped( &scope("workspace-a", "server-a"), scoped_task_request("a", "workspace-a"), ) .unwrap(); let workspace_b = runtime .create_worker_scoped( &scope("workspace-b", "server-b"), scoped_task_request("b", "workspace-b"), ) .unwrap(); let mut subscription = runtime .subscribe_event_selector_scoped( &scope("workspace-a", "server-a"), EventSubscriptionSelector::RuntimeWorkers, ) .unwrap(); let SubscriptionSnapshot::Workers { workers } = subscription.snapshot() else { panic!("runtime_workers must return a Worker snapshot"); }; assert_eq!(workers.len(), 1); assert_eq!( workers[0].worker_id.as_str(), workspace_a.worker_id.to_string() ); runtime.stop_worker(&workspace_b.worker_ref, None).unwrap(); assert!(matches!( subscription.receiver.try_recv(), Err(mpsc::error::TryRecvError::Empty) )); runtime.stop_worker(&workspace_a.worker_ref, None).unwrap(); assert!(receive_subscription_update(&mut subscription).is_ok()); } #[test] fn lagged_runtime_subscription_closes_without_blocking_mutation() { let runtime = runtime_with_backend(); let created = runtime.create_worker(task_request("lag")).unwrap(); let mut subscription = runtime .subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers) .unwrap(); { let mut state = runtime.lock().unwrap(); for _ in 0..=SUBSCRIPTION_QUEUE_CAPACITY { state.publish_worker_upsert(created.worker_id).unwrap(); } } assert_eq!(subscription.receiver.len(), SUBSCRIPTION_QUEUE_CAPACITY); assert!(matches!( receive_subscription_update(&mut subscription), Err(RuntimeSubscriptionRecvError::Lagged) )); } #[test] fn dropping_runtime_subscription_releases_producer_state() { let runtime = runtime_with_backend(); let subscription = runtime .subscribe_event_selector(EventSubscriptionSelector::RuntimeWorkers) .unwrap(); assert_eq!(runtime.lock().unwrap().subscriptions.len(), 1); drop(subscription); assert!(runtime.lock().unwrap().subscriptions.is_empty()); } #[test] fn scoped_worker_access_hides_other_workspace_workers() { let runtime = runtime_with_backend(); let workspace_a = runtime .create_worker_scoped( &scope("workspace-a", "server-a"), scoped_task_request("a", "workspace-a"), ) .unwrap(); let workspace_b = runtime .create_worker_scoped( &scope("workspace-b", "server-b"), scoped_task_request("b", "workspace-b"), ) .unwrap(); assert_eq!(workspace_a.workspace_id.as_deref(), Some("workspace-a")); assert_eq!(workspace_b.workspace_id.as_deref(), Some("workspace-b")); assert_eq!( runtime .list_workers_scoped(&scope("workspace-a", "server-a")) .unwrap() .into_iter() .map(|worker| worker.worker_ref) .collect::>(), vec![workspace_a.worker_ref.clone()] ); assert_eq!( runtime .list_workers_scoped(&scope("workspace-b", "server-b")) .unwrap() .into_iter() .map(|worker| worker.worker_ref) .collect::>(), vec![workspace_b.worker_ref.clone()] ); let detail_error = runtime .worker_detail_scoped(&scope("workspace-a", "server-a"), &workspace_b.worker_ref) .unwrap_err(); assert!(matches!(detail_error, RuntimeError::WorkerNotFound { .. })); let input_error = runtime .send_input_scoped( &scope("workspace-a", "server-a"), &workspace_b.worker_ref, WorkerInput::user("cross workspace"), ) .unwrap_err(); assert!(matches!(input_error, RuntimeError::WorkerNotFound { .. })); let protocol_error = runtime .send_protocol_method_scoped( &scope("workspace-a", "server-a"), &workspace_b.worker_ref, Method::Shutdown { command: test_command(), }, ) .unwrap_err(); assert!(matches!( protocol_error, RuntimeError::WorkerNotFound { .. } )); assert_eq!( runtime .worker_detail(&workspace_b.worker_ref) .unwrap() .status, WorkerStatus::Idle ); } #[test] fn workspace_api_replacement_updates_persisted_request() { let (runtime, _backend) = runtime_and_backend(); let scope = scope("workspace-a", "server-a"); let worker = runtime .create_worker_scoped( &scope, scoped_task_request("repair credential", "workspace-a"), ) .unwrap(); let replacement = WorkspaceApiRef { workspace_id: "workspace-a".to_string(), base_url: "https://workspace.example/workspace-a/".to_string(), }; runtime .replace_worker_workspace_api_scoped(&scope, &worker.worker_ref, replacement.clone()) .unwrap(); let state = runtime.lock().unwrap(); assert_eq!( state .worker(&worker.worker_ref) .unwrap() .request .workspace_api, Some(replacement) ); } #[test] fn workspace_owner_binding_rejects_other_backend_and_forgets_after_last_worker_delete() { let runtime = runtime_with_backend(); let server_a = scope("workspace-a", "server-a"); let server_b = scope("workspace-a", "server-b"); let first = runtime .create_worker_scoped(&server_a, scoped_task_request("first", "workspace-a")) .unwrap(); let second = runtime .create_worker_scoped(&server_a, scoped_task_request("second", "workspace-a")) .unwrap(); let create_error = runtime .create_worker_scoped(&server_b, scoped_task_request("stolen", "workspace-a")) .unwrap_err(); assert!(matches!( create_error, RuntimeError::WorkspaceOwnerMismatch { .. } )); let read_error = runtime .worker_detail_scoped(&server_b, &first.worker_ref) .unwrap_err(); assert!(matches!( read_error, RuntimeError::WorkspaceOwnerMismatch { .. } )); runtime.stop_worker(&first.worker_ref, None).unwrap(); runtime.stop_worker(&second.worker_ref, None).unwrap(); runtime .delete_worker_scoped(&server_a, &first.worker_ref) .unwrap(); let still_owned_error = runtime .create_worker_scoped(&server_b, scoped_task_request("still owned", "workspace-a")) .unwrap_err(); assert!(matches!( still_owned_error, RuntimeError::WorkspaceOwnerMismatch { .. } )); runtime .delete_worker_scoped(&server_a, &second.worker_ref) .unwrap(); let rebound = runtime .create_worker_scoped(&server_b, scoped_task_request("rebound", "workspace-a")) .unwrap(); assert_eq!(rebound.workspace_id.as_deref(), Some("workspace-a")); } #[test] fn scoped_create_rejects_request_workspace_mismatch() { let runtime = runtime_with_backend(); let error = runtime .create_worker_scoped( &scope("workspace-a", "server-a"), scoped_task_request("mismatch", "workspace-b"), ) .unwrap_err(); assert!(matches!(error, RuntimeError::InvalidRequest(_))); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn unscoped_legacy_workers_are_hidden_from_scoped_access() { let runtime = runtime_with_backend(); let legacy = runtime.create_worker(task_request("legacy")).unwrap(); assert!(legacy.workspace_id.is_none()); assert!( runtime .list_workers_scoped(&scope("workspace-a", "server-a")) .unwrap() .is_empty() ); let detail_error = runtime .worker_detail_scoped(&scope("workspace-a", "server-a"), &legacy.worker_ref) .unwrap_err(); assert!(matches!(detail_error, RuntimeError::WorkerNotFound { .. })); } #[test] fn stopped_worker_scoped_list_uses_workspace_boundary() { let runtime = runtime_with_backend(); let workspace_a = runtime .create_worker_scoped( &scope("workspace-a", "server-a"), scoped_task_request("a", "workspace-a"), ) .unwrap(); let workspace_b = runtime .create_worker_scoped( &scope("workspace-b", "server-b"), scoped_task_request("b", "workspace-b"), ) .unwrap(); runtime.stop_worker(&workspace_a.worker_ref, None).unwrap(); runtime.stop_worker(&workspace_b.worker_ref, None).unwrap(); let stopped = runtime .list_stopped_workers_scoped(&scope("workspace-a", "server-a")) .unwrap(); assert_eq!(stopped.len(), 1); assert_eq!(stopped[0].worker_ref, workspace_a.worker_ref); } #[test] fn create_list_and_detail_preserve_runtime_local_worker_authority() { let runtime = runtime_with_backend(); let detail = runtime.create_worker(task_request("implement v0")).unwrap(); assert_eq!(detail.status, WorkerStatus::Idle); assert_eq!(detail.config_bundle.as_ref().unwrap().id, "bundle-1"); let list = runtime.list_workers().unwrap(); assert_eq!(list.len(), 1); assert_eq!(list[0].worker_ref, detail.worker_ref); assert_eq!(list[0].config_bundle, detail.config_bundle); let fetched = runtime.worker_detail(&detail.worker_ref).unwrap(); assert_eq!(fetched.worker_id, detail.worker_id); assert_eq!(fetched.profile, detail.profile); } #[test] fn stopped_worker_list_excludes_alive_and_cancelled_workers() { let runtime = runtime_with_backend(); let alive = runtime.create_worker(task_request("alive")).unwrap(); let stopped = runtime.create_worker(task_request("stopped")).unwrap(); let cancelled = runtime.create_worker(task_request("cancelled")).unwrap(); runtime.stop_worker(&stopped.worker_ref, None).unwrap(); runtime.cancel_worker(&cancelled.worker_ref, None).unwrap(); let candidates = runtime.list_stopped_workers().unwrap(); assert_eq!(candidates.len(), 1); assert_eq!(candidates[0].worker_ref, stopped.worker_ref); assert_eq!(candidates[0].status, WorkerStatus::Stopped); assert_ne!(candidates[0].worker_ref, alive.worker_ref); assert_ne!(candidates[0].worker_ref, cancelled.worker_ref); } #[test] fn synced_config_bundle_is_stored_checked_and_used_for_worker_creation() { let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap(); let mut bundle = test_bundle(); bundle.prompt_catalog = Some( worker::EffectivePromptCatalog::new( BTreeMap::from([("default".to_string(), "workspace prompt".to_string())]), 7, "schema", "toolchain", ) .unwrap(), ); bundle = bundle.with_computed_digest(); let availability = runtime.store_config_bundle(bundle.clone()).unwrap(); assert_eq!(availability.reference.id, "bundle-1"); assert_eq!(availability.reference.digest, bundle.metadata.digest); let mut conflicting_bundle = bundle.clone(); conflicting_bundle.profiles[0].label = Some("conflicting".to_string()); conflicting_bundle = conflicting_bundle.with_computed_digest(); assert!(matches!( runtime.store_config_bundle(conflicting_bundle), Err(RuntimeError::ConfigBundleDigestMismatch { .. }) )); let listed = runtime.list_config_bundles().unwrap(); assert_eq!(listed.len(), 1); assert_eq!(listed[0].id, "bundle-1"); let checked = runtime .check_config_bundle(&availability.reference) .unwrap(); assert_eq!(checked.summary.digest, availability.summary.digest); let detail = runtime .create_worker(bundled_task_request("synced", &bundle)) .unwrap(); assert_eq!(detail.config_bundle, Some(availability.reference)); assert_eq!( backend.config_bundles.lock().unwrap().as_slice(), &[Some(bundle.clone())] ); runtime.stop_worker(&detail.worker_ref, None).unwrap(); runtime.restore_worker(&detail.worker_ref).unwrap(); assert_eq!( backend.config_bundles.lock().unwrap().as_slice(), &[Some(bundle), None] ); } #[test] fn remote_create_refreshes_workspace_config_and_revalidates_cached_etag() { let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_execution_backend(RuntimeOptions::default(), backend.clone()).unwrap(); let mut bundle = test_bundle(); let archive = test_profile_source_archive(); bundle.profile_source_archive = Some(archive.clone()); bundle = bundle.with_computed_digest(); backend.workspace_config_results.lock().unwrap().extend([ WorkspaceConfigFetchResult::Modified(bundle.clone()), WorkspaceConfigFetchResult::NotModified, ]); let archive = archive.reference; let request = |objective: &str| { let mut request = bundled_task_request(objective, &bundle); request.profile_source = ProfileSourceArchiveSource::WorkspaceConfig { archive: archive.clone(), }; request.workspace_api = Some(WorkspaceApiRef { workspace_id: "workspace-test".to_string(), base_url: "https://workspace.example".to_string(), }); request }; let first = request("first refresh"); let second = request("cached refresh"); runtime.create_worker(first).unwrap(); runtime.create_worker(second.clone()).unwrap(); runtime.create_worker(second).unwrap(); let fetches = backend.workspace_config_fetches.lock().unwrap(); assert_eq!(fetches.len(), 2); assert_eq!(fetches[0].cached, None); assert_eq!( fetches[1].cached, Some(ConfigBundleRef { id: bundle.metadata.id.clone(), digest: bundle.metadata.digest.clone(), }) ); } #[test] fn restore_does_not_require_recorded_config_bundle() { let (runtime, backend) = runtime_and_backend(); let bundle = test_bundle(); let detail = runtime .create_worker(bundled_task_request("missing-on-restore", &bundle)) .unwrap(); runtime.stop_worker(&detail.worker_ref, None).unwrap(); runtime.lock().unwrap().config_bundles.clear(); runtime.restore_worker(&detail.worker_ref).unwrap(); assert_eq!( backend.config_bundles.lock().unwrap().as_slice(), &[Some(bundle), None] ); let (runtime, backend) = runtime_and_backend(); let bundle = test_bundle(); let detail = runtime .create_worker(bundled_task_request("mismatch-on-restore", &bundle)) .unwrap(); runtime.stop_worker(&detail.worker_ref, None).unwrap(); let mut replacement = bundle.clone(); replacement.profiles[0].label = Some("replacement".to_string()); replacement = replacement.with_computed_digest(); runtime .lock() .unwrap() .config_bundles .insert(replacement.metadata.id.clone(), replacement); runtime.restore_worker(&detail.worker_ref).unwrap(); assert_eq!( backend.config_bundles.lock().unwrap().as_slice(), &[Some(bundle), None] ); } #[test] fn config_bundle_errors_are_typed() { let runtime = Runtime::new_memory(); let bundle = test_bundle(); runtime.store_config_bundle(bundle.clone()).unwrap(); let mismatch = runtime .check_config_bundle(&ConfigBundleRef { id: bundle.metadata.id.clone(), digest: "0".repeat(64), }) .unwrap_err(); assert!(matches!( mismatch, RuntimeError::ConfigBundleDigestMismatch { .. } )); let mut unsupported = test_bundle(); unsupported.declarations.push(ConfigDeclaration { kind: ConfigDeclarationKind::Unsupported, name: "plugin-registry".to_string(), reference: "plugin-registry:v0".to_string(), }); unsupported = unsupported.with_computed_digest(); let unsupported_err = runtime.store_config_bundle(unsupported).unwrap_err(); assert!(matches!( unsupported_err, RuntimeError::UnsupportedConfigDeclaration { .. } )); } #[test] fn create_worker_idempotency_reuses_worker_and_rejects_different_input() { let runtime = runtime_with_backend(); let mut request = task_request("idempotent"); request.create_fingerprint = "sha256:input-1".to_string(); request.working_directory = Some(WorkingDirectoryClaim { working_directory_id: "workdir-idempotent".to_string(), relative_cwd: None, }); let first = runtime.create_worker(request.clone()).unwrap(); let workdir_count_after_first = runtime.list_working_directories().unwrap().len(); let replayed = runtime.create_worker(request.clone()).unwrap(); assert_eq!(replayed.worker_ref, first.worker_ref); assert_eq!(runtime.list_workers().unwrap().len(), 1); assert_eq!( runtime.list_working_directories().unwrap().len(), workdir_count_after_first ); request.create_fingerprint = "sha256:different".to_string(); let error = runtime.create_worker(request).unwrap_err(); assert!(matches!(error, RuntimeError::InvalidRequest(_))); assert_eq!(runtime.list_workers().unwrap().len(), 1); } #[test] fn create_worker_rejects_system_initial_input_without_persisting_worker() { let runtime = runtime_with_backend(); let mut request = task_request("system initial input"); request.initial_input = Some(WorkerInput::notify("role/system belongs in config bundle")); let error = runtime.create_worker(request).unwrap_err(); assert!(matches!( error, RuntimeError::InvalidInitialInputKind { .. } )); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn create_worker_does_not_infer_state_from_started_submission_ack() { let (runtime, backend) = runtime_and_backend(); backend.set_dispatch_result(WorkerExecutionResult::accepted_submission( WorkerExecutionOperation::Input, "request-test", "test-submission", protocol::SubmissionDisposition::Started, )); let mut request = task_request("committed initial input is already idle"); request.initial_input = Some(WorkerInput::user("start the ticket")); let detail = runtime.create_worker(request).unwrap(); assert_eq!(detail.status, WorkerStatus::Idle); assert_eq!( detail.worker_state.as_ref().map(|snapshot| &snapshot.state), Some(&protocol::WorkerState::Idle) ); } #[test] fn restored_worker_exposes_the_backend_initial_state_snapshot() { let (runtime, _) = runtime_and_backend(); let created = runtime .create_worker(task_request("restore initial state")) .unwrap(); runtime.stop_worker(&created.worker_ref, None).unwrap(); let restored = runtime.restore_worker(&created.worker_ref).unwrap(); let worker_state = restored .worker_state .expect("restored Worker must expose its initial state"); assert_eq!(worker_state.execution_generation, 2); assert_eq!(worker_state.state, protocol::WorkerState::Idle); } #[test] fn runtime_applies_only_newer_worker_state_snapshots() { let (runtime, _) = runtime_and_backend(); let detail = runtime .create_worker(task_request("state ordering")) .unwrap(); let running = protocol::WorkerStateSnapshot { execution_generation: 7, revision: 3, state: protocol::WorkerState::Busy(protocol::WorkerBusyState::Run( protocol::WorkerRunState::Running, )), last_command_id: 2, }; assert!({ let mut state = runtime.lock().unwrap(); state.project_protocol_event_to_worker_state( &detail.worker_ref, &protocol::Event::WorkerState { snapshot: running.clone(), }, ) }); assert_eq!( runtime .worker_detail(&detail.worker_ref) .unwrap() .worker_state, Some(running.clone()) ); assert!({ let mut state = runtime.lock().unwrap(); !state.project_protocol_event_to_worker_state( &detail.worker_ref, &protocol::Event::WorkerState { snapshot: protocol::WorkerStateSnapshot { revision: 2, state: protocol::WorkerState::Idle, ..running.clone() }, }, ) }); assert!({ let mut state = runtime.lock().unwrap(); !state.project_protocol_event_to_worker_state( &detail.worker_ref, &protocol::Event::WorkerState { snapshot: protocol::WorkerStateSnapshot { state: protocol::WorkerState::Idle, ..running.clone() }, }, ) }); let after = runtime.worker_detail(&detail.worker_ref).unwrap(); assert_eq!(after.status, WorkerStatus::Idle); assert_eq!(after.worker_state, Some(running)); } #[test] fn create_worker_rejects_mismatched_submission_acknowledgement() { let (runtime, backend) = runtime_and_backend(); backend.preserve_submission_acknowledgement_id(); backend.set_dispatch_result(WorkerExecutionResult::accepted_submission( WorkerExecutionOperation::Input, "request-test", "forged-submission", protocol::SubmissionDisposition::Started, )); let mut request = task_request("mismatched initial input commit ack"); request.initial_input = Some(WorkerInput::user("start the ticket")); let error = runtime.create_worker(request).unwrap_err(); assert!(matches!( error, RuntimeError::WorkerExecutionRejected { outcome: crate::execution::WorkerExecutionOutcome::Rejected, .. } )); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn create_worker_rejects_initial_input_without_durable_submission_acknowledgement() { let (runtime, backend) = runtime_and_backend(); backend.set_dispatch_result(WorkerExecutionResult::accepted( WorkerExecutionOperation::Input, )); let mut request = task_request("missing initial input commit ack"); request.initial_input = Some(WorkerInput::user("start the ticket")); let error = runtime.create_worker(request).unwrap_err(); assert!(matches!( error, RuntimeError::WorkerExecutionRejected { outcome: crate::execution::WorkerExecutionOutcome::Rejected, .. } )); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn create_worker_without_execution_backend_is_rejected_and_not_persisted() { let runtime = Runtime::new_memory(); runtime.store_config_bundle(test_bundle()).unwrap(); let error = runtime .create_worker(task_request("no backend")) .unwrap_err(); assert!(matches!( error, RuntimeError::ExecutionBackendUnavailable { .. } )); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn create_worker_without_execution_backend_is_rejected_before_persisting_worker() { let runtime = Runtime::new_memory(); let error = runtime .create_worker(task_request("missing backend")) .unwrap_err(); assert!(matches!( error, RuntimeError::ExecutionBackendUnavailable { .. } )); assert!(runtime.list_workers().unwrap().is_empty()); } #[test] fn connected_backend_busy_dispatch_is_typed_and_not_transcribed() { let (runtime, backend) = runtime_and_backend(); backend.set_dispatch_result(WorkerExecutionResult::busy( WorkerExecutionOperation::Input, "worker is already running", )); let detail = runtime.create_worker(task_request("busy")).unwrap(); let err = runtime .send_input(&detail.worker_ref, WorkerInput::user("wait")) .unwrap_err(); assert!(matches!( err, RuntimeError::WorkerExecutionRejected { outcome: crate::execution::WorkerExecutionOutcome::Busy, .. } )); let refreshed = runtime.worker_detail(&detail.worker_ref).unwrap(); assert_eq!(refreshed.status, WorkerStatus::Idle); #[cfg(feature = "ws-server")] assert!( runtime .read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero()) .unwrap() .is_empty() ); } #[cfg(feature = "ws-server")] #[test] fn backend_protocol_publish_hook_writes_observation_bus() { let (runtime, backend) = runtime_and_backend(); let detail = runtime.create_worker(task_request("observe")).unwrap(); let observation = backend .publish_text_delta(&detail.worker_ref, "from backend") .unwrap(); assert_eq!(observation.worker_ref, detail.worker_ref); let observations = runtime .read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero()) .unwrap(); assert_eq!(observations.len(), 1); assert!(matches!( observations[0].payload, protocol::Event::TextDelta { .. } )); } #[cfg(feature = "ws-server")] #[test] fn observation_snapshot_prefers_live_backend_snapshot() { let (runtime, backend) = runtime_and_backend(); let detail = runtime .create_worker(task_request("observe snapshot")) .unwrap(); let expected_entry = serde_json::json!({"kind": "restored-log-entry"}); backend.set_worker_snapshot( &detail.worker_ref, protocol::Event::Snapshot { session: protocol::SessionSnapshot { pending_submissions: protocol::PendingSubmissionsSnapshot::default(), entries: vec![protocol::SessionSnapshotEntry { entry_id: "restored-log-entry".to_owned(), timestamp: 1, provenance: protocol::SessionEntryProvenance::LegacyUnknown, derived_from: Vec::new(), data: protocol::SessionSnapshotEntryData::RunError { message: expected_entry.to_string(), }, }], }, greeting: protocol::Greeting { worker_name: "live-worker".to_string(), cwd: "/tmp/live".to_string(), provider: "test-provider".to_string(), model: "test-model".to_string(), scope_summary: "live snapshot".to_string(), tools: Vec::new(), context_window: 128, context_tokens: 64, }, state: protocol::WorkerStatus::Running.into(), in_flight: protocol::InFlightSnapshot { blocks: Vec::new(), commands: Vec::new(), }, internal_workers: Vec::new(), }, ); let snapshot = runtime .worker_observation_snapshot(&detail.worker_ref) .unwrap(); match snapshot { protocol::Event::Snapshot { session, greeting, state, .. } => { assert_eq!(session.entries.len(), 1); assert_eq!(session.entries[0].entry_id, "restored-log-entry"); assert_eq!(greeting.worker_name, "live-worker"); assert_eq!(state.catalog_status(), protocol::WorkerStatus::Running); } other => panic!("expected snapshot, got {other:?}"), } } #[cfg(feature = "ws-server")] #[test] fn observation_snapshot_fails_closed_when_backend_snapshot_is_unavailable() { let runtime = runtime_with_backend(); let detail = runtime .create_worker(task_request("snapshot unavailable")) .unwrap(); assert!(matches!( runtime .worker_observation_snapshot(&detail.worker_ref) .unwrap_err(), RuntimeError::WorkerExecutionUnavailable { worker_id, message, } if worker_id == detail.worker_ref.worker_id && message == "authoritative Worker snapshot is unavailable" )); } struct InputOnlyBackend; impl WorkerExecutionBackend for InputOnlyBackend { fn backend_id(&self) -> &str { "input-only" } fn spawn_worker(&self, request: WorkerExecutionSpawnRequest) -> WorkerExecutionSpawnResult { WorkerExecutionSpawnResult::Connected { handle: WorkerExecutionHandle::new(request.worker_ref, self.backend_id()), worker_state: protocol::WorkerStateSnapshot { execution_generation: request.run_generation, ..protocol::WorkerStatus::Idle.into() }, working_directory: request .working_directory .as_ref() .map(|binding| binding.status()), } } fn dispatch_input( &self, _handle: &WorkerExecutionHandle, input: WorkerInput, ) -> WorkerExecutionResult { WorkerExecutionResult::accepted_submission( WorkerExecutionOperation::Input, "request-test", input .submission_request_id .expect("Runtime submission request id"), protocol::SubmissionDisposition::Started, ) } } #[test] fn connected_backend_stop_unsupported_is_typed_rejection() { let runtime = Runtime::with_execution_backend(RuntimeOptions::default(), Arc::new(InputOnlyBackend)) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let detail = runtime.create_worker(task_request("no stop")).unwrap(); let err = runtime .stop_worker(&detail.worker_ref, Some("stop".to_string())) .unwrap_err(); assert!(matches!( err, RuntimeError::WorkerExecutionRejected { outcome: crate::execution::WorkerExecutionOutcome::Unsupported, .. } )); assert_eq!( runtime.worker_detail(&detail.worker_ref).unwrap().status, WorkerStatus::Idle ); } #[test] fn protocol_shutdown_transitions_worker_to_stopped() { let runtime = runtime_with_backend(); let detail = runtime .create_worker(task_request("shutdown from client")) .unwrap(); runtime .send_protocol_method( &detail.worker_ref, Method::Shutdown { command: test_command(), }, ) .unwrap(); assert_eq!( runtime.worker_detail(&detail.worker_ref).unwrap().status, WorkerStatus::Stopped ); } #[test] fn stopped_worker_rejects_input_until_explicitly_restored() { let (runtime, backend) = runtime_and_backend(); let detail = runtime .create_worker(task_request("restore explicitly")) .unwrap(); runtime .send_protocol_method( &detail.worker_ref, Method::Shutdown { command: test_command(), }, ) .unwrap(); assert!(matches!( runtime.send_input(&detail.worker_ref, WorkerInput::user("do not wake")), Err(RuntimeError::WorkerExecutionUnavailable { .. }) )); assert_eq!(*backend.restore_count.lock().unwrap(), 0); runtime.restore_worker(&detail.worker_ref).unwrap(); runtime .send_input(&detail.worker_ref, WorkerInput::user("wake up")) .unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 1); assert_eq!(*backend.run_generations.lock().unwrap(), vec![1, 2]); let restored = runtime.worker_detail(&detail.worker_ref).unwrap(); assert_eq!(restored.status, WorkerStatus::Idle); assert_eq!( restored .worker_state .as_ref() .map(|snapshot| (snapshot.execution_generation, &snapshot.state)), Some((2, &protocol::WorkerState::Idle)) ); } #[test] fn restore_does_not_redispatch_spawn_initial_submit() { let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_execution_backend( RuntimeOptions { ..RuntimeOptions::default() }, backend.clone(), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let mut request = task_request("flow restore"); request.initial_input = Some(WorkerInput { kind: WorkerInputKind::User, content: String::new(), submission_request_id: None, segments: Some(vec![ protocol::Segment::Flow { selector: "builtin:coder-review".to_string(), }, protocol::Segment::text("Implement Ticket 00001"), ]), }); let detail = runtime.create_worker(request).unwrap(); assert_eq!(backend.dispatched_inputs.lock().unwrap().len(), 1); runtime .stop_worker(&detail.worker_ref, Some("restore test".to_string())) .unwrap(); runtime.restore_worker(&detail.worker_ref).unwrap(); assert_eq!( backend.dispatched_inputs.lock().unwrap().len(), 1, "restore must continue durable Worker state without replaying spawn initial input" ); } #[test] fn send_input_dispatches_segment_only_flow_submission() { let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_execution_backend( RuntimeOptions { ..RuntimeOptions::default() }, backend.clone(), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let detail = runtime.create_worker(task_request("flow segment")).unwrap(); let input = WorkerInput { kind: WorkerInputKind::User, content: String::new(), submission_request_id: None, segments: Some(vec![protocol::Segment::Flow { selector: "builtin:coder-review".to_string(), }]), }; runtime .send_input(&detail.worker_ref, input.clone()) .unwrap(); let dispatched = backend.dispatched_inputs.lock().unwrap(); assert_eq!(dispatched.len(), 1); assert_eq!(dispatched[0].kind, input.kind); assert_eq!(dispatched[0].content, input.content); assert_eq!(dispatched[0].segments, input.segments); let submission_request_id = dispatched[0] .submission_request_id .as_deref() .expect("Runtime submission request id"); Uuid::parse_str(submission_request_id).expect("submission request id UUID"); } #[cfg(feature = "ws-server")] #[test] fn notify_input_does_not_duplicate_committed_notification_observation() { let runtime = Runtime::with_execution_backend( RuntimeOptions { ..RuntimeOptions::default() }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let detail = runtime.create_worker(task_request("chat")).unwrap(); runtime .send_input(&detail.worker_ref, WorkerInput::user("hello")) .unwrap(); runtime .send_input(&detail.worker_ref, WorkerInput::notify("note")) .unwrap(); let observations = runtime .read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero()) .unwrap(); assert!(observations.is_empty()); runtime .observe_worker_event( &detail.worker_ref, protocol::Event::SystemItem { item: serde_json::json!({ "kind": "notification", "message": "note", "body": "[Notification] note", }), }, ) .unwrap(); let observations = runtime .read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero()) .unwrap(); assert_eq!(observations.len(), 1); let protocol::Event::SystemItem { item } = &observations[0].payload else { panic!("committed notification observation must be a system item"); }; assert_eq!(item["kind"], "notification"); assert!(observations.iter().all(|observation| { !matches!( &observation.payload, protocol::Event::SystemItem { item } if item["kind"] == "embedded_worker_notification" ) })); } #[test] fn stop_and_cancel_workers_keep_four_state_summary() { let runtime = runtime_with_backend(); let stopped = runtime.create_worker(task_request("stop me")).unwrap(); let cancelled = runtime.create_worker(task_request("cancel me")).unwrap(); runtime .send_input(&cancelled.worker_ref, WorkerInput::user("start")) .unwrap(); let stop_ack = runtime .stop_worker(&stopped.worker_ref, Some("done".to_string())) .unwrap(); assert_eq!(stop_ack.status, WorkerStatus::Stopped); let cancel_ack = runtime .cancel_worker(&cancelled.worker_ref, Some("abort".to_string())) .unwrap(); assert_eq!(cancel_ack.status, WorkerStatus::Idle); let summary = runtime.summary().unwrap(); assert_eq!(summary.worker_count, 2); assert_eq!(summary.active_worker_count, 1); assert_eq!(summary.stopped_worker_count, 1); assert!(serde_json::from_value::(serde_json::json!("cancelled")).is_err()); } #[test] fn delete_worker_removes_stopped_worker_from_runtime() { let runtime = runtime_with_backend(); let worker = runtime.create_worker(task_request("delete me")).unwrap(); assert!(runtime.delete_worker(&worker.worker_ref).is_err()); runtime .stop_worker(&worker.worker_ref, Some("done".to_string())) .unwrap(); let result = runtime.delete_worker(&worker.worker_ref).unwrap(); assert!(result.deleted); assert_eq!(result.worker_id, worker.worker_id); assert!(matches!( runtime.worker_detail(&worker.worker_ref), Err(RuntimeError::WorkerNotFound { .. }) )); let summary = runtime.summary().unwrap(); assert_eq!(summary.worker_count, 0); } #[test] fn delete_worker_waits_for_execution_shutdown_before_removing_snapshot() { let root = tempfile::tempdir().unwrap(); let store_root = root.path().join("runtime"); let backend = Arc::new(TestExecutionBackend::default()); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: store_root.clone(), runtime_id: "delete-shutdown-barrier".to_string(), display_name: None, }, backend.clone(), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); backend.set_dispatch_result(WorkerExecutionResult::errored( WorkerExecutionOperation::Input, "initial input failed", )); backend.set_stop_result(WorkerExecutionResult::errored( WorkerExecutionOperation::Stop, "shutdown is still pending", )); let mut request = task_request("delete shutdown barrier"); request.initial_input = Some(WorkerInput::user("start")); let worker_id = request.worker_id; runtime.create_worker(request).unwrap_err(); let worker_ref = WorkerRef::new(worker_id); backend.set_stop_result(WorkerExecutionResult::errored( WorkerExecutionOperation::Stop, "shutdown is still pending", )); let error = runtime.delete_worker(&worker_ref).unwrap_err(); assert!(matches!( error, RuntimeError::WorkerExecutionRejected { operation: WorkerExecutionOperation::Stop, .. } )); assert_eq!( runtime.worker_detail(&worker_ref).unwrap().status, WorkerStatus::Stopped ); assert!( store_root .join("workers") .join(worker_id.to_string()) .exists(), "Worker aggregate must remain until backend shutdown is confirmed" ); assert!(runtime.delete_worker(&worker_ref).unwrap().deleted); assert!( !store_root .join("workers") .join(worker_id.to_string()) .exists() ); } #[test] fn stop_then_cancel_preserves_stopped_terminal_state() { let runtime = runtime_with_backend(); let worker = runtime .create_worker(task_request("stable stopped")) .unwrap(); let stop_ack = runtime .stop_worker(&worker.worker_ref, Some("done".to_string())) .unwrap(); let cancel_ack = runtime .cancel_worker(&worker.worker_ref, Some("late cancel".to_string())) .unwrap(); assert_eq!(stop_ack.status, WorkerStatus::Stopped); assert_eq!(cancel_ack.status, WorkerStatus::Stopped); assert_eq!( runtime.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Stopped ); let summary = runtime.summary().unwrap(); assert_eq!(summary.active_worker_count, 0); assert_eq!(summary.stopped_worker_count, 1); } #[test] fn cancel_then_stop_transitions_idle_session_to_stopped() { let runtime = runtime_with_backend(); let worker = runtime .create_worker(task_request("cancel then stop")) .unwrap(); runtime .send_input(&worker.worker_ref, WorkerInput::user("start")) .unwrap(); let cancel_ack = runtime .cancel_worker(&worker.worker_ref, Some("abort".to_string())) .unwrap(); let stop_ack = runtime .stop_worker(&worker.worker_ref, Some("late stop".to_string())) .unwrap(); assert_eq!(cancel_ack.status, WorkerStatus::Idle); assert_eq!(stop_ack.status, WorkerStatus::Stopped); assert_eq!( runtime.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Stopped ); let summary = runtime.summary().unwrap(); assert_eq!(summary.active_worker_count, 0); assert_eq!(summary.stopped_worker_count, 1); } #[cfg(feature = "fs-store")] static NEXT_FS_TEST_ROOT: AtomicU64 = AtomicU64::new(1); #[cfg(feature = "fs-store")] fn fs_store_root(label: &str) -> std::path::PathBuf { let sequence = NEXT_FS_TEST_ROOT.fetch_add(1, Ordering::Relaxed); let root = std::env::temp_dir().join(format!( "worker-runtime-fs-store-{label}-{}-{sequence}", std::process::id() )); let _ = std::fs::remove_dir_all(&root); root } #[cfg(feature = "fs-store")] fn runtime_store(runtime: &Runtime) -> FsRuntimeStore { let state = runtime.lock().unwrap(); match &state.persistence { RuntimePersistence::Fs(store) => store.clone(), RuntimePersistence::Memory => panic!("expected fs-backed runtime"), } } #[cfg(feature = "fs-store")] #[test] fn fs_store_rejects_schema_older_than_previous_release() { let root = fs_store_root("unsupported-old-schema"); let runtime = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap(); drop(runtime); let runtime_path = root.join("runtime.json"); let mut runtime_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap(); runtime_json["schema_version"] = serde_json::json!(2); std::fs::write( &runtime_path, serde_json::to_vec_pretty(&runtime_json).unwrap(), ) .unwrap(); let error = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap_err(); assert!( error .to_string() .contains("unsupported Runtime store schema version 2; expected 3, 4, 5, or 6") ); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_restores_workers_without_legacy_event_or_protocol_observation_logs() { let root = fs_store_root("restore"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: Some("filesystem runtime".to_string()), }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); assert_eq!( runtime.summary().unwrap().backend, RuntimeBackendKind::FsStore ); let transport_bundle = test_bundle(); runtime .store_config_bundle(transport_bundle.clone()) .unwrap(); let worker = runtime.create_worker(task_request("persist me")).unwrap(); runtime .send_input(&worker.worker_ref, WorkerInput::user("first")) .unwrap(); runtime .send_input(&worker.worker_ref, WorkerInput::notify("second")) .unwrap(); runtime .stop_worker(&worker.worker_ref, Some("finished".to_string())) .unwrap(); let worker_store_dir = root.join("workers").join(worker.worker_id.to_string()); let worker_snapshot: serde_json::Value = serde_json::from_slice(&std::fs::read(worker_store_dir.join("worker.json")).unwrap()) .unwrap(); assert_eq!(worker_snapshot["schema_version"], serde_json::json!(6)); assert_eq!(worker_snapshot["status"], serde_json::json!("stopped")); assert_eq!( worker_snapshot["execution"]["last_run_generation"], serde_json::json!(1) ); assert_eq!( worker_snapshot["execution"]["binding"]["run_generation"], serde_json::json!(1) ); assert_eq!( worker_snapshot["execution"]["restore_intent"], serde_json::json!("explicit") ); assert!(!root.join("events.jsonl").exists()); std::fs::write( worker_store_dir.join("observations.jsonl"), b"{\"legacy\":true}\n", ) .unwrap(); std::fs::write(root.join("events.jsonl"), b"obsolete runtime event\n").unwrap(); drop(runtime); let restored = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap(); let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap(); assert_eq!(restored_worker.status, WorkerStatus::Stopped); assert!(matches!( restored.check_config_bundle(&ConfigBundleRef { id: transport_bundle.metadata.id.clone(), digest: transport_bundle.metadata.digest.clone(), }), Err(RuntimeError::ConfigBundleMissing { .. }) )); assert!(!root.join("events.jsonl").exists()); assert!(!worker_store_dir.join("observations.jsonl").exists()); #[cfg(feature = "ws-server")] { let observations = restored .read_worker_observation_events(&worker.worker_ref, WorkerObservationCursor::zero()) .unwrap(); assert!(observations.is_empty()); } #[cfg(feature = "ws-server")] { let observation = restored .observe_worker_event( &worker.worker_ref, protocol::Event::TextDelta { text: "restored observation bus".to_string(), }, ) .unwrap(); assert_eq!(observation.sequence, 1); let observations = restored .read_worker_observation_events(&worker.worker_ref, WorkerObservationCursor::zero()) .unwrap(); assert_eq!(observations.len(), 1); assert_eq!(observations[0].cursor, observation.cursor); } let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_restores_workspace_scope_and_hides_legacy_workers_from_scoped_access() { let root = fs_store_root("workspace-scope"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let scoped = runtime .create_worker_scoped( &scope("workspace-a", "server-a"), scoped_task_request("persist workspace", "workspace-a"), ) .unwrap(); let legacy = runtime .create_worker(task_request("legacy persist")) .unwrap(); let recoverable_legacy = runtime .create_worker(scoped_task_request( "legacy recoverable persist", "workspace-b", )) .unwrap(); drop(runtime); let restored = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap(); let restored_scoped = restored.worker_detail(&scoped.worker_ref).unwrap(); assert_eq!(restored_scoped.workspace_id.as_deref(), Some("workspace-a")); assert_eq!( restored .list_workers_scoped(&scope("workspace-a", "server-a")) .unwrap() .into_iter() .map(|worker| worker.worker_ref) .collect::>(), vec![scoped.worker_ref.clone()] ); let legacy_error = restored .worker_detail_scoped(&scope("workspace-a", "server-a"), &legacy.worker_ref) .unwrap_err(); assert!(matches!(legacy_error, RuntimeError::WorkerNotFound { .. })); let recovered_legacy = restored .worker_detail_scoped( &scope("workspace-b", "server-b"), &recoverable_legacy.worker_ref, ) .unwrap(); assert_eq!( recovered_legacy.workspace_id.as_deref(), Some("workspace-b") ); let stolen_legacy_error = restored .worker_detail_scoped( &scope("workspace-b", "server-c"), &recoverable_legacy.worker_ref, ) .unwrap_err(); assert!(matches!( stolen_legacy_error, RuntimeError::WorkspaceOwnerMismatch { .. } )); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_restores_active_worker_execution_handles() { let root = fs_store_root("execution-restore"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("restore active worker")) .unwrap(); drop(runtime); let backendless = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap(); let persisted_worker = backendless.worker_detail(&worker.worker_ref).unwrap(); assert_eq!(persisted_worker.status, WorkerStatus::Idle); drop(backendless); let restoring_backend = Arc::new(TestExecutionBackend::default()); let restored = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, restoring_backend.clone(), ) .unwrap(); assert_eq!(*restoring_backend.restore_count.lock().unwrap(), 1); let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap(); assert_eq!(restored_worker.status, WorkerStatus::Idle); restored .send_input(&worker.worker_ref, WorkerInput::user("after restart")) .unwrap(); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_automatically_restores_every_active_lifecycle_state() { for status in ["idle", "running", "paused"] { let root = fs_store_root(&format!("automatic-{status}")); let options = crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }; let runtime = Runtime::with_fs_store_and_execution_backend( options.clone(), Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request(&format!("restore {status}"))) .unwrap(); drop(runtime); let worker_path = root .join("workers") .join(worker.worker_id.to_string()) .join("worker.json"); let mut worker_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap(); worker_json["status"] = serde_json::json!(status); std::fs::write( &worker_path, serde_json::to_vec_pretty(&worker_json).unwrap(), ) .unwrap(); let backend = Arc::new(TestExecutionBackend::default()); let restored = Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 1, "status={status}"); assert_eq!( restored.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Idle, "status={status}" ); let _ = std::fs::remove_dir_all(root); } } #[cfg(feature = "fs-store")] #[test] fn fs_store_current_schema_requires_lifecycle_authority() { let root = fs_store_root("current-schema-requires-lifecycle"); let options = crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }; let runtime = Runtime::with_fs_store_and_execution_backend( options.clone(), Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("missing lifecycle authority")) .unwrap(); drop(runtime); let worker_path = root .join("workers") .join(worker.worker_id.to_string()) .join("worker.json"); let mut worker_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap(); worker_json.as_object_mut().unwrap().remove("status"); std::fs::write( &worker_path, serde_json::to_vec_pretty(&worker_json).unwrap(), ) .unwrap(); let restored = Runtime::with_fs_store_and_execution_backend( options, Arc::new(TestExecutionBackend::default()), ) .unwrap(); assert!(restored.list_workers().unwrap().is_empty()); assert!( restored .diagnostics() .unwrap() .iter() .any(|diagnostic| diagnostic.code == "worker_snapshot_ignored") ); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_shutdown_preserves_active_worker_restore_intent() { let root = fs_store_root("shutdown-preserves-worker"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("preserve active worker")) .unwrap(); runtime.stop_runtime().unwrap(); let snapshot: serde_json::Value = serde_json::from_slice( &std::fs::read( root.join("workers") .join(worker.worker_id.to_string()) .join("worker.json"), ) .unwrap(), ) .unwrap(); assert_eq!(snapshot["status"], serde_json::json!("idle")); assert_eq!( snapshot["execution"]["restore_intent"], serde_json::json!("automatic") ); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_stopped_worker_requires_explicit_restore() { let root = fs_store_root("stopped-explicit-restore"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("explicit restore only")) .unwrap(); runtime .stop_worker(&worker.worker_ref, Some("operator stop".to_string())) .unwrap(); drop(runtime); let backend = Arc::new(TestExecutionBackend::default()); let restored = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, backend.clone(), ) .unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 0); assert_eq!( restored.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Stopped ); assert!(matches!( restored.send_input(&worker.worker_ref, WorkerInput::user("implicit restore")), Err(RuntimeError::WorkerExecutionUnavailable { .. }) )); assert_eq!(*backend.restore_count.lock().unwrap(), 0); restored.restore_worker(&worker.worker_ref).unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 1); assert_eq!( restored.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Idle ); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_migrates_schema_v3_workers_to_stopped_explicit_restore() { let root = fs_store_root("schema-v3-restore-intent"); let options = crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }; let runtime = Runtime::with_fs_store_and_execution_backend( options.clone(), Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("schema v3 worker")) .unwrap(); drop(runtime); let runtime_path = root.join("runtime.json"); let worker_path = root .join("workers") .join(worker.worker_id.to_string()) .join("worker.json"); let mut runtime_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&runtime_path).unwrap()).unwrap(); runtime_json["schema_version"] = serde_json::json!(3); std::fs::write( &runtime_path, serde_json::to_vec_pretty(&runtime_json).unwrap(), ) .unwrap(); let mut worker_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap(); worker_json["schema_version"] = serde_json::json!(3); worker_json.as_object_mut().unwrap().remove("status"); worker_json.as_object_mut().unwrap().remove("execution"); worker_json["run_generation"] = serde_json::json!(7); std::fs::write( &worker_path, serde_json::to_vec_pretty(&worker_json).unwrap(), ) .unwrap(); let backend = Arc::new(TestExecutionBackend::default()); let migrated = Runtime::with_fs_store_and_execution_backend(options, backend.clone()).unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 0); assert_eq!( migrated.worker_detail(&worker.worker_ref).unwrap().status, WorkerStatus::Stopped ); let migrated_json: serde_json::Value = serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap(); assert_eq!(migrated_json["schema_version"], serde_json::json!(6)); assert_eq!(migrated_json["status"], serde_json::json!("stopped")); assert_eq!( migrated_json["execution"]["last_run_generation"], serde_json::json!(7) ); assert_eq!( migrated_json["execution"]["binding"], serde_json::Value::Null ); assert_eq!( migrated_json["execution"]["restore_intent"], serde_json::json!("explicit") ); assert!(matches!( migrated.send_input(&worker.worker_ref, WorkerInput::notify("do not restore")), Err(RuntimeError::WorkerExecutionUnavailable { .. }) )); migrated.restore_worker(&worker.worker_ref).unwrap(); assert_eq!(*backend.restore_count.lock().unwrap(), 1); assert_eq!(backend.run_generations.lock().unwrap().as_slice(), &[8]); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_stops_worker_and_reports_when_execution_restore_fails() { let root = fs_store_root("execution-restore-failed"); let runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); runtime.store_config_bundle(test_bundle()).unwrap(); let worker = runtime .create_worker(task_request("restore failure")) .unwrap(); drop(runtime); let restoring_backend = Arc::new(TestExecutionBackend::default()); *restoring_backend.restore_result.lock().unwrap() = Some(WorkerExecutionSpawnResult::Errored( WorkerExecutionResult::errored(WorkerExecutionOperation::Restore, "restore boom"), )); let restored = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, restoring_backend.clone(), ) .unwrap(); assert_eq!(*restoring_backend.restore_count.lock().unwrap(), 1); let restored_worker = restored.worker_detail(&worker.worker_ref).unwrap(); assert_eq!(restored_worker.status, WorkerStatus::Stopped); assert!( restored .diagnostics() .unwrap() .iter() .any( |diagnostic| diagnostic.code == "worker_execution_restore_failed" && diagnostic.worker_ref.as_ref() == Some(&worker.worker_ref) ) ); let err = restored .send_input( &worker.worker_ref, WorkerInput::user("after failed restore"), ) .unwrap_err(); assert!(matches!( err, RuntimeError::WorkerExecutionUnavailable { .. } )); let _ = std::fs::remove_dir_all(root); } #[cfg(feature = "fs-store")] #[test] fn fs_store_reports_corrupt_and_missing_data() { let corrupt_root = fs_store_root("corrupt"); let corrupt_runtime = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: corrupt_root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap(); let corrupt_store = runtime_store(&corrupt_runtime); std::fs::write( corrupt_store.runtime_dir().join("runtime.json"), b"not json", ) .unwrap(); drop(corrupt_runtime); let err = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: corrupt_root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .unwrap_err(); assert!(matches!(err, RuntimeError::StoreCorrupt { .. })); let _ = std::fs::remove_dir_all(corrupt_root); let missing_root = fs_store_root("missing"); let missing_runtime = Runtime::with_fs_store_and_execution_backend( crate::fs_store::FsRuntimeStoreOptions { root: missing_root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }, Arc::new(TestExecutionBackend::default()), ) .unwrap(); missing_runtime.store_config_bundle(test_bundle()).unwrap(); missing_runtime .create_worker(task_request("missing worker snapshot")) .unwrap(); let missing_store = runtime_store(&missing_runtime); let mut worker_dirs = std::fs::read_dir(missing_store.runtime_dir().join("workers")) .unwrap() .collect::, _>>() .unwrap(); worker_dirs.sort_by_key(|entry| entry.path()); std::fs::remove_file(worker_dirs[0].path().join("worker.json")).unwrap(); drop(missing_runtime); let loaded = Runtime::with_fs_store(crate::fs_store::FsRuntimeStoreOptions { root: missing_root.clone(), runtime_id: "test-runtime".to_string(), display_name: None, }) .expect("invalid worker snapshot should not make runtime store unreadable"); assert!(loaded.list_workers().unwrap().is_empty()); assert!( loaded .diagnostics() .unwrap() .iter() .any(|diagnostic| diagnostic.code == "worker_snapshot_ignored") ); let _ = std::fs::remove_dir_all(missing_root); } }