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_v2::{
38 build_frozen_task, build_run_definition, frozen_definition_messages, RunDefinitionRequest,
39};
40use crate::runtime::state_v2::CheckpointV2;
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};
55use 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, is_runtime_terminal_log,
61 status_string, terminal_event,
62};
63pub(crate) use producer::CheckpointStartOutcome;
64use session_blocking::block_on_session;
65use support::{
66 apply_cancellation_precedence, apply_input_guardrails, apply_output_guardrails, capture_event,
67 effective_event_store, effective_session_id, extract_handoff, initial_budget_usage,
68 insert_context_metadata, merged_tool_policy, 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: CheckpointV2,
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 pub fn run_blocking(
146 &self,
147 agent: &Agent,
148 input: NormalizedInput,
149 config: RunConfig,
150 event_collector: Option<Arc<std::sync::Mutex<Vec<RunEvent>>>>,
151 ) -> Result<RunResult, String> {
152 self.run_blocking_with_event_sender(agent, input, config, event_collector, None, None)
153 }
154
155 fn run_blocking_with_event_sender(
156 &self,
157 agent: &Agent,
158 input: NormalizedInput,
159 config: RunConfig,
160 event_collector: Option<Arc<std::sync::Mutex<Vec<RunEvent>>>>,
161 event_sender: Option<broadcast::Sender<RunEvent>>,
162 checkpoint_admission_sender: Option<CheckpointAdmissionSender>,
163 ) -> Result<RunResult, String> {
164 let event_collector =
165 Some(event_collector.unwrap_or_else(|| Arc::new(std::sync::Mutex::new(Vec::new()))));
166 self.run_agent_chain(
167 agent,
168 input,
169 config,
170 event_collector,
171 event_sender,
172 checkpoint_admission_sender,
173 )
174 }
175
176 fn build_instructions_with_context(
177 &self,
178 request: InstructionBuildRequest<'_>,
179 ) -> Result<(String, Option<ContextBundle>), String> {
180 let InstructionBuildRequest {
181 agent,
182 run_context,
183 input_text,
184 config,
185 model,
186 trace_id,
187 session,
188 workspace,
189 } = request;
190 let providers = self
191 .default_run_config
192 .context_providers
193 .iter()
194 .chain(config.context_providers.iter())
195 .cloned()
196 .collect::<Vec<_>>();
197 let mut request = ContextRequest::new(agent.name(), input_text)
198 .model(model)
199 .trace_id(trace_id)
200 .workspace(workspace);
201 if let Some(session) = session {
202 request = request.session(session);
203 }
204 if let Some(context) = config
205 .app_state
206 .clone()
207 .or_else(|| self.default_run_config.app_state.clone())
208 {
209 request = request.context(context);
210 }
211 request.metadata = agent.metadata().clone();
212 request
213 .metadata
214 .extend(self.default_run_config.metadata.clone());
215 request.metadata.extend(config.metadata.clone());
216 if let Some(max_chars) = config
217 .max_context_chars
218 .or(self.default_run_config.max_context_chars)
219 {
220 request = request.max_prompt_chars(max_chars);
221 }
222 let mut fragments = vec![crate::context_providers::ContextFragment::new(
223 "agent_instructions",
224 agent.resolve_instructions(run_context),
225 )
226 .stable(true)
227 .priority(0)
228 .source("agent.instructions")];
229 if !agent.sub_agents().is_empty() {
230 let available_sub_agents = agent
231 .sub_agents()
232 .iter()
233 .map(|(id, config)| (id.clone(), config.description.clone()))
234 .collect();
235 fragments.push(
236 crate::context_providers::ContextFragment::new(
237 "configured_sub_agents",
238 crate::prompt::templates::render_sub_agents("en-US", &available_sub_agents),
239 )
240 .stable(true)
241 .priority(10)
242 .source("agent.sub_agents"),
243 );
244 }
245 fragments.extend(
246 collect_context_fragments(&request, &providers)
247 .map_err(|error| format!("context provider failed: {error}"))?,
248 );
249 let bundle = assemble_context_fragments(&request, fragments)
250 .map_err(|error| format!("context assembly failed: {error}"))?;
251 Ok((bundle.prompt.clone(), Some(bundle)))
252 }
253}
254
255#[derive(Clone)]
256struct ArcLlmClient(Arc<dyn LlmClient>);
257
258impl LlmClient for ArcLlmClient {
259 fn complete(
260 &self,
261 request: crate::llm::LlmRequest,
262 ) -> Result<crate::types::LLMResponse, crate::llm::LlmError> {
263 self.0.complete(request)
264 }
265
266 fn complete_with_stream(
267 &self,
268 request: crate::llm::LlmRequest,
269 stream_callback: Option<crate::llm::LlmStreamCallback>,
270 ) -> Result<crate::types::LLMResponse, crate::llm::LlmError> {
271 self.0.complete_with_stream(request, stream_callback)
272 }
273}
274
275fn preload_checkpoint(config: Option<&CheckpointConfig>) -> Result<Option<CheckpointV2>, String> {
276 let Some(config) = config else {
277 return Ok(None);
278 };
279 config.validate().map_err(|error| error.to_string())?;
280 if config.resume_policy == ResumePolicy::New {
281 return Ok(None);
282 }
283 let store = config.store.as_ref().ok_or_else(|| {
284 "checkpoint_store_unavailable: process-local Runner resume requires CheckpointConfig.store"
285 .to_string()
286 })?;
287 let key = config.key.as_deref().ok_or_else(|| {
288 "checkpoint_key_required: resume_if_present and require_existing need an explicit key"
289 .to_string()
290 })?;
291 store
292 .load_checkpoint_v2(key)
293 .map_err(|error| error.to_string())
294}
295
296fn format_model_error(error: ModelError) -> String {
297 error.to_string()
298}