From 88b91a2c5599af26c3f19fc8e1aaa3fe8e4cac19 Mon Sep 17 00:00:00 2001 From: Hare Date: Sat, 18 Jul 2026 05:59:02 +0900 Subject: [PATCH] worker: add internal worker runner --- .yoi/tickets/00001KXRXAJFZ/item.md | 4 +- .yoi/tickets/00001KXRXAJFZ/resolution.md | 9 +++ .yoi/tickets/00001KXRXAJFZ/thread.md | 57 +++++++++++++++ crates/worker/src/internal_worker.rs | 92 ++++++++++++++++++++++++ crates/worker/src/lib.rs | 1 + crates/worker/src/worker.rs | 69 ++++++++---------- 6 files changed, 189 insertions(+), 43 deletions(-) create mode 100644 .yoi/tickets/00001KXRXAJFZ/resolution.md create mode 100644 crates/worker/src/internal_worker.rs diff --git a/.yoi/tickets/00001KXRXAJFZ/item.md b/.yoi/tickets/00001KXRXAJFZ/item.md index c64ada02..b6823120 100644 --- a/.yoi/tickets/00001KXRXAJFZ/item.md +++ b/.yoi/tickets/00001KXRXAJFZ/item.md @@ -1,8 +1,8 @@ --- title: 'Internal Worker runnerを実装する' -state: 'inprogress' +state: 'closed' created_at: '2026-07-17T20:47:11Z' -updated_at: '2026-07-17T20:47:48Z' +updated_at: '2026-07-17T20:58:47Z' assignee: null queued_by: 'yoi ticket' queued_at: '2026-07-17T20:47:48Z' diff --git a/.yoi/tickets/00001KXRXAJFZ/resolution.md b/.yoi/tickets/00001KXRXAJFZ/resolution.md new file mode 100644 index 00000000..e166dcba --- /dev/null +++ b/.yoi/tickets/00001KXRXAJFZ/resolution.md @@ -0,0 +1,9 @@ +Implemented reusable Worker-internal runner and migrated memory extract to use it. + +The new runner provides a Worker-owned path for isolated internal LLM jobs with caller-supplied prompt, input, limited tools, cache key, max turns, and usage capture. Memory extract now uses this runner while preserving the existing trigger and staging behavior. + +Validation passed: +- `cargo fmt --check` +- `cargo test -p memory` +- `cargo test -p worker` +- `nix build .#yoi` diff --git a/.yoi/tickets/00001KXRXAJFZ/thread.md b/.yoi/tickets/00001KXRXAJFZ/thread.md index e269ebef..3f2c0a5b 100644 --- a/.yoi/tickets/00001KXRXAJFZ/thread.md +++ b/.yoi/tickets/00001KXRXAJFZ/thread.md @@ -22,4 +22,61 @@ Ticket を `yoi ticket` が queued にしました。 State changed to `inprogress`. +--- + + + +## Implementation report + +Implemented reusable Worker-internal runner and migrated memory extract to use it. + +Changes: +- Added `crates/worker/src/internal_worker.rs` with `InternalWorkerSpec`, `InternalWorkerRunResult`, `InternalWorkerRunError`, and `run_internal_worker`. +- The runner owns isolated Engine creation, limited tool registration, cache key / max turns, usage capture, and foreground-history isolation by construction. +- Updated memory extract path to run through the internal runner with only the extract tool surface supplied by the caller. +- Kept existing extract behavior: post-run threshold trigger, transitional `write_extracted` path, and flat staging output. +- Left persistence / Protocol exposure / session-explore feature for later Tickets. + +Validation: +- `cargo fmt --check` +- `cargo test -p memory` +- `cargo test -p worker` +- `nix build .#yoi` + + +--- + + + +## State changed + +State changed to `done`. + + +--- + + + +## State changed + +Ticket を closed にしました。 + + +--- + + + +## 完了 + +Implemented reusable Worker-internal runner and migrated memory extract to use it. + +The new runner provides a Worker-owned path for isolated internal LLM jobs with caller-supplied prompt, input, limited tools, cache key, max turns, and usage capture. Memory extract now uses this runner while preserving the existing trigger and staging behavior. + +Validation passed: +- `cargo fmt --check` +- `cargo test -p memory` +- `cargo test -p worker` +- `nix build .#yoi` + + --- diff --git a/crates/worker/src/internal_worker.rs b/crates/worker/src/internal_worker.rs new file mode 100644 index 00000000..7eff744e --- /dev/null +++ b/crates/worker/src/internal_worker.rs @@ -0,0 +1,92 @@ +//! Reusable runner for Worker-internal LLM 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. + +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}; + +/// Specification for a single internal worker run. +pub(crate) struct InternalWorkerSpec { + /// Stable purpose label for audit/debug/persistence metadata. + pub purpose: &'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 client: Box, + /// Optional prompt-cache key. + 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, +} + +/// Result metadata for an internal worker run. +#[derive(Debug, Clone, Default)] +pub(crate) struct InternalWorkerRunResult { + pub purpose: &'static str, + /// Last usage event observed for this internal run. + pub usage: Option, +} + +/// Error metadata for an internal worker run. +#[derive(Debug)] +pub(crate) struct InternalWorkerRunError { + pub purpose: &'static str, + pub source: EngineError, + /// Last usage event observed before the failure, if any. + pub usage: Option, +} + +/// Run an isolated internal worker engine. +pub(crate) async fn run_internal_worker( + spec: InternalWorkerSpec, +) -> Result { + let InternalWorkerSpec { + purpose, + system_prompt, + input, + client, + cache_key, + max_turns, + tools, + } = 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); + + 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 run_result = worker.run(input).await; + + let usage = usage_capture + .lock() + .expect("internal worker usage capture poisoned") + .clone(); + + match run_result { + Ok(_output) => Ok(InternalWorkerRunResult { purpose, usage }), + Err(source) => Err(InternalWorkerRunError { + purpose, + source, + usage, + }), + } +} diff --git a/crates/worker/src/lib.rs b/crates/worker/src/lib.rs index cf682fbb..25425ba8 100644 --- a/crates/worker/src/lib.rs +++ b/crates/worker/src/lib.rs @@ -16,6 +16,7 @@ mod shutdown_after_idle; pub mod skill; pub mod spawn; +mod internal_worker; mod interrupt_prep; mod permission; mod ticket_event_notify; diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 6e621858..66492777 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -37,6 +37,7 @@ use crate::hook::{ PreToolCall, }; use crate::in_flight::InFlightEvents; +use crate::internal_worker::{InternalWorkerSpec, run_internal_worker}; const COMPACTION_EXTENSION_DOMAIN: &str = "yoi.compaction"; const COMPACTION_BLOCK_ID: &str = "compact"; @@ -3129,40 +3130,34 @@ impl Worker { return Err(WorkerError::PromptCatalog(err)); } }; - let mut extract_worker = Engine::new(client).system_prompt(extract_system_prompt); - extract_worker.set_cache_key(Some(self.segment_id().to_string())); - - extract_worker.set_max_turns(extract_worker_max_turns); - - let usage_capture = Arc::new(Mutex::new(None)); - let usage_capture_for_worker = usage_capture.clone(); - extract_worker.on_usage(move |event| { - *usage_capture_for_worker - .lock() - .expect("memory extract usage capture poisoned") = - Some(usage_audit_from_event(event)); - }); - let ctx = Arc::new(extract::ExtractWorkerContext::new()); - extract_worker.register_tool(extract::write_extracted_tool(ctx.clone())); - let input_text = extract::build_extract_input(&items_to_extract); - if let Err(err) = extract_worker.run(input_text).await { - let usage = usage_capture - .lock() - .expect("memory extract usage capture poisoned") - .clone(); - audit.emit( - &layout, - event_tx, - lifecycle_status_for_worker_error(&err), - format!("worker_failed: {err}"), - usage, - Some(extract_audit_base), - None, - ); - return Err(WorkerError::Engine(err)); - } + let internal_result = run_internal_worker(InternalWorkerSpec { + purpose: "memory_extract", + system_prompt: extract_system_prompt, + input: input_text, + client, + cache_key: Some(self.segment_id().to_string()), + max_turns: extract_worker_max_turns, + tools: vec![extract::write_extracted_tool(ctx.clone())], + }) + .await; + let usage = match internal_result { + Ok(result) => result.usage.as_ref().map(usage_audit_from_event), + Err(err) => { + let usage = err.usage.as_ref().map(usage_audit_from_event); + audit.emit( + &layout, + event_tx, + lifecycle_status_for_worker_error(&err.source), + format!("worker_failed: {}", err.source), + usage, + Some(extract_audit_base), + None, + ); + return Err(WorkerError::Engine(err.source)); + } + }; let payload = ctx.take_payload().unwrap_or_else(|| { tracing::warn!( @@ -3182,16 +3177,12 @@ impl Worker { match extract::write_staging(&layout, source, payload) { Ok(results) => results, Err(err) => { - let usage = usage_capture - .lock() - .expect("memory extract usage capture poisoned") - .clone(); audit.emit( &layout, event_tx, memory::audit::WorkerLifecycleStatus::Failed, format!("staging_write_failed: {err}"), - usage, + usage.clone(), Some(extract_audit_base), None, ); @@ -3230,10 +3221,6 @@ impl Worker { .staging_paths .push(result.path.display().to_string()); } - let usage = usage_capture - .lock() - .expect("memory extract usage capture poisoned") - .clone(); let reason = if staging_id.is_empty() { "completed_no_staging_output" } else {