From c4814115dee8f88767fefe8026e211b0ca603953 Mon Sep 17 00:00:00 2001 From: Hare Date: Fri, 21 Aug 2026 21:50:12 +0900 Subject: [PATCH 1/4] feat: bind memory language to workspace settings --- crates/manifest/src/config.rs | 2 + crates/manifest/src/defaults.rs | 4 - crates/manifest/src/lib.rs | 43 +- crates/memory/src/audit.rs | 10 + crates/worker-runtime/src/catalog.rs | 3 + crates/worker-runtime/src/http_server.rs | 15 + crates/worker-runtime/src/runtime.rs | 56 ++- crates/worker-runtime/src/worker_backend.rs | 68 ++- crates/worker/src/worker.rs | 254 ++++++++++- crates/workspace-api/src/lib.rs | 17 + crates/workspace-server/src/hosts.rs | 17 + .../src/runtime_subscription_tests.rs | 5 + crates/workspace-server/src/server.rs | 106 ++++- crates/workspace-server/src/store.rs | 399 ++++++++++++++++-- 14 files changed, 927 insertions(+), 72 deletions(-) diff --git a/crates/manifest/src/config.rs b/crates/manifest/src/config.rs index 33f2b295..9671675d 100644 --- a/crates/manifest/src/config.rs +++ b/crates/manifest/src/config.rs @@ -748,6 +748,8 @@ impl MemoryConfig { query_result_limit: upper.query_result_limit.or(self.query_result_limit), query_excerpt_lines: upper.query_excerpt_lines.or(self.query_excerpt_lines), inject_summary: upper.inject_summary.or(self.inject_summary), + workspace_id: upper.workspace_id.or(self.workspace_id), + settings_revision: upper.settings_revision.or(self.settings_revision), language: upper.language.or(self.language), extract_model: upper.extract_model.or(self.extract_model), extract_threshold: upper.extract_threshold.or(self.extract_threshold), diff --git a/crates/manifest/src/defaults.rs b/crates/manifest/src/defaults.rs index c725f1cd..70327972 100644 --- a/crates/manifest/src/defaults.rs +++ b/crates/manifest/src/defaults.rs @@ -95,7 +95,3 @@ pub const COMPACT_DEFAULT_REFERENCE_COUNT: usize = 5; /// Optional maximum extract-worker tool-loop depth. `None` means unlimited. /// See [`crate::MemoryConfig::extract_worker_max_turns`]. pub const MEMORY_EXTRACT_WORKER_MAX_TURNS: Option = Some(8); - -/// Default language used by memory extraction / consolidation workers for -/// durable memory text. See [`crate::MemoryConfig::language`]. -pub const MEMORY_LANGUAGE: &str = "English"; diff --git a/crates/manifest/src/lib.rs b/crates/manifest/src/lib.rs index 7c7c418c..b1265a61 100644 --- a/crates/manifest/src/lib.rs +++ b/crates/manifest/src/lib.rs @@ -449,6 +449,18 @@ pub struct WebFetchConfig { pub allow_private_addresses: Option, } +/// Immutable Workspace Memory settings bound into a Worker launch snapshot. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct WorkspaceMemorySettingsSnapshot { + /// Workspace that owns the settings revision. + pub workspace_id: String, + /// Monotonic Workspace Memory settings revision. + pub settings_revision: u64, + /// Normalized language used for Memory extraction and consolidation output. + pub language: String, +} + /// Memory subsystem configuration. Presence in the manifest enables /// memory; `workspace_root` pins the memory workspace explicitly. When it /// is absent, memory resolution searches upward from the Worker's pwd for a @@ -477,11 +489,14 @@ pub struct MemoryConfig { /// system-prompt section. `None` ⇒ enabled. #[serde(default)] pub inject_summary: Option, - /// Language used by memory extraction / consolidation sub_worker for durable - /// memory text. Free-form so workspaces can use names like - /// `English`, `Japanese`, or locale tags. `None` ⇒ - /// [`defaults::MEMORY_LANGUAGE`]. - #[serde(default)] + /// Workspace that owns the bound Memory settings revision. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace_id: Option, + /// Monotonic revision of the bound Workspace Memory settings. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub settings_revision: Option, + /// Language from the bound Workspace Memory settings revision. + #[serde(default, skip_serializing_if = "Option::is_none")] pub language: Option, /// Optional model for the extract worker. When `None`, /// the main engine model is cloned via `clone_boxed()`. Lightweight @@ -520,6 +535,24 @@ pub struct MemoryConfig { pub consolidation_threshold_bytes: Option, } +impl MemoryConfig { + /// Replace any profile-authored language fields with a trusted Workspace snapshot. + pub fn bind_workspace_settings(&mut self, snapshot: &WorkspaceMemorySettingsSnapshot) { + self.workspace_id = Some(snapshot.workspace_id.clone()); + self.settings_revision = Some(snapshot.settings_revision); + self.language = Some(snapshot.language.clone()); + } + + /// Return the complete bound Workspace settings snapshot, if every field is present. + pub fn workspace_settings(&self) -> Option { + Some(WorkspaceMemorySettingsSnapshot { + workspace_id: self.workspace_id.clone()?, + settings_revision: self.settings_revision?, + language: self.language.clone()?, + }) + } +} + /// Worker metadata. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct WorkerMeta { diff --git a/crates/memory/src/audit.rs b/crates/memory/src/audit.rs index 21d35d25..c27a7990 100644 --- a/crates/memory/src/audit.rs +++ b/crates/memory/src/audit.rs @@ -177,6 +177,13 @@ impl OperationCounts { } } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct MemorySettingsAudit { + pub workspace_id: String, + pub settings_revision: u64, + pub language: String, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct WorkerLifecycleAudit { pub run_id: Uuid, @@ -185,6 +192,8 @@ pub struct WorkerLifecycleAudit { pub trigger: AuditTrigger, pub reason: String, #[serde(default, skip_serializing_if = "Option::is_none")] + pub memory_settings: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] pub model: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub usage: Option, @@ -405,6 +414,7 @@ mod tests { status: WorkerLifecycleStatus::Started, trigger: AuditTrigger::TokenThreshold, reason: "tokens_threshold_reached".to_string(), + memory_settings: None, model: None, usage: None, extract: None, diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index 703cae62..e745e9eb 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -174,6 +174,9 @@ pub struct CreateWorkerRequest { pub worker_observation_grants: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] pub workspace_api: Option, + /// Backend-authored immutable Workspace Memory settings snapshot. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub memory_settings: Option, } /// Worker lifecycle status for the in-memory embedded runtime. diff --git a/crates/worker-runtime/src/http_server.rs b/crates/worker-runtime/src/http_server.rs index 0c0e0dcf..7236f3d3 100644 --- a/crates/worker-runtime/src/http_server.rs +++ b/crates/worker-runtime/src/http_server.rs @@ -1882,6 +1882,11 @@ mod tests { 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 } @@ -2193,6 +2198,11 @@ mod tests { 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(), + }), } } @@ -2828,6 +2838,11 @@ mod ws_tests { 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(), + }), } } diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 7750ecbc..21dc6033 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -2794,6 +2794,27 @@ fn validate_create_workspace_scope( ))); } } + 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 !matches!(snapshot.language.as_str(), "English" | "Japanese") { + return Err(RuntimeError::InvalidRequest( + "Memory settings language must be a normalized supported value".to_string(), + )); + } Ok(()) } @@ -3086,6 +3107,11 @@ mod tests { 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(), + }), } } @@ -3095,9 +3121,37 @@ mod tests { 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 = 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) } @@ -4484,7 +4538,7 @@ mod tests { let worker = runtime .create_worker_scoped( &RuntimeWorkspaceScope::new("workspace-a", "server"), - task_request("legacy"), + scoped_task_request("legacy", "workspace-a"), ) .unwrap(); drop(runtime); diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 7efbc92e..19233e5f 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -642,6 +642,60 @@ fn runtime_local_workdir_session( )) } +fn bind_workspace_memory_settings( + manifest: &mut manifest::WorkerManifest, + request: &CreateWorkerRequest, +) -> Result<(), String> { + let Some(snapshot) = request.memory_settings.as_ref() else { + if request.workspace_api.is_some() { + return Err( + "Workspace Worker request is missing its bound Memory settings snapshot" + .to_string(), + ); + } + return Ok(()); + }; + if let Some(workspace_api) = request.workspace_api.as_ref() + && snapshot.workspace_id != workspace_api.workspace_id + { + return Err(format!( + "Memory settings workspace {} does not match Workspace API scope {}", + snapshot.workspace_id, workspace_api.workspace_id + )); + } + manifest + .memory + .get_or_insert_with(manifest::MemoryConfig::default) + .bind_workspace_settings(snapshot); + Ok(()) +} + +fn validate_worker_memory_settings( + manifest: &manifest::WorkerManifest, + request: &CreateWorkerRequest, +) -> Result<(), String> { + let Some(expected) = request.memory_settings.as_ref() else { + return Ok(()); + }; + let actual = manifest + .memory + .as_ref() + .and_then(manifest::MemoryConfig::workspace_settings) + .ok_or_else(|| { + "Workspace Worker restored without its bound Memory settings snapshot".to_string() + })?; + if &actual != expected { + return Err(format!( + "Workspace Worker Memory settings snapshot mismatch: expected {} revision {}, restored {} revision {}", + expected.workspace_id, + expected.settings_revision, + actual.workspace_id, + actual.settings_revision + )); + } + Ok(()) +} + #[async_trait] impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { fn observe_workspace_prompt_projection( @@ -696,7 +750,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { let archive = self .resolve_profile_source_archive(&request.request.profile_source) .await?; - let (manifest, mut loader) = { + let (mut manifest, mut loader) = { let manifest = archive .resolve_profile(selector, &worker_root, &worker_name) .map_err(|err| format!("failed to resolve profile source archive: {err}"))?; @@ -714,6 +768,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { )? } }; + bind_workspace_memory_settings(&mut manifest, &request.request)?; if let Some(bundle) = request.config_bundle.as_ref() && let Some(resolution) = self.observe_bundle_prompt_projection(bundle, observation_workspace_id.as_deref())? @@ -750,6 +805,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { ) .await .map_err(|err| format!("failed to create Worker from profile: {err}"))?; + validate_worker_memory_settings(worker.manifest(), &request.request)?; if let Some(binding) = request.working_directory.as_ref() { worker.bind_workdir_session(Some(runtime_local_workdir_session( &binding.working_directory.id, @@ -856,7 +912,8 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { self.embedded_worker_mutation_dispatcher.as_ref(), Some(self.prompt_projection_cache.clone()), ); - let (manifest, loader) = Self::restore_fallback_manifest(&worker_name)?; + let (mut manifest, loader) = Self::restore_fallback_manifest(&worker_name)?; + bind_workspace_memory_settings(&mut manifest, &request.request)?; let worker_aggregate_dir = self.worker_aggregate_dir(&request.worker_ref)?; let session_dir = worker_aggregate_dir.join("session"); @@ -924,6 +981,7 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { } Err(err) => return Err(format!("failed to restore Worker from metadata: {err}")), }; + validate_worker_memory_settings(worker.manifest(), &request.request)?; let flow_transition_enabled = worker.manifest().feature.flow.enabled; if let Some(binding) = request.working_directory.as_ref() { worker.bind_workdir_session(Some(runtime_local_workdir_session( @@ -2504,6 +2562,7 @@ mod tests { worker_observation_enabled: false, worker_observation_grants: Vec::new(), workspace_api: None, + memory_settings: None, } } @@ -2816,6 +2875,11 @@ mod tests { workspace_id: "workspace-restore".to_string(), base_url: "http://workspace.invalid".to_string(), }); + request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: "workspace-restore".to_string(), + settings_revision: 1, + language: "English".to_string(), + }); let identity = RuntimeIdentityMaterial::generate("runtime-restore").unwrap(); let error = match ProfileRuntimeWorkerFactory::new(root.path()) .with_runtime_store_dir(&runtime_store_dir) diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index fdf45c71..ea00b0b6 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -2163,6 +2163,33 @@ impl Worker { let Some(template) = self.system_prompt_template.take() else { return Ok(()); }; + let is_memory_consolidation = self.manifest.profile.as_ref().is_some_and(|snapshot| { + matches!( + &snapshot.source, + manifest::ProfileSource::Registry { + source: manifest::ProfileRegistrySource::Builtin, + name, + .. + } if name == "memory-consolidation" + ) + }); + if is_memory_consolidation { + let memory_config = self.manifest.memory.as_ref().ok_or_else(|| { + WorkerError::InvalidState( + "Memory consolidation Worker has no Memory configuration".to_string(), + ) + })?; + let language = memory_language(memory_config)?; + let rendered = self + .prompts + .load_full() + .memory_consolidation_system(&language)?; + self.engine + .as_mut() + .expect("worker present") + .set_system_prompt(rendered); + return Ok(()); + } let alerter = self.alerter.clone(); let tool_names: Vec = { let worker = self.engine.as_mut().expect("worker present"); @@ -3750,6 +3777,7 @@ impl Worker { memory::audit::AuditTrigger::TokenThreshold, Some(model_audit_from_manifest(model)), ) + .with_memory_settings(&memory_cfg) .emit( self.workspace_client(), self.event_tx.as_ref(), @@ -3780,6 +3808,7 @@ impl Worker { memory::audit::AuditTrigger::TokenThreshold, Some(model_audit_from_manifest(model)), ) + .with_memory_settings(&memory_cfg) .emit( self.workspace_client(), self.event_tx.as_ref(), @@ -3845,7 +3874,8 @@ impl Worker { memory::audit::AuditWorker::MemoryExtract, memory::audit::AuditTrigger::TokenThreshold, Some(model_audit_from_manifest(model)), - ); + ) + .with_memory_settings(memory_cfg); let event_tx = self.event_tx.as_ref(); let pointer_snapshot = self @@ -3997,11 +4027,11 @@ impl Worker { return Err(err); } }; - let memory_language = memory_language(memory_cfg); + let memory_language = memory_language(memory_cfg)?; let extract_system_prompt = match self .prompts .load_full() - .memory_extract_system(memory_language) + .memory_extract_system(&memory_language) { Ok(prompt) => prompt, Err(err) => { @@ -4191,6 +4221,7 @@ impl Worker { memory::audit::AuditTrigger::StagingBacklog, Some(model_audit_from_manifest(model)), ) + .with_memory_settings(&memory_cfg) .emit( self.workspace_client(), self.event_tx.as_ref(), @@ -4232,6 +4263,7 @@ impl Worker { memory::audit::AuditTrigger::StagingBacklog, Some(model_audit_from_manifest(model)), ) + .with_memory_settings(&memory_cfg) .emit( self.workspace_client(), self.event_tx.as_ref(), @@ -4312,6 +4344,7 @@ struct WorkerAuditBase { worker: memory::audit::AuditWorker, trigger: memory::audit::AuditTrigger, model: Option, + memory_settings: Option, } impl WorkerAuditBase { @@ -4325,9 +4358,22 @@ impl WorkerAuditBase { worker, trigger, model, + memory_settings: None, } } + fn with_memory_settings(mut self, memory_config: &manifest::MemoryConfig) -> Self { + self.memory_settings = + memory_config + .workspace_settings() + .map(|snapshot| memory::audit::MemorySettingsAudit { + workspace_id: snapshot.workspace_id, + settings_revision: snapshot.settings_revision, + language: snapshot.language, + }); + self + } + async fn emit( &self, workspace_client: &dyn WorkspaceClient, @@ -4345,6 +4391,7 @@ impl WorkerAuditBase { status, trigger: self.trigger, reason: reason.clone(), + memory_settings: self.memory_settings.clone(), model: self.model.clone(), usage, extract, @@ -4391,12 +4438,14 @@ fn is_idle_consolidation_skip_reason(reason: &str) -> bool { || reason.starts_with("threshold_not_reached") } -fn memory_language(cfg: &manifest::MemoryConfig) -> &str { - cfg.language - .as_deref() - .map(str::trim) - .filter(|language| !language.is_empty()) - .unwrap_or(manifest::defaults::MEMORY_LANGUAGE) +fn memory_language(cfg: &manifest::MemoryConfig) -> Result { + cfg.workspace_settings() + .map(|snapshot| snapshot.language) + .ok_or_else(|| { + WorkerError::InvalidState( + "Memory operation requires a bound Workspace Memory settings snapshot".to_string(), + ) + }) } fn worker_language(cfg: &manifest::EngineManifest) -> &str { @@ -4453,6 +4502,7 @@ where workspace_context: WorkerWorkspaceContext, filesystem_authority: WorkerFilesystemAuthority, ) -> Result { + validate_workspace_memory_snapshot(&manifest.worker.name, &manifest, &workspace_context)?; let common = prepare_worker_common_with_context( &manifest, &loader, @@ -4654,6 +4704,7 @@ where workspace_context: WorkerWorkspaceContext, filesystem_authority: WorkerFilesystemAuthority, ) -> Result { + validate_workspace_memory_snapshot(&manifest.worker.name, &manifest, &workspace_context)?; let common = prepare_worker_common_with_context( &manifest, &loader, @@ -5156,8 +5207,48 @@ fn worker_metadata_for_manifest( metadata } +fn validate_workspace_memory_snapshot( + worker_name: &str, + manifest: &WorkerManifest, + workspace_context: &WorkerWorkspaceContext, +) -> Result<(), WorkerError> { + let Some(workspace_id) = workspace_context.workspace_id() else { + return Ok(()); + }; + let snapshot = manifest + .memory + .as_ref() + .and_then(manifest::MemoryConfig::workspace_settings) + .ok_or_else(|| { + WorkerError::InvalidState(format!( + "Workspace Worker {worker_name} has no complete persisted Memory settings snapshot" + )) + })?; + if snapshot.workspace_id != workspace_id.as_str() { + return Err(WorkerError::InvalidState(format!( + "Workspace Worker {worker_name} Memory settings belong to {} instead of {}", + snapshot.workspace_id, + workspace_id.as_str() + ))); + } + if snapshot.settings_revision == 0 + || !matches!(snapshot.language.as_str(), "English" | "Japanese") + { + return Err(WorkerError::InvalidState(format!( + "Workspace Worker {worker_name} has corrupt Memory settings snapshot metadata" + ))); + } + Ok(()) +} + fn should_persist_resolved_manifest_snapshot(manifest: &WorkerManifest) -> bool { - manifest.profile.is_some() || manifest.plugins.has_resolved_plan() + manifest.profile.is_some() + || manifest.plugins.has_resolved_plan() + || manifest + .memory + .as_ref() + .and_then(manifest::MemoryConfig::workspace_settings) + .is_some() } fn restore_manifest_from_worker_metadata_snapshot( @@ -6179,6 +6270,92 @@ permission = "write" ); } + #[test] + fn workspace_memory_settings_snapshot_is_persisted_and_scope_checked() { + let mut manifest = WorkerManifest::from_toml( + r#" +[worker] +name = "memory-snapshot" + +[model] +scheme = "anthropic" +model_id = "claude-sonnet-4-20250514" + +[engine] +instruction = "default" + +[[scope.allow]] +target = "/workspace" +permission = "read" + +[[delegation_scope.allow]] +target = "/workspace" +permission = "read" +"#, + ) + .unwrap(); + manifest.memory = Some(manifest::MemoryConfig::default()); + manifest.memory.as_mut().unwrap().bind_workspace_settings( + &manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: "workspace-a".to_string(), + settings_revision: 7, + language: "Japanese".to_string(), + }, + ); + + let metadata = worker_metadata_for_manifest(&manifest, None, None, None); + let restored: WorkerManifest = serde_json::from_value( + metadata + .resolved_manifest_snapshot + .expect("Memory settings require a resolved manifest snapshot"), + ) + .unwrap(); + assert_eq!( + restored.memory.unwrap().workspace_settings(), + Some(manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: "workspace-a".to_string(), + settings_revision: 7, + language: "Japanese".to_string(), + }) + ); + assert!( + validate_workspace_memory_snapshot( + "memory-snapshot", + &manifest, + &WorkerWorkspaceContext::unavailable( + Some(WorkspaceId::new("workspace-a").unwrap()), + "test", + ) + ) + .is_ok() + ); + assert!( + validate_workspace_memory_snapshot( + "memory-snapshot", + &manifest, + &WorkerWorkspaceContext::unavailable( + Some(WorkspaceId::new("workspace-b").unwrap()), + "test", + ) + ) + .is_err() + ); + + let mut missing = manifest.clone(); + missing.memory.as_mut().unwrap().settings_revision = None; + assert!( + validate_workspace_memory_snapshot( + "memory-snapshot", + &missing, + &WorkerWorkspaceContext::unavailable( + Some(WorkspaceId::new("workspace-a").unwrap()), + "test", + ) + ) + .is_err() + ); + } + #[test] fn plugin_resolved_manifest_snapshot_is_persisted_without_profile() { let mut manifest = WorkerManifest::from_toml( @@ -7064,6 +7241,52 @@ mod build_summary_prompt_tests { } } + #[tokio::test] + async fn memory_consolidation_prompt_uses_bound_workspace_language() { + let dir = tempfile::tempdir().unwrap(); + let cwd = dir.path().join("workspace"); + std::fs::create_dir_all(&cwd).unwrap(); + let store = session_store::FsStore::new(dir.path().join("sessions")).unwrap(); + let mut manifest = minimal_manifest(); + manifest.profile = Some(manifest::ProfileManifestSnapshot { + source: manifest::ProfileSource::Registry { + source: manifest::ProfileRegistrySource::Builtin, + name: "memory-consolidation".to_string(), + path: None, + provenance: None, + }, + profile: None, + }); + let mut memory = manifest::MemoryConfig::default(); + memory.bind_workspace_settings(&manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: "workspace-test".to_string(), + settings_revision: 3, + language: "Japanese".to_string(), + }); + manifest.memory = Some(memory); + let mut worker = Worker::new( + manifest, + Engine::new(NoopClient), + store, + WorkerWorkspaceContext::no_workspace(), + WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()), + Scope::writable(&cwd).unwrap(), + ) + .await + .unwrap(); + worker.set_system_prompt_template( + SystemPromptTemplate::parse( + "default", + crate::prompt::source::PromptCatalogSource::builtins_only(), + ) + .unwrap(), + ); + worker.ensure_system_prompt_materialized().await.unwrap(); + let prompt = worker.engine().get_system_prompt().unwrap(); + assert!(prompt.contains("`language`: `Japanese`")); + assert!(!prompt.contains("`language`: `English`")); + } + async fn render_system_prompt_with_summary( summary_doc: Option<&str>, memory_config: Option, @@ -7358,6 +7581,9 @@ mod build_summary_prompt_tests { let mut manifest = minimal_manifest(); manifest.memory = Some(manifest::MemoryConfig { extract_threshold: Some(1), + workspace_id: Some("workspace-test".to_string()), + settings_revision: Some(1), + language: Some("English".to_string()), ..Default::default() }); let memory_config = manifest.memory.clone().unwrap(); @@ -7454,6 +7680,14 @@ mod build_summary_prompt_tests { assert_eq!(audits.len(), 2); assert_eq!(audits[0].run_id, audits[1].run_id); assert_eq!(audits[0].worker, memory::audit::AuditWorker::MemoryExtract); + assert!(audits.iter().all(|audit| { + audit.memory_settings + == Some(memory::audit::MemorySettingsAudit { + workspace_id: "workspace-test".to_string(), + settings_revision: 1, + language: "English".to_string(), + }) + })); assert_eq!( audits.iter().map(|audit| audit.status).collect::>(), vec![ diff --git a/crates/workspace-api/src/lib.rs b/crates/workspace-api/src/lib.rs index a24cd814..f52eea8e 100644 --- a/crates/workspace-api/src/lib.rs +++ b/crates/workspace-api/src/lib.rs @@ -169,6 +169,23 @@ pub struct WorkerRestoreResponse { pub result: WorkerRestoreResult, } +/// Workspace-owned Memory settings returned by the shared Server API. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct WorkspaceMemorySettings { + pub workspace_id: String, + pub settings_revision: u64, + pub language: String, +} + +/// Compare-and-swap update for Workspace-owned Memory settings. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct UpdateWorkspaceMemorySettingsRequest { + pub expected_revision: u64, + pub language: String, +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 29b610ba..87cf4261 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -518,6 +518,9 @@ pub struct WorkerSpawnRequest { pub resolved_config_bundle: Option, #[serde(skip, default)] pub resolved_workspace_api: Option, + /// Backend-authored immutable Workspace Memory settings snapshot. + #[serde(skip, default)] + pub resolved_memory_settings: Option, /// Backend-owned feature enablement; client input cannot set it. #[serde(skip, default)] pub resolved_worker_observation_enabled: bool, @@ -2184,6 +2187,7 @@ impl WorkspaceWorkerRuntime for EmbeddedWorkerRuntime { worker_observation_enabled: request.resolved_worker_observation_enabled, worker_observation_grants: request.resolved_worker_observation_grants.clone(), workspace_api: Some(workspace_api), + memory_settings: request.resolved_memory_settings.clone(), }; let workspace_scope = RuntimeWorkspaceScope::new(workspace_id, "embedded-backend"); match self @@ -3331,6 +3335,7 @@ impl WorkspaceWorkerRuntime for RemoteWorkerRuntime { worker_observation_enabled: request.resolved_worker_observation_enabled, worker_observation_grants: request.resolved_worker_observation_grants.clone(), workspace_api: Some(workspace_api), + memory_settings: request.resolved_memory_settings.clone(), }; match self.post_json::<_, RuntimeHttpWorkerResponse>("/v1/workers", &create) { Ok(response) => WorkerSpawnResult { @@ -4408,6 +4413,14 @@ mod tests { } } + fn test_memory_settings() -> manifest::WorkspaceMemorySettingsSnapshot { + manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: "workspace-test".to_string(), + settings_revision: 1, + language: "English".to_string(), + } + } + #[test] fn worker_summary_keeps_flat_wire_identity_while_using_structured_internal_identity() { let summary = placeholder_worker("placeholder"); @@ -5007,6 +5020,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: Some(test_workspace_api()), + resolved_memory_settings: Some(test_memory_settings()), } } @@ -5246,6 +5260,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: Some(test_workspace_api()), + resolved_memory_settings: Some(test_memory_settings()), }, ) .unwrap(); @@ -5345,6 +5360,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: Some(test_workspace_api()), + resolved_memory_settings: Some(test_memory_settings()), }, ) .unwrap(); @@ -5383,6 +5399,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: Some(test_workspace_api()), + resolved_memory_settings: Some(test_memory_settings()), }, ) .unwrap(); diff --git a/crates/workspace-server/src/runtime_subscription_tests.rs b/crates/workspace-server/src/runtime_subscription_tests.rs index 60b587a4..a42c0c02 100644 --- a/crates/workspace-server/src/runtime_subscription_tests.rs +++ b/crates/workspace-server/src/runtime_subscription_tests.rs @@ -86,6 +86,11 @@ fn create_request(name: &str) -> CreateWorkerRequest { 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(), + }), } } diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 47c0aa6e..c00d985b 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1182,8 +1182,11 @@ impl WorkspaceApi { &now_registry_timestamp(), )?; } - let create_fingerprint = worker_spawn_create_fingerprint(&request) + let request_fingerprint = worker_spawn_create_fingerprint(&request) .map_err(|message| Error::Config(message.to_string()))?; + let current_memory_settings = self + .config_store + .get_workspace_memory_settings(&self.config.workspace_id)?; let allocation_key = request .resolved_control_operation .as_ref() @@ -1195,22 +1198,25 @@ impl WorkspaceApi { .map(|assignment| assignment.operation_id.clone()) }) .unwrap_or_else(|| format!("manual:{}", WorkerId::now_v7())); - let worker_id = self + let reservation = self .config_store .reserve_worker_create( &self.config.workspace_id, runtime_id, &allocation_key, - &create_fingerprint, + &request_fingerprint, + ¤t_memory_settings, ) .map_err(|error| Error::RuntimeOperationFailed { runtime_id: runtime_id.to_string(), code: "workspace_worker_allocation_conflict".to_string(), message: error.to_string(), })?; + let worker_id = reservation.worker_id; + request.resolved_memory_settings = Some(reservation.memory_settings); let create_binding = WorkerCreateBinding { worker_id, - create_fingerprint, + create_fingerprint: reservation.create_fingerprint, }; let result = match self .runtime @@ -1601,6 +1607,11 @@ pub fn build_router(api: WorkspaceApi) -> Router { "/api/w/{workspace_id}/settings/workspace", get(scoped_get_workspace_settings).put(scoped_update_workspace_settings), ) + .route( + "/api/w/{workspace_id}/settings/memory", + get(scoped_get_workspace_memory_settings) + .put(scoped_update_workspace_memory_settings), + ) .route( "/api/w/{workspace_id}/config/source-tree", get(scoped_get_workspace_config_tree), @@ -2998,6 +3009,57 @@ struct WorkspaceConfigTreeResponse { projection_digest: String, } +async fn scoped_get_workspace_memory_settings( + State(api): State, + AxumPath(workspace_id): AxumPath, +) -> ApiResult> { + if api + .config_store + .get_workspace(&workspace_id) + .await? + .is_none() + { + return Err(ApiError::from(Error::WorkspaceIdMismatch)); + } + let settings = api + .config_store + .get_workspace_memory_settings(&workspace_id) + .map_err(ApiError::from)?; + Ok(Json(workspace_api::WorkspaceMemorySettings { + workspace_id: settings.workspace_id, + settings_revision: settings.settings_revision, + language: settings.language, + })) +} + +async fn scoped_update_workspace_memory_settings( + State(api): State, + AxumPath(workspace_id): AxumPath, + Json(request): Json, +) -> ApiResult> { + if api + .config_store + .get_workspace(&workspace_id) + .await? + .is_none() + { + return Err(ApiError::from(Error::WorkspaceIdMismatch)); + } + let settings = api + .config_store + .update_workspace_memory_settings( + &workspace_id, + request.expected_revision, + &request.language, + ) + .map_err(ApiError::from)?; + Ok(Json(workspace_api::WorkspaceMemorySettings { + workspace_id: settings.workspace_id, + settings_revision: settings.settings_revision, + language: settings.language, + })) +} + async fn scoped_get_workspace_config_tree( State(api): State, AxumPath(path): AxumPath, @@ -6126,6 +6188,7 @@ fn start_memory_staging_consolidation( resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }, )?; if result.state != WorkerOperationState::Accepted { @@ -7192,6 +7255,7 @@ async fn scoped_start_workspace_orchestrator( resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }, )?; if result.state != WorkerOperationState::Accepted || result.worker.is_none() { @@ -9736,6 +9800,7 @@ async fn create_workspace_worker( resolved_worker_observation_grants: Vec::new(), resolved_control_operation, resolved_workspace_api: None, + resolved_memory_settings: None, }; validate_ticket_assignment_spawn(&api, &runtime_id, &request)?; let assignment = request.ticket_assignment.clone(); @@ -13164,6 +13229,14 @@ mod tests { } } + fn test_worker_memory_settings() -> manifest::WorkspaceMemorySettingsSnapshot { + manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: TEST_WORKSPACE_ID.to_string(), + settings_revision: 1, + language: "English".to_string(), + } + } + #[test] fn ticket_api_errors_preserve_http_status() { let not_found = ApiError::from(Error::Ticket(ticket::TicketError::NotFound( @@ -13432,6 +13505,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }; assert!( api.validate_worker_spawn_repository_scope(&workdir_flow_launch) @@ -13604,6 +13678,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }; assert!( @@ -15069,6 +15144,7 @@ mod tests { resolved_workspace_api: Some(test_worker_workspace_api( EMBEDDED_WORKER_RUNTIME_ID, )), + resolved_memory_settings: Some(test_worker_memory_settings()), resolved_control_operation: None, }, ) @@ -15282,6 +15358,7 @@ mod tests { resolved_workspace_api: Some(test_worker_workspace_api( EMBEDDED_WORKER_RUNTIME_ID, )), + resolved_memory_settings: Some(test_worker_memory_settings()), resolved_control_operation: None, }, ) @@ -15492,6 +15569,7 @@ mod tests { resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: Some(test_worker_workspace_api(EMBEDDED_WORKER_RUNTIME_ID)), + resolved_memory_settings: Some(test_worker_memory_settings()), resolved_control_operation: None, }; let source_worker = api @@ -15756,6 +15834,7 @@ mod tests { resolved_workspace_api: Some(test_worker_workspace_api( EMBEDDED_WORKER_RUNTIME_ID, )), + resolved_memory_settings: Some(test_worker_memory_settings()), resolved_control_operation: None, }, ) @@ -15881,6 +15960,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }; let Json(first) = scoped_create_runtime_worker( State(api.clone()), @@ -16045,22 +16125,28 @@ mod tests { TEST_CREATED_AT, ) .unwrap(); - let reserved_worker_id = api + let current_memory_settings = api + .config_store + .get_workspace_memory_settings(TEST_WORKSPACE_ID) + .unwrap(); + let reservation = api .config_store .reserve_worker_create( TEST_WORKSPACE_ID, EMBEDDED_WORKER_RUNTIME_ID, "pending-spawn-operation", &pending_fingerprint, + ¤t_memory_settings, ) .unwrap(); + pending_request.resolved_memory_settings = Some(reservation.memory_settings.clone()); let spawned_before_backend_failure = api .runtime .spawn_worker( EMBEDDED_WORKER_RUNTIME_ID, WorkerCreateBinding { - worker_id: reserved_worker_id, - create_fingerprint: pending_fingerprint.clone(), + worker_id: reservation.worker_id, + create_fingerprint: reservation.create_fingerprint, }, pending_request.clone(), ) @@ -16139,6 +16225,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }; let Json(created) = scoped_create_runtime_worker( State(api.clone()), @@ -16699,6 +16786,7 @@ mod tests { resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, + resolved_memory_settings: None, resolved_control_operation: None, }, ) @@ -16755,6 +16843,7 @@ mod tests { resolved_worker_observation_enabled: false, resolved_worker_observation_grants: Vec::new(), resolved_workspace_api: None, + resolved_memory_settings: None, resolved_control_operation: None, }, ) @@ -17680,6 +17769,7 @@ mod tests { worker_observation_enabled: false, worker_observation_grants: Vec::new(), workspace_api: None, + memory_settings: None, } } @@ -18974,6 +19064,7 @@ mod tests { resolved_workspace_api: Some(test_worker_workspace_api( "embedded-worker-runtime", )), + resolved_memory_settings: Some(test_worker_memory_settings()), resolved_control_operation: None, }, ) @@ -19493,6 +19584,7 @@ mod tests { resolved_worker_observation_grants: Vec::new(), resolved_control_operation: None, resolved_workspace_api: None, + resolved_memory_settings: None, }; let spawned = api .spawn_workspace_worker(EMBEDDED_WORKER_RUNTIME_ID, spawn_request) diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 339a0782..4456bc78 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -8,6 +8,7 @@ use rusqlite::{ Connection, OpenFlags, OptionalExtension, TransactionBehavior, backup::Backup, params, }; use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; use uuid::Uuid; use worker_runtime::identity::{ @@ -230,6 +231,11 @@ const MIGRATIONS: &[Migration] = &[ name: "rename Workspace resource keys", apply: verify_workspace_resource_key_schema, }, + Migration { + version: 42, + name: "create Workspace Memory settings authority", + apply: create_workspace_memory_settings_authority, + }, ]; struct Migration { @@ -261,6 +267,22 @@ pub struct WorkspaceRecord { pub updated_at: String, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkspaceMemorySettingsRecord { + pub workspace_id: String, + pub settings_revision: u64, + pub language: String, + pub created_at: String, + pub updated_at: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkerCreateReservation { + pub worker_id: WorkerId, + pub create_fingerprint: String, + pub memory_settings: manifest::WorkspaceMemorySettingsSnapshot, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct RepositoryRecord { pub workspace_id: String, @@ -1140,14 +1162,97 @@ impl SqliteWorkspaceStore { f(&mut conn) } + pub(crate) fn get_workspace_memory_settings( + &self, + workspace_id: &str, + ) -> Result { + self.with_conn(|conn| { + conn.query_row( + "SELECT workspace_id, settings_revision, language, created_at, updated_at \ + FROM workspace_memory_settings WHERE workspace_id = ?1", + params![workspace_id], + |row| { + let revision = row.get::<_, i64>(1)?; + Ok(WorkspaceMemorySettingsRecord { + workspace_id: row.get(0)?, + settings_revision: revision + .try_into() + .map_err(|_| rusqlite::Error::IntegralValueOutOfRange(1, revision))?, + language: row.get(2)?, + created_at: row.get(3)?, + updated_at: row.get(4)?, + }) + }, + ) + .optional()? + .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string())) + }) + } + + pub(crate) fn update_workspace_memory_settings( + &self, + workspace_id: &str, + expected_revision: u64, + language: &str, + ) -> Result { + let language = normalize_workspace_memory_language(language)?; + self.with_conn_mut(|conn| { + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + let current_revision = tx + .query_row( + "SELECT settings_revision FROM workspace_memory_settings WHERE workspace_id = ?1", + params![workspace_id], + |row| row.get::<_, i64>(0), + ) + .optional()? + .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?; + let current_revision: u64 = current_revision.try_into().map_err(|_| { + rusqlite::Error::IntegralValueOutOfRange(0, current_revision) + })?; + if current_revision != expected_revision { + return Err(Error::WorkspaceConfigConflict(format!( + "Workspace Memory settings revision changed: expected {expected_revision}, current {current_revision}" + ))); + } + let next_revision = current_revision.checked_add(1).ok_or_else(|| { + Error::InvalidInput("Workspace Memory settings revision overflow".to_string()) + })?; + let now = chrono::Utc::now().to_rfc3339(); + tx.execute( + "UPDATE workspace_memory_settings \ + SET settings_revision = ?2, language = ?3, updated_at = ?4 \ + WHERE workspace_id = ?1", + params![workspace_id, next_revision as i64, language, now], + )?; + let record = tx.query_row( + "SELECT workspace_id, settings_revision, language, created_at, updated_at \ + FROM workspace_memory_settings WHERE workspace_id = ?1", + params![workspace_id], + |row| { + let revision = row.get::<_, i64>(1)?; + Ok(WorkspaceMemorySettingsRecord { + workspace_id: row.get(0)?, + settings_revision: revision as u64, + language: row.get(2)?, + created_at: row.get(3)?, + updated_at: row.get(4)?, + }) + }, + )?; + tx.commit()?; + Ok(record) + }) + } + pub(crate) fn reserve_worker_create( &self, workspace_id: &str, runtime_id: &str, allocation_key: &str, - create_fingerprint: &str, - ) -> Result { - if allocation_key.trim().is_empty() || create_fingerprint.trim().is_empty() { + request_fingerprint: &str, + current_memory_settings: &WorkspaceMemorySettingsRecord, + ) -> Result { + if allocation_key.trim().is_empty() || request_fingerprint.trim().is_empty() { return Err(Error::InvalidInput( "Worker create allocation key and fingerprint must be non-empty".to_string(), )); @@ -1156,7 +1261,8 @@ impl SqliteWorkspaceStore { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let existing = tx .query_row( - "SELECT worker_id, runtime_id, create_fingerprint \ + "SELECT worker_id, runtime_id, request_fingerprint, create_fingerprint, \ + memory_settings_revision, memory_language \ FROM worker_create_reservations \ WHERE workspace_id = ?1 AND allocation_key = ?2", params![workspace_id, allocation_key], @@ -1164,24 +1270,77 @@ impl SqliteWorkspaceStore { Ok(( row.get::<_, String>(0)?, row.get::<_, String>(1)?, - row.get::<_, String>(2)?, + row.get::<_, Option>(2)?, + row.get::<_, String>(3)?, + row.get::<_, Option>(4)?, + row.get::<_, Option>(5)?, )) }, ) .optional()?; - if let Some((worker_id, reserved_runtime_id, reserved_fingerprint)) = existing { - if reserved_runtime_id != runtime_id || reserved_fingerprint != create_fingerprint { + if let Some((worker_id, reserved_runtime_id, stored_request_fingerprint, create_fingerprint, revision, language)) = existing { + if reserved_runtime_id != runtime_id + || stored_request_fingerprint.as_deref() != Some(request_fingerprint) + { return Err(Error::InvalidInput(format!( - "Worker create allocation `{allocation_key}` was already used with different input" + "Worker create allocation {allocation_key} was already used with different input" ))); } - return worker_id.parse::().map_err(|_| { - Error::Store(format!( - "Worker create allocation `{allocation_key}` has a non-UUIDv7 worker id" + let revision = revision.ok_or_else(|| { + Error::InvalidInput(format!( + "Worker create allocation {allocation_key} has no persisted Memory settings snapshot" )) + })?; + let language = language.ok_or_else(|| { + Error::InvalidInput(format!( + "Worker create allocation {allocation_key} has no persisted Memory language" + )) + })?; + let worker_id = worker_id.parse::().map_err(|_| { + Error::Store(format!( + "Worker create allocation {allocation_key} has a non-UUIDv7 worker id" + )) + })?; + return Ok(WorkerCreateReservation { + worker_id, + create_fingerprint, + memory_settings: manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: workspace_id.to_string(), + settings_revision: revision.try_into().map_err(|_| { + rusqlite::Error::IntegralValueOutOfRange(4, revision) + })?, + language, + }, }); } + let (authoritative_revision, authoritative_language) = tx + .query_row( + "SELECT settings_revision, language FROM workspace_memory_settings WHERE workspace_id = ?1", + params![workspace_id], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)), + ) + .optional()? + .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?; + let authoritative_revision: u64 = authoritative_revision.try_into().map_err(|_| { + rusqlite::Error::IntegralValueOutOfRange(0, authoritative_revision) + })?; + if authoritative_revision != current_memory_settings.settings_revision + || authoritative_language != current_memory_settings.language + { + return Err(Error::WorkspaceConfigConflict( + "Workspace Memory settings changed while the Worker create reservation was being accepted" + .to_string(), + )); + } + + let snapshot = manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: workspace_id.to_string(), + settings_revision: authoritative_revision, + language: authoritative_language, + }; + let create_fingerprint = + bound_worker_create_fingerprint(request_fingerprint, &snapshot); let worker_id = WorkerId::now_v7(); let now = chrono::Utc::now().to_rfc3339(); allocate_resource_key( @@ -1194,19 +1353,27 @@ impl SqliteWorkspaceStore { tx.execute( "INSERT INTO worker_create_reservations(\ workspace_id, allocation_key, worker_id, runtime_id, create_fingerprint,\ - state, created_at, updated_at\ - ) VALUES (?1, ?2, ?3, ?4, ?5, 'reserved', ?6, ?6)", + state, created_at, updated_at, request_fingerprint,\ + memory_settings_revision, memory_language\ + ) VALUES (?1, ?2, ?3, ?4, ?5, 'reserved', ?6, ?6, ?7, ?8, ?9)", params![ workspace_id, allocation_key, worker_id.to_string(), runtime_id, create_fingerprint, - now + now, + request_fingerprint, + snapshot.settings_revision as i64, + snapshot.language, ], )?; tx.commit()?; - Ok(worker_id) + Ok(WorkerCreateReservation { + worker_id, + create_fingerprint, + memory_settings: snapshot, + }) }) } @@ -1375,8 +1542,9 @@ impl ControlPlaneStore for SqliteWorkspaceStore { } async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()> { - self.with_conn(|conn| { - conn.execute( + self.with_conn_mut(|conn| { + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + tx.execute( r#"INSERT INTO workspaces ( workspace_id, owner_account_id, display_name, state, created_at, updated_at ) VALUES (?1, ?2, ?3, ?4, ?5, ?6) @@ -1394,6 +1562,13 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.updated_at, ], )?; + tx.execute( + r#"INSERT OR IGNORE INTO workspace_memory_settings ( + workspace_id, settings_revision, language, created_at, updated_at + ) VALUES (?1, 1, 'English', ?2, ?3)"#, + params![record.workspace_id, record.created_at, record.updated_at], + )?; + tx.commit()?; Ok(()) })?; self.materialize_workspace_config(&record.workspace_id, &record.created_at) @@ -1518,6 +1693,16 @@ impl ControlPlaneStore for SqliteWorkspaceStore { record.workspace.updated_at, ], )?; + tx.execute( + r#"INSERT INTO workspace_memory_settings ( + workspace_id, settings_revision, language, created_at, updated_at + ) VALUES (?1, 1, 'English', ?2, ?3)"#, + params![ + record.workspace.workspace_id, + record.workspace.created_at, + record.workspace.updated_at, + ], + )?; tx.execute( r#"INSERT INTO repositories ( workspace_id, repository_id, name, kind, provider, uri, default_ref, @@ -5730,6 +5915,37 @@ fn collect_reference_diagnostics( Ok(()) } +fn normalize_workspace_memory_language(language: &str) -> Result { + match language.trim().to_ascii_lowercase().as_str() { + "english" | "en" => Ok("English".to_string()), + "japanese" | "ja" => Ok("Japanese".to_string()), + _ => Err(Error::InvalidInput( + "Workspace Memory language must be one of: English, Japanese".to_string(), + )), + } +} + +fn bound_worker_create_fingerprint( + request_fingerprint: &str, + snapshot: &manifest::WorkspaceMemorySettingsSnapshot, +) -> String { + let mut digest = Sha256::new(); + digest.update(b"workspace-worker-create-v2\0"); + digest.update(request_fingerprint.as_bytes()); + digest.update(b"\0"); + digest.update(snapshot.workspace_id.as_bytes()); + digest.update(b"\0"); + digest.update(snapshot.settings_revision.to_be_bytes()); + digest.update(b"\0"); + digest.update(snapshot.language.as_bytes()); + let encoded = digest + .finalize() + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + format!("sha256:{encoded}") +} + fn current_schema_version(conn: &Connection) -> Result { conn.query_row( "SELECT COALESCE(MAX(version), 0) FROM __yoi_schema_migrations", @@ -6373,6 +6589,42 @@ fn add_workspace_resource_human_keys(conn: &Connection) -> Result<()> { Ok(()) } +fn create_workspace_memory_settings_authority(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" + CREATE TABLE workspace_memory_settings ( + workspace_id TEXT PRIMARY KEY NOT NULL, + settings_revision INTEGER NOT NULL CHECK(settings_revision >= 1), + language TEXT NOT NULL CHECK(length(trim(language)) > 0), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE + ); + + INSERT INTO workspace_memory_settings ( + workspace_id, + settings_revision, + language, + created_at, + updated_at + ) + SELECT workspace_id, 1, 'English', created_at, updated_at + FROM workspaces; + + ALTER TABLE worker_create_reservations + ADD COLUMN request_fingerprint TEXT; + ALTER TABLE worker_create_reservations + ADD COLUMN memory_settings_revision INTEGER; + ALTER TABLE worker_create_reservations + ADD COLUMN memory_language TEXT; + + UPDATE worker_create_reservations + SET request_fingerprint = create_fingerprint; + "#, + )?; + Ok(()) +} + fn verify_workspace_resource_key_schema(conn: &Connection) -> Result<()> { ticket::migrate_sqlite_ticket_resource_key_schema_in_transaction(conn).map_err(|error| { Error::Store(format!( @@ -7682,7 +7934,7 @@ mod tests { let before = std::fs::read(&path).unwrap(); let plan = SqliteWorkspaceStore::migration_plan(&path).unwrap(); assert_eq!(plan.current_schema_version, 36); - assert_eq!(plan.target_schema_version, 41); + assert_eq!(plan.target_schema_version, 42); assert!(plan.migration_required); assert_eq!(plan.worker_count, 1); assert_eq!(plan.mappings[0].legacy_worker_id, 7); @@ -7696,7 +7948,7 @@ mod tests { store .with_conn(|conn| { assert!(table_exists(conn, "worker_diagnostics_archives")?); - assert_eq!(current_schema_version(conn)?, 41); + assert_eq!(current_schema_version(conn)?, 42); Ok(()) }) .unwrap(); @@ -7775,7 +8027,7 @@ mod tests { ), ] ); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); let foreign_key_error: Option = conn .query_row("PRAGMA foreign_key_check", [], |row| row.get(0)) .optional() @@ -7904,7 +8156,7 @@ INSERT INTO worker_orphan_diagnostics ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); assert!(!table_exists(&conn, "worker_control_delegation_operations").unwrap()); let controller_worker_id: String = conn .query_row( @@ -8022,10 +8274,36 @@ INSERT INTO worker_orphan_diagnostics ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); assert!(table_exists(&conn, "worker_workdir_attachment_reservations").unwrap()); } + #[test] + fn schema_v42_initializes_existing_workspaces_with_explicit_english_memory_settings() { + let conn = Connection::open_in_memory().unwrap(); + configure_sqlite(&conn).unwrap(); + apply_migrations_through(&conn, 41).unwrap(); + conn.execute( + "INSERT INTO workspaces (workspace_id, display_name, state, created_at, updated_at) \ + VALUES ('workspace-existing', 'Existing', 'active', '1', '1')", + [], + ) + .unwrap(); + + apply_migrations(&conn).unwrap(); + + assert_eq!(current_schema_version(&conn).unwrap(), 42); + let settings = conn + .query_row( + "SELECT settings_revision, language FROM workspace_memory_settings \ + WHERE workspace_id = 'workspace-existing'", + [], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)), + ) + .unwrap(); + assert_eq!(settings, (1, "English".to_string())); + } + #[test] fn schema_v26_removes_legacy_backend_flow_runtime_tables() { let conn = Connection::open_in_memory().unwrap(); @@ -8055,7 +8333,7 @@ CREATE TABLE flow_events (event_id TEXT PRIMARY KEY); apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); assert!(table_exists(&conn, "flow_sources").unwrap()); assert!(table_exists(&conn, "flow_source_revisions").unwrap()); assert!(!table_exists(&conn, "flow_instances").unwrap()); @@ -8122,7 +8400,7 @@ INSERT INTO worker_workdir_attachment_reservations ( apply_migrations(&conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); let repositories_sql: String = conn .query_row( "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'repositories'", @@ -8300,7 +8578,7 @@ INSERT INTO workdir_registry ( let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 41); + assert_eq!(store.schema_version().await.unwrap(), 42); assert!( !store .with_conn(|conn| table_exists(conn, "worker_workspace_credentials")) @@ -8317,7 +8595,7 @@ INSERT INTO workdir_registry ( store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 41); + assert_eq!(reopened.schema_version().await.unwrap(), 42); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -8386,22 +8664,50 @@ INSERT INTO workdir_registry ( .await .unwrap(); + let memory_settings = store.get_workspace_memory_settings("workspace-a").unwrap(); + assert_eq!(memory_settings.settings_revision, 1); + assert_eq!(memory_settings.language, "English"); let reserved = store - .reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:one") + .reserve_worker_create( + "workspace-a", + "arcadia", + "operation-1", + "sha256:one", + &memory_settings, + ) .unwrap(); assert_eq!( - reserved.as_uuid().get_version(), + reserved.worker_id.as_uuid().get_version(), Some(uuid::Version::SortRand) ); + assert_eq!(reserved.memory_settings.settings_revision, 1); + assert_eq!(reserved.memory_settings.language, "English"); + let updated_memory_settings = store + .update_workspace_memory_settings("workspace-a", 1, "ja") + .unwrap(); + assert_eq!(updated_memory_settings.settings_revision, 2); + assert_eq!(updated_memory_settings.language, "Japanese"); assert_eq!( store - .reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:one") + .reserve_worker_create( + "workspace-a", + "arcadia", + "operation-1", + "sha256:one", + &updated_memory_settings, + ) .unwrap(), reserved ); assert!( store - .reserve_worker_create("workspace-a", "arcadia", "operation-1", "sha256:different") + .reserve_worker_create( + "workspace-a", + "arcadia", + "operation-1", + "sha256:different", + &updated_memory_settings, + ) .is_err() ); assert_eq!( @@ -8409,7 +8715,7 @@ INSERT INTO workdir_registry ( .resource_key( "workspace-a", WorkspaceResourceKind::Worker, - &reserved.to_string() + &reserved.worker_id.to_string() ) .unwrap() .as_deref(), @@ -8419,31 +8725,38 @@ INSERT INTO workdir_registry ( store .resolve_resource_reference("workspace-a", WorkspaceResourceKind::Worker, "W-1") .unwrap(), - Some(reserved.to_string()) + Some(reserved.worker_id.to_string()) ); let second = store - .reserve_worker_create("workspace-a", "arcadia", "operation-2", "sha256:two") + .reserve_worker_create( + "workspace-a", + "arcadia", + "operation-2", + "sha256:two", + &updated_memory_settings, + ) .unwrap(); + assert_eq!(second.memory_settings.settings_revision, 2); assert_eq!( store .resource_key( "workspace-a", WorkspaceResourceKind::Worker, - &second.to_string() + &second.worker_id.to_string() ) .unwrap() .as_deref(), Some("W-2") ); store - .complete_worker_create_reservation("workspace-a", reserved) + .complete_worker_create_reservation("workspace-a", reserved.worker_id) .unwrap(); let state: String = store .with_conn(|conn| { conn.query_row( "SELECT state FROM worker_create_reservations \ WHERE workspace_id = 'workspace-a' AND worker_id = ?1", - [reserved.to_string()], + [reserved.worker_id.to_string()], |row| row.get(0), ) .map_err(Error::from) @@ -8816,13 +9129,13 @@ INSERT INTO worker_registry ( configure_sqlite(&conn).unwrap(); apply_migrations(&conn).unwrap(); conn.execute( - "INSERT INTO __yoi_schema_migrations (version, name) VALUES (42, 'future')", + "INSERT INTO __yoi_schema_migrations (version, name) VALUES (43, 'future')", [], ) .unwrap(); let error = apply_migrations(&conn).unwrap_err().to_string(); - assert!(error.contains("schema version 42 is newer"), "{error}"); + assert!(error.contains("schema version 43 is newer"), "{error}"); assert!(error.contains("refusing to serve"), "{error}"); } @@ -9043,7 +9356,7 @@ VALUES ('workspace-b', 'ticket-b', 'related', 'ticket-a', NULL, 'tester', '2026- apply_migrations(&mut conn).unwrap(); - assert_eq!(current_schema_version(&conn).unwrap(), 41); + assert_eq!(current_schema_version(&conn).unwrap(), 42); let workspace_id: Option = conn .query_row( "SELECT workspace_id FROM trusted_runtime_records WHERE runtime_id = 'runtime-a'", @@ -9660,7 +9973,7 @@ WHERE workspace_id = 'workspace-a' .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 41); + assert_eq!(store.schema_version().await.unwrap(), 42); store .with_conn(|conn| { @@ -9849,7 +10162,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn repository_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 41); + assert_eq!(store.schema_version().await.unwrap(), 42); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -9915,7 +10228,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn memory_authority_records_round_trip_and_close_staging() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 41); + assert_eq!(store.schema_version().await.unwrap(), 42); let workspace = WorkspaceRecord { workspace_id: "local-dev".to_string(), owner_account_id: None, @@ -10317,7 +10630,7 @@ CREATE TABLE ticket_assignment_operations ( #[tokio::test] async fn account_and_login_records_round_trip() { let store = SqliteWorkspaceStore::in_memory().unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 41); + assert_eq!(store.schema_version().await.unwrap(), 42); let now = "2026-07-22T00:00:00Z".to_string(); let account = AccountRecord { account_id: "acct-user-alice".to_string(), From 61d174b17461863c009029190a4a8ea4db2b1ae9 Mon Sep 17 00:00:00 2001 From: Hare Date: Fri, 21 Aug 2026 22:05:53 +0900 Subject: [PATCH 2/4] fix: enforce workspace memory settings contract --- crates/manifest/src/lib.rs | 25 ++++++++- crates/manifest/src/profile.rs | 77 +++++++++++++++++++++++++++- crates/worker-runtime/src/runtime.rs | 8 +-- crates/worker/src/worker.rs | 2 +- crates/workspace-server/src/store.rs | 47 ++++++++++++----- 5 files changed, 138 insertions(+), 21 deletions(-) diff --git a/crates/manifest/src/lib.rs b/crates/manifest/src/lib.rs index b1265a61..3d057c05 100644 --- a/crates/manifest/src/lib.rs +++ b/crates/manifest/src/lib.rs @@ -449,6 +449,17 @@ pub struct WebFetchConfig { pub allow_private_addresses: Option, } +/// Maximum Unicode scalar values accepted in a normalized Workspace Memory language name. +pub const MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS: usize = 64; + +/// Return whether a Workspace Memory language is already normalized and safe to persist. +pub fn is_normalized_workspace_memory_language(language: &str) -> bool { + !language.is_empty() + && language == language.trim() + && language.chars().count() <= MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS + && !language.chars().any(char::is_control) +} + /// Immutable Workspace Memory settings bound into a Worker launch snapshot. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] @@ -536,7 +547,7 @@ pub struct MemoryConfig { } impl MemoryConfig { - /// Replace any profile-authored language fields with a trusted Workspace snapshot. + /// Replace any untrusted manifest values with a trusted Workspace snapshot. pub fn bind_workspace_settings(&mut self, snapshot: &WorkspaceMemorySettingsSnapshot) { self.workspace_id = Some(snapshot.workspace_id.clone()); self.settings_revision = Some(snapshot.settings_revision); @@ -1256,6 +1267,18 @@ model_id = "claude-sonnet-4-20250514" ); } + #[test] + fn workspace_memory_language_validation_is_bounded_free_form_utf8() { + assert!(is_normalized_workspace_memory_language("Français")); + assert!(is_normalized_workspace_memory_language("日本語")); + assert!(!is_normalized_workspace_memory_language("")); + assert!(!is_normalized_workspace_memory_language(" English ")); + assert!(!is_normalized_workspace_memory_language("English\n")); + assert!(!is_normalized_workspace_memory_language( + &"x".repeat(MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS + 1) + )); + } + #[test] fn memory_section_with_language() { let toml = format!("{MINIMAL_REQUIRED}\n[memory]\nlanguage = \"Japanese\"\n"); diff --git a/crates/manifest/src/profile.rs b/crates/manifest/src/profile.rs index 63605b0d..83576db3 100644 --- a/crates/manifest/src/profile.rs +++ b/crates/manifest/src/profile.rs @@ -562,7 +562,7 @@ fn resolve_profile_value( mcp: profile.mcp, compaction, web: profile.web, - memory: profile.memory, + memory: profile.memory.map(Into::into), skills: profile.skills, }; let config = WorkerManifestConfig::builtin_defaults().merge(config.resolve_paths(profile_dir)); @@ -582,6 +582,51 @@ fn resolve_profile_value( }) } +#[derive(Debug, Default, Deserialize)] +#[serde(deny_unknown_fields)] +struct ProfileMemoryConfig { + #[serde(default)] + workspace_root: Option, + #[serde(default)] + query_result_limit: Option, + #[serde(default)] + query_excerpt_lines: Option, + #[serde(default)] + inject_summary: Option, + #[serde(default)] + extract_model: Option, + #[serde(default)] + extract_threshold: Option, + #[serde(default)] + extract_worker_max_turns: Option, + #[serde(default)] + consolidation_model: Option, + #[serde(default)] + consolidation_threshold_files: Option, + #[serde(default)] + consolidation_threshold_bytes: Option, +} + +impl From for MemoryConfig { + fn from(profile: ProfileMemoryConfig) -> Self { + Self { + workspace_root: profile.workspace_root, + query_result_limit: profile.query_result_limit, + query_excerpt_lines: profile.query_excerpt_lines, + inject_summary: profile.inject_summary, + workspace_id: None, + settings_revision: None, + language: None, + extract_model: profile.extract_model, + extract_threshold: profile.extract_threshold, + extract_worker_max_turns: profile.extract_worker_max_turns, + consolidation_model: profile.consolidation_model, + consolidation_threshold_files: profile.consolidation_threshold_files, + consolidation_threshold_bytes: profile.consolidation_threshold_bytes, + } + } +} + #[derive(Debug, Default, Deserialize)] #[serde(deny_unknown_fields)] struct ProfileConfig { @@ -612,7 +657,7 @@ struct ProfileConfig { #[serde(default)] web: Option, #[serde(default)] - memory: Option, + memory: Option, #[serde(default)] skills: Option, } @@ -1334,6 +1379,34 @@ mod tests { } } + #[test] + fn profile_rejects_workspace_memory_snapshot_authority_fields() { + let tmp = TempDir::new().unwrap(); + for (field, value) in [ + ("workspace_id", serde_json::json!("workspace-a")), + ("settings_revision", serde_json::json!(2)), + ("language", serde_json::json!("Japanese")), + ] { + let artifact = serde_json::json!({ "memory": { (field): value } }); + let error = resolve_profile_artifact_value( + artifact, + ProfileSource::Registry { + source: ProfileRegistrySource::Builtin, + name: "test".to_string(), + path: None, + provenance: None, + }, + tmp.path(), + "test-worker", + ) + .unwrap_err(); + assert!( + error.to_string().contains("unknown field"), + "unexpected error for {field}: {error}" + ); + } + } + #[test] fn builtin_companion_can_manage_workdirs() { let tmp = TempDir::new().unwrap(); diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 21dc6033..97ad176c 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -2810,9 +2810,9 @@ fn validate_create_workspace_scope( "Memory settings revision must be at least 1".to_string(), )); } - if !matches!(snapshot.language.as_str(), "English" | "Japanese") { + if !manifest::is_normalized_workspace_memory_language(&snapshot.language) { return Err(RuntimeError::InvalidRequest( - "Memory settings language must be a normalized supported value".to_string(), + "Memory settings language must be a normalized bounded UTF-8 value".to_string(), )); } Ok(()) @@ -3133,6 +3133,8 @@ mod tests { 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()); @@ -3147,7 +3149,7 @@ mod tests { request.memory_settings = Some(manifest::WorkspaceMemorySettingsSnapshot { workspace_id: "workspace-a".to_string(), settings_revision: 2, - language: "english".to_string(), + language: " english ".to_string(), }); assert!(validate_create_workspace_scope(&request, Some("workspace-a")).is_err()); } diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index ea00b0b6..01f9c0c9 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -5232,7 +5232,7 @@ fn validate_workspace_memory_snapshot( ))); } if snapshot.settings_revision == 0 - || !matches!(snapshot.language.as_str(), "English" | "Japanese") + || !manifest::is_normalized_workspace_memory_language(&snapshot.language) { return Err(WorkerError::InvalidState(format!( "Workspace Worker {worker_name} has corrupt Memory settings snapshot metadata" diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index 4456bc78..b3ba4fdb 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -1198,22 +1198,36 @@ impl SqliteWorkspaceStore { let language = normalize_workspace_memory_language(language)?; self.with_conn_mut(|conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; - let current_revision = tx + let current = tx .query_row( - "SELECT settings_revision FROM workspace_memory_settings WHERE workspace_id = ?1", + "SELECT workspace_id, settings_revision, language, created_at, updated_at \ + FROM workspace_memory_settings WHERE workspace_id = ?1", params![workspace_id], - |row| row.get::<_, i64>(0), + |row| { + let revision = row.get::<_, i64>(1)?; + Ok(WorkspaceMemorySettingsRecord { + workspace_id: row.get(0)?, + settings_revision: revision.try_into().map_err(|_| { + rusqlite::Error::IntegralValueOutOfRange(1, revision) + })?, + language: row.get(2)?, + created_at: row.get(3)?, + updated_at: row.get(4)?, + }) + }, ) .optional()? .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?; - let current_revision: u64 = current_revision.try_into().map_err(|_| { - rusqlite::Error::IntegralValueOutOfRange(0, current_revision) - })?; + let current_revision = current.settings_revision; if current_revision != expected_revision { return Err(Error::WorkspaceConfigConflict(format!( "Workspace Memory settings revision changed: expected {expected_revision}, current {current_revision}" ))); } + if current.language == language { + tx.commit()?; + return Ok(current); + } let next_revision = current_revision.checked_add(1).ok_or_else(|| { Error::InvalidInput("Workspace Memory settings revision overflow".to_string()) })?; @@ -5916,13 +5930,14 @@ fn collect_reference_diagnostics( } fn normalize_workspace_memory_language(language: &str) -> Result { - match language.trim().to_ascii_lowercase().as_str() { - "english" | "en" => Ok("English".to_string()), - "japanese" | "ja" => Ok("Japanese".to_string()), - _ => Err(Error::InvalidInput( - "Workspace Memory language must be one of: English, Japanese".to_string(), - )), + let language = language.trim(); + if !manifest::is_normalized_workspace_memory_language(language) { + return Err(Error::InvalidInput(format!( + "Workspace Memory language must be a non-empty UTF-8 string of at most {} characters without control characters", + manifest::MAX_WORKSPACE_MEMORY_LANGUAGE_CHARS + ))); } + Ok(language.to_string()) } fn bound_worker_create_fingerprint( @@ -8682,11 +8697,15 @@ INSERT INTO workdir_registry ( ); assert_eq!(reserved.memory_settings.settings_revision, 1); assert_eq!(reserved.memory_settings.language, "English"); + let unchanged_memory_settings = store + .update_workspace_memory_settings("workspace-a", 1, " English ") + .unwrap(); + assert_eq!(unchanged_memory_settings, memory_settings); let updated_memory_settings = store - .update_workspace_memory_settings("workspace-a", 1, "ja") + .update_workspace_memory_settings("workspace-a", 1, " Français ") .unwrap(); assert_eq!(updated_memory_settings.settings_revision, 2); - assert_eq!(updated_memory_settings.language, "Japanese"); + assert_eq!(updated_memory_settings.language, "Français"); assert_eq!( store .reserve_worker_create( From 2a3ece036411783dbc3c25380486673b29e6f819 Mon Sep 17 00:00:00 2001 From: Hare Date: Fri, 21 Aug 2026 22:16:41 +0900 Subject: [PATCH 3/4] fix: validate workspace memory settings authority --- crates/workspace-server/src/server.rs | 44 ++++++++----- crates/workspace-server/src/store.rs | 95 ++++++++++++++++++++++++--- 2 files changed, 114 insertions(+), 25 deletions(-) diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index c00d985b..2d9d9931 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -3013,14 +3013,7 @@ async fn scoped_get_workspace_memory_settings( State(api): State, AxumPath(workspace_id): AxumPath, ) -> ApiResult> { - if api - .config_store - .get_workspace(&workspace_id) - .await? - .is_none() - { - return Err(ApiError::from(Error::WorkspaceIdMismatch)); - } + validate_workspace_scope(&api, &workspace_id)?; let settings = api .config_store .get_workspace_memory_settings(&workspace_id) @@ -3037,14 +3030,7 @@ async fn scoped_update_workspace_memory_settings( AxumPath(workspace_id): AxumPath, Json(request): Json, ) -> ApiResult> { - if api - .config_store - .get_workspace(&workspace_id) - .await? - .is_none() - { - return Err(ApiError::from(Error::WorkspaceIdMismatch)); - } + validate_workspace_scope(&api, &workspace_id)?; let settings = api .config_store .update_workspace_memory_settings( @@ -16591,6 +16577,32 @@ mod tests { test_api_with_recording_backend(workspace_root).await.0 } + #[tokio::test] + async fn memory_settings_handlers_reject_foreign_workspace_path_scope() { + let temp = tempfile::tempdir().unwrap(); + let api = test_api(temp.path()).await; + assert!( + scoped_get_workspace_memory_settings( + State(api.clone()), + AxumPath("workspace-foreign".to_string()), + ) + .await + .is_err() + ); + assert!( + scoped_update_workspace_memory_settings( + State(api), + AxumPath("workspace-foreign".to_string()), + Json(workspace_api::UpdateWorkspaceMemorySettingsRequest { + expected_revision: 1, + language: "English".to_string(), + }), + ) + .await + .is_err() + ); + } + #[tokio::test] async fn destructive_worker_remove_rejects_browser_and_legacy_source_headers() { let headers = HeaderMap::new(); diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index b3ba4fdb..4f3b48f8 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -1166,7 +1166,7 @@ impl SqliteWorkspaceStore { &self, workspace_id: &str, ) -> Result { - self.with_conn(|conn| { + let record = self.with_conn(|conn| { conn.query_row( "SELECT workspace_id, settings_revision, language, created_at, updated_at \ FROM workspace_memory_settings WHERE workspace_id = ?1", @@ -1186,7 +1186,9 @@ impl SqliteWorkspaceStore { ) .optional()? .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string())) - }) + })?; + validate_workspace_memory_settings_record(&record, workspace_id)?; + Ok(record) } pub(crate) fn update_workspace_memory_settings( @@ -1218,6 +1220,7 @@ impl SqliteWorkspaceStore { ) .optional()? .ok_or_else(|| Error::Store("Workspace Memory settings are missing".to_string()))?; + validate_workspace_memory_settings_record(¤t, workspace_id)?; let current_revision = current.settings_revision; if current_revision != expected_revision { return Err(Error::WorkspaceConfigConflict(format!( @@ -1271,6 +1274,7 @@ impl SqliteWorkspaceStore { "Worker create allocation key and fingerprint must be non-empty".to_string(), )); } + validate_workspace_memory_settings_record(current_memory_settings, workspace_id)?; self.with_conn_mut(|conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let existing = tx @@ -1315,16 +1319,18 @@ impl SqliteWorkspaceStore { "Worker create allocation {allocation_key} has a non-UUIDv7 worker id" )) })?; + let snapshot = manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: workspace_id.to_string(), + settings_revision: revision.try_into().map_err(|_| { + rusqlite::Error::IntegralValueOutOfRange(4, revision) + })?, + language, + }; + validate_workspace_memory_settings_snapshot(&snapshot, workspace_id)?; return Ok(WorkerCreateReservation { worker_id, create_fingerprint, - memory_settings: manifest::WorkspaceMemorySettingsSnapshot { - workspace_id: workspace_id.to_string(), - settings_revision: revision.try_into().map_err(|_| { - rusqlite::Error::IntegralValueOutOfRange(4, revision) - })?, - language, - }, + memory_settings: snapshot, }); } @@ -1353,6 +1359,7 @@ impl SqliteWorkspaceStore { settings_revision: authoritative_revision, language: authoritative_language, }; + validate_workspace_memory_settings_snapshot(&snapshot, workspace_id)?; let create_fingerprint = bound_worker_create_fingerprint(request_fingerprint, &snapshot); let worker_id = WorkerId::now_v7(); @@ -5929,6 +5936,35 @@ fn collect_reference_diagnostics( Ok(()) } +fn validate_workspace_memory_settings_snapshot( + snapshot: &manifest::WorkspaceMemorySettingsSnapshot, + expected_workspace_id: &str, +) -> Result<()> { + if snapshot.workspace_id != expected_workspace_id + || snapshot.settings_revision == 0 + || !manifest::is_normalized_workspace_memory_language(&snapshot.language) + { + return Err(Error::Store( + "Workspace Memory settings are corrupt or belong to another Workspace".to_string(), + )); + } + Ok(()) +} + +fn validate_workspace_memory_settings_record( + record: &WorkspaceMemorySettingsRecord, + expected_workspace_id: &str, +) -> Result<()> { + validate_workspace_memory_settings_snapshot( + &manifest::WorkspaceMemorySettingsSnapshot { + workspace_id: record.workspace_id.clone(), + settings_revision: record.settings_revision, + language: record.language.clone(), + }, + expected_workspace_id, + ) +} + fn normalize_workspace_memory_language(language: &str) -> Result { let language = language.trim(); if !manifest::is_normalized_workspace_memory_language(language) { @@ -8782,6 +8818,47 @@ INSERT INTO workdir_registry ( }) .unwrap(); assert_eq!(state, "created"); + + store + .with_conn(|conn| { + conn.execute( + "UPDATE workspace_memory_settings SET language = ' English ' \ + WHERE workspace_id = 'workspace-a'", + [], + )?; + Ok(()) + }) + .unwrap(); + assert!(store.get_workspace_memory_settings("workspace-a").is_err()); + assert!( + store + .update_workspace_memory_settings("workspace-a", 2, "Spanish") + .is_err() + ); + let mut corrupt = updated_memory_settings.clone(); + corrupt.language = " English ".to_string(); + assert!( + store + .reserve_worker_create( + "workspace-a", + "arcadia", + "operation-corrupt", + "sha256:corrupt", + &corrupt, + ) + .is_err() + ); + + store + .with_conn(|conn| { + conn.execute( + "DELETE FROM workspace_memory_settings WHERE workspace_id = 'workspace-a'", + [], + )?; + Ok(()) + }) + .unwrap(); + assert!(store.get_workspace_memory_settings("workspace-a").is_err()); } #[tokio::test] From 71c906e04dd2682159bccf2a3ed646fd42dba65d Mon Sep 17 00:00:00 2001 From: Hare Date: Sat, 22 Aug 2026 00:28:33 +0900 Subject: [PATCH 4/4] fix: reject legacy workspace worker restore --- crates/worker-runtime/src/worker_backend.rs | 34 +++------------------ crates/worker/src/worker.rs | 18 ++++++++++- 2 files changed, 21 insertions(+), 31 deletions(-) diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 19233e5f..2b1123c8 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -2824,7 +2824,7 @@ mod tests { #[tokio::test] #[serial_test::serial(worker_allocation)] - async fn restore_pending_workspace_worker_without_system_prompt_fails_closed() { + async fn restore_legacy_workspace_worker_without_manifest_snapshot_requires_replacement() { let root = tempfile::tempdir().unwrap(); let runtime_store_dir = root.path().join("runtime"); let worker_ref = WorkerRef::new(crate::identity::WorkerId::from_legacy_u64(1)); @@ -2833,32 +2833,6 @@ mod tests { .join(worker_ref.worker_id.to_string()); let worker_name = ProfileRuntimeWorkerFactory::runtime_worker_name_for_ref(&worker_ref); let session_id = session_store::new_session_id(); - let manifest = manifest::WorkerManifest::from_toml(&format!( - r#" - [worker] - name = "{}" - pwd = "{}" - - [model] - scheme = "anthropic" - model_id = "test-model" - auth = {{ kind = "none" }} - - [engine] - max_tokens = 100 - - [feature.flow] - enabled = true - - [[scope.allow]] - target = "{}" - permission = "write" - "#, - worker_name, - root.path().display(), - root.path().display(), - )) - .unwrap(); WorkerAggregateStore::new(&worker_aggregate_dir, &worker_name) .unwrap() .set_active( @@ -2866,7 +2840,7 @@ mod tests { Some(session_store::WorkerActiveSegmentRef::pending_segment( session_id, )), - Some(serde_json::to_value(&manifest).unwrap()), + None, ) .unwrap(); @@ -2899,10 +2873,10 @@ mod tests { }) .await { - Ok(_) => panic!("pending Workspace Worker restore unexpectedly succeeded"), + Ok(_) => panic!("legacy Workspace Worker restore unexpectedly succeeded"), Err(error) => error, }; - assert!(error.contains("requires operation-owned launch material")); + assert!(error.contains("replacement Worker is required"), "{error}"); } #[tokio::test] diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 01f9c0c9..2ebac04a 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -4819,6 +4819,13 @@ where .ok_or_else(|| WorkerError::WorkerMetadataMissing { worker_name: worker_name.to_string(), })?; + if workspace_context.workspace_id().is_some() + && metadata.resolved_manifest_snapshot.is_none() + { + return Err(WorkerError::WorkerMetadataManifestSnapshotMissing { + worker_name: worker_name.to_string(), + }); + } let active = metadata .active .ok_or_else(|| WorkerError::WorkerMetadataInactive { @@ -4868,6 +4875,13 @@ where .ok_or_else(|| WorkerError::WorkerMetadataMissing { worker_name: worker_name.to_string(), })?; + if workspace_context.workspace_id().is_some() + && metadata.resolved_manifest_snapshot.is_none() + { + return Err(WorkerError::WorkerMetadataManifestSnapshotMissing { + worker_name: worker_name.to_string(), + }); + } let active = metadata .active .ok_or_else(|| WorkerError::WorkerMetadataInactive { @@ -5769,7 +5783,9 @@ pub enum WorkerError { session_id: SessionId, }, - #[error("worker metadata for {worker_name} does not include a resolved manifest snapshot")] + #[error( + "worker metadata for {worker_name} does not include a trusted resolved manifest snapshot; a replacement Worker is required" + )] WorkerMetadataManifestSnapshotMissing { worker_name: String }, #[error(