1use super::checkpointing::{
4 AgentProgress, save_checkpoint, validate_exact_usage, validate_usage_floor,
5};
6use super::completion::TerminalCompletionContext;
7use super::observability::{consume_budget, emit_usage, record_domain, terminal_event};
8use super::{
9 Agent, AgentCheckpoint, AgentCheckpointPhase, AgentCheckpointState, AgentError,
10 AgentEventStream, AgentFuture, AgentObserver, AgentOutcome, AgentStreamEvent, Arc,
11 BufferedObserver, CheckpointCursor, ContentPart, DurableConversationCheckpoint, Either,
12 EventId, Instant, InvocationId, LifecycleEvent, Message, ModelCallContext, ModelError,
13 ModelErrorKind, ModelRequest, ModelResponse, ModelStreamAccumulator, NoopObserver,
14 ResumePolicy, Role, RunContext, RunEventKind, StreamExt, TOOL_RESULT_EXECUTION_ID_METADATA,
15 ToolCall, ToolChoice, Usage, emit_agent_event, select,
16};
17use crate::conversation::{
18 AgentConversationError, AgentConversationOutcome, AutomaticConversationSummary,
19 ConversationAppend, ConversationContextPolicy, ConversationId, ConversationStore,
20 ConversationSummaryCommit, ConversationSummaryRequest, DurableConversationCommit,
21 DurableConversationRequest, DurableConversationStore, MemoryNamespace, SemanticMemoryQuery,
22 is_transient_context, semantic_memory_message, summary_message,
23};
24use runifold_core::{CheckpointId, CheckpointStore};
25use runifold_retrieval::RetrievalContext;
26
27impl Agent {
28 pub fn prompt<'a>(
35 &'a self,
36 input: impl Into<String> + Send + 'a,
37 ) -> AgentFuture<'a, Result<AgentOutcome, AgentError>> {
38 let input = input.into();
39 Box::pin(async move {
40 let run = self.default_run_context();
41 self.run(input, &run).await
42 })
43 }
44
45 pub fn prompt_text<'a>(
51 &'a self,
52 input: impl Into<String> + Send + 'a,
53 ) -> AgentFuture<'a, Result<String, AgentError>> {
54 let input = input.into();
55 Box::pin(async move { self.prompt(input).await.map(AgentOutcome::into_text) })
56 }
57
58 pub fn run<'a>(
60 &'a self,
61 input: impl Into<String> + Send + 'a,
62 run: &'a RunContext,
63 ) -> AgentFuture<'a, Result<AgentOutcome, AgentError>> {
64 let input = input.into();
65 let state = self.initial_state(input, InvocationId::new().to_string());
66 Box::pin(async move {
67 self.execute_state(state, run, None, Arc::new(NoopObserver), true, true)
68 .await
69 })
70 }
71
72 pub fn run_conversation<'a>(
78 &'a self,
79 input: impl Into<String> + Send + 'a,
80 run: &'a RunContext,
81 store: &'a dyn ConversationStore,
82 conversation_id: ConversationId,
83 namespace: MemoryNamespace,
84 policy: ConversationContextPolicy,
85 ) -> AgentFuture<'a, Result<AgentConversationOutcome, AgentConversationError>> {
86 let input = input.into();
87 Box::pin(async move {
88 store.create(conversation_id, namespace.clone()).await?;
89 let view = store
90 .load_view(
91 conversation_id,
92 namespace.clone(),
93 policy.window,
94 policy.summary_batch,
95 )
96 .await?;
97 if view.requires_summary() {
98 return Err(AgentConversationError::SummaryRequired {
99 conversation_id,
100 buffered_entries: u64::try_from(view.summary_buffer.len())
101 .unwrap_or(u64::MAX)
102 .saturating_add(view.summary_backlog),
103 });
104 }
105 let mut transcript = self.instructions.clone();
106 if let Some(summary) = &view.summary {
107 transcript.push(summary_message(summary));
108 }
109 if let Some(limit) = policy.semantic_memory_limit {
110 let query =
111 SemanticMemoryQuery::new(namespace.clone(), input.clone(), limit.get())?;
112 let search = store
113 .search_memory_scoped(query, RetrievalContext::for_run(run))
114 .await?;
115 if search.usage != Usage::default() {
116 consume_budget(run, search.usage, None).map_err(AgentConversationError::Run)?;
117 }
118 if let Some(message) = semantic_memory_message(&search.memories) {
119 transcript.push(message);
120 }
121 }
122 transcript.extend(view.window.iter().map(|entry| entry.message.clone()));
123 let persisted_prefix_len = transcript.len();
124 transcript.push(Message::user(input));
125 let state =
126 self.initial_state_from_transcript(transcript, InvocationId::new().to_string());
127 let outcome = self
128 .execute_state(state, run, None, Arc::new(NoopObserver), true, true)
129 .await
130 .map_err(AgentConversationError::Run)?;
131 let messages = outcome
132 .transcript
133 .iter()
134 .skip(persisted_prefix_len)
135 .filter(|message| !is_transient_context(message))
136 .cloned()
137 .collect();
138 let append = ConversationAppend {
139 conversation_id,
140 expected_version: view.version,
141 messages,
142 };
143 match store.append(namespace, append).await {
144 Ok(conversation_version) => Ok(AgentConversationOutcome {
145 outcome,
146 conversation_version,
147 }),
148 Err(source) => Err(AgentConversationError::Commit {
149 source,
150 outcome: Box::new(outcome),
151 }),
152 }
153 })
154 }
155
156 pub fn run_conversation_with_summary<'a>(
162 &'a self,
163 input: impl Into<String> + Send + 'a,
164 run: &'a RunContext,
165 store: &'a dyn ConversationStore,
166 conversation_id: ConversationId,
167 namespace: MemoryNamespace,
168 automatic_summary: AutomaticConversationSummary<'a>,
169 ) -> AgentFuture<'a, Result<AgentConversationOutcome, AgentConversationError>> {
170 let input = input.into();
171 Box::pin(async move {
172 let policy = automatic_summary.context;
173 store.create(conversation_id, namespace.clone()).await?;
174 for pass in 0..automatic_summary.max_passes.get() {
175 let view = store
176 .load_view(
177 conversation_id,
178 namespace.clone(),
179 policy.window,
180 policy.summary_batch,
181 )
182 .await?;
183 let Some(through_sequence) = view.summary_buffer.last().map(|entry| entry.sequence)
184 else {
185 break;
186 };
187 let summary_backlog = view.summary_backlog;
188 let summary = automatic_summary
189 .summarizer
190 .summarize(
191 ConversationSummaryRequest {
192 transcript_version: view.version,
193 previous_summary: view.summary,
194 entries: view.summary_buffer,
195 },
196 run,
197 )
198 .await?;
199 store
200 .commit_summary(
201 namespace.clone(),
202 ConversationSummaryCommit {
203 conversation_id,
204 expected_version: view.version,
205 through_sequence,
206 content: summary,
207 },
208 )
209 .await?;
210 if summary_backlog == 0 {
211 break;
212 }
213 if pass + 1 == automatic_summary.max_passes.get() {
214 return Err(AgentConversationError::SummaryPassLimitExceeded {
215 conversation_id,
216 remaining_entries: summary_backlog,
217 });
218 }
219 }
220 self.run_conversation(input, run, store, conversation_id, namespace, policy)
221 .await
222 })
223 }
224
225 pub fn stream<'a>(
227 &'a self,
228 input: impl Into<String> + Send + 'a,
229 run: &'a RunContext,
230 ) -> AgentEventStream<'a> {
231 let state = self.initial_state(input.into(), InvocationId::new().to_string());
232 let observer = BufferedObserver::default();
233 let events = observer.events();
234 let execution =
235 Box::pin(self.execute_state(state, run, None, Arc::new(observer), true, true));
236 AgentEventStream::new(execution, events)
237 }
238
239 pub fn run_durable_conversation<'a>(
245 &'a self,
246 input: impl Into<String> + Send + 'a,
247 run: &'a RunContext,
248 store: Arc<dyn DurableConversationStore>,
249 request: DurableConversationRequest,
250 ) -> AgentFuture<'a, Result<AgentConversationOutcome, AgentConversationError>> {
251 let input = input.into();
252 Box::pin(async move {
253 let DurableConversationRequest {
254 checkpoint_id,
255 conversation_id,
256 namespace,
257 policy,
258 } = request;
259 store.create(conversation_id, namespace.clone()).await?;
260 let view = store
261 .load_view(
262 conversation_id,
263 namespace.clone(),
264 policy.window,
265 policy.summary_batch,
266 )
267 .await?;
268 if view.requires_summary() {
269 return Err(AgentConversationError::SummaryRequired {
270 conversation_id,
271 buffered_entries: u64::try_from(view.summary_buffer.len())
272 .unwrap_or(u64::MAX)
273 .saturating_add(view.summary_backlog),
274 });
275 }
276 let mut transcript = self.instructions.clone();
277 if let Some(summary) = &view.summary {
278 transcript.push(summary_message(summary));
279 }
280 if let Some(limit) = policy.semantic_memory_limit {
281 let query =
282 SemanticMemoryQuery::new(namespace.clone(), input.clone(), limit.get())?;
283 let search = store
284 .search_memory_scoped(query, RetrievalContext::for_run(run))
285 .await?;
286 if search.usage != Usage::default() {
287 consume_budget(run, search.usage, None).map_err(AgentConversationError::Run)?;
288 }
289 if let Some(message) = semantic_memory_message(&search.memories) {
290 transcript.push(message);
291 }
292 }
293 transcript.extend(view.window.iter().map(|entry| entry.message.clone()));
294 let persisted_prefix_len = u64::try_from(transcript.len()).map_err(|_| {
295 AgentConversationError::Run(checkpoint_payload_error(
296 "conversation context length exceeds durable checkpoint range",
297 ))
298 })?;
299 transcript.push(Message::user(input));
300 let durable = DurableConversationCheckpoint {
301 conversation_id,
302 namespace,
303 expected_version: view.version,
304 persisted_prefix_len,
305 };
306 let mut state =
307 self.initial_state_from_transcript(transcript, checkpoint_id.to_string());
308 state.durable_conversation = Some(durable.clone());
309 state.usage = run.budget().usage();
310 let checkpoint_store: Arc<dyn CheckpointStore> = store.clone();
311 let checkpoint = AgentCheckpoint::existing(checkpoint_id, checkpoint_store);
312 let mut cursor = CheckpointCursor::create(&checkpoint, run, &state)
313 .map_err(AgentConversationError::Run)?;
314 let outcome = self
315 .execute_state(
316 state,
317 run,
318 Some(&mut cursor),
319 Arc::new(NoopObserver),
320 true,
321 false,
322 )
323 .await
324 .map_err(AgentConversationError::Run)?;
325 self.commit_durable_outcome(store.as_ref(), run, &cursor, durable, outcome)
326 .await
327 })
328 }
329
330 pub fn resume_durable_conversation<'a>(
332 &'a self,
333 store: Arc<dyn DurableConversationStore>,
334 checkpoint_id: CheckpointId,
335 run: &'a RunContext,
336 policy: ResumePolicy,
337 ) -> AgentFuture<'a, Result<AgentConversationOutcome, AgentConversationError>> {
338 Box::pin(async move {
339 let checkpoint_store: Arc<dyn CheckpointStore> = store.clone();
340 let checkpoint = AgentCheckpoint::existing(checkpoint_id, checkpoint_store);
341 let (envelope, mut state) = checkpoint
342 .load()
343 .map_err(AgentError::from)
344 .map_err(AgentConversationError::Run)?;
345 self.validate_checkpoint_identity(&state)
346 .map_err(AgentConversationError::Run)?;
347 let durable = state.durable_conversation.clone().ok_or_else(|| {
348 AgentConversationError::Run(checkpoint_payload_error(
349 "checkpoint is not a durable conversation turn",
350 ))
351 })?;
352 if let Some(outcome) = state.outcome() {
353 let conversation_version = durable
354 .expected_version
355 .get()
356 .checked_add(1)
357 .map(crate::ConversationVersion::new)
358 .ok_or_else(|| {
359 AgentConversationError::Run(checkpoint_payload_error(
360 "durable conversation version overflow",
361 ))
362 })?;
363 return Ok(AgentConversationOutcome {
364 outcome,
365 conversation_version,
366 });
367 }
368 if let Some(error) = state.terminal_failure() {
369 validate_exact_usage(state.usage, run.budget().usage())
370 .map_err(AgentConversationError::Run)?;
371 return Err(AgentConversationError::Run(error));
372 }
373 Self::prepare_resume_state(&mut state, run, policy)
374 .map_err(AgentConversationError::Run)?;
375 let mut cursor = CheckpointCursor::loaded(&checkpoint, envelope);
376 let outcome = self
377 .execute_state(
378 state,
379 run,
380 Some(&mut cursor),
381 Arc::new(NoopObserver),
382 false,
383 false,
384 )
385 .await
386 .map_err(AgentConversationError::Run)?;
387 self.commit_durable_outcome(store.as_ref(), run, &cursor, durable, outcome)
388 .await
389 })
390 }
391
392 pub fn run_checkpointed<'a>(
394 &'a self,
395 input: impl Into<String> + Send + 'a,
396 run: &'a RunContext,
397 checkpoint: &'a AgentCheckpoint,
398 ) -> AgentFuture<'a, Result<AgentOutcome, AgentError>> {
399 let input = input.into();
400 Box::pin(async move {
401 let mut state = self.initial_state(input, checkpoint.id().to_string());
402 state.usage = run.budget().usage();
403 let mut cursor = CheckpointCursor::create(checkpoint, run, &state)?;
404 self.execute_state(
405 state,
406 run,
407 Some(&mut cursor),
408 Arc::new(NoopObserver),
409 true,
410 true,
411 )
412 .await
413 })
414 }
415
416 pub fn resume<'a>(
418 &'a self,
419 checkpoint: &'a AgentCheckpoint,
420 run: &'a RunContext,
421 policy: ResumePolicy,
422 ) -> AgentFuture<'a, Result<AgentOutcome, AgentError>> {
423 Box::pin(async move {
424 let (envelope, mut state) = checkpoint.load()?;
425 self.validate_checkpoint_identity(&state)?;
426 if let Some(outcome) = state.outcome() {
427 validate_exact_usage(state.usage, run.budget().usage())?;
428 return Ok(outcome);
429 }
430 if let Some(error) = state.terminal_failure() {
431 validate_exact_usage(state.usage, run.budget().usage())?;
432 return Err(error);
433 }
434 Self::prepare_resume_state(&mut state, run, policy)?;
435 let mut cursor = CheckpointCursor::loaded(checkpoint, envelope);
436 self.execute_state(
437 state,
438 run,
439 Some(&mut cursor),
440 Arc::new(NoopObserver),
441 false,
442 true,
443 )
444 .await
445 })
446 }
447
448 fn initial_state(&self, input: String, execution_id: String) -> AgentCheckpointState {
449 let mut transcript = self.instructions.clone();
450 transcript.push(Message::user(input));
451 self.initial_state_from_transcript(transcript, execution_id)
452 }
453
454 fn initial_state_from_transcript(
455 &self,
456 transcript: Vec<Message>,
457 execution_id: String,
458 ) -> AgentCheckpointState {
459 AgentCheckpointState {
460 execution_id,
461 agent: self.name.clone(),
462 model: self.model_ref.clone(),
463 transcript,
464 turns: 0,
465 tool_calls: 0,
466 delegations: 0,
467 usage: Usage::default(),
468 turn_reviewer: self
469 .turn_review
470 .as_ref()
471 .map(|review| review.descriptor.clone()),
472 turn_review_policy: self.turn_review.as_ref().map(|review| review.policy),
473 turn_reviewer_capabilities: self.turn_review.as_ref().map_or_else(Vec::new, |review| {
474 review
475 .capabilities
476 .iter()
477 .map(|capability| capability.id)
478 .collect()
479 }),
480 terminal_reviewer: self
481 .terminal_review
482 .as_ref()
483 .map(|review| review.descriptor.clone()),
484 terminal_review_policy: self.terminal_review.as_ref().map(|review| review.policy),
485 terminal_reviewer_capabilities: self.terminal_review.as_ref().map_or_else(
486 Vec::new,
487 |review| {
488 review
489 .capabilities
490 .iter()
491 .map(|capability| capability.id)
492 .collect()
493 },
494 ),
495 phase: AgentCheckpointPhase::ReadyForTurn,
496 durable_conversation: None,
497 }
498 }
499
500 async fn execute_state(
501 &self,
502 state: AgentCheckpointState,
503 run: &RunContext,
504 mut checkpoint: Option<&mut CheckpointCursor>,
505 observer: Arc<dyn AgentObserver>,
506 retrieve_context: bool,
507 persist_terminal_checkpoint: bool,
508 ) -> Result<AgentOutcome, AgentError> {
509 let started = run
510 .record(
511 RunEventKind::Lifecycle(LifecycleEvent::Started),
512 run.caused_by(),
513 )?
514 .map(|event| event.meta.event_id);
515 emit_agent_event(
516 observer.as_ref(),
517 AgentStreamEvent::Started {
518 agent: self.name.clone(),
519 },
520 )
521 .await;
522 let result = async {
523 let has_context = !self.context.is_empty() || !self.dynamic_context.is_empty();
524 let state = if retrieve_context && has_context {
525 let mut prepared = self
526 .prepare_context(state, run, started, observer.as_ref())
527 .await?;
528 prepared.usage = run.budget().usage();
529 save_checkpoint(&mut checkpoint, &prepared)?;
530 prepared
531 } else {
532 state
533 };
534 self.run_loop(
535 state,
536 run,
537 started,
538 checkpoint,
539 observer.as_ref(),
540 persist_terminal_checkpoint,
541 )
542 .await
543 }
544 .await;
545 let terminal = terminal_event(&self.name, &result);
546 run.record(terminal, started)?;
547 if let Ok(outcome) = &result {
548 emit_agent_event(
549 observer.as_ref(),
550 AgentStreamEvent::Completed {
551 outcome: outcome.clone(),
552 },
553 )
554 .await;
555 }
556 result
557 }
558
559 async fn run_loop(
560 &self,
561 state: AgentCheckpointState,
562 run: &RunContext,
563 caused_by: Option<EventId>,
564 mut checkpoint: Option<&mut CheckpointCursor>,
565 observer: &dyn AgentObserver,
566 persist_terminal_checkpoint: bool,
567 ) -> Result<AgentOutcome, AgentError> {
568 self.validate_config()?;
569 let completion_context = TerminalCompletionContext {
570 caused_by,
571 observer,
572 persist_terminal_checkpoint,
573 };
574 let (mut progress, resumed, mut approved_response) = self
575 .prepare_run_loop_progress(state, run, &mut checkpoint, &completion_context)
576 .await?;
577 if let Some(outcome) = resumed {
578 return Ok(outcome);
579 }
580
581 loop {
582 Self::check_lifecycle(run)?;
583 let Some((response, requires_tool)) = self
584 .next_reviewed_response(
585 approved_response.take(),
586 &mut progress,
587 run,
588 &mut checkpoint,
589 &completion_context,
590 )
591 .await?
592 else {
593 continue;
594 };
595 let calls = tool_calls_from(&response.content);
596 if calls.is_empty() {
597 if self.continue_provider_turn(
598 response.clone(),
599 &mut progress,
600 run,
601 &mut checkpoint,
602 )? {
603 continue;
604 }
605 if matches!(
606 response.finish_reason,
607 runifold_model::FinishReason::ToolCalls
608 ) && !response.content.is_empty()
609 {
610 return Err(AgentError::Protocol(
611 "model stopped for tool calls without emitting a tool call".into(),
612 ));
613 }
614 if requires_tool {
615 return Err(AgentError::ToolRequirementUnsatisfied {
616 required: self.min_successful_tool_calls,
617 successful: self.successful_local_tool_calls(&progress)?,
618 });
619 }
620 if let Some(outcome) = self
621 .complete_terminal_candidate(
622 response,
623 run,
624 &mut progress,
625 &mut checkpoint,
626 TerminalCompletionContext {
627 caused_by,
628 observer,
629 persist_terminal_checkpoint,
630 },
631 )
632 .await?
633 {
634 return Ok(outcome);
635 }
636 continue;
637 }
638
639 let assistant = Message::new(Role::Assistant, response.content.clone())
640 .map_err(|error| AgentError::Protocol(error.to_string()))?;
641 progress.transcript.push(assistant);
642
643 self.execute_calls(calls, run, caused_by, &mut progress, observer)
644 .await?;
645 save_checkpoint(
646 &mut checkpoint,
647 &self.checkpoint_state(&progress, run, AgentCheckpointPhase::ReadyForTurn),
648 )?;
649 }
650 }
651
652 async fn next_reviewed_response(
653 &self,
654 approved_response: Option<ModelResponse>,
655 progress: &mut AgentProgress,
656 run: &RunContext,
657 checkpoint: &mut Option<&mut CheckpointCursor>,
658 context: &TerminalCompletionContext<'_>,
659 ) -> Result<Option<(ModelResponse, bool)>, AgentError> {
660 let (response, requires_tool, already_reviewed) = if let Some(response) = approved_response
661 {
662 let requires_tool =
663 self.successful_local_tool_calls(progress)? < self.min_successful_tool_calls;
664 (response, requires_tool, true)
665 } else {
666 let tool_choice = self
667 .begin_turn(
668 progress,
669 run,
670 checkpoint,
671 context.caused_by,
672 context.observer,
673 )
674 .await?;
675 let requires_tool = matches!(tool_choice, ToolChoice::Required);
676 let response = self
677 .invoke_model(
678 &progress.transcript,
679 run,
680 progress.turns,
681 tool_choice,
682 context.caused_by,
683 context.observer,
684 )
685 .await?;
686 (response, requires_tool, false)
687 };
688
689 let calls = tool_calls_from(&response.content);
690 validate_tool_call_completion(&calls, &response.finish_reason)?;
691 if already_reviewed
692 || !self
693 .turn_review
694 .as_ref()
695 .is_some_and(|review| review.policy.scope().includes(&response))
696 {
697 return Ok(Some((response, requires_tool)));
698 }
699
700 save_checkpoint(
701 checkpoint,
702 &self.checkpoint_state(
703 progress,
704 run,
705 AgentCheckpointPhase::TurnReviewReady {
706 response: Box::new(response.clone()),
707 turn: progress.turns,
708 },
709 ),
710 )?;
711 let response = self
712 .review_turn_candidate(response, run, progress, checkpoint, context)
713 .await?;
714 Ok(response.map(|response| (response, requires_tool)))
715 }
716
717 async fn begin_turn(
718 &self,
719 progress: &mut AgentProgress,
720 run: &RunContext,
721 checkpoint: &mut Option<&mut CheckpointCursor>,
722 caused_by: Option<EventId>,
723 observer: &dyn AgentObserver,
724 ) -> Result<ToolChoice, AgentError> {
725 let tool_choice = self.next_tool_choice(progress, run)?;
726 save_checkpoint(
727 checkpoint,
728 &self.checkpoint_state(
729 progress,
730 run,
731 AgentCheckpointPhase::TurnInFlight {
732 turn: progress.turns + 1,
733 },
734 ),
735 )?;
736 consume_budget(
737 run,
738 Usage {
739 turns: 1,
740 ..Usage::default()
741 },
742 caused_by,
743 )?;
744 progress.turns += 1;
745 emit_agent_event(
746 observer,
747 AgentStreamEvent::TurnStarted {
748 turn: progress.turns,
749 },
750 )
751 .await;
752 emit_usage(observer, run).await;
753 self.record_turn_started(run, progress.turns, caused_by)?;
754 Ok(tool_choice)
755 }
756
757 fn record_turn_started(
758 &self,
759 run: &RunContext,
760 turn: u32,
761 caused_by: Option<EventId>,
762 ) -> Result<(), AgentError> {
763 record_domain(
764 run,
765 "turn.started",
766 serde_json::json!({"agent": self.name, "turn": turn}),
767 caused_by,
768 )
769 }
770
771 fn continue_provider_turn(
772 &self,
773 response: ModelResponse,
774 progress: &mut AgentProgress,
775 run: &RunContext,
776 checkpoint: &mut Option<&mut CheckpointCursor>,
777 ) -> Result<bool, AgentError> {
778 if !matches!(
779 &response.finish_reason,
780 runifold_model::FinishReason::Other(reason) if reason == "pause_turn"
781 ) {
782 return Ok(false);
783 }
784 let assistant = Message::new(Role::Assistant, response.content)
785 .map_err(|error| AgentError::Protocol(error.to_string()))?;
786 progress.transcript.push(assistant);
787 save_checkpoint(
788 checkpoint,
789 &self.checkpoint_state(progress, run, AgentCheckpointPhase::ReadyForTurn),
790 )?;
791 Ok(true)
792 }
793
794 async fn invoke_model(
795 &self,
796 transcript: &[Message],
797 run: &RunContext,
798 turn: u32,
799 tool_choice: ToolChoice,
800 caused_by: Option<EventId>,
801 observer: &dyn AgentObserver,
802 ) -> Result<ModelResponse, AgentError> {
803 record_domain(
804 run,
805 "model.started",
806 serde_json::json!({
807 "agent": self.name,
808 "turn": turn,
809 "provider": self.model_ref.provider,
810 "model": self.model_ref.name,
811 }),
812 caused_by,
813 )?;
814 let response = match self
815 .stream_model_response(self.request(transcript, tool_choice)?, run, turn, observer)
816 .await
817 {
818 Ok(response) => response,
819 Err(error) => {
820 record_domain(
821 run,
822 "model.failed",
823 serde_json::json!({
824 "agent": self.name,
825 "turn": turn,
826 "kind": format!("{:?}", error.kind),
827 }),
828 caused_by,
829 )?;
830 return Err(error.into());
831 }
832 };
833 record_domain(
834 run,
835 "model.completed",
836 serde_json::json!({
837 "agent": self.name,
838 "turn": turn,
839 "finish_reason": response.finish_reason,
840 "usage": response.usage,
841 }),
842 caused_by,
843 )?;
844 consume_budget(run, response.usage.into(), caused_by)?;
845 emit_usage(observer, run).await;
846 Ok(response)
847 }
848
849 async fn stream_model_response(
850 &self,
851 request: ModelRequest,
852 run: &RunContext,
853 turn: u32,
854 observer: &dyn AgentObserver,
855 ) -> Result<ModelResponse, ModelError> {
856 let context = ModelCallContext::for_run(run);
857 let cancellation = context.cancellation().clone();
858 let opening = self.model.stream(request, context);
859 let mut stream = match select(Box::pin(cancellation.cancelled()), Box::pin(opening)).await {
860 Either::Left(_) => return Err(cancelled_model_error()),
861 Either::Right((result, _)) => result?,
862 };
863 let mut accumulator = ModelStreamAccumulator::new();
864 loop {
865 let next = stream.next();
866 let event = match select(Box::pin(cancellation.cancelled()), Box::pin(next)).await {
867 Either::Left(_) => return Err(cancelled_model_error()),
868 Either::Right((Some(event), _)) => event?,
869 Either::Right((None, _)) => {
870 return Err(ModelError::local(
871 ModelErrorKind::Protocol,
872 "model stream ended before a terminal response event",
873 ));
874 }
875 };
876 let response = accumulator.push(event.clone())?;
877 emit_agent_event(observer, AgentStreamEvent::Model { turn, event }).await;
878 if let Some(response) = response {
879 return Ok(response);
880 }
881 }
882 }
883
884 fn validate_config(&self) -> Result<(), AgentError> {
885 if self.name.trim().is_empty() {
886 return Err(AgentError::InvalidConfig(
887 "agent name cannot be empty".into(),
888 ));
889 }
890 if self.config.max_turns == 0 {
891 return Err(AgentError::InvalidConfig(
892 "max_turns must be greater than zero".into(),
893 ));
894 }
895 if self.min_successful_tool_calls > 0 && self.tools.is_empty() {
896 return Err(AgentError::InvalidConfig(format!(
897 "min_successful_tool_calls={} requires at least one registered local Tool",
898 self.min_successful_tool_calls
899 )));
900 }
901 if let Some(collision) = self
902 .agents
903 .model_specs()
904 .into_iter()
905 .find(|spec| self.tools.contains(&spec.name))
906 {
907 return Err(AgentError::InvalidConfig(format!(
908 "callable name `{}` is registered as both a tool and an agent",
909 collision.name
910 )));
911 }
912 Ok(())
913 }
914
915 fn validate_checkpoint_identity(&self, state: &AgentCheckpointState) -> Result<(), AgentError> {
916 let terminal_reviewer = self
917 .terminal_review
918 .as_ref()
919 .map(|review| &review.descriptor);
920 let turn_reviewer = self.turn_review.as_ref().map(|review| &review.descriptor);
921 let terminal_review_policy = self.terminal_review.as_ref().map(|review| review.policy);
922 let turn_review_policy = self.turn_review.as_ref().map(|review| review.policy);
923 let terminal_reviewer_capabilities =
924 self.terminal_review
925 .as_ref()
926 .map_or_else(Vec::new, |review| {
927 review
928 .capabilities
929 .iter()
930 .map(|capability| capability.id)
931 .collect()
932 });
933 let turn_reviewer_capabilities =
934 self.turn_review.as_ref().map_or_else(Vec::new, |review| {
935 review
936 .capabilities
937 .iter()
938 .map(|capability| capability.id)
939 .collect()
940 });
941 if state.agent != self.name
942 || state.model != self.model_ref
943 || state.terminal_reviewer.as_ref() != terminal_reviewer
944 || state.turn_reviewer.as_ref() != turn_reviewer
945 || state.terminal_review_policy != terminal_review_policy
946 || state.turn_review_policy != turn_review_policy
947 || state.terminal_reviewer_capabilities != terminal_reviewer_capabilities
948 || state.turn_reviewer_capabilities != turn_reviewer_capabilities
949 {
950 return Err(runifold_core::CheckpointError::new(
951 runifold_core::CheckpointErrorKind::InvalidPayload,
952 "checkpoint Agent, model, reviewer identity, policy, or reviewer capabilities do not match",
953 )
954 .into());
955 }
956 Ok(())
957 }
958
959 async fn prepare_run_loop_progress(
960 &self,
961 state: AgentCheckpointState,
962 run: &RunContext,
963 checkpoint: &mut Option<&mut CheckpointCursor>,
964 context: &TerminalCompletionContext<'_>,
965 ) -> Result<(AgentProgress, Option<AgentOutcome>, Option<ModelResponse>), AgentError> {
966 enum Pending {
967 Terminal(ModelResponse, u32),
968 Turn(ModelResponse, u32),
969 Approved(ModelResponse, u32),
970 }
971
972 let pending = match &state.phase {
973 AgentCheckpointPhase::TerminalReviewReady { response, attempt } => {
974 Some(Pending::Terminal(response.as_ref().clone(), *attempt))
975 }
976 AgentCheckpointPhase::TurnReviewReady { response, turn } => {
977 Some(Pending::Turn(response.as_ref().clone(), *turn))
978 }
979 AgentCheckpointPhase::TurnReviewApproved { response, turn } => {
980 Some(Pending::Approved(response.as_ref().clone(), *turn))
981 }
982 AgentCheckpointPhase::ReadyForTurn => None,
983 _ => {
984 return Err(checkpoint_payload_error(
985 "checkpoint phase is not ready for Agent execution",
986 ));
987 }
988 };
989 let mut progress = AgentProgress::from(state);
990 let mut outcome = None;
991 let mut approved_response = None;
992 match pending {
993 Some(Pending::Terminal(response, attempt)) => {
994 outcome = self
995 .review_terminal_candidate(
996 response,
997 attempt,
998 run,
999 &mut progress,
1000 checkpoint,
1001 context,
1002 )
1003 .await?;
1004 }
1005 Some(Pending::Turn(response, turn)) => {
1006 validate_review_turn(turn, progress.turns)?;
1007 approved_response = self
1008 .review_turn_candidate(response, run, &mut progress, checkpoint, context)
1009 .await?;
1010 }
1011 Some(Pending::Approved(response, turn)) => {
1012 validate_review_turn(turn, progress.turns)?;
1013 approved_response = Some(response);
1014 }
1015 None => {}
1016 }
1017 Ok((progress, outcome, approved_response))
1018 }
1019
1020 fn prepare_resume_state(
1021 state: &mut AgentCheckpointState,
1022 run: &RunContext,
1023 policy: ResumePolicy,
1024 ) -> Result<(), AgentError> {
1025 match state.phase.clone() {
1026 AgentCheckpointPhase::TurnInFlight { turn } => {
1027 if policy == ResumePolicy::RejectAmbiguous {
1028 return Err(AgentError::AmbiguousCheckpoint { turn });
1029 }
1030 validate_usage_floor(state.usage, run.budget().usage())?;
1031 state.usage = run.budget().usage();
1032 state.phase = AgentCheckpointPhase::ReadyForTurn;
1033 }
1034 AgentCheckpointPhase::TerminalReviewInFlight { response, attempt } => {
1035 if policy == ResumePolicy::RejectAmbiguous {
1036 return Err(AgentError::AmbiguousTerminalReview { attempt });
1037 }
1038 validate_usage_floor(state.usage, run.budget().usage())?;
1039 state.usage = run.budget().usage();
1040 state.phase = AgentCheckpointPhase::TerminalReviewReady { response, attempt };
1041 }
1042 AgentCheckpointPhase::TurnReviewInFlight { response, turn } => {
1043 if policy == ResumePolicy::RejectAmbiguous {
1044 return Err(AgentError::AmbiguousTurnReview { turn });
1045 }
1046 validate_usage_floor(state.usage, run.budget().usage())?;
1047 state.usage = run.budget().usage();
1048 state.phase = AgentCheckpointPhase::TurnReviewReady { response, turn };
1049 }
1050 AgentCheckpointPhase::TurnReviewApproved { ref response, turn }
1051 if !tool_calls_from(&response.content).is_empty() =>
1052 {
1053 if policy == ResumePolicy::RejectAmbiguous {
1054 return Err(AgentError::AmbiguousCheckpoint { turn });
1055 }
1056 validate_usage_floor(state.usage, run.budget().usage())?;
1057 state.usage = run.budget().usage();
1058 }
1059 _ => validate_exact_usage(state.usage, run.budget().usage())?,
1060 }
1061 Ok(())
1062 }
1063
1064 pub(super) fn checkpoint_state(
1065 &self,
1066 progress: &AgentProgress,
1067 run: &RunContext,
1068 phase: AgentCheckpointPhase,
1069 ) -> AgentCheckpointState {
1070 AgentCheckpointState {
1071 execution_id: progress.execution_id.clone(),
1072 agent: self.name.clone(),
1073 model: self.model_ref.clone(),
1074 transcript: progress.transcript.clone(),
1075 turns: progress.turns,
1076 tool_calls: progress.tool_calls,
1077 delegations: progress.delegations,
1078 usage: run.budget().usage(),
1079 turn_reviewer: self
1080 .turn_review
1081 .as_ref()
1082 .map(|review| review.descriptor.clone()),
1083 turn_review_policy: self.turn_review.as_ref().map(|review| review.policy),
1084 turn_reviewer_capabilities: self.turn_review.as_ref().map_or_else(Vec::new, |review| {
1085 review
1086 .capabilities
1087 .iter()
1088 .map(|capability| capability.id)
1089 .collect()
1090 }),
1091 terminal_reviewer: self
1092 .terminal_review
1093 .as_ref()
1094 .map(|review| review.descriptor.clone()),
1095 terminal_review_policy: self.terminal_review.as_ref().map(|review| review.policy),
1096 terminal_reviewer_capabilities: self.terminal_review.as_ref().map_or_else(
1097 Vec::new,
1098 |review| {
1099 review
1100 .capabilities
1101 .iter()
1102 .map(|capability| capability.id)
1103 .collect()
1104 },
1105 ),
1106 phase,
1107 durable_conversation: progress.durable_conversation.clone(),
1108 }
1109 }
1110
1111 async fn commit_durable_outcome(
1112 &self,
1113 store: &dyn DurableConversationStore,
1114 run: &RunContext,
1115 cursor: &CheckpointCursor,
1116 durable: DurableConversationCheckpoint,
1117 outcome: AgentOutcome,
1118 ) -> Result<AgentConversationOutcome, AgentConversationError> {
1119 let persisted_prefix_len = usize::try_from(durable.persisted_prefix_len).map_err(|_| {
1120 AgentConversationError::Run(checkpoint_payload_error(
1121 "durable conversation prefix does not fit this platform",
1122 ))
1123 })?;
1124 if persisted_prefix_len >= outcome.transcript.len() {
1125 return Err(AgentConversationError::Run(checkpoint_payload_error(
1126 "durable conversation checkpoint has an invalid transcript prefix",
1127 )));
1128 }
1129 let messages = outcome
1130 .transcript
1131 .iter()
1132 .skip(persisted_prefix_len)
1133 .filter(|message| !is_transient_context(message))
1134 .cloned()
1135 .collect();
1136 let state = AgentCheckpointState {
1137 execution_id: cursor.id().to_string(),
1138 agent: self.name.clone(),
1139 model: self.model_ref.clone(),
1140 transcript: outcome.transcript.clone(),
1141 turns: outcome.turns,
1142 tool_calls: outcome.tool_calls,
1143 delegations: outcome.delegations,
1144 usage: run.budget().usage(),
1145 turn_reviewer: self
1146 .turn_review
1147 .as_ref()
1148 .map(|review| review.descriptor.clone()),
1149 turn_review_policy: self.turn_review.as_ref().map(|review| review.policy),
1150 turn_reviewer_capabilities: self.turn_review.as_ref().map_or_else(Vec::new, |review| {
1151 review
1152 .capabilities
1153 .iter()
1154 .map(|capability| capability.id)
1155 .collect()
1156 }),
1157 terminal_reviewer: self
1158 .terminal_review
1159 .as_ref()
1160 .map(|review| review.descriptor.clone()),
1161 terminal_review_policy: self.terminal_review.as_ref().map(|review| review.policy),
1162 terminal_reviewer_capabilities: self.terminal_review.as_ref().map_or_else(
1163 Vec::new,
1164 |review| {
1165 review
1166 .capabilities
1167 .iter()
1168 .map(|capability| capability.id)
1169 .collect()
1170 },
1171 ),
1172 phase: AgentCheckpointPhase::Completed {
1173 response: Box::new(outcome.response.clone()),
1174 },
1175 durable_conversation: Some(durable.clone()),
1176 };
1177 let checkpoint = cursor.next(&state).map_err(AgentConversationError::Run)?;
1178 let command = DurableConversationCommit {
1179 namespace: durable.namespace,
1180 append: ConversationAppend {
1181 conversation_id: durable.conversation_id,
1182 expected_version: durable.expected_version,
1183 messages,
1184 },
1185 checkpoint,
1186 expected_checkpoint_revision: cursor.revision(),
1187 };
1188 match store.commit_durable_turn(command).await {
1189 Ok(conversation_version) => Ok(AgentConversationOutcome {
1190 outcome,
1191 conversation_version,
1192 }),
1193 Err(source) => Err(AgentConversationError::Commit {
1194 source,
1195 outcome: Box::new(outcome),
1196 }),
1197 }
1198 }
1199
1200 pub(super) fn check_lifecycle(run: &RunContext) -> Result<(), AgentError> {
1201 let error = if run.cancellation().is_cancelled() {
1202 Some((
1203 runifold_model::ModelErrorKind::Cancelled,
1204 "agent run was cancelled",
1205 ))
1206 } else if run
1207 .deadline()
1208 .is_some_and(|deadline| deadline <= Instant::now())
1209 {
1210 Some((
1211 runifold_model::ModelErrorKind::DeadlineExceeded,
1212 "agent run deadline elapsed",
1213 ))
1214 } else {
1215 None
1216 };
1217 if let Some((kind, message)) = error {
1218 return Err(runifold_model::ModelError::local(kind, message).into());
1219 }
1220 Ok(())
1221 }
1222
1223 fn request(
1224 &self,
1225 transcript: &[Message],
1226 tool_choice: ToolChoice,
1227 ) -> Result<ModelRequest, AgentError> {
1228 let (first, rest) = transcript
1229 .split_first()
1230 .ok_or_else(|| AgentError::Protocol("agent transcript is empty".into()))?;
1231 let mut request = ModelRequest::new(self.model_ref.clone(), first.clone());
1232 request.messages.extend_from_slice(rest);
1233 request.tools = self.tools.model_specs();
1234 request.tools.extend(self.agents.model_specs());
1235 request.tool_choice = tool_choice;
1236 for tool in &self.provider_tools {
1237 request = request.provider_tool(tool.clone());
1238 }
1239 request.generation.clone_from(&self.generation);
1240 request = request.response_mode(self.response_mode);
1241 request.provider_options.clone_from(&self.provider_options);
1242 request.feature_policy = self.config.feature_policy;
1243 request.output_format.clone_from(&self.output_format);
1244 Ok(request)
1245 }
1246
1247 fn successful_local_tool_calls(&self, progress: &AgentProgress) -> Result<u32, AgentError> {
1248 let count = progress
1249 .transcript
1250 .iter()
1251 .filter(|message| {
1252 message
1253 .metadata
1254 .get(TOOL_RESULT_EXECUTION_ID_METADATA)
1255 .and_then(serde_json::Value::as_str)
1256 == Some(progress.execution_id.as_str())
1257 })
1258 .flat_map(|message| &message.content)
1259 .filter(|part| {
1260 matches!(
1261 part,
1262 ContentPart::ToolResult(result)
1263 if !result.is_error
1264 && result
1265 .name
1266 .as_deref()
1267 .is_some_and(|name| self.tools.contains(name))
1268 )
1269 })
1270 .count();
1271 u32::try_from(count)
1272 .map_err(|_| AgentError::Protocol("successful Tool-call counter overflow".into()))
1273 }
1274
1275 fn next_tool_choice(
1276 &self,
1277 progress: &AgentProgress,
1278 run: &RunContext,
1279 ) -> Result<ToolChoice, AgentError> {
1280 let successful = self.successful_local_tool_calls(progress)?;
1281 let remaining_required = self.min_successful_tool_calls.saturating_sub(successful);
1282 Self::validate_tool_requirement_budget(remaining_required, run)?;
1283 if progress.turns >= self.config.max_turns {
1284 if remaining_required > 0 {
1285 return Err(AgentError::ToolRequirementUnsatisfied {
1286 required: self.min_successful_tool_calls,
1287 successful,
1288 });
1289 }
1290 return Err(AgentError::MaxTurns {
1291 max_turns: self.config.max_turns,
1292 });
1293 }
1294 Ok(if remaining_required > 0 {
1295 ToolChoice::Required
1296 } else {
1297 ToolChoice::Auto
1298 })
1299 }
1300
1301 fn validate_tool_requirement_budget(
1302 remaining_required: u32,
1303 run: &RunContext,
1304 ) -> Result<(), AgentError> {
1305 let Some(limit) = run.budget().limit().tool_calls else {
1306 return Ok(());
1307 };
1308 let remaining = limit.saturating_sub(run.budget().usage().tool_calls);
1309 if u64::from(remaining_required) > remaining {
1310 return Err(AgentError::ToolRequirementExceedsBudget {
1311 required: remaining_required,
1312 remaining,
1313 });
1314 }
1315 Ok(())
1316 }
1317}
1318
1319fn validate_tool_call_completion(
1320 calls: &[ToolCall],
1321 finish_reason: &runifold_model::FinishReason,
1322) -> Result<(), AgentError> {
1323 if !calls.is_empty() && !matches!(finish_reason, runifold_model::FinishReason::ToolCalls) {
1324 return Err(AgentError::Protocol(format!(
1325 "refusing to execute tool calls from a {finish_reason:?} model response"
1326 )));
1327 }
1328 Ok(())
1329}
1330
1331fn checkpoint_payload_error(message: &str) -> AgentError {
1332 runifold_core::CheckpointError::new(runifold_core::CheckpointErrorKind::InvalidPayload, message)
1333 .into()
1334}
1335
1336fn validate_review_turn(expected: u32, actual: u32) -> Result<(), AgentError> {
1337 if expected != actual {
1338 return Err(checkpoint_payload_error(
1339 "turn review checkpoint does not match the completed model turn count",
1340 ));
1341 }
1342 Ok(())
1343}
1344
1345fn cancelled_model_error() -> ModelError {
1346 ModelError::local(ModelErrorKind::Cancelled, "model invocation was cancelled")
1347}
1348
1349fn tool_calls_from(content: &[ContentPart]) -> Vec<ToolCall> {
1350 content
1351 .iter()
1352 .filter_map(|part| match part {
1353 ContentPart::ToolCall(call) => Some(call.clone()),
1354 _ => None,
1355 })
1356 .collect()
1357}