Skip to main content

vv_agent/
runner.rs

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