Skip to main content

runifold_agent/agent/
execution.rs

1//! Canonical Agent execution engine and its private runtime helpers.
2
3use 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    /// Runs a user text turn with a default root context.
29    ///
30    /// This is the ergonomic surface for one-off prompts. It grants only
31    /// callables registered on this Agent and applies no hard budget limit.
32    /// Use [`Self::run`] when the caller must provide explicit authority,
33    /// budget, deadline, observability, or run-tree identity.
34    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    /// Runs an ergonomic prompt and returns only model-visible text.
46    ///
47    /// Rich content, usage, warnings, the canonical transcript, and provider
48    /// events are intentionally discarded. Use [`Self::prompt`] when that
49    /// information matters.
50    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    /// Runs a user text turn inside an existing runtime context.
59    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    /// Runs and atomically commits one bounded multi-turn conversation.
73    ///
74    /// Transcript messages remain append-only. Execution-journal events stay
75    /// in [`runifold_core::Journal`], summaries remain lossy derived views,
76    /// and semantic memory is injected only as explicitly untrusted context.
77    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    /// Summarizes an overflowing prefix before running a conversational turn.
157    ///
158    /// Summary generation uses the supplied [`AutomaticConversationSummary`]
159    /// and the same [`RunContext`], preserving cancellation, deadline, budget,
160    /// and journal behavior. The immutable transcript is never rewritten.
161    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    /// Streams real-time events while driving the canonical Agent loop.
226    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    /// Runs one conversational turn with atomic transcript and checkpoint commit.
240    ///
241    /// Intermediate checkpoints are written ahead of model and callable work.
242    /// The terminal checkpoint and transcript append are committed together by
243    /// [`DurableConversationStore`].
244    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    /// Resumes a durable conversational turn from its write-ahead checkpoint.
331    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    /// Runs with write-ahead checkpoint persistence.
393    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    /// Resumes a persisted Agent execution.
417    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}