worker: run internal extraction through Worker
This commit is contained in:
@@ -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
|
//! Internal jobs are intentionally not Runtime-catalogued Workers. Each run still owns a
|
||||||
//! [`Worker`](crate::Worker). They do not append their harness prompt/input to
|
//! distinct Worker identity and executes through [`Worker`], including feature installation,
|
||||||
//! the foreground session history. Callers provide a deliberately limited tool
|
//! Workspace authority, session history, lifecycle records, usage accounting, cancellation,
|
||||||
//! surface and decide how to apply the result.
|
//! 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 std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
use llm_engine::llm_client::client::LlmClient;
|
use llm_engine::timeline::event::UsageEvent;
|
||||||
use llm_engine::llm_client::event::UsageEvent;
|
use llm_engine::{Engine, llm_client::LlmClient};
|
||||||
use llm_engine::tool::ToolDefinition;
|
use manifest::{Scope, WorkerManifest};
|
||||||
use llm_engine::{Engine, EngineError};
|
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 {
|
pub(crate) struct InternalWorkerSpec {
|
||||||
/// Stable slug for audit/debug/future persistence metadata.
|
pub identity: InternalWorkerIdentity,
|
||||||
pub slug: &'static str,
|
pub manifest: WorkerManifest,
|
||||||
/// 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<dyn LlmClient>,
|
pub client: Box<dyn LlmClient>,
|
||||||
/// Optional prompt-cache key.
|
pub system_prompt: String,
|
||||||
|
pub input: String,
|
||||||
pub cache_key: Option<String>,
|
pub cache_key: Option<String>,
|
||||||
/// Optional turn limit for the internal engine.
|
|
||||||
pub max_turns: Option<u32>,
|
pub max_turns: Option<u32>,
|
||||||
/// Deliberately limited tool surface for this internal run.
|
pub features: FeatureRegistryBuilder,
|
||||||
pub tools: Vec<ToolDefinition>,
|
pub required_tools: &'static [&'static str],
|
||||||
|
pub authority: InternalWorkerAuthority,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Result metadata for an internal worker run.
|
pub(crate) struct InternalWorkerResult {
|
||||||
#[derive(Debug, Clone, Default)]
|
|
||||||
pub(crate) struct InternalWorkerRunResult {
|
|
||||||
pub slug: &'static str,
|
|
||||||
/// Last usage event observed for this internal run.
|
|
||||||
pub usage: Option<UsageEvent>,
|
pub usage: Option<UsageEvent>,
|
||||||
|
pub identity: InternalWorkerIdentity,
|
||||||
|
pub lifecycle: WorkerRunResult,
|
||||||
|
pub history_entries: usize,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Error metadata for an internal worker run.
|
pub(crate) struct InternalWorkerError {
|
||||||
#[derive(Debug)]
|
pub source: WorkerError,
|
||||||
pub(crate) struct InternalWorkerRunError {
|
|
||||||
pub slug: &'static str,
|
|
||||||
pub source: EngineError,
|
|
||||||
/// Last usage event observed before the failure, if any.
|
|
||||||
pub usage: Option<UsageEvent>,
|
pub usage: Option<UsageEvent>,
|
||||||
|
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(
|
pub(crate) async fn run_internal_worker(
|
||||||
spec: InternalWorkerSpec,
|
spec: InternalWorkerSpec,
|
||||||
) -> Result<InternalWorkerRunResult, InternalWorkerRunError> {
|
) -> Result<InternalWorkerResult, InternalWorkerError> {
|
||||||
let InternalWorkerSpec {
|
let InternalWorkerSpec {
|
||||||
slug,
|
identity,
|
||||||
|
mut manifest,
|
||||||
|
client,
|
||||||
system_prompt,
|
system_prompt,
|
||||||
input,
|
input,
|
||||||
client,
|
|
||||||
cache_key,
|
cache_key,
|
||||||
max_turns,
|
max_turns,
|
||||||
tools,
|
features,
|
||||||
|
required_tools,
|
||||||
|
authority,
|
||||||
} = spec;
|
} = spec;
|
||||||
|
|
||||||
let mut worker = Engine::new(client).system_prompt(system_prompt);
|
// Internal identities are run-scoped and never enter the public Runtime Worker catalog.
|
||||||
worker.set_cache_key(cache_key);
|
manifest.worker.name = format!("internal-{}-{}", identity.kind, identity.run_id);
|
||||||
worker.set_max_turns(max_turns);
|
// Internal jobs only receive features supplied below. A parent manifest must not accidentally
|
||||||
worker.register_tools(tools);
|
// 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 last_usage = Arc::new(Mutex::new(None::<UsageEvent>));
|
||||||
let usage_capture_for_worker = usage_capture.clone();
|
let usage_slot = last_usage.clone();
|
||||||
worker.on_usage(move |event| {
|
let mut engine = Engine::new(client).system_prompt(system_prompt);
|
||||||
*usage_capture_for_worker
|
engine.on_usage(move |usage| {
|
||||||
.lock()
|
if let Ok(mut slot) = usage_slot.lock() {
|
||||||
.expect("internal worker usage capture poisoned") = Some(event.clone());
|
*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::<Vec<_>>()
|
||||||
|
.join("; ");
|
||||||
|
let missing = required_tools
|
||||||
|
.iter()
|
||||||
|
.filter(|required| {
|
||||||
|
!installed_tools
|
||||||
|
.iter()
|
||||||
|
.any(|installed| installed == **required)
|
||||||
|
})
|
||||||
|
.copied()
|
||||||
|
.collect::<Vec<_>>()
|
||||||
|
.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
|
match worker.run_text(&input).await {
|
||||||
.lock()
|
Ok(lifecycle) => Ok(InternalWorkerResult {
|
||||||
.expect("internal worker usage capture poisoned")
|
usage: last_usage.lock().ok().and_then(|slot| slot.clone()),
|
||||||
.clone();
|
identity,
|
||||||
|
lifecycle,
|
||||||
match run_result {
|
history_entries: store.entries_count(session_id, segment_id),
|
||||||
Ok(_output) => Ok(InternalWorkerRunResult { slug, usage }),
|
}),
|
||||||
Err(source) => Err(InternalWorkerRunError {
|
Err(source) => Err(InternalWorkerError {
|
||||||
slug,
|
|
||||||
source,
|
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<Mutex<HashMap<(SessionId, SegmentId), Vec<LogEntry>>>>,
|
||||||
|
traces: Arc<Mutex<HashMap<(SessionId, SegmentId), Vec<TraceEntry>>>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Vec<LogEntry>, 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<Vec<SessionId>, StoreError> {
|
||||||
|
let mut sessions = self
|
||||||
|
.entries
|
||||||
|
.lock()
|
||||||
|
.expect("ephemeral session store mutex poisoned")
|
||||||
|
.keys()
|
||||||
|
.map(|(session_id, _)| *session_id)
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
sessions.sort_unstable();
|
||||||
|
sessions.dedup();
|
||||||
|
sessions.reverse();
|
||||||
|
Ok(sessions)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn list_segments(&self, session_id: SessionId) -> Result<Vec<SegmentId>, 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::<Vec<_>>();
|
||||||
|
segments.sort_unstable();
|
||||||
|
segments.reverse();
|
||||||
|
Ok(segments)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn lookup_session_of(&self, segment_id: SegmentId) -> Result<Option<SessionId>, 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<bool, StoreError> {
|
||||||
|
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<usize, StoreError> {
|
||||||
|
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<AtomicUsize>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl LlmClient for OneTurnClient {
|
||||||
|
fn clone_boxed(&self) -> Box<dyn LlmClient> {
|
||||||
|
Box::new(self.clone())
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn stream(
|
||||||
|
&self,
|
||||||
|
_request: Request,
|
||||||
|
) -> Result<Pin<Box<dyn Stream<Item = Result<LlmEvent, ClientError>> + 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<AtomicUsize>,
|
||||||
|
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(_)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+44
-40
@@ -43,7 +43,9 @@ use crate::hook::{
|
|||||||
PreToolCall,
|
PreToolCall,
|
||||||
};
|
};
|
||||||
use crate::in_flight::InFlightEvents;
|
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_EXTENSION_DOMAIN: &str = "yoi.compaction";
|
||||||
const COMPACTION_BLOCK_ID: &str = "compact";
|
const COMPACTION_BLOCK_ID: &str = "compact";
|
||||||
@@ -3435,51 +3437,53 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
let session_explore_state =
|
let session_explore_state =
|
||||||
SessionExploreState::new(session_view, self.workspace_client_handle(), source);
|
SessionExploreState::new(session_view, self.workspace_client_handle(), source);
|
||||||
let input_text = render_extract_input(session_explore_state.view());
|
let input_text = render_extract_input(session_explore_state.view());
|
||||||
let mut internal_tools = Vec::new();
|
let features = FeatureRegistryBuilder::new()
|
||||||
let mut internal_hook_builder = HookRegistryBuilder::new();
|
.with_module(SessionExploreFeature::new(session_explore_state.clone()));
|
||||||
let feature_report = FeatureRegistryBuilder::new()
|
let mut internal_manifest = self.manifest.clone();
|
||||||
.with_module(SessionExploreFeature::new(session_explore_state.clone()))
|
internal_manifest.model = model.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 internal_result = run_internal_worker(InternalWorkerSpec {
|
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,
|
system_prompt: extract_system_prompt,
|
||||||
input: input_text,
|
input: input_text,
|
||||||
client,
|
|
||||||
cache_key: Some(self.segment_id().to_string()),
|
cache_key: Some(self.segment_id().to_string()),
|
||||||
max_turns: extract_worker_max_turns,
|
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;
|
.await;
|
||||||
let usage = match internal_result {
|
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) => {
|
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);
|
let usage = err.usage.as_ref().map(usage_audit_from_event);
|
||||||
audit
|
audit
|
||||||
.emit(
|
.emit(
|
||||||
@@ -3492,7 +3496,7 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
None,
|
None,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
return Err(WorkerError::Engine(err.source));
|
return Err(err.source);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -3625,8 +3629,8 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn lifecycle_status_for_worker_error(err: &EngineError) -> memory::audit::WorkerLifecycleStatus {
|
fn lifecycle_status_for_worker_error(err: &WorkerError) -> memory::audit::WorkerLifecycleStatus {
|
||||||
if matches!(err, EngineError::Cancelled) {
|
if matches!(err, WorkerError::Engine(EngineError::Cancelled)) {
|
||||||
memory::audit::WorkerLifecycleStatus::Cancelled
|
memory::audit::WorkerLifecycleStatus::Cancelled
|
||||||
} else {
|
} else {
|
||||||
memory::audit::WorkerLifecycleStatus::Failed
|
memory::audit::WorkerLifecycleStatus::Failed
|
||||||
|
|||||||
Reference in New Issue
Block a user