From 2c8e617b2af1228319c7c7c44f5050e9bd98bf79 Mon Sep 17 00:00:00 2001 From: Hare Date: Thu, 6 Aug 2026 07:54:32 +0900 Subject: [PATCH] worker: run internal extraction through Worker --- crates/worker/src/internal_worker.rs | 452 +++++++++++++++++++++++---- crates/worker/src/worker.rs | 84 ++--- 2 files changed, 438 insertions(+), 98 deletions(-) diff --git a/crates/worker/src/internal_worker.rs b/crates/worker/src/internal_worker.rs index 3342f3d4..027a2ee7 100644 --- a/crates/worker/src/internal_worker.rs +++ b/crates/worker/src/internal_worker.rs @@ -1,92 +1,428 @@ -//! Reusable runner for Worker-internal LLM jobs. +//! Reusable execution substrate for Worker-backed internal jobs. //! -//! Internal workers are isolated model/tool runs owned by the foreground -//! [`Worker`](crate::Worker). They do not append their harness prompt/input to -//! the foreground session history. Callers provide a deliberately limited tool -//! surface and decide how to apply the result. +//! Internal jobs are intentionally not Runtime-catalogued Workers. Each run still owns a +//! distinct Worker identity and executes through [`Worker`], including feature installation, +//! Workspace authority, session history, lifecycle records, usage accounting, cancellation, +//! and error handling. The session store is in-memory and dropped with the run; callers that +//! need durable domain audit must keep using their domain authority (for example Memory audit). +use std::collections::HashMap; use std::sync::{Arc, Mutex}; -use llm_engine::llm_client::client::LlmClient; -use llm_engine::llm_client::event::UsageEvent; -use llm_engine::tool::ToolDefinition; -use llm_engine::{Engine, EngineError}; +use llm_engine::timeline::event::UsageEvent; +use llm_engine::{Engine, llm_client::LlmClient}; +use manifest::{Scope, WorkerManifest}; +use session_store::{LogEntry, SegmentId, SessionId, Store, StoreError, TraceEntry}; +use uuid::Uuid; + +use crate::feature::FeatureRegistryBuilder; +use crate::worker::{ + Worker, WorkerError, WorkerFilesystemAuthority, WorkerRunResult, WorkerWorkspaceContext, +}; + +/// Per-run identity for an internal Worker that is not registered in the public Runtime catalog. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct InternalWorkerIdentity { + pub kind: &'static str, + pub run_id: Uuid, +} + +/// Explicit authority granted to one internal Worker run. +/// +/// Extraction currently receives Workspace authority but no filesystem authority. Workdir +/// capabilities for Flow verifiers are deliberately left to the downstream Flow ticket. +pub(crate) struct InternalWorkerAuthority { + pub workspace: WorkerWorkspaceContext, + pub filesystem: WorkerFilesystemAuthority, + pub scope: Scope, +} -/// Specification for a single internal worker run. pub(crate) struct InternalWorkerSpec { - /// Stable slug for audit/debug/future persistence metadata. - pub slug: &'static str, - /// System prompt for the isolated engine. - pub system_prompt: String, - /// Initial user input for the isolated engine. - pub input: String, - /// Client used by the internal engine. + pub identity: InternalWorkerIdentity, + pub manifest: WorkerManifest, pub client: Box, - /// Optional prompt-cache key. + pub system_prompt: String, + pub input: String, pub cache_key: Option, - /// Optional turn limit for the internal engine. pub max_turns: Option, - /// Deliberately limited tool surface for this internal run. - pub tools: Vec, + pub features: FeatureRegistryBuilder, + pub required_tools: &'static [&'static str], + pub authority: InternalWorkerAuthority, } -/// Result metadata for an internal worker run. -#[derive(Debug, Clone, Default)] -pub(crate) struct InternalWorkerRunResult { - pub slug: &'static str, - /// Last usage event observed for this internal run. +pub(crate) struct InternalWorkerResult { pub usage: Option, + pub identity: InternalWorkerIdentity, + pub lifecycle: WorkerRunResult, + pub history_entries: usize, } -/// Error metadata for an internal worker run. -#[derive(Debug)] -pub(crate) struct InternalWorkerRunError { - pub slug: &'static str, - pub source: EngineError, - /// Last usage event observed before the failure, if any. +pub(crate) struct InternalWorkerError { + pub source: WorkerError, pub usage: Option, + pub identity: InternalWorkerIdentity, + pub history_entries: usize, } -/// Run an isolated internal worker engine. +/// Execute an internal job through the normal Worker substrate. +/// +/// The caller supplies the effective model client, an explicitly restricted feature set, and +/// explicit authority. No tools are registered directly on `Engine`, and no ambient filesystem +/// authority is inferred. pub(crate) async fn run_internal_worker( spec: InternalWorkerSpec, -) -> Result { +) -> Result { let InternalWorkerSpec { - slug, + identity, + mut manifest, + client, system_prompt, input, - client, cache_key, max_turns, - tools, + features, + required_tools, + authority, } = spec; - let mut worker = Engine::new(client).system_prompt(system_prompt); - worker.set_cache_key(cache_key); - worker.set_max_turns(max_turns); - worker.register_tools(tools); + // Internal identities are run-scoped and never enter the public Runtime Worker catalog. + manifest.worker.name = format!("internal-{}-{}", identity.kind, identity.run_id); + // Internal jobs only receive features supplied below. A parent manifest must not accidentally + // grant its normal public tool surface or recursively schedule memory work. + manifest.feature = Default::default(); + manifest.plugins = Default::default(); + manifest.mcp = Default::default(); + manifest.skills = None; + manifest.compaction = None; + manifest.memory = None; - let usage_capture = Arc::new(Mutex::new(None)); - let usage_capture_for_worker = usage_capture.clone(); - worker.on_usage(move |event| { - *usage_capture_for_worker - .lock() - .expect("internal worker usage capture poisoned") = Some(event.clone()); + let last_usage = Arc::new(Mutex::new(None::)); + let usage_slot = last_usage.clone(); + let mut engine = Engine::new(client).system_prompt(system_prompt); + engine.on_usage(move |usage| { + if let Ok(mut slot) = usage_slot.lock() { + *slot = Some(usage.clone()); + } }); + engine.set_cache_key(cache_key); + engine.set_max_turns(max_turns); + let store = EphemeralSessionStore::default(); + let mut worker = Worker::new( + manifest, + engine, + store.clone(), + authority.workspace, + authority.filesystem, + authority.scope, + ) + .await + .map_err(|source| InternalWorkerError { + source, + usage: last_usage.lock().ok().and_then(|slot| slot.clone()), + identity: identity.clone(), + history_entries: 0, + })?; - let run_result = worker.run(input).await; + let install_report = worker.install_features(features); + let installed_tools = install_report.installed_tool_names(); + let required_tools_missing = required_tools.iter().any(|required| { + !installed_tools + .iter() + .any(|installed| installed == required) + }); + let install_failed = install_report + .reports + .iter() + .any(|report| !report.installed); + if install_failed || required_tools_missing { + let diagnostics = install_report + .reports + .iter() + .flat_map(|report| report.diagnostics.iter()) + .map(|diagnostic| diagnostic.message.as_str()) + .collect::>() + .join("; "); + let missing = required_tools + .iter() + .filter(|required| { + !installed_tools + .iter() + .any(|installed| installed == **required) + }) + .copied() + .collect::>() + .join(", "); + return Err(InternalWorkerError { + source: WorkerError::FeatureInstall(format!( + "internal Worker feature installation failed: {diagnostics}; missing tools: {missing}" + )), + usage: last_usage.lock().ok().and_then(|slot| slot.clone()), + identity, + history_entries: 0, + }); + } + let session_id = worker.session_id(); + let segment_id = worker.segment_id(); - let usage = usage_capture - .lock() - .expect("internal worker usage capture poisoned") - .clone(); - - match run_result { - Ok(_output) => Ok(InternalWorkerRunResult { slug, usage }), - Err(source) => Err(InternalWorkerRunError { - slug, + match worker.run_text(&input).await { + Ok(lifecycle) => Ok(InternalWorkerResult { + usage: last_usage.lock().ok().and_then(|slot| slot.clone()), + identity, + lifecycle, + history_entries: store.entries_count(session_id, segment_id), + }), + Err(source) => Err(InternalWorkerError { source, - usage, + usage: last_usage.lock().ok().and_then(|slot| slot.clone()), + identity, + history_entries: store.entries_count(session_id, segment_id), }), } } + +/// Session history for an ephemeral internal Worker. +/// +/// Keeping the normal Store contract makes history/lifecycle/error records identical to a normal +/// Worker while avoiding a second public persistence/catalog policy for helper executions. +#[derive(Clone, Default)] +struct EphemeralSessionStore { + entries: Arc>>>, + traces: Arc>>>, +} + +impl EphemeralSessionStore { + fn entries_count(&self, session_id: SessionId, segment_id: SegmentId) -> usize { + self.entries + .lock() + .ok() + .and_then(|entries| entries.get(&(session_id, segment_id)).map(Vec::len)) + .unwrap_or_default() + } +} + +impl Store for EphemeralSessionStore { + fn append( + &self, + session_id: SessionId, + segment_id: SegmentId, + entry: &LogEntry, + ) -> Result<(), StoreError> { + self.entries + .lock() + .expect("ephemeral session store mutex poisoned") + .entry((session_id, segment_id)) + .or_default() + .push(entry.clone()); + Ok(()) + } + + fn read_all( + &self, + session_id: SessionId, + segment_id: SegmentId, + ) -> Result, StoreError> { + Ok(self + .entries + .lock() + .expect("ephemeral session store mutex poisoned") + .get(&(session_id, segment_id)) + .cloned() + .unwrap_or_default()) + } + + fn list_sessions(&self) -> Result, StoreError> { + let mut sessions = self + .entries + .lock() + .expect("ephemeral session store mutex poisoned") + .keys() + .map(|(session_id, _)| *session_id) + .collect::>(); + sessions.sort_unstable(); + sessions.dedup(); + sessions.reverse(); + Ok(sessions) + } + + fn list_segments(&self, session_id: SessionId) -> Result, StoreError> { + let mut segments = self + .entries + .lock() + .expect("ephemeral session store mutex poisoned") + .keys() + .filter_map(|(entry_session_id, segment_id)| { + (*entry_session_id == session_id).then_some(*segment_id) + }) + .collect::>(); + segments.sort_unstable(); + segments.reverse(); + Ok(segments) + } + + fn lookup_session_of(&self, segment_id: SegmentId) -> Result, StoreError> { + Ok(self + .entries + .lock() + .expect("ephemeral session store mutex poisoned") + .keys() + .find_map(|(session_id, entry_segment_id)| { + (*entry_segment_id == segment_id).then_some(*session_id) + })) + } + + fn create_segment( + &self, + session_id: SessionId, + segment_id: SegmentId, + entries: &[LogEntry], + ) -> Result<(), StoreError> { + self.entries + .lock() + .expect("ephemeral session store mutex poisoned") + .insert((session_id, segment_id), entries.to_vec()); + Ok(()) + } + + fn exists(&self, session_id: SessionId, segment_id: SegmentId) -> Result { + Ok(self + .entries + .lock() + .expect("ephemeral session store mutex poisoned") + .contains_key(&(session_id, segment_id))) + } + + fn read_entry_count( + &self, + session_id: SessionId, + segment_id: SegmentId, + ) -> Result { + Ok(self.entries_count(session_id, segment_id)) + } + + fn append_trace( + &self, + session_id: SessionId, + segment_id: SegmentId, + entry: &TraceEntry, + ) -> Result<(), StoreError> { + self.traces + .lock() + .expect("ephemeral session trace store mutex poisoned") + .entry((session_id, segment_id)) + .or_default() + .push(entry.clone()); + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use std::pin::Pin; + use std::sync::atomic::{AtomicUsize, Ordering}; + + use async_trait::async_trait; + use futures::Stream; + use llm_engine::llm_client::event::{Event as LlmEvent, ResponseStatus, StatusEvent}; + use llm_engine::llm_client::{ClientError, Request}; + + use super::*; + + #[derive(Clone)] + struct OneTurnClient { + calls: Arc, + } + + #[async_trait] + impl LlmClient for OneTurnClient { + fn clone_boxed(&self) -> Box { + Box::new(self.clone()) + } + + async fn stream( + &self, + _request: Request, + ) -> Result> + Send>>, ClientError> + { + self.calls.fetch_add(1, Ordering::SeqCst); + Ok(Box::pin(futures::stream::iter(vec![ + Ok(LlmEvent::text_block_start(0)), + Ok(LlmEvent::text_delta(0, "done")), + Ok(LlmEvent::text_block_stop(0, None)), + Ok(LlmEvent::Status(StatusEvent { + status: ResponseStatus::Completed, + })), + ]))) + } + } + + fn manifest() -> WorkerManifest { + WorkerManifest::from_toml( + r#" +[worker] +name = "parent" + +[model] +scheme = "anthropic" +model_id = "test-model" + +[engine] + +[[scope.allow]] +target = "/abs/scope" +permission = "write" +"#, + ) + .unwrap() + } + + fn spec( + calls: Arc, + required_tools: &'static [&'static str], + ) -> InternalWorkerSpec { + InternalWorkerSpec { + identity: InternalWorkerIdentity { + kind: "test", + run_id: Uuid::from_u128(1), + }, + manifest: manifest(), + client: Box::new(OneTurnClient { calls }), + system_prompt: "system".to_string(), + input: "input".to_string(), + cache_key: Some("internal-test".to_string()), + max_turns: Some(1), + features: FeatureRegistryBuilder::new(), + required_tools, + authority: InternalWorkerAuthority { + workspace: WorkerWorkspaceContext::no_workspace(), + filesystem: WorkerFilesystemAuthority::None, + scope: Scope::empty(), + }, + } + } + + #[tokio::test] + async fn executes_through_worker_and_records_ephemeral_history() { + let calls = Arc::new(AtomicUsize::new(0)); + let result = match run_internal_worker(spec(calls.clone(), &[])).await { + Ok(result) => result, + Err(error) => panic!("internal Worker should complete: {}", error.source), + }; + + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert!(matches!(result.lifecycle, WorkerRunResult::Finished)); + assert!(result.history_entries >= 4); + assert_eq!(result.identity.kind, "test"); + } + + #[tokio::test] + async fn rejects_missing_explicit_tools_before_model_execution() { + let calls = Arc::new(AtomicUsize::new(0)); + let error = match run_internal_worker(spec(calls.clone(), &["missing_tool"])).await { + Err(error) => error, + Ok(_) => panic!("required tool must be installed through the Worker registry"), + }; + + assert_eq!(calls.load(Ordering::SeqCst), 0); + assert!(matches!(error.source, WorkerError::FeatureInstall(_))); + } +} diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index d1a191cc..b3e2ccec 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -43,7 +43,9 @@ use crate::hook::{ PreToolCall, }; use crate::in_flight::InFlightEvents; -use crate::internal_worker::{InternalWorkerSpec, run_internal_worker}; +use crate::internal_worker::{ + InternalWorkerAuthority, InternalWorkerIdentity, InternalWorkerSpec, run_internal_worker, +}; const COMPACTION_EXTENSION_DOMAIN: &str = "yoi.compaction"; const COMPACTION_BLOCK_ID: &str = "compact"; @@ -3435,51 +3437,53 @@ impl Worker { let session_explore_state = SessionExploreState::new(session_view, self.workspace_client_handle(), source); let input_text = render_extract_input(session_explore_state.view()); - let mut internal_tools = Vec::new(); - let mut internal_hook_builder = HookRegistryBuilder::new(); - let feature_report = FeatureRegistryBuilder::new() - .with_module(SessionExploreFeature::new(session_explore_state.clone())) - .install_into_pending(&mut internal_tools, &mut internal_hook_builder); - let installed_tool_names = feature_report.installed_tool_names(); - let expected_extract_tools = [ - "search_evidence", - "read_evidence", - "stage_candidate", - "finish_extraction", - ]; - if !expected_extract_tools.iter().all(|name| { - installed_tool_names - .iter() - .any(|installed| installed == name) - }) { - audit - .emit( - self.workspace_client(), - event_tx, - memory::audit::WorkerLifecycleStatus::Failed, - "session_explore_feature_install_failed", - None, - Some(extract_audit_base), - None, - ) - .await; - return Err(WorkerError::FeatureInstall( - "session-explore feature install failed".to_string(), - )); - } + let features = FeatureRegistryBuilder::new() + .with_module(SessionExploreFeature::new(session_explore_state.clone())); + let mut internal_manifest = self.manifest.clone(); + internal_manifest.model = model.clone(); let internal_result = run_internal_worker(InternalWorkerSpec { - slug: "memory-extract", + identity: InternalWorkerIdentity { + kind: "memory-extract", + run_id: audit.run_id, + }, + manifest: internal_manifest, + client, system_prompt: extract_system_prompt, input: input_text, - client, cache_key: Some(self.segment_id().to_string()), max_turns: extract_worker_max_turns, - tools: internal_tools, + features, + required_tools: &[ + "search_evidence", + "read_evidence", + "stage_candidate", + "finish_extraction", + ], + authority: InternalWorkerAuthority { + workspace: self.workspace_context.clone(), + filesystem: WorkerFilesystemAuthority::None, + scope: Scope::empty(), + }, }) .await; let usage = match internal_result { - Ok(result) => result.usage.as_ref().map(usage_audit_from_event), + Ok(result) => { + tracing::debug!( + internal_worker_kind = result.identity.kind, + internal_worker_run_id = %result.identity.run_id, + history_entries = result.history_entries, + lifecycle = ?result.lifecycle, + "internal Worker execution completed" + ); + result.usage.as_ref().map(usage_audit_from_event) + } Err(err) => { + tracing::debug!( + internal_worker_kind = err.identity.kind, + internal_worker_run_id = %err.identity.run_id, + history_entries = err.history_entries, + "internal Worker execution failed" + ); let usage = err.usage.as_ref().map(usage_audit_from_event); audit .emit( @@ -3492,7 +3496,7 @@ impl Worker { None, ) .await; - return Err(WorkerError::Engine(err.source)); + return Err(err.source); } }; @@ -3625,8 +3629,8 @@ impl Worker { } } -fn lifecycle_status_for_worker_error(err: &EngineError) -> memory::audit::WorkerLifecycleStatus { - if matches!(err, EngineError::Cancelled) { +fn lifecycle_status_for_worker_error(err: &WorkerError) -> memory::audit::WorkerLifecycleStatus { + if matches!(err, WorkerError::Engine(EngineError::Cancelled)) { memory::audit::WorkerLifecycleStatus::Cancelled } else { memory::audit::WorkerLifecycleStatus::Failed