From b50b94612a1185dc5bdc935d83d63f9d6ace826b Mon Sep 17 00:00:00 2001 From: Hare Date: Sat, 11 Jul 2026 07:30:17 +0900 Subject: [PATCH] feat: add explicit worker filesystem authority --- crates/worker-runtime/src/worker_backend.rs | 90 +++++-- crates/worker/src/controller.rs | 93 ++++--- crates/worker/src/discovery.rs | 26 +- crates/worker/src/entrypoint.rs | 14 +- crates/worker/src/lib.rs | 5 +- crates/worker/src/prompt/system.rs | 14 +- crates/worker/src/ticket_event_notify.rs | 2 +- crates/worker/src/worker.rs | 250 +++++++++++++----- crates/worker/tests/compact_events_test.rs | 13 +- crates/worker/tests/consolidation_test.rs | 13 +- crates/worker/tests/controller_test.rs | 8 +- crates/worker/tests/session_metrics_test.rs | 39 ++- .../tests/system_prompt_template_test.rs | 10 +- 13 files changed, 422 insertions(+), 155 deletions(-) diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index c62d860c..36c9ed50 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -35,7 +35,7 @@ use session_store::{CombinedStore, FsWorkerStore}; use tokio::runtime::Runtime; use tokio::sync::broadcast; -use worker::{Worker, WorkerController, WorkerHandle}; +use worker::{Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle}; const DEFAULT_BACKEND_ID: &str = "worker-crate"; const RUNTIME_TASK_TIMEOUT: Duration = Duration::from_secs(10); @@ -298,11 +298,16 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { .as_ref() .map(|binding| binding.root().to_path_buf()) .unwrap_or_else(|| self.profile_base_dir.clone()); - let cwd = request + let filesystem_authority = request .working_directory .as_ref() - .map(|binding| binding.cwd().to_path_buf()) - .unwrap_or_else(|| self.cwd.clone()); + .map(|binding| { + WorkerFilesystemAuthority::local( + binding.root().to_path_buf(), + binding.cwd().to_path_buf(), + ) + }) + .unwrap_or(WorkerFilesystemAuthority::None); let selector = profile.as_deref().unwrap_or("builtin:default"); let archive = self .resolve_profile_source_archive(&request.request.profile_source) @@ -335,9 +340,15 @@ impl RuntimeWorkerFactory for ProfileRuntimeWorkerFactory { })?; let store = CombinedStore::new(session_store, worker_metadata_store); - let worker = Worker::from_manifest_with_context(manifest, store, loader, worker_root, cwd) - .await - .map_err(|err| format!("failed to create Worker from profile: {err}"))?; + let worker = Worker::from_manifest_with_context( + manifest, + store, + loader, + worker_root, + filesystem_authority, + ) + .await + .map_err(|err| format!("failed to create Worker from profile: {err}"))?; let runtime_base = self.runtime_base_dir()?; let (handle, _shutdown_rx) = WorkerController::spawn(worker, &runtime_base) @@ -806,7 +817,7 @@ where #[cfg(test)] mod tests { use super::*; - use std::collections::BTreeMap; + use std::collections::{BTreeMap, BTreeSet}; use std::fs; use std::pin::Pin; use std::process::Command; @@ -902,18 +913,27 @@ mod tests { FsStore::new(&self.store_dir).map_err(|err| err.to_string())?, FsWorkerStore::new(&self.worker_metadata_dir).map_err(|err| err.to_string())?, ); - let cwd = request + let filesystem_authority = request .working_directory .as_ref() - .map(|binding| binding.cwd().to_path_buf()) + .map(|binding| { + let cwd = binding.cwd().to_path_buf(); + self.observed_cwds.lock().unwrap().push(cwd.clone()); + WorkerFilesystemAuthority::local(binding.root().to_path_buf(), cwd) + }) + .unwrap_or(WorkerFilesystemAuthority::None); + let scope_root = request + .working_directory + .as_ref() + .map(|binding| binding.root().to_path_buf()) .unwrap_or_else(|| self.cwd.clone()); - self.observed_cwds.lock().unwrap().push(cwd.clone()); - let scope = Scope::writable(&cwd).map_err(|err| err.to_string())?; + let scope = Scope::writable(&scope_root).map_err(|err| err.to_string())?; let worker = Worker::new( manifest, Engine::new(self.client.clone()), store, - cwd, + self.cwd.clone(), + filesystem_authority, scope, ) .await @@ -925,6 +945,20 @@ mod tests { } } + fn core_filesystem_tool_names() -> BTreeSet<&'static str> { + ["Read", "Write", "Edit", "Glob", "Grep", "Bash"] + .into_iter() + .collect() + } + + fn captured_tool_names(client: &MockClient, index: usize) -> BTreeSet { + client.captured.lock().unwrap()[index] + .tools + .iter() + .map(|tool| tool.name.clone()) + .collect() + } + fn simple_text_events() -> Vec { vec![ LlmEvent::text_block_start(0), @@ -1109,7 +1143,7 @@ mod tests { cwd: cwd.path().to_path_buf(), store_dir: store.path().join("sessions"), worker_metadata_dir: store.path().join("workers"), - observed_cwds, + observed_cwds: observed_cwds.clone(), }; let backend = WorkerRuntimeExecutionBackend::new(factory).unwrap(); let runtime = EmbeddedRuntime::with_execution_backend( @@ -1148,6 +1182,14 @@ mod tests { } assert_eq!(client.captured.lock().unwrap().len(), 1); + assert!(observed_cwds.lock().unwrap().is_empty()); + let names = captured_tool_names(&client, 0); + for forbidden in core_filesystem_tool_names() { + assert!( + !names.contains(forbidden), + "no-workdir Worker unexpectedly exposed {forbidden}; tools={names:?}" + ); + } let observations = runtime .read_worker_observation_events(&detail.worker_ref, WorkerObservationCursor::zero()) .unwrap(); @@ -1166,7 +1208,7 @@ mod tests { let store = tempfile::tempdir().unwrap(); let observed_cwds = Arc::new(Mutex::new(Vec::new())); let factory = MockFactory { - client, + client: client.clone(), runtime_base: runtime_base.path().to_path_buf(), cwd: repo.path().to_path_buf(), store_dir: store.path().join("sessions"), @@ -1191,6 +1233,24 @@ mod tests { request.working_directory_request = Some(working_directory_request(repo.path())); let detail = runtime.create_worker(request).unwrap(); + runtime + .send_input(&detail.worker_ref, WorkerInput::user("inspect tools")) + .unwrap(); + let deadline = std::time::Instant::now() + Duration::from_secs(5); + while client.captured.lock().unwrap().is_empty() { + assert!( + std::time::Instant::now() < deadline, + "timed out waiting for materialized-worker request" + ); + std::thread::sleep(Duration::from_millis(20)); + } + let names = captured_tool_names(&client, 0); + for expected in core_filesystem_tool_names() { + assert!( + names.contains(expected), + "local Worker did not expose {expected}; tools={names:?}" + ); + } assert!(detail.execution.working_directory.is_some()); let cwds = observed_cwds.lock().unwrap(); diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index 5a98ee54..4741e5f3 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -274,7 +274,9 @@ impl WorkerController { manifest_toml.clone(), greeting, )); - shared_state.set_fs_view(crate::fs_view::WorkerFsView::new(fs_for_view)); + if let Some(fs_for_view) = fs_for_view { + shared_state.set_fs_view(crate::fs_view::WorkerFsView::new(fs_for_view)); + } shared_state.set_workflows( worker .workflow_completions() @@ -528,7 +530,10 @@ fn install_ticket_event_companion_notify_hook( return; } - let Ok(ticket_config) = TicketConfig::load_workspace(worker.cwd()) else { + let Some(local) = worker.local_working_directory() else { + return; + }; + let Ok(ticket_config) = TicketConfig::load_workspace(&local.cwd) else { return; }; let backend_root = ticket_config.backend_root().to_path_buf(); @@ -540,7 +545,7 @@ fn install_ticket_event_companion_notify_hook( worker.worker_metadata_store(), worker.manifest().worker.name.clone(), runtime_base, - worker.cwd().to_path_buf(), + Some(local.cwd.clone()), spawned_registry, ); match discovery.ensure_existing_peer(&companion_worker_name) { @@ -589,7 +594,7 @@ async fn register_worker_tools( spawner_socket: PathBuf, runtime_base: PathBuf, spawned_registry: Arc, -) -> std::io::Result +) -> std::io::Result> where C: LlmClient + Clone + 'static, St: Store + WorkerMetadataStore + Clone + 'static, @@ -597,7 +602,7 @@ where // Worker-immutable snapshots taken before the mutable worker borrow // below so the worker borrow doesn't conflict with reads on `worker`. let scope_handle = worker.scope().clone(); - let cwd = worker.cwd().to_path_buf(); + let local_filesystem = worker.local_working_directory().cloned(); let workspace_root = worker.workspace_root().to_path_buf(); let task_feature = worker.task_feature(); let session_id_for_usage = worker.segment_id().to_string(); @@ -611,24 +616,24 @@ where let worker_metadata_store = worker.store().clone(); let self_parent_socket = worker.callback_socket().cloned(); - // The Worker's SharedScope (already augmented with the bash-output - // Read rule by the caller) is the single source of truth — every - // ScopedFs (builtin tools, fs_view, compact worker) reads from it, - // and any future scope mutation (SpawnWorker-style revoke, future - // GrantScope) propagates through it. - let fs = tools::ScopedFs::with_shared_scope(scope_handle.clone(), cwd.clone()); - let tracker = tools::Tracker::new(); - // Same ScopedFs also powers the IPC `ListCompletions` query — keep - // a clone for the FS view we attach below, since the tools consume - // `fs` itself. - let fs_for_view = fs.clone(); - worker - .engine_mut() - .register_tools(tools::core_builtin_tools( - fs, - tracker.clone(), - bash_output_dir, - )); + // The Worker's SharedScope is the single source of truth for every + // ScopedFs when local filesystem authority exists. No-workdir Workers + // deliberately skip constructing/registering filesystem and Bash tools. + let (fs_for_view, tracker) = if let Some(local) = local_filesystem.as_ref() { + let fs = tools::ScopedFs::with_shared_scope(scope_handle.clone(), local.cwd.clone()); + let tracker = tools::Tracker::new(); + let fs_for_view = fs.clone(); + worker + .engine_mut() + .register_tools(tools::core_builtin_tools( + fs, + tracker.clone(), + bash_output_dir, + )); + (Some(fs_for_view), Some(tracker)) + } else { + (None, None) + }; if feature_config.web.enabled { worker .engine_mut() @@ -649,11 +654,20 @@ where } }; // Ticket tools are typed operations over the currently checked-out work - // tree. Use the Worker cwd rather than the runtime workspace root so a - // dedicated Orchestrator worktree gets its own `.yoi/tickets` backend. + // tree. They require explicit local filesystem authority; workspace_root + // is context only and must not be used as a cwd fallback. + let ticket_cwd = local_filesystem + .as_ref() + .map(|local| &local.cwd) + .ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "ticket tools require local Worker filesystem authority", + ) + })?; feature_registry.add_module( crate::feature::builtin::ticket::ticket_tools_feature_with_options( - &cwd, + ticket_cwd, feature_config.ticket.enabled.then_some(ticket_access), feature_config.ticket_orchestration.enabled, ), @@ -709,12 +723,21 @@ where "[feature.workers].enabled = true requires non-empty [[delegation_scope.allow]]", )); } + let spawner_cwd = local_filesystem + .as_ref() + .map(|local| local.cwd.clone()) + .ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "worker spawn tools require local Worker filesystem authority", + ) + })?; worker.register_tool(spawn_worker_tool( spawner_name.clone(), spawner_socket, runtime_base.clone(), workspace_root.clone(), - cwd.clone(), + spawner_cwd.clone(), spawned_registry.clone(), self_parent_socket, spawner_manifest, @@ -728,7 +751,7 @@ where worker_metadata_store, spawner_name, runtime_base, - cwd, + Some(spawner_cwd), spawned_registry, ); worker.register_tool(list_workers_tool(discovery.clone())); @@ -737,7 +760,9 @@ where } } let _feature_install_report = worker.install_features(feature_registry); - worker.attach_tracker(tracker); + if let Some(tracker) = tracker { + worker.attach_tracker(tracker); + } Ok(fs_for_view) } @@ -774,11 +799,14 @@ async fn controller_loop( .parent() .map(PathBuf::from) .unwrap_or_else(|| runtime_dir.path().to_path_buf()); + let discovery_cwd = worker + .local_working_directory() + .map(|local| local.cwd.clone()); let discovery = WorkerDiscovery::new( worker.store().clone(), spawner_name.clone(), discovery_runtime_base, - worker.cwd().to_path_buf(), + discovery_cwd, spawned_registry.clone(), ); let mut pending: Option = None; @@ -1445,7 +1473,10 @@ where .collect(); protocol::Greeting { worker_name: manifest.worker.name.clone(), - cwd: worker.cwd().display().to_string(), + cwd: worker + .local_working_directory() + .map(|local| local.cwd.display().to_string()) + .unwrap_or_default(), provider: provider_name, model: model_id, scope_summary: worker.scope_snapshot().summary(), diff --git a/crates/worker/src/discovery.rs b/crates/worker/src/discovery.rs index ad182f72..0d7e0a28 100644 --- a/crates/worker/src/discovery.rs +++ b/crates/worker/src/discovery.rs @@ -42,7 +42,7 @@ pub struct WorkerDiscovery { store: St, self_worker_name: String, runtime_base: PathBuf, - cwd: PathBuf, + cwd: Option, store_dir: Option, spawned_registry: Arc, } @@ -55,7 +55,7 @@ where store: St, self_worker_name: String, runtime_base: PathBuf, - cwd: PathBuf, + cwd: Option, spawned_registry: Arc, ) -> Self { let store_dir = store.root_dir(); @@ -432,13 +432,19 @@ where ) -> Result<(), WorkerDiscoveryError> { let runtime_command = WorkerRuntimeCommand::resolve().map_err(WorkerDiscoveryError::RestoreSpawn)?; + let Some(cwd) = &self.cwd else { + return Err(WorkerDiscoveryError::NotRestorable { + worker_name: worker_name.to_string(), + reason: "restore requires local Worker filesystem authority".into(), + }); + }; let mut command = Command::new(runtime_command.program()); command .args(runtime_command.prefix_args()) .arg("--worker") .arg(worker_name) .arg("--require-worker-state") - .current_dir(&self.cwd) + .current_dir(cwd) .stdin(Stdio::null()) .stdout(Stdio::null()) .stderr(Stdio::null()) @@ -1228,7 +1234,7 @@ mod tests { store.clone(), "parent".into(), runtime_base.clone(), - root.path().to_path_buf(), + Some(root.path().to_path_buf()), registry, ); @@ -1355,7 +1361,7 @@ mod tests { store.clone(), "source".into(), runtime_base.clone(), - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); let result = discovery.register_peer("target").unwrap(); @@ -1390,7 +1396,7 @@ mod tests { store, "source".into(), runtime_base, - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); @@ -1430,7 +1436,7 @@ mod tests { store.clone(), "source".into(), runtime_base, - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); @@ -1481,7 +1487,7 @@ mod tests { store, "source".into(), runtime_base.clone(), - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); @@ -1599,7 +1605,7 @@ mod tests { store, "source".into(), runtime_base.clone(), - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); @@ -1701,7 +1707,7 @@ mod tests { store, "source".into(), runtime_base, - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ); diff --git a/crates/worker/src/entrypoint.rs b/crates/worker/src/entrypoint.rs index 7a650074..d23a66de 100644 --- a/crates/worker/src/entrypoint.rs +++ b/crates/worker/src/entrypoint.rs @@ -2,7 +2,7 @@ use std::ffi::OsString; use std::path::{Path, PathBuf}; use std::process::ExitCode; -use crate::{PromptLoader, Worker, WorkerController}; +use crate::{PromptLoader, Worker, WorkerController, WorkerFilesystemAuthority}; use clap::{CommandFactory, FromArgMatches, Parser}; use manifest::{ Permission, ProfileResolveOptions, ProfileResolver, ProfileSelector, ScopeConfig, ScopeRule, @@ -511,6 +511,8 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { } }; let store = CombinedStore::new(session_store, worker_metadata_store); + let filesystem_authority = + WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()); let mut worker = if cli.adopt { let callback = match cli.callback.clone() { @@ -526,7 +528,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { loader, callback, workspace_root.clone(), - cwd.clone(), + filesystem_authority.clone(), ) .await { @@ -557,7 +559,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { store, loader, workspace_root.clone(), - cwd.clone(), + filesystem_authority.clone(), ) .await { @@ -577,7 +579,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { store, loader, workspace_root.clone(), - cwd.clone(), + filesystem_authority.clone(), ) .await { @@ -598,7 +600,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { store, loader, workspace_root.clone(), - cwd.clone(), + filesystem_authority.clone(), ) .await { @@ -620,7 +622,7 @@ async fn run_cli_inner(cli: Cli) -> ExitCode { store, loader, workspace_root.clone(), - cwd.clone(), + filesystem_authority.clone(), ) .await { diff --git a/crates/worker/src/lib.rs b/crates/worker/src/lib.rs index 399dfc4c..592cca40 100644 --- a/crates/worker/src/lib.rs +++ b/crates/worker/src/lib.rs @@ -38,4 +38,7 @@ pub use provider::{ProviderError, build_client}; pub use runtime::dir::RuntimeDir; pub use segment_log_sink::SegmentLogSink; pub use shared_state::WorkerSharedState; -pub use worker::{Worker, WorkerError, WorkerRunResult, apply_worker_manifest}; +pub use worker::{ + LocalWorkingDirectory, Worker, WorkerError, WorkerFilesystemAuthority, WorkerRunResult, + apply_worker_manifest, +}; diff --git a/crates/worker/src/prompt/system.rs b/crates/worker/src/prompt/system.rs index 7dfc930b..89fd6574 100644 --- a/crates/worker/src/prompt/system.rs +++ b/crates/worker/src/prompt/system.rs @@ -13,7 +13,9 @@ //! `set_system_prompt`. Subsequent turns and compactions reuse that //! materialised string verbatim. +use std::borrow::Cow; use std::collections::BTreeMap; +#[cfg(test)] use std::path::Path; use std::sync::Arc; @@ -146,7 +148,7 @@ impl std::fmt::Debug for SystemPromptTemplate { /// templates cannot drop them on the floor. pub struct SystemPromptContext<'a> { pub now: DateTime, - pub cwd: &'a Path, + pub cwd: Cow<'a, str>, /// Language policy exposed to instruction templates as `{{ language }}`. pub language: &'a str, pub scope: &'a Scope, @@ -189,7 +191,7 @@ impl<'a> SystemPromptContext<'a> { "datetime".into(), Value::from(self.now.to_rfc3339_opts(SecondsFormat::Secs, true)), ); - root.insert("cwd".into(), Value::from(self.cwd.display().to_string())); + root.insert("cwd".into(), Value::from(self.cwd.as_ref())); root.insert("language".into(), Value::from(self.language)); root.insert( "tools".into(), @@ -442,7 +444,7 @@ mod tests { ) -> SystemPromptContext<'a> { SystemPromptContext { now: fixed_now(), - cwd, + cwd: cwd.display().to_string().into(), language: manifest::defaults::WORKER_LANGUAGE, scope, tool_names: tools, @@ -461,7 +463,7 @@ mod tests { ) -> SystemPromptContext<'a> { SystemPromptContext { now: fixed_now(), - cwd, + cwd: cwd.display().to_string().into(), language: manifest::defaults::WORKER_LANGUAGE, scope, tool_names: Vec::new(), @@ -480,7 +482,7 @@ mod tests { ) -> SystemPromptContext<'a> { SystemPromptContext { now: fixed_now(), - cwd, + cwd: cwd.display().to_string().into(), language: manifest::defaults::WORKER_LANGUAGE, scope, tool_names: Vec::new(), @@ -499,7 +501,7 @@ mod tests { ) -> SystemPromptContext<'a> { SystemPromptContext { now: fixed_now(), - cwd, + cwd: cwd.display().to_string().into(), language: manifest::defaults::WORKER_LANGUAGE, scope, tool_names: Vec::new(), diff --git a/crates/worker/src/ticket_event_notify.rs b/crates/worker/src/ticket_event_notify.rs index bb719d29..1f0a898f 100644 --- a/crates/worker/src/ticket_event_notify.rs +++ b/crates/worker/src/ticket_event_notify.rs @@ -414,7 +414,7 @@ mod tests { store, "orchestrator".into(), runtime_base.clone(), - root.path().to_path_buf(), + Some(root.path().to_path_buf()), SpawnedWorkerRegistry::new(runtime_dir), ), "companion", diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 1673c621..972e455f 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -57,6 +57,41 @@ use tokio::task::JoinHandle; const RESTORE_RECONCILIATION_REACHABILITY_TIMEOUT: Duration = Duration::from_millis(500); +/// Explicit filesystem authority held by a Worker. +/// +/// `None` means the Worker has no local filesystem authority: no cwd, no +/// filesystem view, and no filesystem/Bash tool surface. Workspace context may +/// still exist separately for memory, workflows, and project records. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum WorkerFilesystemAuthority { + None, + Local(LocalWorkingDirectory), +} + +impl WorkerFilesystemAuthority { + pub fn local(root: PathBuf, cwd: PathBuf) -> Self { + Self::Local(LocalWorkingDirectory { root, cwd }) + } + + pub fn as_local(&self) -> Option<&LocalWorkingDirectory> { + match self { + Self::None => None, + Self::Local(local) => Some(local), + } + } +} + +/// Local filesystem authority for a Worker. +/// +/// `root` is the authority root retained for control-plane semantics; +/// `cwd` is the default working directory used by filesystem tools, Bash, +/// file references, and local worktree-scoped features. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct LocalWorkingDirectory { + pub root: PathBuf, + pub cwd: PathBuf, +} + /// `(SessionId, SegmentId)` pair the Worker is currently writing to. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct SegmentLocation { @@ -249,8 +284,9 @@ pub struct Worker { /// `segment_id` and append tally. `self.segment_id()` is a thin /// wrapper over `segment_state.segment_id()`. segment_state: Arc, - /// Absolute tool/process working directory of the Worker. - cwd: PathBuf, + /// Explicit local filesystem authority, or `None` for Workers with no + /// local cwd and no filesystem/Bash tool surface. + filesystem_authority: WorkerFilesystemAuthority, /// Absolute runtime workspace root used for project records, workflow, /// memory, Ticket config, Profile context, and spawned-child inheritance. workspace_root: PathBuf, @@ -446,7 +482,7 @@ impl Worker store: self.store.clone(), worker_metadata_writer: None, segment_state: self.segment_state.clone(), - cwd: self.cwd.clone(), + filesystem_authority: self.filesystem_authority.clone(), workspace_root: self.workspace_root.clone(), scope: self.scope.clone(), delegation_scope: self.delegation_scope.clone(), @@ -607,9 +643,10 @@ impl Worker impl Worker { /// Create a new Worker from a pre-built Engine and store. /// - /// Callers must pre-resolve `cwd` (absolute) and build a [`Scope`] + /// Callers must pass explicit filesystem authority and build a [`Scope`] /// — typically via [`Scope::from_config`] when coming from a - /// manifest, or [`Scope::writable`] in tests. + /// manifest, or [`Scope::writable`] in tests. Use + /// [`WorkerFilesystemAuthority::None`] for no-workdir Workers. /// /// Note: this constructor does **not** parse `manifest.worker.system_prompt` /// as a template. `Worker::from_manifest` is the production path for @@ -619,7 +656,8 @@ impl Worker { manifest: WorkerManifest, worker: Engine, store: St, - cwd: PathBuf, + workspace_root: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, scope: Scope, ) -> Result { // Segment creation is deferred to `ensure_segment_head` at first @@ -636,8 +674,8 @@ impl Worker { store, worker_metadata_writer: None, segment_state: SegmentState::new(session_id, segment_id, 0), - workspace_root: cwd.clone(), - cwd, + filesystem_authority, + workspace_root, scope: SharedScope::new(scope), delegation_scope, hook_builder: HookRegistryBuilder::new(), @@ -749,9 +787,14 @@ impl Worker { self.runtime_ticket_role = role; } - /// The Worker's tool/process working directory. - pub fn cwd(&self) -> &Path { - &self.cwd + /// Explicit filesystem authority held by this Worker. + pub fn filesystem_authority(&self) -> &WorkerFilesystemAuthority { + &self.filesystem_authority + } + + /// Local working directory when this Worker has local filesystem authority. + pub fn local_working_directory(&self) -> Option<&LocalWorkingDirectory> { + self.filesystem_authority.as_local() } /// The Worker's runtime workspace root. This stays separate from `cwd` for @@ -1376,9 +1419,13 @@ impl Worker { self.resident_exposure_snapshots(&resident, &resident_workflows); let worker_language = worker_language(&self.manifest.engine); let scope_snapshot = self.scope.snapshot(); + let cwd_for_prompt = self + .local_working_directory() + .map(|local| local.cwd.display().to_string()) + .unwrap_or_else(|| "no local working directory".to_string()); let ctx = SystemPromptContext { now: chrono::Utc::now(), - cwd: &self.cwd, + cwd: cwd_for_prompt.into(), language: worker_language, scope: &scope_snapshot, tool_names, @@ -1622,9 +1669,21 @@ impl Worker { /// unresolved placeholder stays in the flattened user message so the LLM /// still sees the intent. fn resolve_file_refs(&self, segments: &[Segment]) -> Vec { + let Some(local) = self.local_working_directory() else { + for seg in segments { + if let Segment::FileRef { path } = seg { + self.alert( + AlertLevel::Warn, + AlertSource::Worker, + format!("file ref @{path} could not be resolved: Worker has no local filesystem authority"), + ); + } + } + return Vec::new(); + }; let view = crate::fs_view::WorkerFsView::new(tools::ScopedFs::with_shared_scope( self.scope.clone(), - self.cwd.clone(), + local.cwd.clone(), )); let mut out = Vec::new(); for seg in segments { @@ -2549,11 +2608,13 @@ impl Worker { auto_read_budget, ))); - // Build an independent compact worker. Scope and cwd are shared - // with the main Worker (reads go through the same policy) but the - // Tracker is fresh — compact-time reads must not pollute the - // main session's recency list, which feeds `default_refs` above. - let scoped_fs = tools::ScopedFs::with_shared_scope(self.scope.clone(), self.cwd.clone()); + // Build an independent compact worker. When the main Worker has local + // filesystem authority, compact-time reads go through the same scope + // and cwd policy. No-workdir Workers deliberately omit compact-time + // filesystem tools as well. + let scoped_fs = self + .local_working_directory() + .map(|local| tools::ScopedFs::with_shared_scope(self.scope.clone(), local.cwd.clone())); let summary_tracker = tools::Tracker::new(); let summary_client: Box = self.build_compactor_client()?; let summary_system_prompt = self @@ -2591,10 +2652,12 @@ impl Worker { // Tools: read_file (shared scope, fresh tracker), bounded session // history exploration, and compact-specific tools that populate `ctx`. let compact_target_items = Arc::new(items_to_summarise.clone()); - summary_worker.register_tool(tools::read_tool(scoped_fs.clone(), summary_tracker)); + if let Some(scoped_fs) = scoped_fs.clone() { + summary_worker.register_tool(tools::read_tool(scoped_fs.clone(), summary_tracker)); + summary_worker.register_tool(mark_read_required_tool(scoped_fs, ctx.clone())); + } summary_worker.register_tool(search_session_log_tool(compact_target_items.clone())); summary_worker.register_tool(read_session_items_tool(compact_target_items)); - summary_worker.register_tool(mark_read_required_tool(scoped_fs.clone(), ctx.clone())); summary_worker.register_tool(add_reference_tool(ctx.clone())); summary_worker.register_tool(write_summary_tool(ctx.clone())); @@ -2671,8 +2734,12 @@ impl Worker { // logged and skipped inside `render_auto_read` rather than // aborting compaction — a missing / moved file should not fail // the whole compact. - let auto_read_messages = - WorkerFsView::new(scoped_fs.clone()).render_auto_read(&final_ctx.read_required); + let auto_read_messages = scoped_fs + .clone() + .map(|scoped_fs| { + WorkerFsView::new(scoped_fs).render_auto_read(&final_ctx.read_required) + }) + .unwrap_or_default(); // Reference list as a single system message; omitted when empty. let reference_message = (!final_ctx.references.is_empty()).then(|| { @@ -3824,7 +3891,8 @@ where loader: PromptLoader, ) -> Result { let cwd = current_cwd()?; - Self::from_manifest_with_context(manifest, store, loader, cwd.clone(), cwd).await + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); + Self::from_manifest_with_context(manifest, store, loader, cwd, authority).await } pub async fn from_manifest_with_context( @@ -3832,14 +3900,14 @@ where store: St, loader: PromptLoader, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, ) -> Result { let mut common = prepare_worker_common_with_context( &manifest, &loader, /* parse_template */ true, workspace_root, - cwd, + filesystem_authority, manifest.scope.clone(), )?; let skill_shadows = std::mem::take(&mut common.skill_shadows); @@ -3878,7 +3946,7 @@ where store, worker_metadata_writer, segment_state: SegmentState::new(session_id, segment_id, 0), - cwd: common.cwd, + filesystem_authority: common.filesystem_authority, workspace_root: common.workspace_root, scope: SharedScope::new(common.scope), delegation_scope: common.delegation_scope, @@ -3938,13 +4006,14 @@ where callback_socket: PathBuf, ) -> Result { let cwd = current_cwd()?; + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); Self::from_manifest_spawned_with_context( manifest, store, loader, callback_socket, - cwd.clone(), cwd, + authority, ) .await } @@ -3955,14 +4024,14 @@ where loader: PromptLoader, callback_socket: PathBuf, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, ) -> Result { let mut common = prepare_worker_common_with_context( &manifest, &loader, /* parse_template */ true, workspace_root, - cwd, + filesystem_authority, manifest.scope.clone(), )?; let skill_shadows = std::mem::take(&mut common.skill_shadows); @@ -3988,7 +4057,7 @@ where store, worker_metadata_writer, segment_state: SegmentState::new(session_id, segment_id, 0), - cwd: common.cwd, + filesystem_authority: common.filesystem_authority, workspace_root: common.workspace_root, scope: SharedScope::new(common.scope), delegation_scope: common.delegation_scope, @@ -4044,13 +4113,14 @@ where loader: PromptLoader, ) -> Result { let cwd = current_cwd()?; + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); Self::restore_from_worker_metadata_with_context( worker_name, manifest, store, loader, - cwd.clone(), cwd, + authority, ) .await } @@ -4061,7 +4131,7 @@ where store: St, loader: PromptLoader, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, ) -> Result { let metadata = store @@ -4092,7 +4162,7 @@ where store, loader, workspace_root, - cwd, + filesystem_authority, ) .await } @@ -4122,14 +4192,9 @@ where loader: PromptLoader, ) -> Result { let cwd = current_cwd()?; + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); Self::restore_from_manifest_with_context( - session_id, - segment_id, - manifest, - store, - loader, - cwd.clone(), - cwd, + session_id, segment_id, manifest, store, loader, cwd, authority, ) .await } @@ -4141,7 +4206,7 @@ where store: St, loader: PromptLoader, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, ) -> Result { // Read raw entries once so we can both reconstruct state and // seed the broadcast sink's mirror with the same prefix that @@ -4159,7 +4224,7 @@ where &loader, /* parse_template */ false, workspace_root, - cwd, + filesystem_authority, scope_config, )?; let skill_shadows = std::mem::take(&mut common.skill_shadows); @@ -4223,7 +4288,7 @@ where store, worker_metadata_writer, segment_state: SegmentState::new(session_id, segment_id, state.entries_count), - cwd: common.cwd, + filesystem_authority: common.filesystem_authority, workspace_root: common.workspace_root, scope: SharedScope::new(common.scope), delegation_scope: common.delegation_scope, @@ -4914,7 +4979,7 @@ pub enum WorkerError { /// [`prepare_worker_common_with_context`] from the resolved manifest and then split into Worker /// fields. struct WorkerCommon { - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, workspace_root: PathBuf, scope: Scope, delegation_scope: DelegationScope, @@ -5003,7 +5068,7 @@ fn prepare_worker_common_with_context( loader: &PromptLoader, parse_template: bool, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, scope_config: ScopeConfig, ) -> Result { let workspace_root = std::fs::canonicalize(&workspace_root).map_err(|source| { @@ -5012,10 +5077,23 @@ fn prepare_worker_common_with_context( source, } })?; - let cwd = std::fs::canonicalize(&cwd).map_err(|source| WorkerError::InvalidCwd { - cwd: cwd.clone(), - source, - })?; + let filesystem_authority = match filesystem_authority { + WorkerFilesystemAuthority::None => WorkerFilesystemAuthority::None, + WorkerFilesystemAuthority::Local(local) => { + let root = std::fs::canonicalize(&local.root).map_err(|source| { + WorkerError::InvalidWorkspaceRoot { + workspace_root: local.root.clone(), + source, + } + })?; + let cwd = + std::fs::canonicalize(&local.cwd).map_err(|source| WorkerError::InvalidCwd { + cwd: local.cwd.clone(), + source, + })?; + WorkerFilesystemAuthority::Local(LocalWorkingDirectory { root, cwd }) + } + }; let mut scope_config = scope_config; if let Some(mem) = manifest.memory.as_ref() { let layout = memory::WorkspaceLayout::resolve(mem, &workspace_root); @@ -5026,7 +5104,14 @@ fn prepare_worker_common_with_context( } scope_config.allow.extend(skill_dir_read_rules(manifest)); let scope = Scope::from_config(&scope_config).map_err(WorkerError::Scope)?; - prepare_worker_common_from_scope(manifest, loader, parse_template, workspace_root, cwd, scope) + prepare_worker_common_from_scope( + manifest, + loader, + parse_template, + workspace_root, + filesystem_authority, + scope, + ) } fn prepare_worker_common_from_scope( @@ -5034,14 +5119,18 @@ fn prepare_worker_common_from_scope( loader: &PromptLoader, parse_template: bool, workspace_root: PathBuf, - cwd: PathBuf, + filesystem_authority: WorkerFilesystemAuthority, scope: Scope, ) -> Result { if !scope.is_readable(&workspace_root) { return Err(WorkerError::WorkspaceRootOutsideScope { workspace_root }); } - if !scope.is_readable(&cwd) { - return Err(WorkerError::CwdOutsideScope { cwd }); + if let Some(local) = filesystem_authority.as_local() { + if !scope.is_readable(&local.cwd) { + return Err(WorkerError::CwdOutsideScope { + cwd: local.cwd.clone(), + }); + } } let delegation_scope = DelegationScope::from_config(&manifest.delegation_scope).map_err(WorkerError::Scope)?; @@ -5070,7 +5159,7 @@ fn prepare_worker_common_from_scope( }; Ok(WorkerCommon { - cwd, + filesystem_authority, workspace_root, scope, delegation_scope, @@ -5174,7 +5263,7 @@ mod spawned_context_tests { &PromptLoader::builtins_only(), false, workspace_root.clone(), - cwd.clone(), + WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()), manifest.scope.clone(), ) .unwrap(); @@ -5183,7 +5272,10 @@ mod spawned_context_tests { common.workspace_root, workspace_root.canonicalize().unwrap() ); - assert_eq!(common.cwd, cwd.canonicalize().unwrap()); + assert_eq!( + common.filesystem_authority.as_local().unwrap().cwd, + cwd.canonicalize().unwrap() + ); assert_eq!( common.memory_layout.as_ref().unwrap().root(), workspace_root.canonicalize().unwrap() @@ -5204,7 +5296,7 @@ mod spawned_context_tests { &PromptLoader::builtins_only(), false, workspace_root.clone(), - cwd.clone(), + WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()), ScopeConfig { allow: vec![ScopeRule { target: cwd.clone(), @@ -5242,7 +5334,7 @@ mod spawned_context_tests { &PromptLoader::builtins_only(), false, workspace_root.clone(), - cwd.clone(), + WorkerFilesystemAuthority::local(workspace_root.clone(), cwd.clone()), ScopeConfig { allow: vec![ScopeRule { target: workspace_root.clone(), @@ -5703,9 +5795,17 @@ mod build_summary_prompt_tests { let cwd = dir.path().join("workspace"); std::fs::create_dir_all(&cwd).unwrap(); let scope = Scope::writable(&cwd).unwrap(); - let mut worker = Worker::new(manifest, Engine::new(NoopClient), store, cwd, scope) - .await - .unwrap(); + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); + let mut worker = Worker::new( + manifest, + Engine::new(NoopClient), + store, + cwd.clone(), + authority, + scope, + ) + .await + .unwrap(); worker.ensure_segment_head().unwrap(); (dir, worker) } @@ -5851,9 +5951,17 @@ mod build_summary_prompt_tests { let cwd = dir.path().join("workspace"); std::fs::create_dir_all(&cwd).unwrap(); let scope = Scope::writable(&cwd).unwrap(); - let mut worker = Worker::new(manifest, Engine::new(NoopClient), store, cwd, scope) - .await - .unwrap(); + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); + let mut worker = Worker::new( + manifest, + Engine::new(NoopClient), + store, + cwd.clone(), + authority, + scope, + ) + .await + .unwrap(); worker.ensure_segment_head().unwrap(); worker.wire_history_persistence(); @@ -5978,9 +6086,17 @@ mod build_summary_prompt_tests { let mut manifest = minimal_manifest_with_skills(vec![]); manifest.memory = memory_config; let scope = Scope::writable(&cwd).unwrap(); - let mut worker = Worker::new(manifest, Engine::new(NoopClient), store, cwd.clone(), scope) - .await - .unwrap(); + let authority = WorkerFilesystemAuthority::local(cwd.clone(), cwd.clone()); + let mut worker = Worker::new( + manifest, + Engine::new(NoopClient), + store, + cwd.clone(), + authority, + scope, + ) + .await + .unwrap(); worker.memory_layout = worker .manifest .memory diff --git a/crates/worker/tests/compact_events_test.rs b/crates/worker/tests/compact_events_test.rs index ab6d1fba..2e3da105 100644 --- a/crates/worker/tests/compact_events_test.rs +++ b/crates/worker/tests/compact_events_test.rs @@ -164,9 +164,16 @@ async fn make_worker_with_manifest( std::mem::forget(pwd_tmp); let worker = Engine::new(client); - let mut worker = Worker::new(manifest, worker, store, pwd, scope) - .await - .unwrap(); + let mut worker = Worker::new( + manifest, + worker, + store, + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await + .unwrap(); worker.enable_worker_metadata_write_through().unwrap(); worker } diff --git a/crates/worker/tests/consolidation_test.rs b/crates/worker/tests/consolidation_test.rs index 4b68c64e..1d71fe8d 100644 --- a/crates/worker/tests/consolidation_test.rs +++ b/crates/worker/tests/consolidation_test.rs @@ -170,9 +170,16 @@ async fn make_worker_with( let scope = worker::Scope::writable(&pwd).unwrap(); let worker = Engine::new(client); - Worker::new(manifest, worker, store, pwd, scope) - .await - .unwrap() + Worker::new( + manifest, + worker, + store, + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await + .unwrap() } fn write_n_staging(layout: &WorkspaceLayout, n: usize) -> Vec { diff --git a/crates/worker/tests/controller_test.rs b/crates/worker/tests/controller_test.rs index 097967e2..429209a3 100644 --- a/crates/worker/tests/controller_test.rs +++ b/crates/worker/tests/controller_test.rs @@ -12,7 +12,10 @@ use llm_engine::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput}; use session_store::{CombinedStore, FsWorkerStore}; use session_store::{FsStore, LogEntry}; -use worker::{Event, Method, Worker, WorkerController, WorkerHandle, WorkerManifest, WorkerStatus}; +use worker::{ + Event, Method, Worker, WorkerController, WorkerFilesystemAuthority, WorkerHandle, + WorkerManifest, WorkerStatus, +}; type TestStore = CombinedStore; @@ -186,7 +189,8 @@ async fn make_worker_with_pwd_and_manifest( std::mem::forget(pwd_tmp); let worker = Engine::new(client); - let worker = Worker::new(manifest, worker, store, pwd.clone(), scope) + let authority = WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()); + let worker = Worker::new(manifest, worker, store, pwd.clone(), authority, scope) .await .unwrap(); (worker, pwd) diff --git a/crates/worker/tests/session_metrics_test.rs b/crates/worker/tests/session_metrics_test.rs index d932a1a2..950b90f5 100644 --- a/crates/worker/tests/session_metrics_test.rs +++ b/crates/worker/tests/session_metrics_test.rs @@ -190,9 +190,16 @@ async fn make_worker( let mut worker = Engine::new(client); worker.register_tool(big_content_tool_definition(tool_name)); - let worker = Worker::new(manifest, worker, store, pwd, scope) - .await - .unwrap(); + let worker = Worker::new( + manifest, + worker, + store, + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await + .unwrap(); (worker, store_tmp, pwd_tmp) } @@ -453,9 +460,16 @@ async fn metric_write_failure_emits_warn_alert_and_does_not_abort_run() { // the failure path: at least one metric attempts to write. let client = MockClient::new(vec![text_response_with_cache("hi", 0, 0)]); let worker = Engine::new(client); - let mut worker = Worker::new(manifest, worker, store.clone(), pwd, scope) - .await - .unwrap(); + let mut worker = Worker::new( + manifest, + worker, + store.clone(), + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await + .unwrap(); let (tx, mut rx) = broadcast::channel::(64); let alerter = worker::Alerter::new(tx); @@ -522,9 +536,16 @@ permission = "write" let pwd = pwd_tmp.path().to_path_buf(); let scope = worker::Scope::writable(&pwd).unwrap(); let worker = Engine::new(client); - let mut worker = Worker::new(manifest, worker, store.clone(), pwd, scope) - .await - .unwrap(); + let mut worker = Worker::new( + manifest, + worker, + store.clone(), + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await + .unwrap(); let session_id = worker.session_id(); let segment_id = worker.segment_id(); worker.run_text("hello").await.unwrap(); diff --git a/crates/worker/tests/system_prompt_template_test.rs b/crates/worker/tests/system_prompt_template_test.rs index f4422a88..0094b6be 100644 --- a/crates/worker/tests/system_prompt_template_test.rs +++ b/crates/worker/tests/system_prompt_template_test.rs @@ -123,7 +123,15 @@ async fn make_worker_with_body( std::mem::forget(user_prompts_tmp); let worker = Engine::new(client); - let mut worker = Worker::new(manifest, worker, store, pwd.clone(), scope).await?; + let mut worker = Worker::new( + manifest, + worker, + store, + pwd.clone(), + worker::WorkerFilesystemAuthority::local(pwd.clone(), pwd.clone()), + scope, + ) + .await?; let template = SystemPromptTemplate::parse("$user/test", loader) .map_err(|source| WorkerError::InvalidSystemPromptTemplate { source })?;