1use 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#[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#[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#[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#[derive(Debug, Clone)]
104pub enum TurnLifecycleEffect {
105 TurnCompleted {
107 input_message_id: MessageId,
108 data: TurnCompletedData,
109 },
110 ResolveAskUserUnattended {
118 turn_id: Option<TurnId>,
121 input_message_id: MessageId,
122 calls: Vec<(String, serde_json::Value)>,
124 },
125 SessionIdled {
127 turn_id: TurnId,
128 input_message_id: MessageId,
129 iterations: Option<u32>,
130 usage: Option<TokenUsage>,
131 },
132 TurnFailedWithDisclosure {
135 turn_id: TurnId,
136 input_message_id: MessageId,
137 text: String,
138 user_error: Option<UserFacingError>,
139 disclosure: Option<ErrorDisclosure>,
140 },
141 FireTurnEndHooks {
143 harness_id: HarnessId,
144 agent_id: Option<AgentId>,
145 turn_id: TurnId,
146 success: bool,
147 },
148 WaitingForToolResults,
150}
151
152#[derive(Debug, Clone, Copy, Default)]
154pub struct ActOutcome {
155 pub blocked: bool,
156 pub waiting_for_tool_results: bool,
157 pub waiting_for_url_elicitation: bool,
160 pub waiting_for_ask_user: bool,
163 pub waiting_for_tool_approval: bool,
167}
168
169#[derive(Debug, Clone, Default)]
177pub struct ActSchedulingFacts {
178 pub blueprint_id: Option<String>,
179 pub workspace_id: Option<WorkspaceId>,
180}
181
182pub enum ActivityOutcome {
188 ProcessInput { turn_id: Option<TurnId> },
189 Reason(Box<ReasonResult>),
191 Act(ActOutcome),
192}
193
194#[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 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 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 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
295pub 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
307pub 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
343pub 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
360pub 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 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 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 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
508pub 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
530pub 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 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 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}