From d7e35ea9eedc8ef9f711f6c67e209423997e7c66 Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 19 Aug 2026 04:41:38 +0900 Subject: [PATCH] fix: keep restore independent of live prompts --- crates/worker-runtime/src/worker_backend.rs | 120 ++++---------------- crates/worker/src/prompt/catalog.rs | 5 + 2 files changed, 25 insertions(+), 100 deletions(-) diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 2fdcabe4..c67e6af8 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -246,33 +246,14 @@ impl WorkspacePromptProjectionCache { Ok(projection) } - fn fetch_current( + fn current( &self, - workspace_client: &Arc, workspace_id: &worker::WorkspaceId, - ) -> Result, String> { - let response = workspace_client - .execute(worker::WorkspaceRequest { - method: worker::WorkspaceRequestMethod::Get, - path: format!( - "/api/w/{}/config/projections/prompts", - workspace_id.as_str() - ), - body: None, - }) - .map_err(|error| error.to_string())?; - if !(200..300).contains(&response.status) { - return Err(format!( - "Workspace Prompt projection request failed with status {}", - response.status - )); - } - let projection: worker::WorkspacePromptProjection = - serde_json::from_str(&response.body).map_err(|error| error.to_string())?; - if projection.workspace_id != workspace_id.as_str() { - return Err("Workspace Prompt projection returned a mismatched workspace".to_string()); - } - self.observe(projection) + ) -> Result>, String> { + self.active + .lock() + .map_err(|_| "Workspace Prompt projection cache lock was poisoned".to_string()) + .map(|active| active.get(workspace_id.as_str()).cloned()) } } @@ -816,10 +797,9 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { self.embedded_worker_mutation_dispatcher.as_ref(), ); let (manifest, mut loader) = Self::restore_fallback_manifest(&worker_name)?; - if let Some(workspace_id) = workspace_context.workspace_id() { - let projection = self - .prompt_projection_cache - .fetch_current(&workspace_context.client_handle(), workspace_id)?; + if let Some(workspace_id) = workspace_context.workspace_id() + && let Some(projection) = self.prompt_projection_cache.current(workspace_id)? + { loader = loader.with_effective_catalog(projection.catalog.clone()); } @@ -1948,75 +1928,6 @@ mod tests { use manifest::{Scope, WorkerManifest}; use session_store::{LogEntry, WorkerMetadataStore}; - #[derive(Debug)] - struct PromptProjectionWorkspaceClient { - response: String, - calls: AtomicUsize, - } - - impl worker::WorkspaceClient for PromptProjectionWorkspaceClient { - fn kind(&self) -> &'static str { - "prompt-projection-test" - } - - fn workspace_id(&self) -> Option<&str> { - Some("workspace-a") - } - - fn is_available(&self) -> bool { - true - } - - fn execute( - &self, - request: worker::WorkspaceRequest, - ) -> Result { - assert_eq!( - request.path, - "/api/w/workspace-a/config/projections/prompts" - ); - self.calls.fetch_add(1, Ordering::SeqCst); - Ok(worker::WorkspaceResponse { - status: 200, - body: self.response.clone(), - }) - } - } - - #[test] - fn workspace_prompt_projection_cache_fetches_and_validates_current_projection() { - let catalog = worker::EffectivePromptCatalog::new( - BTreeMap::from([("default".to_string(), "prompt".to_string())]), - 8, - "schema", - "toolchain", - ) - .unwrap(); - let projection = worker::WorkspacePromptProjection::new( - "workspace-a", - "source-digest", - catalog.catalog_digest.clone(), - catalog, - ) - .unwrap(); - let client = Arc::new(PromptProjectionWorkspaceClient { - response: serde_json::to_string(&projection).unwrap(), - calls: AtomicUsize::new(0), - }); - let cache = WorkspacePromptProjectionCache::default(); - let workspace_id = worker::WorkspaceId::new("workspace-a".to_string()).unwrap(); - - let resolved = cache - .fetch_current( - &(client.clone() as Arc), - &workspace_id, - ) - .unwrap(); - - assert_eq!(resolved.as_ref(), &projection); - assert_eq!(client.calls.load(Ordering::SeqCst), 1); - } - #[test] fn workspace_prompt_projection_cache_rejects_same_revision_source_drift() { let catalog = worker::EffectivePromptCatalog::new( @@ -2704,14 +2615,23 @@ mod tests { ) .unwrap(); - let request = create_request("restore"); + let mut request = create_request("restore"); + request.workspace_api = Some(crate::catalog::WorkspaceApiRef { + workspace_id: "workspace-restore".to_string(), + base_url: "http://workspace.invalid".to_string(), + }); + let identity = RuntimeIdentityMaterial::generate("runtime-restore").unwrap(); let controller = ProfileRuntimeWorkerFactory::new(root.path()) .with_runtime_store_dir(&runtime_store_dir) + .with_remote_worker_mutation_identity(identity) .restore_controller(WorkerExecutionRestoreRequest { worker_ref: worker_ref.clone(), run_generation: 1, request, - workspace_scope: None, + workspace_scope: Some(crate::runtime::RuntimeWorkspaceScope::new( + "workspace-restore", + "server-main", + )), context: test_execution_context(worker_ref), previous_working_directory: None, working_directory: None, diff --git a/crates/worker/src/prompt/catalog.rs b/crates/worker/src/prompt/catalog.rs index 9c5d187a..6cf95753 100644 --- a/crates/worker/src/prompt/catalog.rs +++ b/crates/worker/src/prompt/catalog.rs @@ -211,6 +211,11 @@ impl WorkspacePromptProjection { "Workspace Prompt projection digest must not be empty".to_string(), )); } + if projection_digest != catalog.catalog_digest { + return Err(CatalogError::InvalidTemplateCatalog( + "Workspace Prompt projection digest does not match its catalog".to_string(), + )); + } if !catalog.source_digest.is_empty() && catalog.source_digest != source_digest { return Err(CatalogError::InvalidTemplateCatalog( "Workspace Prompt projection source digest does not match its catalog".to_string(),