diff --git a/crates/worker/src/ipc/interceptor.rs b/crates/worker/src/ipc/interceptor.rs index 4f388853..55266d31 100644 --- a/crates/worker/src/ipc/interceptor.rs +++ b/crates/worker/src/ipc/interceptor.rs @@ -11,6 +11,7 @@ use std::borrow::Cow; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; +use arc_swap::ArcSwap; use async_trait::async_trait; use llm_engine::Item; use llm_engine::UsageRecord; @@ -64,7 +65,7 @@ pub(crate) struct WorkerInterceptor { pending_attachments: Arc>>, /// Prompt catalog used to render pending notification entries into the /// same system-message text that will be persisted in history. - prompts: Arc, + prompts: Arc>, /// Type-erased commit handle. The interceptor uses it to commit /// `LogEntry::SystemItem` entries directly (sync) before /// returning the corresponding `Item::system_message`s up to the @@ -84,7 +85,7 @@ impl WorkerInterceptor { usage_history: Option>>>, pending_notifies: NotifyBuffer, pending_attachments: Arc>>, - prompts: Arc, + prompts: Arc>, log_writer: Option>, ) -> Self { Self { @@ -208,10 +209,11 @@ impl Interceptor for WorkerInterceptor { return Ok(Vec::new()); } + let prompts = self.prompts.load_full(); let mut system_items: Vec = Vec::with_capacity(drained.len()); let mut items: Vec = Vec::with_capacity(drained.len()); for entry in drained { - match build_system_item(&entry, &self.prompts) { + match build_system_item(&entry, &prompts) { Ok(system_item) => { items.push(system_item.to_history_item()); system_items.push(system_item); @@ -440,6 +442,10 @@ mod tests { HookTurnEndAction, OnTurnEnd, PostToolCall, PreLlmRequest, PreToolCall, }; + fn test_prompts() -> Arc> { + Arc::new(ArcSwap::from(PromptCatalog::builtins_only().unwrap())) + } + struct CountingHook(Arc); #[async_trait] @@ -541,7 +547,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx = ctx_items; @@ -571,7 +577,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), Some(Arc::new(RecordingSystemItemCommitter { committed: Arc::clone(&committed), })), @@ -609,7 +615,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ) .with_usage_tracker(usage_tracker); @@ -634,7 +640,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx = ctx_items; @@ -675,7 +681,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx = ctx_items; @@ -702,7 +708,7 @@ mod tests { Some(history), NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx = ctx_items; @@ -723,7 +729,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx: Vec = Vec::new(); @@ -751,7 +757,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), Some(committer), ); @@ -798,7 +804,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); @@ -855,7 +861,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut info = task_tool_call_info("TaskList", serde_json::json!({"scope": "all"})); @@ -902,7 +908,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let info = task_tool_call_info("TaskList", serde_json::json!({})); @@ -953,7 +959,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let history = vec![Item::user_message("hi"), Item::assistant_message("done")]; @@ -985,7 +991,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), Some(Arc::new(RecordingSystemItemCommitter { committed: Arc::clone(&committed), })), @@ -1033,6 +1039,45 @@ mod tests { assert!(body.contains("track active work")); } + #[tokio::test] + async fn pending_notifications_use_the_latest_prompt_projection() { + let prompts = test_prompts(); + let buffer = NotifyBuffer::new(); + let interceptor = WorkerInterceptor::new( + Arc::new(HookRegistryBuilder::new().build()), + None, + None, + buffer.clone(), + Arc::new(Mutex::new(Vec::new())), + prompts.clone(), + None, + ); + + let current = prompts.load_full(); + let projection = current.projection(); + let mut templates = projection.templates.clone(); + templates.insert( + "internal.notify_wrapper".to_string(), + "CURRENT-PROJECTION {{ message }}".to_string(), + ); + let mut projection = crate::prompt::catalog::EffectivePromptCatalog::new( + templates, + 2, + projection.schema_fingerprint.clone(), + projection.toolchain_fingerprint.clone(), + ) + .unwrap(); + projection.source_digest = "source-2".to_string(); + prompts.store(Arc::new( + PromptCatalog::from_projection(projection).unwrap(), + )); + + buffer.push_notify("updated".to_string(), false); + let appends = interceptor.pending_history_appends().await.unwrap(); + assert_eq!(appends.len(), 1); + assert!(format!("{:?}", appends[0]).contains("CURRENT-PROJECTION updated")); + } + #[tokio::test] async fn pending_history_appends_drains_buffer_into_items() { let registry = Arc::new(HookRegistryBuilder::new().build()); @@ -1046,7 +1091,7 @@ mod tests { None, buffer.clone(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); @@ -1083,7 +1128,7 @@ mod tests { None, buffer.clone(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx: Vec = vec![Item::user_message("hi")]; @@ -1113,7 +1158,7 @@ mod tests { None, NotifyBuffer::new(), Arc::new(Mutex::new(Vec::new())), - PromptCatalog::builtins_only().unwrap(), + test_prompts(), None, ); let mut ctx: Vec = Vec::new(); diff --git a/crates/worker/src/spawn/tool.rs b/crates/worker/src/spawn/tool.rs index 17096356..26dae685 100644 --- a/crates/worker/src/spawn/tool.rs +++ b/crates/worker/src/spawn/tool.rs @@ -8,6 +8,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; +use arc_swap::ArcSwap; use async_trait::async_trait; use llm_engine::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput}; use manifest::{ @@ -946,7 +947,7 @@ pub(crate) fn sub_worker_spawn_tool( registry: Arc, spawner_manifest: WorkerManifest, spawner_scope: SharedScope, - prompts: Arc, + prompts: Arc>, ) -> ToolDefinition { sub_worker_spawn_tool_impl( spawner_name, @@ -972,13 +973,14 @@ fn sub_worker_spawn_tool_impl( registry: Arc, spawner_manifest: WorkerManifest, spawner_scope: SharedScope, - prompts: Arc, + prompts: Arc>, ) -> ToolDefinition { Arc::new(move || { let schema = schemars::schema_for!(SubWorkerSpawnInput); let schema_value = serde_json::to_value(schema).unwrap_or(serde_json::json!({})); let available_profiles = AvailableProfiles::discover(&workspace_root); let description = prompts + .load_full() .sub_worker_spawn_tool_description( &available_profiles.compact_list(), &available_profiles.default_label(), @@ -1002,7 +1004,7 @@ fn sub_worker_spawn_tool_impl( spawner_cwd.clone(), registry.clone(), spawner_manifest.clone(), - prompts.source(), + prompts.load_full().source(), available_profiles, spawner_scope.clone(), DelegationScope::from_config(&spawner_manifest.delegation_scope) diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 56be1b96..f7606b20 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -67,7 +67,7 @@ use crate::ipc::alerter::Alerter; use crate::ipc::interceptor::WorkerInterceptor; use crate::ipc::notify_buffer::NotifyBuffer; use crate::prompt::agents_md::read_agents_md; -use crate::prompt::catalog::{CatalogError, PromptCatalog}; +use crate::prompt::catalog::{CatalogError, PromptCatalog, WorkspacePromptProjection}; use crate::prompt::source::PromptCatalogSource; use crate::prompt::system::{SystemPromptContext, SystemPromptError, SystemPromptTemplate}; use crate::runtime::dir; @@ -224,6 +224,15 @@ pub trait WorkspaceClient: std::fmt::Debug + Send + Sync { fn execute(&self, request: WorkspaceRequest) -> Result; + /// Resolve the Workspace's current immutable Prompt projection for future + /// operation boundaries. Creation and restore continue to use persisted + /// launch/session state; this hook never reconstructs historical prompts. + fn current_prompt_projection( + &self, + ) -> Result, WorkspaceClientError> { + Ok(None) + } + /// Executes the destructive WorkerRemove operation through Runtime-owned source proof. /// Target identity is operation data; source identity and permission are never caller inputs. fn execute_worker_remove( @@ -285,6 +294,12 @@ impl WorkspaceClient for ReviewerChildWorkspaceClient { Some(&self.context) } + fn current_prompt_projection( + &self, + ) -> Result, WorkspaceClientError> { + self.inner.current_prompt_projection() + } + fn execute( &self, mut request: WorkspaceRequest, @@ -872,7 +887,7 @@ pub struct Worker { /// sections, ...). Built from the 4-layer overlay in /// [`Self::from_manifest`], or defaults to the builtin pack when a /// Worker is constructed through lower-level paths that have no loader. - prompts: Arc, + prompts: Arc>, /// When true (default), the system-prompt assembler may append resident /// context from the workspace Memory document. Internal disposable /// workers disable this so resident memory exposure is opt-in per Worker. @@ -1146,7 +1161,7 @@ impl Worker { // `set_system_prompt_template`) can be captured by `SegmentStart`. let session_id = session_store::new_session_id(); let segment_id = session_store::new_segment_id(); - let prompts = PromptCatalog::builtins_only()?; + let prompts = Arc::new(ArcSwap::from(PromptCatalog::builtins_only()?)); let delegation_scope = DelegationScope::from_config(&manifest.delegation_scope).map_err(WorkerError::Scope)?; let scope = SharedScope::new(scope); @@ -1232,10 +1247,49 @@ impl Worker { self.inject_resident_summary = enabled; } - pub fn prompts(&self) -> Arc { + pub fn prompts(&self) -> Arc> { Arc::clone(&self.prompts) } + fn refresh_prompt_projection_for_future_operations(&self) -> Result<(), WorkerError> { + // The launch catalog remains authoritative until the initial system + // Prompt has been rendered and committed. Later operation boundaries + // may adopt the Workspace's current immutable projection. + if self.system_prompt_template.is_some() { + return Ok(()); + } + let Some(projection) = self + .workspace_context + .client() + .current_prompt_projection() + .map_err(|source| WorkerError::WorkspacePromptProjection { + message: source.to_string(), + })? + else { + return Ok(()); + }; + projection + .validate() + .map_err(|source| WorkerError::WorkspacePromptProjection { + message: source.to_string(), + })?; + let current = self.prompts.load(); + if current.projection().config_revision == projection.config_revision + && current.projection().source_digest == projection.source_digest + && current.projection().catalog_digest == projection.projection_digest + { + return Ok(()); + } + let catalog = PromptCatalog::load( + &PromptCatalogSource::builtins_only().with_effective_catalog(projection.catalog), + ) + .map_err(|source| WorkerError::WorkspacePromptProjection { + message: source.to_string(), + })?; + self.prompts.store(catalog); + Ok(()) + } + /// The current segment ID. Read lock-free from the shared session /// pointer so fork-time swaps are observed immediately. pub fn segment_id(&self) -> SegmentId { @@ -2020,6 +2074,7 @@ impl Worker { .local_working_directory() .map(|local| local.cwd.display().to_string()) .unwrap_or_else(|| "no local working directory".to_string()); + let prompt_catalog = self.prompts.load_full(); let ctx = SystemPromptContext { now: chrono::Utc::now(), cwd: cwd_for_prompt.into(), @@ -2029,7 +2084,7 @@ impl Worker { feature_instructions: &self.feature_instructions, agents_md: agents_md_read.and_then(|read| read.body), resident_summary: resident_summary.as_deref(), - prompts: &self.prompts, + prompts: &prompt_catalog, }; let rendered = template .render(&ctx) @@ -2085,6 +2140,7 @@ impl Worker { /// store, and runs pre-run compact (joining any in-flight memory task /// first so extract sees a stable history range). async fn prepare_for_run(&mut self) -> Result<(), WorkerError> { + self.refresh_prompt_projection_for_future_operations()?; self.ensure_interceptor_installed(); self.ensure_system_prompt_materialized().await?; self.cleanup_finished_memory_task(); @@ -2430,10 +2486,12 @@ impl Worker { fn apply_interrupt_prep(&mut self) -> Result<(), WorkerError> { let tool_result_summary = self .prompts() + .load_full() .interrupt_tool_result_summary() .map_err(WorkerError::from)?; let system_note = self .prompts() + .load_full() .interrupt_system_note() .map_err(WorkerError::from)?; @@ -3173,6 +3231,7 @@ impl Worker { let summary_client: Box = self.build_compactor_client()?; let summary_system_prompt = self .prompts + .load_full() .compact_system() .map_err(WorkerError::PromptCatalog)?; let mut summary_worker = Engine::new(summary_client).system_prompt(summary_system_prompt); @@ -3802,7 +3861,11 @@ impl Worker { } }; let memory_language = memory_language(memory_cfg); - let extract_system_prompt = match self.prompts.memory_extract_system(memory_language) { + let extract_system_prompt = match self + .prompts + .load_full() + .memory_extract_system(memory_language) + { Ok(prompt) => prompt, Err(err) => { audit @@ -5446,6 +5509,9 @@ pub enum WorkerError { #[error(transparent)] PromptCatalog(#[from] CatalogError), + #[error("failed to resolve current Workspace Prompt projection: {message}")] + WorkspacePromptProjection { message: String }, + #[error(transparent)] Skill(#[from] SkillClientError), @@ -5513,7 +5579,7 @@ struct WorkerCommon { scope: Scope, delegation_scope: DelegationScope, client: Box, - prompts: Arc, + prompts: Arc>, system_prompt_template: Option, feature_instructions: Vec, } @@ -5655,7 +5721,7 @@ fn prepare_worker_common_from_scope( DelegationScope::from_config(&manifest.delegation_scope).map_err(WorkerError::Scope)?; let client = crate::model_client::build_client(&manifest.model)?; - let prompts = PromptCatalog::load(loader)?; + let prompts = Arc::new(ArcSwap::from(PromptCatalog::load(loader)?)); let system_prompt_template = if parse_template { Some( SystemPromptTemplate::parse(&manifest.engine.instruction, loader.clone())