From 2fd043b6341dcb9020ab86ee1ed26c27b3bc2000 Mon Sep 17 00:00:00 2001 From: Hare Date: Mon, 7 Sep 2026 19:43:26 +0900 Subject: [PATCH] fix: expand worker creation timeout budgets --- crates/worker-runtime/src/runtime.rs | 34 +++++++++++++++++++++ crates/worker-runtime/src/worker_backend.rs | 32 +++++++++++-------- crates/workspace-server/src/hosts.rs | 11 ++++++- 3 files changed, 63 insertions(+), 14 deletions(-) diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 3ed3f76f..adc6584c 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -3408,6 +3408,25 @@ fn validate_create_worker_request(request: &CreateWorkerRequest) -> Result<(), R 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) { @@ -3453,6 +3472,12 @@ fn write_runtime_worker_create_failure( error: &RuntimeError, ) { 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", @@ -3462,6 +3487,7 @@ fn write_runtime_worker_create_failure( error_kind, operation = operation.as_deref().unwrap_or(""), outcome = outcome.as_deref().unwrap_or(""), + execution_failure_code, "Worker creation failed" ); } @@ -3597,6 +3623,14 @@ mod tests { 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 { diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 5697e7d7..374fd6ea 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -89,11 +89,12 @@ use worker::{ const DEFAULT_BACKEND_ID: &str = "worker-crate"; const RUNTIME_TASK_TIMEOUT: Duration = Duration::from_secs(10); +const SPAWN_RESTORE_TASK_TIMEOUT: Duration = Duration::from_secs(60); +const USER_INPUT_TASK_TIMEOUT: Duration = Duration::from_secs(35); const WORKSPACE_CONFIG_HTTP_TIMEOUT: Duration = Duration::from_secs(8); const MAX_WORKSPACE_CONFIG_RESPONSE_BYTES: usize = 72 * 1024 * 1024; -// Keep this below the adapter task timeout so a failed acknowledgement task -// returns a typed execution error instead of leaving the outer waiter to time out. -const USER_INPUT_COMMIT_TIMEOUT: Duration = Duration::from_secs(9); +// Leave adapter cancellation margin after the durable submission deadline. +const USER_INPUT_COMMIT_TIMEOUT: Duration = Duration::from_secs(30); pub struct RuntimeWorkerController { pub handle: WorkerHandle, @@ -1244,7 +1245,7 @@ where working_directory_materializer: None, runtime: Mutex::new(Some(runtime)), workers: Mutex::new(HashMap::new()), - spawn_restore_timeout: RUNTIME_TASK_TIMEOUT, + spawn_restore_timeout: SPAWN_RESTORE_TASK_TIMEOUT, }) } @@ -1305,12 +1306,15 @@ where Self::wait_for_runtime_task(rx) } - fn run_spawn_restore_on_adapter_runtime(&self, task: Fut) -> Result + fn run_cancellable_on_adapter_runtime( + &self, + timeout: Duration, + 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); @@ -1406,7 +1410,7 @@ where submission_request_id: String, ) -> WorkerExecutionResult { let request_id = submission_request_id.clone(); - self.run_on_adapter_runtime(async move { + self.run_cancellable_on_adapter_runtime(USER_INPUT_TASK_TIMEOUT, async move { // Subscribe before enqueueing so a fast durable acceptance cannot // race the Runtime acknowledgement. let mut events = worker.subscribe(); @@ -1797,9 +1801,10 @@ where let factory = self.factory.clone(); let bridge_context = request.context.clone(); let worker_ref = request.worker_ref.clone(); - let spawn_result = self.run_spawn_restore_on_adapter_runtime(async move { - factory.spawn_controller(request).await - }); + let spawn_result = self + .run_cancellable_on_adapter_runtime(self.spawn_restore_timeout, async move { + factory.spawn_controller(request).await + }); let controller = match spawn_result { Ok(controller) => controller, @@ -1901,9 +1906,10 @@ where let factory = self.factory.clone(); let bridge_context = request.context.clone(); let worker_ref = request.worker_ref.clone(); - let restore_result = self.run_spawn_restore_on_adapter_runtime(async move { - factory.restore_controller(request).await - }); + let restore_result = self + .run_cancellable_on_adapter_runtime(self.spawn_restore_timeout, async move { + factory.restore_controller(request).await + }); let controller = match restore_result { Ok(controller) => controller, diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 43142bcf..f1459e41 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -69,6 +69,9 @@ const EMBEDDED_HOST_KIND: &str = "embedded-worker-runtime-host"; const REMOTE_HOST_KIND: &str = "remote-worker-runtime-host"; const MAX_DIAGNOSTICS: usize = 16; const MAX_RUNTIME_PING_RESPONSE_BYTES: usize = 8 * 1024; +// Runtime creation can spend up to 60s bootstrapping, 35s waiting for the +// durable initial-input acknowledgement, and 5s confirming shutdown. +const REMOTE_WORKER_CREATE_TIMEOUT: Duration = Duration::from_secs(105); const MAX_HOST_SCAN: usize = 256; const MAX_IDENTIFIER_LEN: usize = 120; const ID_DIGEST_HEX_LEN: usize = 16; @@ -3769,7 +3772,13 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { workspace_api: Some(workspace_api), memory_settings: request.resolved_memory_settings.clone(), }; - match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) { + match self.send_json::( + "/v1/workers", + self.http + .post(self.endpoint("/v1/workers")) + .timeout(REMOTE_WORKER_CREATE_TIMEOUT) + .json(&create), + ) { Ok(response) => WorkerSpawnResult { state: WorkerOperationState::Accepted, worker: Some(self.map_worker_detail(response.worker)),