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}