Skip to main content

vv_agent/
runner.rs

1mod builder;
2mod checkpoint_runtime;
3mod event_stream;
4mod handoff;
5mod helpers;
6mod producer;
7mod resume;
8mod run_single;
9mod session_blocking;
10mod support;
11mod trace_lifecycle;
12
13use std::path::PathBuf;
14use std::sync::{Arc, Mutex};
15
16use serde_json::Value;
17use tokio::sync::broadcast;
18
19use crate::agent::Agent;
20use crate::approval::ApprovalBroker;
21use crate::checkpoint::{CheckpointConfig, ResumePolicy};
22use crate::config::apply_resolved_model_limits;
23use crate::context::RunContext;
24use crate::context_providers::{
25    assemble_context_fragments, collect_context_fragments, ContextBundle, ContextRequest,
26};
27use crate::events::{AgentErrorPayload, RunEvent};
28use crate::guardrails::GuardrailOutcome;
29use crate::llm::LlmClient;
30use crate::model::{ModelError, ModelProvider, ModelRef};
31use crate::result::{RunResult, RunResumeContext};
32use crate::run_config::{validate_max_cycles, RunConfig, INITIAL_BUDGET_USAGE_METADATA_KEY};
33use crate::runtime::checkpoint_resume::{
34    CheckpointController, CheckpointControllerRequest, CheckpointEventSink,
35    CheckpointResumeController,
36};
37use crate::runtime::run_definition::{
38    build_frozen_task, build_run_definition, frozen_definition_messages, RunDefinitionRequest,
39};
40use crate::runtime::state::Checkpoint;
41use crate::runtime::tool_planner::project_tool_policy;
42use crate::runtime::{
43    AgentRuntime, CheckpointRuntimeControl, ExecutionContext, RuntimeRunControls,
44};
45use crate::sessions::{checkpoint_session_commit_id, session_commit_payload_digest, SessionItem};
46use crate::tools::{ToolEnablementContext, ToolRegistry};
47use crate::types::{AgentResult, AgentStatus, AgentTask, MessageRole};
48use crate::workspace::LocalWorkspaceBackend;
49
50pub use builder::RunnerBuilder;
51use checkpoint_runtime::{
52    prepare_checkpoint_resume, prepare_checkpoint_runtime, prepare_checkpoint_terminal,
53    replay_checkpoint_terminal, CheckpointRuntimeRequest, CheckpointRuntimeState,
54};
55pub(crate) use event_stream::map_stream_event;
56pub use event_stream::RunEventStream;
57#[doc(hidden)]
58pub use event_stream::{map_runtime_event, RuntimeEventContext};
59use helpers::{
60    effective_model_ref, effective_trace_id, effective_workflow_name, status_string, terminal_event,
61};
62pub(crate) use producer::CheckpointStartOutcome;
63use session_blocking::block_on_session;
64use support::{
65    apply_cancellation_precedence, apply_input_guardrails, apply_optional_output_validation,
66    apply_output_guardrails, capture_event, effective_event_store, effective_session_id,
67    extract_handoff, initial_budget_usage, insert_context_metadata, merged_tool_policy,
68    ApprovalHook, SingleRunOutcome,
69};
70use trace_lifecycle::RunTrace;
71
72#[derive(Clone, Debug, PartialEq)]
73pub struct NormalizedInput {
74    pub text: String,
75}
76
77impl From<&str> for NormalizedInput {
78    fn from(value: &str) -> Self {
79        Self {
80            text: value.to_string(),
81        }
82    }
83}
84
85impl From<String> for NormalizedInput {
86    fn from(value: String) -> Self {
87        Self { text: value }
88    }
89}
90
91struct InstructionBuildRequest<'a> {
92    agent: &'a Agent,
93    run_context: &'a RunContext,
94    input_text: &'a str,
95    config: &'a RunConfig,
96    model: &'a str,
97    trace_id: &'a str,
98    session: Option<Arc<dyn crate::sessions::Session>>,
99    workspace: &'a std::path::Path,
100}
101
102pub(super) struct CheckpointAdmission {
103    pub checkpoint: Checkpoint,
104    pub terminal_replayed: bool,
105}
106
107pub(super) type CheckpointAdmissionSender = tokio::sync::oneshot::Sender<CheckpointAdmission>;
108
109#[derive(Clone)]
110pub struct Runner {
111    model_provider: Arc<dyn ModelProvider>,
112    workspace: PathBuf,
113    tool_registry: ToolRegistry,
114    default_run_config: RunConfig,
115}
116
117impl Runner {
118    pub fn builder() -> RunnerBuilder {
119        RunnerBuilder::default()
120    }
121
122    pub async fn run(
123        &self,
124        agent: &Agent,
125        input: impl Into<NormalizedInput>,
126    ) -> Result<RunResult, String> {
127        self.run_with_config(agent, input, RunConfig::default())
128            .await
129    }
130
131    pub async fn run_with_config(
132        &self,
133        agent: &Agent,
134        input: impl Into<NormalizedInput>,
135        config: RunConfig,
136    ) -> Result<RunResult, String> {
137        let runner = self.clone();
138        let agent = agent.clone();
139        let input = input.into();
140        tokio::task::spawn_blocking(move || runner.run_blocking(&agent, input, config, None))
141            .await
142            .map_err(|error| format!("runner task failed: {error}"))?
143    }
144
145    async fn run_with_config_and_run_id(
146        &self,
147        agent: &Agent,
148        input: NormalizedInput,
149        config: RunConfig,
150        run_id: String,
151    ) -> Result<RunResult, String> {
152        let runner = self.clone();
153        let agent = agent.clone();
154        tokio::task::spawn_blocking(move || {
155            runner.run_agent_chain_with_initial(
156                &agent,
157                input,
158                config,
159                Some(Arc::new(Mutex::new(Vec::new()))),
160                None,
161                None,
162                None,
163                Some(run_id),
164            )
165        })
166        .await
167        .map_err(|error| format!("runner task failed: {error}"))?
168    }
169
170    pub fn run_blocking(
171        &self,
172        agent: &Agent,
173        input: NormalizedInput,
174        config: RunConfig,
175        event_collector: Option<Arc<std::sync::Mutex<Vec<RunEvent>>>>,
176    ) -> Result<RunResult, String> {
177        self.run_blocking_with_event_sender(agent, input, config, event_collector, None, None)
178    }
179
180    fn run_blocking_with_event_sender(
181        &self,
182        agent: &Agent,
183        input: NormalizedInput,
184        config: RunConfig,
185        event_collector: Option<Arc<std::sync::Mutex<Vec<RunEvent>>>>,
186        event_sender: Option<broadcast::Sender<RunEvent>>,
187        checkpoint_admission_sender: Option<CheckpointAdmissionSender>,
188    ) -> Result<RunResult, String> {
189        let event_collector =
190            Some(event_collector.unwrap_or_else(|| Arc::new(std::sync::Mutex::new(Vec::new()))));
191        self.run_agent_chain(
192            agent,
193            input,
194            config,
195            event_collector,
196            event_sender,
197            checkpoint_admission_sender,
198        )
199    }
200
201    fn build_instructions_with_context(
202        &self,
203        request: InstructionBuildRequest<'_>,
204    ) -> Result<(String, Option<ContextBundle>), String> {
205        let InstructionBuildRequest {
206            agent,
207            run_context,
208            input_text,
209            config,
210            model,
211            trace_id,
212            session,
213            workspace,
214        } = request;
215        let providers = self
216            .default_run_config
217            .context_providers
218            .iter()
219            .chain(config.context_providers.iter())
220            .cloned()
221            .collect::<Vec<_>>();
222        let mut request = ContextRequest::new(agent.name(), input_text)
223            .model(model)
224            .trace_id(trace_id)
225            .workspace(workspace);
226        if let Some(session) = session {
227            request = request.session(session);
228        }
229        if let Some(context) = config
230            .app_state
231            .clone()
232            .or_else(|| self.default_run_config.app_state.clone())
233        {
234            request = request.context(context);
235        }
236        request.metadata = agent.metadata().clone();
237        request
238            .metadata
239            .extend(self.default_run_config.metadata.clone());
240        request.metadata.extend(config.metadata.clone());
241        if let Some(max_chars) = config
242            .max_context_chars
243            .or(self.default_run_config.max_context_chars)
244        {
245            request = request.max_prompt_chars(max_chars);
246        }
247        let mut fragments = vec![crate::context_providers::ContextFragment::new(
248            "agent_instructions",
249            agent.resolve_instructions(run_context),
250        )
251        .stable(true)
252        .priority(0)
253        .source("agent.instructions")];
254        if !agent.sub_agents().is_empty() {
255            let available_sub_agents = agent
256                .sub_agents()
257                .iter()
258                .map(|(id, config)| (id.clone(), config.description.clone()))
259                .collect();
260            fragments.push(
261                crate::context_providers::ContextFragment::new(
262                    "configured_sub_agents",
263                    crate::prompt::templates::render_sub_agents("en-US", &available_sub_agents),
264                )
265                .stable(true)
266                .priority(10)
267                .source("agent.sub_agents"),
268            );
269        }
270        fragments.extend(
271            collect_context_fragments(&request, &providers)
272                .map_err(|error| format!("context provider failed: {error}"))?,
273        );
274        let bundle = assemble_context_fragments(&request, fragments)
275            .map_err(|error| format!("context assembly failed: {error}"))?;
276        Ok((bundle.prompt.clone(), Some(bundle)))
277    }
278}
279
280#[derive(Clone)]
281struct ArcLlmClient(Arc<dyn LlmClient>);
282
283impl LlmClient for ArcLlmClient {
284    fn complete(
285        &self,
286        request: crate::llm::LlmRequest,
287    ) -> Result<crate::types::LLMResponse, crate::llm::LlmError> {
288        self.0.complete(request)
289    }
290
291    fn complete_with_stream(
292        &self,
293        request: crate::llm::LlmRequest,
294        stream_callback: Option<crate::llm::LlmStreamCallback>,
295    ) -> Result<crate::types::LLMResponse, crate::llm::LlmError> {
296        self.0.complete_with_stream(request, stream_callback)
297    }
298}
299
300fn preload_checkpoint(config: Option<&CheckpointConfig>) -> Result<Option<Checkpoint>, String> {
301    let Some(config) = config else {
302        return Ok(None);
303    };
304    config.validate().map_err(|error| error.to_string())?;
305    if config.resume_policy == ResumePolicy::New {
306        return Ok(None);
307    }
308    let store = config.store.as_ref().ok_or_else(|| {
309        "checkpoint_store_unavailable: process-local Runner resume requires CheckpointConfig.store"
310            .to_string()
311    })?;
312    let key = config.key.as_deref().ok_or_else(|| {
313        "checkpoint_key_required: resume_if_present and require_existing need an explicit key"
314            .to_string()
315    })?;
316    store
317        .load_checkpoint(key)
318        .map_err(|error| error.to_string())
319}
320
321fn format_model_error(error: ModelError) -> String {
322    error.to_string()
323}