Skip to main content

everruns_engine/
turn.rs

1// The sans-IO turn planner (EVE-840, Sans-IO Turn State epic).
2//
3// This is the authoritative turn-planning brain. Every function here is pure and
4// deterministic: it reads only its arguments and returns a `TurnPlan` (plus, for
5// terminal outcomes, a list of `TurnLifecycleEffect`s the host must perform). It
6// never touches a store, socket, process, event bus, or `Utc::now()` — the host
7// resolves those facts, passes `now` in, and performs the returned effects.
8
9use crate::{ActInput, ExecutionContext, ReasonResult};
10use chrono::{DateTime, Utc};
11use everruns_contracts::typed_id::{
12    AgentId, ExecId, HarnessId, MessageId, SessionId, TurnId, WorkspaceId,
13};
14use everruns_contracts::user_facing_error::codes as user_facing_error_codes;
15use everruns_contracts::user_facing_error::{
16    ErrorDisclosure, UserFacingError, UserFacingErrorContext, classify_runtime_error_message,
17};
18use everruns_core::events::{TokenUsage, TurnCompletedData};
19use everruns_core::turn::TurnStopReason;
20use serde::{Deserialize, Serialize};
21use tracing::{debug, info};
22
23/// Host-owned state carried across turn phases.
24///
25/// Durable hosts can persist this between activities; in-memory hosts can hold
26/// it directly in memory. The type itself is engine-level and has no host,
27/// store, or durable-engine coupling.
28///
29/// Hosts are expected to serialize this however they want. `everruns-engine`
30/// only defines the fields required to resume the next semantic step.
31#[derive(Debug, Clone, Serialize, Deserialize)]
32pub struct TurnState {
33    pub org_id: i64,
34    pub session_id: SessionId,
35    pub harness_id: HarnessId,
36    pub agent_id: Option<AgentId>,
37    pub input_message_id: MessageId,
38    #[serde(skip_serializing_if = "Option::is_none")]
39    pub turn_id: Option<TurnId>,
40    #[serde(skip_serializing_if = "Option::is_none", default)]
41    pub previous_response_id: Option<String>,
42    #[serde(default = "default_iteration")]
43    pub iteration: u32,
44    #[serde(skip_serializing_if = "Option::is_none", default)]
45    pub request_id: Option<String>,
46    #[serde(skip_serializing_if = "Option::is_none", default)]
47    pub started_at: Option<DateTime<Utc>>,
48    #[serde(skip_serializing_if = "Option::is_none", default)]
49    pub cumulative_usage: Option<TokenUsage>,
50    #[serde(default)]
51    pub tool_call_count: u32,
52    #[serde(default)]
53    pub llm_call_count: u32,
54    #[serde(skip_serializing_if = "Option::is_none", default)]
55    pub time_to_first_token_ms: Option<u64>,
56    #[serde(skip_serializing_if = "Option::is_none", default)]
57    pub final_message_id: Option<MessageId>,
58    #[serde(skip_serializing_if = "Option::is_none", default)]
59    pub final_answer_preview: Option<String>,
60}
61
62fn default_iteration() -> u32 {
63    1
64}
65
66/// Engine-owned act scheduling payload.
67///
68/// Hosts enqueue or execute this immediately using their own worker model.
69#[derive(Debug, Clone)]
70pub struct ActPlan {
71    pub input: ActInput,
72    pub previous_response_id: Option<String>,
73    pub iteration: u32,
74    pub request_id: Option<String>,
75    pub resume_state: Box<TurnState>,
76}
77
78/// Generic next-step decision for a host turn.
79///
80/// This intentionally stops at the semantic boundary:
81/// - the engine decides what should happen next
82/// - the host decides how to persist, enqueue, retry, or resume it
83#[derive(Debug, Clone)]
84pub enum TurnPlan {
85    ScheduleReason(TurnState),
86    ScheduleAct(ActPlan),
87    Complete {
88        stop_reason: TurnStopReason,
89        error: Option<String>,
90    },
91    WaitForToolResults {
92        resume: TurnState,
93    },
94}
95
96/// A lifecycle side effect the engine decided must be recorded, described as
97/// data rather than performed.
98///
99/// The engine never emits events or fires hooks — it *returns* these, and the
100/// host applies them (in list order) through its own lifecycle machinery. This
101/// keeps planning deterministic while preserving the exact event stream. These
102/// are NOT part of the public [`TurnPlan`]; they travel alongside it.
103#[derive(Debug, Clone)]
104pub enum TurnLifecycleEffect {
105    /// Emit `turn.completed` with the summarized turn fields.
106    TurnCompleted {
107        input_message_id: MessageId,
108        data: TurnCompletedData,
109    },
110    /// Answer an `ask_user` call the client structurally cannot be asked.
111    ///
112    /// Emitted instead of parking when the session never declared the
113    /// `ask_user` hint — a scheduled run, a trigger, an SDK caller. The host
114    /// completes each call with the declared defaults and
115    /// `answered_by: "unattended"`, and the turn carries on in the same
116    /// iteration (EVE-1057).
117    ResolveAskUserUnattended {
118        /// The turn that asked, so the completion is filed against it rather
119        /// than floating free in the session.
120        turn_id: Option<TurnId>,
121        input_message_id: MessageId,
122        /// The `ask_user` calls to answer, as `(tool_call_id, arguments)`.
123        calls: Vec<(String, serde_json::Value)>,
124    },
125    /// Idle the session and emit `session.idled`.
126    SessionIdled {
127        turn_id: TurnId,
128        input_message_id: MessageId,
129        iterations: Option<u32>,
130        usage: Option<TokenUsage>,
131    },
132    /// Fail the turn with the already-disclosure-filtered error, emitting
133    /// `turn.failed` + `session.idled`.
134    TurnFailedWithDisclosure {
135        turn_id: TurnId,
136        input_message_id: MessageId,
137        text: String,
138        user_error: Option<UserFacingError>,
139        disclosure: Option<ErrorDisclosure>,
140    },
141    /// Fire the advisory `turn_end` lifecycle hooks.
142    FireTurnEndHooks {
143        harness_id: HarnessId,
144        agent_id: Option<AgentId>,
145        turn_id: TurnId,
146        success: bool,
147    },
148    /// Mark the session `waiting_for_tool_results`.
149    WaitingForToolResults,
150}
151
152/// Parsed `act` activity output the planner decides over.
153#[derive(Debug, Clone, Copy, Default)]
154pub struct ActOutcome {
155    pub blocked: bool,
156    pub waiting_for_tool_results: bool,
157    /// The pause is a URL mode elicitation consent card, which only a client
158    /// that declared `url_elicitation` can answer.
159    pub waiting_for_url_elicitation: bool,
160    /// The pause is an `ask_user` question set, which only a client that
161    /// declared `ask_user` can answer (EVE-1057).
162    pub waiting_for_ask_user: bool,
163    /// The pause includes a hard tool-approval request (EVE-1140). Held
164    /// whatever the client declared: the gated call already failed closed, and
165    /// a person can answer it through the API even when no card is drawn.
166    pub waiting_for_tool_approval: bool,
167}
168
169/// Session facts the host pre-resolves for the reason→act scheduling case.
170///
171/// The host fetches these (from its session store) only when
172/// [`reason_schedules_act`] is true, mirroring the original conditional fetch:
173/// `blueprint_id` scopes blueprint tool resolution, and `workspace_id` points
174/// tool file I/O at the (possibly shared) workspace rather than the session's
175/// own keyspace.
176#[derive(Debug, Clone, Default)]
177pub struct ActSchedulingFacts {
178    pub blueprint_id: Option<String>,
179    pub workspace_id: Option<WorkspaceId>,
180}
181
182/// Typed, parsed activity output the engine plans the next step from.
183///
184/// The host parses the raw serialized activity output into this before calling
185/// [`plan_next_turn`]; unknown activity kinds are rejected by the host, so the
186/// engine stays total.
187pub enum ActivityOutcome {
188    ProcessInput { turn_id: Option<TurnId> },
189    // Boxed: `ReasonResult` dwarfs the other variants (clippy::large_enum_variant).
190    Reason(Box<ReasonResult>),
191    Act(ActOutcome),
192}
193
194/// Host-resolved facts the engine needs but cannot fetch itself.
195///
196/// The host populates only the field relevant to the completed activity, doing
197/// I/O in exactly the same conditions as the original planner:
198/// `act_scheduling` only when [`reason_schedules_act`] is true, and
199/// `setup_connection_hint_enabled` only when the act paused for tool results.
200#[derive(Debug, Clone, Default)]
201pub struct HostFacts {
202    pub act_scheduling: Option<ActSchedulingFacts>,
203    pub setup_connection_hint_enabled: bool,
204    pub url_elicitation_hint_enabled: bool,
205    pub ask_user_hint_enabled: bool,
206    /// The `ask_user` calls the act left pending, as `(tool_call_id,
207    /// arguments)`. Carried here rather than on [`ActOutcome`], which is `Copy`
208    /// and cannot hold them. Empty unless the act paused on `ask_user`.
209    pub ask_user_calls: Vec<(String, serde_json::Value)>,
210}
211
212fn preview_final_answer(text: &str) -> Option<String> {
213    if text.is_empty() {
214        return None;
215    }
216
217    Some(text.chars().take(2000).collect())
218}
219
220fn add_usage(current: &mut Option<TokenUsage>, next: &TokenUsage) {
221    match current {
222        Some(current) => current.add(next),
223        None => *current = Some(next.clone()),
224    }
225}
226
227impl TurnState {
228    pub(crate) fn with_reason_summary(&self, reason_result: &ReasonResult) -> Self {
229        let mut next = self.clone();
230        next.llm_call_count = next.llm_call_count.saturating_add(
231            reason_result
232                .native_counts
233                .as_ref()
234                .map_or(1, |counts| counts.llm_calls),
235        );
236        next.tool_call_count = next.tool_call_count.saturating_add(
237            reason_result
238                .native_counts
239                .as_ref()
240                .map_or(reason_result.tool_calls.len() as u32, |counts| {
241                    counts.tool_calls
242                }),
243        );
244        if let Some(usage) = &reason_result.usage {
245            add_usage(&mut next.cumulative_usage, usage);
246        }
247        if next.time_to_first_token_ms.is_none() {
248            next.time_to_first_token_ms = reason_result.time_to_first_token_ms;
249        }
250        next.final_message_id = reason_result.output_message_id;
251        next.final_answer_preview = preview_final_answer(&reason_result.text);
252        next
253    }
254
255    /// Wall-clock duration since `started_at`, measured against the host-supplied
256    /// `now` so the calculation stays deterministic.
257    fn duration_ms(&self, now: DateTime<Utc>) -> Option<u64> {
258        self.started_at
259            .map(|started_at| now.signed_duration_since(started_at))
260            .and_then(|duration| u64::try_from(duration.num_milliseconds()).ok())
261    }
262}
263
264fn classify_reason_failure(reason_result: &ReasonResult) -> UserFacingError {
265    // The reason atom already classified and disclosure-filtered the failure.
266    // Reuse it so the turn.failed event matches what the session message
267    // showed; re-classifying strings here could leak past a generic mode.
268    if let Some(user_error) = &reason_result.user_facing_error {
269        return user_error.clone();
270    }
271
272    let from_text =
273        classify_runtime_error_message(&reason_result.text, &UserFacingErrorContext::default());
274
275    let Some(error) = reason_result.error.as_deref() else {
276        return from_text;
277    };
278
279    let from_error = classify_runtime_error_message(error, &UserFacingErrorContext::default());
280
281    if from_error.code == user_facing_error_codes::PROCESSING_ERROR {
282        return from_text;
283    }
284
285    if from_error.code == from_text.code
286        && from_error.fields.is_empty()
287        && !from_text.fields.is_empty()
288    {
289        return from_text;
290    }
291
292    from_error
293}
294
295/// Does this reason outcome schedule an act phase?
296///
297/// The host consults this predicate to decide whether to resolve
298/// [`ActSchedulingFacts`] (a session fetch) before calling [`plan_after_reason`]
299/// — the same condition under which the original planner fetched the session.
300/// The reason planner branches on this same function, so the rule has exactly
301/// one definition.
302pub fn reason_schedules_act(state: &TurnState, reason_result: &ReasonResult) -> bool {
303    let max_turn_requests_reached = state.iteration >= reason_result.max_iterations as u32;
304    reason_result.has_tool_calls && reason_result.success && !max_turn_requests_reached
305}
306
307/// Plan the next host step after an activity finishes.
308///
309/// The authoritative, deterministic turn-planning entry point. Given the carried
310/// [`TurnState`], the parsed [`ActivityOutcome`], the count of queued steering
311/// messages, the host-supplied `now`, and any [`HostFacts`] the host pre-resolved,
312/// it returns the [`TurnPlan`] together with the [`TurnLifecycleEffect`]s the
313/// host must perform (in order). It performs no I/O of its own.
314pub fn plan_next_turn(
315    state: &TurnState,
316    outcome: ActivityOutcome,
317    pending_user_message_count: usize,
318    now: DateTime<Utc>,
319    facts: HostFacts,
320) -> (TurnPlan, Vec<TurnLifecycleEffect>) {
321    match outcome {
322        ActivityOutcome::ProcessInput { turn_id } => {
323            (plan_after_process_input(state, turn_id, now), Vec::new())
324        }
325        ActivityOutcome::Reason(reason_result) => plan_after_reason(
326            state,
327            *reason_result,
328            pending_user_message_count,
329            now,
330            facts.act_scheduling,
331        ),
332        ActivityOutcome::Act(outcome) => plan_after_act(
333            state,
334            outcome,
335            facts.setup_connection_hint_enabled,
336            facts.url_elicitation_hint_enabled,
337            facts.ask_user_hint_enabled,
338            facts.ask_user_calls.clone(),
339        ),
340    }
341}
342
343/// Plan the reason step that follows a completed `process_input` activity.
344pub fn plan_after_process_input(
345    state: &TurnState,
346    turn_id: Option<TurnId>,
347    now: DateTime<Utc>,
348) -> TurnPlan {
349    let next = TurnState {
350        turn_id,
351        previous_response_id: None,
352        iteration: 1,
353        started_at: state.started_at.or(Some(now)),
354        ..state.clone()
355    };
356    debug!(session_id = %state.session_id, turn_id = ?turn_id, "planned reason step");
357    TurnPlan::ScheduleReason(next)
358}
359
360/// Plan the next step after a `reason` activity finishes.
361///
362/// When [`reason_schedules_act`] holds, `act_scheduling` supplies the session
363/// facts the host resolved for the act phase; it is ignored otherwise. A
364/// terminal reason outcome returns the lifecycle effects the host must perform;
365/// the continuing outcomes return an empty effect list.
366pub fn plan_after_reason(
367    state: &TurnState,
368    reason_result: ReasonResult,
369    pending_user_message_count: usize,
370    now: DateTime<Utc>,
371    act_scheduling: Option<ActSchedulingFacts>,
372) -> (TurnPlan, Vec<TurnLifecycleEffect>) {
373    let response_id = reason_result.response_id.clone();
374    let summarized_state = state.with_reason_summary(&reason_result);
375    let max_turn_requests_reached = state.iteration >= reason_result.max_iterations as u32;
376
377    // A remote tool loop paused on a tool call (EVE-1124). Checked before
378    // anything else: the provider still holds that call open, so neither a
379    // steering message nor the iteration cap may end or advance the turn.
380    if reason_result.success && reason_result.waiting_for_tool_results {
381        let next = TurnState {
382            previous_response_id: response_id,
383            iteration: state.iteration.saturating_add(1),
384            ..summarized_state
385        };
386        return (
387            TurnPlan::WaitForToolResults { resume: next },
388            vec![TurnLifecycleEffect::WaitingForToolResults],
389        );
390    }
391
392    if reason_schedules_act(state, &reason_result) {
393        let facts = act_scheduling.unwrap_or_default();
394        let plan = ActPlan {
395            input: ActInput {
396                org_id: Some(state.org_id),
397                context: ExecutionContext {
398                    session_id: state.session_id,
399                    turn_id: state.turn_id.unwrap_or_default(),
400                    input_message_id: state.input_message_id,
401                    exec_id: ExecId::new(),
402                    workspace_id: facts.workspace_id,
403                },
404                harness_id: state.harness_id,
405                agent_id: state.agent_id,
406                tool_calls: reason_result.tool_calls,
407                tool_definitions: reason_result.tool_definitions,
408                locale: reason_result.locale,
409                blueprint_id: facts.blueprint_id,
410                network_access: reason_result.network_access,
411                // Request-level parallel tool calling preference, carried
412                // from agent config through the reason path (EVE-598).
413                parallel_tool_calls: reason_result.parallel_tool_calls,
414            },
415            previous_response_id: response_id,
416            iteration: state.iteration,
417            request_id: state.request_id.clone(),
418            resume_state: Box::new(summarized_state),
419        };
420        return (TurnPlan::ScheduleAct(plan), Vec::new());
421    }
422
423    if reason_result.success && pending_user_message_count > 0 && !max_turn_requests_reached {
424        if pending_user_message_count > 1 {
425            info!(
426                session_id = %state.session_id,
427                pending_user_message_count,
428                "multiple steering messages arrived during turn"
429            );
430        }
431
432        let next = TurnState {
433            previous_response_id: response_id,
434            iteration: state.iteration.saturating_add(1),
435            ..summarized_state
436        };
437        return (TurnPlan::ScheduleReason(next), Vec::new());
438    }
439
440    let turn_id = state.turn_id.unwrap_or_default();
441    let mut effects = Vec::new();
442
443    if reason_result.success {
444        effects.push(TurnLifecycleEffect::TurnCompleted {
445            input_message_id: state.input_message_id,
446            data: TurnCompletedData {
447                turn_id,
448                iterations: state.iteration,
449                duration_ms: summarized_state.duration_ms(now),
450                usage: summarized_state.cumulative_usage.clone(),
451                input_content: None,
452                final_message_id: summarized_state.final_message_id,
453                final_answer_preview: summarized_state.final_answer_preview.clone(),
454                time_to_first_token_ms: summarized_state.time_to_first_token_ms,
455                tool_call_count: Some(summarized_state.tool_call_count),
456                llm_call_count: Some(summarized_state.llm_call_count),
457                status: Some("completed".to_string()),
458            },
459        });
460        effects.push(TurnLifecycleEffect::SessionIdled {
461            turn_id,
462            input_message_id: state.input_message_id,
463            iterations: Some(state.iteration),
464            usage: summarized_state.cumulative_usage.clone(),
465        });
466    } else {
467        let user_error = classify_reason_failure(&reason_result);
468        effects.push(TurnLifecycleEffect::TurnFailedWithDisclosure {
469            turn_id,
470            input_message_id: state.input_message_id,
471            text: reason_result.text.clone(),
472            user_error: Some(user_error),
473            disclosure: reason_result.error_disclosure,
474        });
475    }
476
477    // turn_end lifecycle hooks (advisory). Fired once the turn reaches a
478    // terminal reason outcome on the durable/strategy path.
479    effects.push(TurnLifecycleEffect::FireTurnEndHooks {
480        harness_id: state.harness_id,
481        agent_id: state.agent_id,
482        turn_id,
483        success: reason_result.success,
484    });
485
486    let stop_reason = if !reason_result.success {
487        match TurnStopReason::from_provider_finish_reason(reason_result.finish_reason.as_deref()) {
488            TurnStopReason::Refusal => TurnStopReason::Refusal,
489            _ => TurnStopReason::Error,
490        }
491    } else if max_turn_requests_reached
492        && (reason_result.has_tool_calls || pending_user_message_count > 0)
493    {
494        TurnStopReason::MaxTurnRequests
495    } else {
496        TurnStopReason::from_provider_finish_reason(reason_result.finish_reason.as_deref())
497    };
498
499    (
500        TurnPlan::Complete {
501            stop_reason,
502            error: reason_result.error,
503        },
504        effects,
505    )
506}
507
508/// Whether an act outcome parks the turn, given the session's pause hints.
509///
510/// One definition shared by [`plan_after_act`] and hosts that run tools
511/// inside a remote reason loop (the OpenAI Agents API backend), so both pause
512/// on exactly the same conditions.
513pub fn act_pauses_turn(
514    outcome: ActOutcome,
515    setup_connection_hint_enabled: bool,
516    url_elicitation_hint_enabled: bool,
517    ask_user_hint_enabled: bool,
518) -> bool {
519    outcome.waiting_for_tool_results
520        && if outcome.waiting_for_tool_approval {
521            true
522        } else if outcome.waiting_for_ask_user {
523            ask_user_hint_enabled
524        } else {
525            setup_connection_hint_enabled
526                || (outcome.waiting_for_url_elicitation && url_elicitation_hint_enabled)
527        }
528}
529
530/// Plan the next step after an `act` activity finishes.
531///
532/// `setup_connection_hint_enabled` and `url_elicitation_hint_enabled` are the
533/// resolved session hints; the host reads them only when the act reported
534/// `waiting_for_tool_results`, so passing `false` otherwise matches the original
535/// short-circuit exactly.
536pub fn plan_after_act(
537    state: &TurnState,
538    outcome: ActOutcome,
539    setup_connection_hint_enabled: bool,
540    url_elicitation_hint_enabled: bool,
541    ask_user_hint_enabled: bool,
542    ask_user_calls: Vec<(String, serde_json::Value)>,
543) -> (TurnPlan, Vec<TurnLifecycleEffect>) {
544    if outcome.blocked {
545        return (
546            TurnPlan::Complete {
547                stop_reason: TurnStopReason::EndTurn,
548                error: None,
549            },
550            Vec::new(),
551        );
552    }
553
554    // A pause is only useful if the client on the other end can answer it. A
555    // URL elicitation waits on a consent card, so it needs a client that
556    // declared it renders one; everything else rides the `setup_connection`
557    // hint as before. Without the matching hint the turn continues and the
558    // elicitation reaches the user as an ordinary tool result instead.
559    // An `ask_user` pause needs its own hint for the same reason: a question
560    // nobody can render is not worth holding a turn for (EVE-1057). It does not
561    // ride `setup_connection`, because a client can be able to finish a
562    // connection setup and still have no way to draw a question.
563    //
564    // A hard tool-approval request is the exception that pauses regardless of
565    // hints (EVE-1140). Continuing would not let the gated call through — it
566    // already failed closed — but it would throw away the only chance a person
567    // has to approve it, and an API caller can answer without any card.
568    let should_pause_for_tool_results = act_pauses_turn(
569        outcome,
570        setup_connection_hint_enabled,
571        url_elicitation_hint_enabled,
572        ask_user_hint_enabled,
573    );
574
575    let next = TurnState {
576        iteration: state.iteration.saturating_add(1),
577        ..state.clone()
578    };
579
580    if should_pause_for_tool_results {
581        return (
582            TurnPlan::WaitForToolResults { resume: next },
583            vec![TurnLifecycleEffect::WaitingForToolResults],
584        );
585    }
586
587    if outcome.waiting_for_tool_results {
588        info!(
589            session_id = %state.session_id,
590            waiting_for_url_elicitation = outcome.waiting_for_url_elicitation,
591            waiting_for_ask_user = outcome.waiting_for_ask_user,
592            waiting_for_tool_approval = outcome.waiting_for_tool_approval,
593            "no hint declares this client can answer the pause, continuing turn instead"
594        );
595    }
596
597    // A URL elicitation continuing unanswered reaches the model as an ordinary
598    // tool result, which reads correctly. An `ask_user` call would instead sit
599    // in the transcript with nothing answering it, so the defaults are applied
600    // here and the model is told plainly that nobody could be asked (EVE-1057).
601    let effects = if outcome.waiting_for_ask_user && !ask_user_calls.is_empty() {
602        vec![TurnLifecycleEffect::ResolveAskUserUnattended {
603            turn_id: state.turn_id,
604            input_message_id: state.input_message_id,
605            calls: ask_user_calls,
606        }]
607    } else {
608        Vec::new()
609    };
610
611    (TurnPlan::ScheduleReason(next), effects)
612}