From e88b0cda18e4db813aa47af9cc41b527b53aceb5 Mon Sep 17 00:00:00 2001 From: Hare Date: Tue, 11 Aug 2026 04:58:34 +0900 Subject: [PATCH] fix: restore production orchestrator spawn --- crates/workspace-server/src/hosts.rs | 67 +++++++++++++- crates/workspace-server/src/server.rs | 127 ++++++++++++++++++++++++++ resources/profiles/orchestrator.dcdl | 2 +- 3 files changed, 192 insertions(+), 4 deletions(-) diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index d86ad12b..b2046416 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -3602,6 +3602,59 @@ fn embedded_workdir_unsupported_diagnostic() -> RuntimeDiagnostic { ) } +fn sanitize_embedded_execution_message( + message: &str, + operation: &impl std::fmt::Debug, + outcome: &impl std::fmt::Debug, +) -> String { + let summary = + format!("Embedded Worker execution backend rejected {operation:?} with {outcome:?}"); + let mut redact_next = false; + let detail = message + .split_whitespace() + .map(|part| { + if redact_next { + redact_next = false; + return "[redacted]"; + } + let lowercase = part.to_ascii_lowercase(); + let label = + lowercase.trim_matches(|character: char| !character.is_ascii_alphanumeric()); + if matches!( + label, + "bearer" | "credential" | "key" | "password" | "secret" | "session" | "token" + ) { + redact_next = true; + } + if part.contains('/') + || part.contains('\\') + || lowercase.contains("credential=") + || lowercase.contains("key=") + || lowercase.contains("password=") + || lowercase.contains("secret=") + || lowercase.contains("session=") + || lowercase.contains("session_id=") + || lowercase.contains("token=") + { + "[redacted]" + } else { + part + } + }) + .collect::>() + .join(" "); + let detail = detail.trim(); + if detail.is_empty() { + return summary; + } + let truncated = detail.chars().count() > 512; + let mut detail = detail.chars().take(512).collect::(); + if truncated { + detail.push('…'); + } + format!("{summary}: {detail}") +} + fn embedded_runtime_diagnostic(error: &EmbeddedRuntimeError) -> RuntimeDiagnostic { match error { EmbeddedRuntimeError::RuntimeStopped => diagnostic( @@ -3621,11 +3674,14 @@ fn embedded_runtime_diagnostic(error: &EmbeddedRuntimeError) -> RuntimeDiagnosti "Embedded Worker has no execution backend attached".to_string(), ), EmbeddedRuntimeError::WorkerExecutionRejected { - operation, outcome, .. + operation, + outcome, + message, + .. } => diagnostic( "embedded_worker_execution_rejected", DiagnosticSeverity::Warning, - format!("Embedded Worker execution backend rejected {operation:?} with {outcome:?}"), + sanitize_embedded_execution_message(message, operation, outcome), ), EmbeddedRuntimeError::LimitTooLarge { requested, max } => diagnostic( "embedded_runtime_limit_too_large", @@ -4263,7 +4319,7 @@ mod tests { worker_runtime::execution::WorkerExecutionSpawnResult::Errored( worker_runtime::execution::WorkerExecutionResult::errored( worker_runtime::execution::WorkerExecutionOperation::Spawn, - "provider setup failed at /tmp/secret-provider-config", + "provider setup failed at /tmp/secret-provider-config token=secret-value session_id=session-42", ), ) } @@ -4643,7 +4699,12 @@ mod tests { assert!(spawned.acceptance_evidence.is_empty()); assert!(spawned.diagnostics.iter().any(|diagnostic| { diagnostic.code == "embedded_worker_execution_rejected" + && diagnostic + .message + .contains("provider setup failed at [redacted]") && !diagnostic.message.contains("/tmp/secret-provider-config") + && !diagnostic.message.contains("secret-value") + && !diagnostic.message.contains("session-42") })); assert!(spawned.worker.is_none()); } diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 3661cf86..6652821e 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -10421,6 +10421,115 @@ mod tests { assert!(matches!(error, Error::WorkerSourceIdentity(_))); } + #[tokio::test] + async fn production_profile_backend_launches_and_restores_workspace_orchestrator() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let config = test_server_config(workspace.path()); + let store = SqliteWorkspaceStore::open(config.database_path.clone()).unwrap(); + let api = WorkspaceApi::new(config, Arc::new(store)).await.unwrap(); + let workspace_id = api.config.workspace_id.clone(); + + let result = scoped_start_workspace_orchestrator( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + ) + .await; + + let Json(started) = result.unwrap_or_else(|error| { + panic!( + "production Workspace Orchestrator launch failed: error={:?}, diagnostics={:?}", + error.error, error.diagnostics + ) + }); + assert_eq!(started.disposition, "created"); + assert!(started.online); + let worker = started + .worker + .expect("production Workspace Orchestrator Worker") + .worker; + + let stopped = api + .runtime + .stop_worker( + &worker, + WorkerLifecycleRequest { + reason: Some("production Orchestrator restore regression test".to_string()), + ticket_assignment: None, + }, + ) + .unwrap(); + assert_eq!(stopped.state, WorkerOperationState::Accepted); + let Json(restored) = scoped_start_workspace_orchestrator( + State(api), + AxumPath(ScopedWorkspacePath { workspace_id }), + ) + .await + .unwrap(); + assert_eq!(restored.disposition, "restored"); + assert!(restored.online); + assert_eq!( + restored.worker.expect("restored Orchestrator").worker, + worker + ); + } + + #[tokio::test] + async fn rejected_orchestrator_spawn_stays_offline_and_can_be_retried() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let config = test_server_config(workspace.path()); + let store = SqliteWorkspaceStore::open(config.database_path.clone()).unwrap(); + let api = WorkspaceApi::new_with_execution_backend( + config, + Arc::new(store), + Arc::new(DeterministicExecutionBackend::fail_first_spawn( + "provider setup rejected safe-root /tmp/private token=private-token session_id=private-session", + )), + ) + .await + .unwrap(); + let workspace_id = api.config.workspace_id.clone(); + + let error = scoped_start_workspace_orchestrator( + State(api.clone()), + AxumPath(ScopedWorkspacePath { + workspace_id: workspace_id.clone(), + }), + ) + .await + .unwrap_err(); + assert!(matches!( + error.error, + Error::RuntimeOperationFailed { + ref code, + .. + } if code == "workspace_orchestrator_spawn_rejected" + )); + assert!(error.diagnostics.iter().any(|diagnostic| { + diagnostic.code == "embedded_worker_execution_rejected" + && diagnostic + .message + .contains("provider setup rejected safe-root") + && !diagnostic.message.contains("/tmp/private") + && !diagnostic.message.contains("private-token") + && !diagnostic.message.contains("private-session") + })); + assert!(find_workspace_orchestrator(&api).is_none()); + assert!(!workspace_orchestrator_response(&api, "failed").online); + + let Json(retried) = scoped_start_workspace_orchestrator( + State(api), + AxumPath(ScopedWorkspacePath { workspace_id }), + ) + .await + .unwrap(); + assert_eq!(retried.disposition, "created"); + assert!(retried.online); + } + #[tokio::test] async fn explicit_orchestrator_launch_marks_only_the_dedicated_worker() { let workspace = tempfile::tempdir().unwrap(); @@ -11036,6 +11145,7 @@ mod tests { >, >, materializer: worker_runtime::working_directory::LocalGitWorktreeMaterializer, + spawn_failure: std::sync::Mutex>, } impl Default for DeterministicExecutionBackend { @@ -11053,10 +11163,19 @@ mod tests { materializer: worker_runtime::working_directory::LocalGitWorktreeMaterializer::new( std::env::temp_dir().join(unique), ), + spawn_failure: std::sync::Mutex::new(None), } } } + impl DeterministicExecutionBackend { + fn fail_first_spawn(message: impl Into) -> Self { + let backend = Self::default(); + *backend.spawn_failure.lock().unwrap() = Some(message.into()); + backend + } + } + impl worker_runtime::execution::WorkerExecutionBackend for DeterministicExecutionBackend { fn backend_id(&self) -> &str { "deterministic-workspace-server-test" @@ -11104,6 +11223,14 @@ mod tests { &self, request: worker_runtime::execution::WorkerExecutionSpawnRequest, ) -> worker_runtime::execution::WorkerExecutionSpawnResult { + if let Some(message) = self.spawn_failure.lock().unwrap().take() { + return worker_runtime::execution::WorkerExecutionSpawnResult::Errored( + worker_runtime::execution::WorkerExecutionResult::errored( + worker_runtime::execution::WorkerExecutionOperation::Spawn, + message, + ), + ); + } let working_directory = match request.request.working_directory.as_ref() { Some(claim) => match self.materializer.bind_working_directory( &claim.working_directory_id, diff --git a/resources/profiles/orchestrator.dcdl b/resources/profiles/orchestrator.dcdl index 7fc07c19..5c82075a 100644 --- a/resources/profiles/orchestrator.dcdl +++ b/resources/profiles/orchestrator.dcdl @@ -7,7 +7,7 @@ import "./base.dcdl" // { task = { enabled = true; }; memory = { enabled = true; }; web = { enabled = true; }; - sub_worker = { enabled = true; }; + sub_worker = { enabled = false; }; worker = { enabled = true; }; manage_workdir = { enabled = true; }; ticket = { enabled = true; thread = true; orchestration_control = true; };