diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 4db8d8d3..9722a3ea 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -153,6 +153,7 @@ impl Drop for RuntimeEventSelectorSubscription { #[derive(Clone, Debug)] pub struct Runtime { inner: Arc>, + worker_operations: Arc>>>>, } impl Runtime { @@ -166,6 +167,7 @@ impl Runtime { let state = RuntimeState::new(options.display_name); Self { inner: Arc::new(Mutex::new(state)), + worker_operations: Arc::new(Mutex::new(BTreeMap::new())), } } @@ -222,6 +224,7 @@ impl Runtime { 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) @@ -781,6 +784,10 @@ impl Runtime { 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); } @@ -893,8 +900,7 @@ impl Runtime { 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() { - let _ = backend.stop_worker(&handle); - self.rollback_failed_create(&worker_ref)?; + self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; return Err(RuntimeError::WorkerExecutionRejected { worker_id: worker_ref.worker_id.clone(), operation: dispatch_result.operation, @@ -908,8 +914,7 @@ impl Runtime { .as_ref() .is_some_and(|ack| ack.submission_request_id == expected_submission_id); if !has_commit_ack { - let _ = backend.stop_worker(&handle); - self.rollback_failed_create(&worker_ref)?; + self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; let result = WorkerExecutionResult::rejected( WorkerExecutionOperation::Input, "execution backend accepted initial input without a durable session commit acknowledgement", @@ -922,21 +927,36 @@ impl Runtime { result, }); } - let detail = self.commit_created_worker( + let detail = match self.commit_created_worker( &worker_ref, - handle, + handle.clone(), working_directory, dispatch_result, - )?; - self.record_input_observation(&worker_ref, initial_input)?; + ) { + 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 { - self.commit_created_worker( + match self.commit_created_worker( &worker_ref, - handle, + handle.clone(), working_directory, WorkerExecutionResult::accepted(WorkerExecutionOperation::Spawn), - ) + ) { + Ok(detail) => Ok(detail), + Err(error) => { + self.cleanup_connected_failed_create(&backend, &worker_ref, &handle)?; + Err(error) + } + } } } @@ -1639,14 +1659,43 @@ impl Runtime { 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 let Some(record) = state.workers.remove(&worker_ref.worker_id) { + 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(()) } @@ -1800,6 +1849,10 @@ impl Runtime { &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 mut state = self.lock()?; state.ensure_running()?; state.ensure_worker_ref(worker_ref)?; @@ -2260,6 +2313,17 @@ impl Runtime { 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) } diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 860496f9..5697e7d7 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -1216,6 +1216,7 @@ pub struct WorkerRuntimeExecutionBackend { working_directory_materializer: Option>, runtime: Mutex>, workers: Mutex>, + spawn_restore_timeout: Duration, } impl WorkerRuntimeExecutionBackend { @@ -1243,6 +1244,7 @@ where working_directory_materializer: None, runtime: Mutex::new(Some(runtime)), workers: Mutex::new(HashMap::new()), + spawn_restore_timeout: RUNTIME_TASK_TIMEOUT, }) } @@ -1259,6 +1261,12 @@ where self } + #[cfg(test)] + fn with_spawn_restore_timeout(mut self, timeout: Duration) -> Self { + self.spawn_restore_timeout = timeout; + self + } + fn wait_for_runtime_task(receiver: mpsc::Receiver>) -> Result { receiver .recv_timeout(RUNTIME_TASK_TIMEOUT) @@ -1297,6 +1305,39 @@ where Self::wait_for_runtime_task(rx) } + fn run_spawn_restore_on_adapter_runtime(&self, task: Fut) -> Result + where + T: Send + 'static, + Fut: Future> + Send + 'static, + { + let timeout = self.spawn_restore_timeout; + let (tx, rx) = mpsc::sync_channel(1); + self.spawn_on_adapter_runtime(async move { + let mut handle = tokio::spawn(task); + let result = tokio::select! { + biased; + result = &mut handle => match result { + Ok(result) => result, + Err(err) => Err(format!("worker adapter task failed: {err}")), + }, + _ = tokio::time::sleep(timeout) => { + handle.abort(); + match handle.await { + Ok(result) => result, + Err(err) if err.is_cancelled() => Err(format!( + "worker adapter task did not complete within {} seconds and was cancelled", + timeout.as_secs_f64() + )), + Err(err) => Err(format!("worker adapter task failed: {err}")), + } + } + }; + let _ = tx.send(result); + })?; + rx.recv() + .map_err(|err| format!("worker adapter task did not complete: {err}"))? + } + fn get_execution( &self, handle: &WorkerExecutionHandle, @@ -1756,8 +1797,9 @@ where let factory = self.factory.clone(); let bridge_context = request.context.clone(); let worker_ref = request.worker_ref.clone(); - let spawn_result = - self.run_on_adapter_runtime(async move { factory.spawn_controller(request).await }); + let spawn_result = self.run_spawn_restore_on_adapter_runtime(async move { + factory.spawn_controller(request).await + }); let controller = match spawn_result { Ok(controller) => controller, @@ -1859,8 +1901,9 @@ where let factory = self.factory.clone(); let bridge_context = request.context.clone(); let worker_ref = request.worker_ref.clone(); - let restore_result = - self.run_on_adapter_runtime(async move { factory.restore_controller(request).await }); + let restore_result = self.run_spawn_restore_on_adapter_runtime(async move { + factory.restore_controller(request).await + }); let controller = match restore_result { Ok(controller) => controller, @@ -2100,21 +2143,24 @@ where if let Err(message) = shutdown_wait { return WorkerExecutionResult::errored(WorkerExecutionOperation::Stop, message); } - if let Err(error) = artifact_cleanup.delete_uncommitted_uploaded_files() { - return WorkerExecutionResult::errored( - WorkerExecutionOperation::Stop, - format!("uploaded_file_cleanup_failed: {error}"), - ); - } + let artifact_cleanup_error = artifact_cleanup + .delete_uncommitted_uploaded_files() + .err() + .map(|error| format!("uploaded_file_cleanup_failed: {error}")); match self.workers.lock() { Ok(mut workers) => { workers.remove(handle.worker_ref()); - result } - Err(_) => WorkerExecutionResult::errored( - WorkerExecutionOperation::Stop, - "worker adapter registry lock is poisoned after shutdown", - ), + Err(poisoned) => { + poisoned.into_inner().remove(handle.worker_ref()); + } + } + if let Some(message) = artifact_cleanup_error { + let mut result = result; + result.message = Some(message); + result + } else { + result } } @@ -2178,7 +2224,7 @@ mod tests { use std::fs; use std::pin::Pin; use std::process::Command; - use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use crate::Runtime as EmbeddedRuntime; use crate::catalog::{ @@ -2489,6 +2535,87 @@ mod tests { WorkerExecutionContext::new(worker_ref) } + struct DelayedFactory { + completed: Arc, + delay: Duration, + } + + #[async_trait] + impl RuntimeWorkerFactory for DelayedFactory { + async fn spawn_controller( + &self, + _request: WorkerExecutionSpawnRequest, + ) -> Result { + tokio::time::sleep(self.delay).await; + self.completed.store(true, Ordering::SeqCst); + Err("delayed factory completed".to_string()) + } + + async fn restore_controller( + &self, + _request: WorkerExecutionRestoreRequest, + ) -> Result { + tokio::time::sleep(self.delay).await; + self.completed.store(true, Ordering::SeqCst); + Err("delayed factory completed".to_string()) + } + } + + #[test] + fn create_timeout_cancels_factory_and_removes_persisted_worker() { + let root = tempfile::tempdir().unwrap(); + let runtime_store_dir = root.path().join("runtime"); + let completed = Arc::new(AtomicBool::new(false)); + let backend = Arc::new( + WorkerRuntimeExecutionBackend::new(DelayedFactory { + completed: completed.clone(), + delay: Duration::from_millis(200), + }) + .unwrap() + .with_spawn_restore_timeout(Duration::from_millis(20)), + ); + let runtime = EmbeddedRuntime::with_fs_store_and_execution_backend( + crate::fs_store::FsRuntimeStoreOptions { + root: runtime_store_dir.clone(), + runtime_id: "create-timeout-runtime".to_string(), + display_name: None, + }, + backend.clone(), + ) + .unwrap(); + runtime.store_config_bundle(test_bundle()).unwrap(); + let request = create_request("create timeout"); + let worker_id = request.worker_id; + let create_runtime = runtime.clone(); + let create = std::thread::spawn(move || create_runtime.create_worker(request)); + std::thread::sleep(Duration::from_millis(5)); + + let delete_error = runtime + .delete_worker(&crate::identity::WorkerRef::new(worker_id)) + .unwrap_err(); + let error = create.join().unwrap().unwrap_err(); + + assert!(matches!( + delete_error, + crate::error::RuntimeError::WorkerNotFound { worker_id: missing } if missing == worker_id + )); + assert!(error.to_string().contains("was cancelled")); + std::thread::sleep(Duration::from_millis(250)); + assert!( + !completed.load(Ordering::SeqCst), + "timed out factory future must not resume after create returns" + ); + assert!(runtime.list_workers().unwrap().is_empty()); + assert!(backend.workers.lock().unwrap().is_empty()); + assert!( + !runtime_store_dir + .join("workers") + .join(worker_id.to_string()) + .exists(), + "failed create must remove its persisted Worker aggregate" + ); + } + struct MockFactory { client: MockClient, runtime_base: PathBuf, diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 69db8362..5e73ecd4 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -2338,14 +2338,6 @@ impl WorkspaceApi { Uuid::new_v4().to_string(), ) }); - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - self.store.reserve_worker_workdir_attachment( - &self.config.workspace_id, - workdir_id, - reservation_id, - &now_registry_timestamp(), - )?; - } let request_fingerprint = worker_spawn_create_fingerprint(&request) .map_err(|message| Error::Config(message.to_string()))?; let current_memory_settings = self @@ -2378,54 +2370,127 @@ impl WorkspaceApi { })?; let worker_id = reservation.worker_id; request.resolved_memory_settings = Some(reservation.memory_settings); + let reservation_fingerprint = reservation.create_fingerprint.clone(); let create_binding = WorkerCreateBinding { worker_id, create_fingerprint: reservation.create_fingerprint, }; - let result = match self + let compensation_context = WorkerSpawnCompensationContext { + assignment: None, + prepared_workdir_id: attachment_reservation + .as_ref() + .map(|(workdir_id, _)| workdir_id.as_str()), + cleanup_spawned_workdir: false, + }; + if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() + && let Err(error) = self.store.reserve_worker_workdir_attachment( + &self.config.workspace_id, + workdir_id, + reservation_id, + &now_registry_timestamp(), + ) + { + let mut diagnostics = Vec::new(); + if let Err(cleanup_error) = self.config_store.fail_worker_create_reservation( + &self.config.workspace_id, + runtime_id, + worker_id, + &reservation_fingerprint, + ) { + diagnostics.push(spawn_compensation_diagnostic( + "worker_spawn_compensation_create_reservation_release_failed", + format!( + "Failed to terminalize Worker create reservation {} and release its resource key: {}", + worker_id, + sanitize_backend_error(&cleanup_error.to_string()) + ), + )); + } + return Err(ApiError::with_diagnostics(error, diagnostics)); + } + let mut result = match self .runtime .spawn_worker(runtime_id, create_binding, request) { Ok(result) => result, Err(error) => { - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - ); - } - return Err(error.into_error().into()); + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + None, + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) + }), + ); + return Err(ApiError::with_diagnostics(error.into_error(), diagnostics)); } }; let Some(worker) = result.worker.as_ref() else { - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - )?; - } + result + .diagnostics + .extend(compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + None, + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) + }), + )); return Ok(result); }; let worker_ref = worker.worker.clone(); if worker_ref.worker_id != worker_id.to_string() { - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - ); - } - return Err(Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), - code: "workspace_worker_identity_mismatch".to_string(), - message: format!( - "Runtime returned Worker {} for reserved Workspace Worker {}", - worker_ref.worker_id, worker_id - ), - } - .into()); + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) + }), + ); + return Err(ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "workspace_worker_identity_mismatch".to_string(), + message: format!( + "Runtime returned Worker {} for reserved Workspace Worker {}", + worker_ref.worker_id, worker_id + ), + }, + diagnostics, + )); + } + if result.state != WorkerOperationState::Accepted { + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) + }), + ); + result.diagnostics.extend(diagnostics); + return Ok(result); } let replacement = match self .runtime @@ -2433,48 +2498,63 @@ impl WorkspaceApi { { Ok(replacement) => replacement, Err(error) => { - let _ = self.runtime.delete_worker(&worker_ref); - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - ); - } - return Err(error.into_error().into()); + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) + }), + ); + return Err(ApiError::with_diagnostics(error.into_error(), diagnostics)); } }; if replacement.state != WorkerOperationState::Accepted { - let _ = self.runtime.delete_worker(&worker_ref); - if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, - ); - } - return Err(Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), - code: "worker_workspace_api_replace_failed".to_string(), - message: replacement - .diagnostics - .first() - .map(|diagnostic| diagnostic.message.clone()) - .unwrap_or_else(|| { - "Runtime rejected Workspace API replacement after spawn".to_string() + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + attachment_reservation + .as_ref() + .map(|(workdir_id, reservation_id)| { + (workdir_id.as_str(), reservation_id.as_str()) }), - } - .into()); + ); + return Err(ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "worker_workspace_api_replace_failed".to_string(), + message: replacement + .diagnostics + .first() + .map(|diagnostic| diagnostic.message.clone()) + .unwrap_or_else(|| { + "Runtime rejected Workspace API replacement after spawn".to_string() + }), + }, + diagnostics, + )); } if let Some((workdir_id, reservation_id)) = attachment_reservation.as_ref() { if let Err(error) = parse_runtime_worker_id_for_registry(&worker.worker.worker_id) { - let _ = self.runtime.delete_worker(&worker_ref); - let _ = self.store.release_worker_workdir_attachment_reservation( - &self.config.workspace_id, - workdir_id, - reservation_id, + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + Some((workdir_id.as_str(), reservation_id.as_str())), ); - return Err(error); + return Err(api_error_with_additional_diagnostics(error, diagnostics)); } let compensation_context = WorkerSpawnCompensationContext { assignment: None, @@ -2489,20 +2569,23 @@ impl WorkspaceApi { WorkerRegistryDisplayNamePolicy::PreserveExisting, ) .map(|_| ()); - if let Err(mut error) = finalize_worker_spawn_stage( + if let Err(error) = finalize_worker_spawn_stage( self, worker, &compensation_context, WorkerSpawnFinalizeStage::WorkerRegistry, registry_result, ) { - append_attachment_reservation_release_diagnostic( + let diagnostics = compensate_failed_workspace_worker_create( self, - workdir_id, - reservation_id, - &mut error, + runtime_id, + worker_id, + &reservation_fingerprint, + None, + &compensation_context, + Some((workdir_id.as_str(), reservation_id.as_str())), ); - return Err(error); + return Err(api_error_with_additional_diagnostics(error, diagnostics)); } let attachment = WorkerWorkdirLinkRecord { @@ -2517,25 +2600,43 @@ impl WorkspaceApi { .store .finalize_reserved_worker_workdir_attachment(&attachment, reservation_id) .map_err(ApiError::from); - if let Err(mut error) = finalize_worker_spawn_stage( + if let Err(error) = finalize_worker_spawn_stage( self, worker, &compensation_context, WorkerSpawnFinalizeStage::WorkdirAttachment, attachment_result, ) { - append_attachment_reservation_release_diagnostic( + let diagnostics = compensate_failed_workspace_worker_create( self, - workdir_id, - reservation_id, - &mut error, + runtime_id, + worker_id, + &reservation_fingerprint, + None, + &compensation_context, + Some((workdir_id.as_str(), reservation_id.as_str())), ); - return Err(error); + return Err(api_error_with_additional_diagnostics(error, diagnostics)); } } - self.config_store + if let Err(error) = self + .config_store .complete_worker_create_reservation(&self.config.workspace_id, worker_id) - .map_err(|error| Error::Config(error.to_string()))?; + { + let diagnostics = compensate_failed_workspace_worker_create( + self, + runtime_id, + worker_id, + &reservation_fingerprint, + Some(worker), + &compensation_context, + None, + ); + return Err(ApiError::with_diagnostics( + Error::Config(error.to_string()), + diagnostics, + )); + } Ok(result) } @@ -14174,25 +14275,87 @@ fn finalize_worker_spawn_stage( )) } -fn append_attachment_reservation_release_diagnostic( +fn api_error_with_additional_diagnostics( + mut error: ApiError, + diagnostics: Vec, +) -> ApiError { + error.diagnostics.extend(diagnostics); + error +} + +fn compensate_failed_workspace_worker_create( api: &WorkspaceApi, - workdir_id: &str, - reservation_id: &str, - error: &mut ApiError, -) { - if let Err(release_error) = api.store.release_worker_workdir_attachment_reservation( - &api.config.workspace_id, - workdir_id, - reservation_id, - ) { - error.diagnostics.push(spawn_compensation_diagnostic( + runtime_id: &str, + reservation_worker_id: WorkerId, + create_fingerprint: &str, + worker: Option<&WorkerSummary>, + context: &WorkerSpawnCompensationContext<'_>, + attachment_reservation: Option<(&str, &str)>, +) -> Vec { + let mut diagnostics = Vec::new(); + let reserved_worker_ref = + RuntimeWorkerRef::new(runtime_id.to_string(), reservation_worker_id.to_string()); + let mut reserved_worker_absent = false; + if let Some(worker) = worker { + let (returned_worker_absent, delete_diagnostics) = + delete_runtime_worker_for_spawn_compensation(api, &worker.worker); + diagnostics.extend(delete_diagnostics); + if returned_worker_absent { + diagnostics.extend(finalize_spawn_compensation_after_worker_delete( + api, worker, context, + )); + } + if worker.worker == reserved_worker_ref { + reserved_worker_absent = returned_worker_absent; + } + } + if !reserved_worker_absent { + let (worker_absent, delete_diagnostics) = + delete_runtime_worker_for_spawn_compensation(api, &reserved_worker_ref); + reserved_worker_absent = worker_absent; + diagnostics.extend(delete_diagnostics); + } + if let Some((workdir_id, reservation_id)) = attachment_reservation + && let Err(error) = api.store.release_worker_workdir_attachment_reservation( + &api.config.workspace_id, + workdir_id, + reservation_id, + ) + { + diagnostics.push(spawn_compensation_diagnostic( "worker_spawn_compensation_attachment_reservation_release_failed", format!( "Failed to release Workdir `{workdir_id}` attachment reservation `{reservation_id}`: {}", - sanitize_backend_error(&release_error.to_string()) + sanitize_backend_error(&error.to_string()) ), )); } + if reserved_worker_absent { + if let Err(error) = api.config_store.fail_worker_create_reservation( + &api.config.workspace_id, + runtime_id, + reservation_worker_id, + create_fingerprint, + ) { + diagnostics.push(spawn_compensation_diagnostic( + "worker_spawn_compensation_create_reservation_release_failed", + format!( + "Failed to terminalize Worker create reservation {} and release its resource key: {}", + reservation_worker_id, + sanitize_backend_error(&error.to_string()) + ), + )); + } + } else { + diagnostics.push(spawn_compensation_diagnostic( + "worker_spawn_compensation_create_reservation_retained", + format!( + "Retained Worker create reservation {} and its resource key because Runtime Worker absence could not be confirmed", + reservation_worker_id + ), + )); + } + diagnostics } fn compensate_failed_worker_spawn( @@ -14200,6 +14363,21 @@ fn compensate_failed_worker_spawn( worker: &WorkerSummary, context: &WorkerSpawnCompensationContext<'_>, ) -> Vec { + let (runtime_deleted, mut diagnostics) = + delete_runtime_worker_for_spawn_compensation(api, &worker.worker); + if !runtime_deleted { + return diagnostics; + } + diagnostics.extend(finalize_spawn_compensation_after_worker_delete( + api, worker, context, + )); + diagnostics +} + +fn delete_runtime_worker_for_spawn_compensation( + api: &WorkspaceApi, + worker_ref: &RuntimeWorkerRef, +) -> (bool, Vec) { let mut diagnostics = Vec::new(); let lifecycle_request = WorkerLifecycleRequest { reason: Some("Backend spawn finalize failed; compensating Runtime Worker".to_string()), @@ -14207,8 +14385,8 @@ fn compensate_failed_worker_spawn( }; let cancellation = api .runtime - .cancel_worker(&worker.worker, lifecycle_request.clone()); - let stop = api.runtime.stop_worker(&worker.worker, lifecycle_request); + .cancel_worker(worker_ref, lifecycle_request.clone()); + let stop = api.runtime.stop_worker(worker_ref, lifecycle_request); let stop_accepted = stop .as_ref() .is_ok_and(|result| result.state == WorkerOperationState::Accepted); @@ -14218,12 +14396,12 @@ fn compensate_failed_worker_spawn( format!("{cancellation}; {stop}") }); - let runtime_deleted = match api.runtime.delete_worker(&worker.worker) { + let runtime_deleted = match api.runtime.delete_worker(worker_ref) { Ok(result) if result.state == WorkerOperationState::Accepted && result.deleted => true, Ok(result) => { let mut message = format!( "Runtime did not delete Worker {}:{}: state={:?}, deleted={}", - worker.worker.runtime_id, worker.worker.worker_id, result.state, result.deleted + worker_ref.runtime_id, worker_ref.worker_id, result.state, result.deleted ); if let Some(detail) = termination_detail.as_deref() { message.push_str(&format!("; cancellation: {detail}")); @@ -14244,8 +14422,8 @@ fn compensate_failed_worker_spawn( Err(error) => { let mut message = format!( "Failed to delete Runtime Worker {}:{}: {}", - worker.worker.runtime_id, - worker.worker.worker_id, + worker_ref.runtime_id, + worker_ref.worker_id, error.message() ); if let Some(detail) = termination_detail.as_deref() { @@ -14258,13 +14436,7 @@ fn compensate_failed_worker_spawn( false } }; - if !runtime_deleted { - return diagnostics; - } - diagnostics.extend(finalize_spawn_compensation_after_worker_delete( - api, worker, context, - )); - diagnostics + (runtime_deleted, diagnostics) } fn finalize_spawn_compensation_after_worker_delete( @@ -14337,6 +14509,20 @@ fn finalize_spawn_compensation_after_worker_delete( } } } + if let Ok(worker_id) = parse_runtime_worker_id_for_registry(&worker.worker.worker_id) + && let Err(error) = api + .config_store + .release_removed_worker_resource_key(&api.config.workspace_id, worker_id) + { + diagnostics.push(spawn_compensation_diagnostic( + "worker_spawn_compensation_resource_key_release_failed", + format!( + "Failed to release removed Workspace Worker resource key {}: {}", + worker.worker.worker_id, + sanitize_backend_error(&error.to_string()) + ), + )); + } diagnostics } @@ -18886,6 +19072,25 @@ mod tests { })); assert!(find_workspace_orchestrator(&api).is_none()); assert!(!workspace_orchestrator_response(&api, "failed").online); + let (reserved_creates, worker_resource_keys): (i64, i64) = api + .config_store + .with_conn(|conn| { + Ok(( + conn.query_row( + "SELECT COUNT(*) FROM worker_create_reservations WHERE state = 'reserved'", + [], + |row| row.get(0), + )?, + conn.query_row( + "SELECT COUNT(*) FROM workspace_resource_keys WHERE resource_kind = 'worker'", + [], + |row| row.get(0), + )?, + )) + }) + .unwrap(); + assert_eq!(reserved_creates, 0); + assert_eq!(worker_resource_keys, 0); let Json(retried) = scoped_start_workspace_orchestrator( State(api), diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 4b22b1a8..cdf6c2bb 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -1468,6 +1468,146 @@ impl SqliteWorkspaceStore { }) } + pub(crate) fn fail_worker_create_reservation( + &self, + workspace_id: &str, + runtime_id: &str, + worker_id: WorkerId, + create_fingerprint: &str, + ) -> Result<()> { + self.with_conn_mut(|conn| { + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + let reservation = tx + .query_row( + "SELECT runtime_id, create_fingerprint, state \ + FROM worker_create_reservations \ + WHERE workspace_id = ?1 AND worker_id = ?2", + params![workspace_id, worker_id.to_string()], + |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + )) + }, + ) + .optional()? + .ok_or_else(|| { + Error::Store(format!( + "Worker create reservation {} was not found", + worker_id + )) + })?; + if reservation.0 != runtime_id || reservation.1 != create_fingerprint { + return Err(Error::InvalidInput(format!( + "Worker create reservation {} does not match the failed create", + worker_id + ))); + } + let registry_exists = tx + .query_row( + "SELECT 1 FROM worker_registry \ + WHERE workspace_id = ?1 AND worker_id = ?2 LIMIT 1", + params![workspace_id, worker_id.to_string()], + |_| Ok(()), + ) + .optional()? + .is_some(); + if registry_exists { + return Err(Error::InvalidInput(format!( + "Worker create reservation {} cannot fail while its registry row exists", + worker_id + ))); + } + match reservation.2.as_str() { + "reserved" | "created" => { + tx.execute( + "UPDATE worker_create_reservations \ + SET state = 'removed', updated_at = ?3 \ + WHERE workspace_id = ?1 AND worker_id = ?2 AND state IN ('reserved', 'created')", + params![ + workspace_id, + worker_id.to_string(), + chrono::Utc::now().to_rfc3339() + ], + )?; + } + "removed" => {} + state => { + return Err(Error::InvalidInput(format!( + "Worker create reservation {} cannot fail from state {state}", + worker_id + ))); + } + } + tx.execute( + "DELETE FROM workspace_resource_keys \ + WHERE workspace_id = ?1 AND resource_kind = ?2 AND resource_id = ?3", + params![ + workspace_id, + WorkspaceResourceKind::Worker.as_str(), + worker_id.to_string() + ], + )?; + tx.commit()?; + Ok(()) + }) + } + + pub(crate) fn release_removed_worker_resource_key( + &self, + workspace_id: &str, + worker_id: WorkerId, + ) -> Result { + self.with_conn_mut(|conn| { + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + let reservation_state = tx + .query_row( + "SELECT state FROM worker_create_reservations \ + WHERE workspace_id = ?1 AND worker_id = ?2", + params![workspace_id, worker_id.to_string()], + |row| row.get::<_, String>(0), + ) + .optional()?; + let Some(reservation_state) = reservation_state else { + tx.commit()?; + return Ok(false); + }; + if reservation_state != "removed" { + return Err(Error::InvalidInput(format!( + "Worker resource key {} cannot be released from create reservation state {reservation_state}", + worker_id + ))); + } + let registry_exists = tx + .query_row( + "SELECT 1 FROM worker_registry \ + WHERE workspace_id = ?1 AND worker_id = ?2 LIMIT 1", + params![workspace_id, worker_id.to_string()], + |_| Ok(()), + ) + .optional()? + .is_some(); + if registry_exists { + return Err(Error::InvalidInput(format!( + "Worker resource key {} cannot be released while its registry row exists", + worker_id + ))); + } + let released = tx.execute( + "DELETE FROM workspace_resource_keys \ + WHERE workspace_id = ?1 AND resource_kind = ?2 AND resource_id = ?3", + params![ + workspace_id, + WorkspaceResourceKind::Worker.as_str(), + worker_id.to_string() + ], + )? > 0; + tx.commit()?; + Ok(released) + }) + } + fn materialize_workspace_config(&self, workspace_id: &str, created_at: &str) -> Result<()> { self.with_conn_mut(|conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; @@ -8680,6 +8820,58 @@ mod tests { .as_deref(), Some("W-2") ); + store + .fail_worker_create_reservation( + "workspace-a", + "arcadia", + second.worker_id, + &second.create_fingerprint, + ) + .unwrap(); + store + .fail_worker_create_reservation( + "workspace-a", + "arcadia", + second.worker_id, + &second.create_fingerprint, + ) + .unwrap(); + assert!( + !store + .has_active_worker_create_reservation( + "workspace-a", + &RuntimeWorkerRef::new("arcadia", second.worker_id.to_string()) + ) + .unwrap() + ); + assert_eq!( + store + .resource_key( + "workspace-a", + WorkspaceResourceKind::Worker, + &second.worker_id.to_string() + ) + .unwrap(), + None + ); + assert_eq!( + store + .resolve_resource_reference("workspace-a", WorkspaceResourceKind::Worker, "W-2") + .unwrap(), + None + ); + let failed_state: String = store + .with_conn(|conn| { + conn.query_row( + "SELECT state FROM worker_create_reservations \ + WHERE workspace_id = 'workspace-a' AND worker_id = ?1", + [second.worker_id.to_string()], + |row| row.get(0), + ) + .map_err(Error::from) + }) + .unwrap(); + assert_eq!(failed_state, "removed"); store .complete_worker_create_reservation("workspace-a", reserved.worker_id) .unwrap(); @@ -8712,6 +8904,16 @@ mod tests { Ok(()) }) .unwrap(); + assert!( + store + .fail_worker_create_reservation( + "workspace-a", + "arcadia", + reserved.worker_id, + &reserved.create_fingerprint, + ) + .is_err() + ); assert!( store .delete_worker_registry("workspace-a", &reserved_worker) @@ -8729,6 +8931,29 @@ mod tests { }) .unwrap(); assert_eq!(removed_state, "removed"); + assert!( + store + .release_removed_worker_resource_key("workspace-a", reserved.worker_id) + .unwrap() + ); + store + .fail_worker_create_reservation( + "workspace-a", + "arcadia", + reserved.worker_id, + &reserved.create_fingerprint, + ) + .unwrap(); + assert_eq!( + store + .resource_key( + "workspace-a", + WorkspaceResourceKind::Worker, + &reserved.worker_id.to_string(), + ) + .unwrap(), + None + ); assert!( store .reserve_worker_create(