Skip to main content

lash_core/runtime/
turn_loop.rs

1#[cfg(test)]
2use super::logical_turn::agent_frame_follow_turn_id;
3use super::logical_turn::{
4    LogicalTurnClaims, LogicalTurnStart, PhysicalTurnExecution, PreparedLogicalTurn,
5};
6use super::turn_control::ActiveTurnControl;
7use super::*;
8
9fn trace_fields_from_outcome(
10    outcome: &TurnOutcome,
11) -> (
12    &'static str,
13    &'static str,
14    Option<lash_trace::TraceAgentFrameSwitch>,
15) {
16    match outcome {
17        TurnOutcome::Finished(TurnFinish::AssistantMessage { .. }) => {
18            ("completed", "assistant_message", None)
19        }
20        TurnOutcome::Finished(TurnFinish::FinalValue { .. }) => ("completed", "final_value", None),
21        TurnOutcome::Finished(TurnFinish::ToolValue { .. }) => ("completed", "tool_value", None),
22        TurnOutcome::AgentFrameSwitch { frame_id, .. } => (
23            "completed",
24            "agent_frame_switch",
25            Some(lash_trace::TraceAgentFrameSwitch {
26                frame_id: frame_id.clone(),
27            }),
28        ),
29        TurnOutcome::Stopped(stop) => ("failed", trace_stop_reason(stop), None),
30    }
31}
32
33fn trace_stop_reason(stop: &TurnStop) -> &'static str {
34    match stop {
35        TurnStop::Cancelled => "cancelled",
36        TurnStop::Incomplete => "incomplete",
37        TurnStop::InvalidInput => "invalid_input",
38        TurnStop::MaxTurns => "max_turns",
39        TurnStop::ToolFailure => "tool_failure",
40        TurnStop::ProviderError => "provider_error",
41        TurnStop::PluginAbort => "plugin_abort",
42        TurnStop::RuntimeError => "runtime_error",
43        TurnStop::SubmittedError { .. } => "submitted_error",
44        TurnStop::ToolError { .. } => "tool_error",
45    }
46}
47
48fn session_head_refresh_error(err: SessionError) -> RuntimeError {
49    RuntimeError::new(
50        RuntimeErrorCode::Other("session_head_refresh".to_string()),
51        err.to_string(),
52    )
53}
54
55#[derive(Clone, Copy)]
56pub(super) enum SessionExecutionLeaseReleasePolicy {
57    KeepOnAgentFrameSwitch,
58}
59
60impl SessionExecutionLeaseReleasePolicy {
61    fn should_release(self, outcome: &TurnOutcome) -> bool {
62        match self {
63            Self::KeepOnAgentFrameSwitch => {
64                !matches!(outcome, TurnOutcome::AgentFrameSwitch { .. })
65            }
66        }
67    }
68}
69
70fn queued_work_payload_type(payload: &crate::QueuedWorkPayload) -> &'static str {
71    match payload {
72        crate::QueuedWorkPayload::ProcessWake { .. } => "process_wake",
73        crate::QueuedWorkPayload::AgentFrameTask { .. } => "agent_frame_task",
74        crate::QueuedWorkPayload::SessionCommand { command } => command.kind(),
75    }
76}
77
78fn queued_work_batch_ids(claim: &crate::QueuedWorkClaim) -> Vec<String> {
79    claim
80        .batches
81        .iter()
82        .map(|batch| batch.batch_id.clone())
83        .collect()
84}
85
86/// Measures the whole host-visible turn.
87///
88/// Opened before the runtime claims the turn (session-execution lease and
89/// queued-work/turn-input claims) and stamped onto the assembled turn after
90/// the final commit and post-persist hooks complete, so
91/// [`ExecutionSummary`](crate::ExecutionSummary) timing covers
92/// claim → final commit. Reads only the injected [`Clock`](crate::Clock):
93/// `started_at_ms` comes from the wall-clock source and the duration from the
94/// monotonic source, so deterministic clocks produce deterministic timing.
95#[derive(Clone, Copy)]
96pub(super) struct TurnStopwatch {
97    started: std::time::Instant,
98    started_at_ms: u64,
99}
100
101impl TurnStopwatch {
102    pub(super) fn start(clock: &dyn crate::Clock) -> Self {
103        Self {
104            started: clock.now(),
105            started_at_ms: clock.timestamp_ms(),
106        }
107    }
108
109    pub(super) fn stamp(&self, turn: &mut AssembledTurn, clock: &dyn crate::Clock) {
110        turn.execution.started_at_ms = self.started_at_ms;
111        turn.execution.duration_ms = clock
112            .now()
113            .saturating_duration_since(self.started)
114            .as_millis() as u64;
115    }
116}
117
118fn turn_phase_id(parent_turn_id: &str, phase: &str) -> String {
119    format!("{parent_turn_id}:{phase}")
120}
121
122fn scoped_child_turn_controller<'run>(
123    scoped_effect_controller: &'run ScopedEffectController<'_>,
124    session_id: &str,
125    turn_id: &str,
126) -> Result<ScopedEffectController<'run>, RuntimeError> {
127    ScopedEffectController::borrowed(
128        scoped_effect_controller.controller(),
129        ExecutionScope::turn(session_id, turn_id),
130    )
131}
132
133pub(in crate::runtime) fn queued_work_trace_payload(
134    boundary: crate::QueuedWorkClaimBoundary,
135    claim: &crate::QueuedWorkClaim,
136    causes: &[crate::TurnCause],
137) -> serde_json::Value {
138    serde_json::json!({
139        "boundary": boundary,
140        "claim_id": claim.claim_id,
141        "owner_id": claim.owner.owner_id,
142        "incarnation_id": claim.owner.incarnation_id,
143        "batch_ids": queued_work_batch_ids(claim),
144        "payload_types": claim.batches.iter()
145            .flat_map(|batch| batch.items.iter())
146            .map(|item| queued_work_payload_type(&item.payload))
147            .collect::<Vec<_>>(),
148        "causes": causes,
149    })
150}
151
152pub(in crate::runtime) fn queued_work_completion_trace_payload(
153    completions: &[crate::QueuedWorkCompletion],
154) -> serde_json::Value {
155    serde_json::json!({
156        "claims": completions.iter().map(|completion| {
157            serde_json::json!({
158                "session_id": completion.session_id,
159                "claim_id": completion.claim_id,
160                "batch_ids": completion.batch_ids,
161            })
162        }).collect::<Vec<_>>(),
163    })
164}
165
166pub(in crate::runtime) fn turn_input_completion_trace_payload(
167    completions: &[crate::TurnInputCompletion],
168) -> serde_json::Value {
169    serde_json::json!({
170        "claims": completions.iter().map(|completion| {
171            serde_json::json!({
172                "session_id": completion.session_id,
173                "claim_id": completion.claim_id,
174                "input_ids": completion.input_ids,
175            })
176        }).collect::<Vec<_>>(),
177    })
178}
179
180async fn emit_queued_work_started_to_sink(
181    events: &dyn TurnActivitySink,
182    boundary: crate::QueuedWorkClaimBoundary,
183    claim: &crate::QueuedWorkClaim,
184    causes: Vec<crate::TurnCause>,
185) {
186    emit_turn_activity_to_sink(
187        events,
188        TurnActivity::independent(TurnEvent::QueuedWorkStarted {
189            boundary,
190            batch_ids: queued_work_batch_ids(claim),
191            causes,
192        }),
193    )
194    .await;
195}
196
197pub(in crate::runtime) async fn send_queued_work_started_event(
198    event_tx: &mpsc::Sender<RuntimeStreamEvent>,
199    boundary: crate::QueuedWorkClaimBoundary,
200    claim: &crate::QueuedWorkClaim,
201    causes: Vec<crate::TurnCause>,
202) {
203    send_turn_activity(
204        event_tx,
205        TurnActivityId::fresh(),
206        TurnEvent::QueuedWorkStarted {
207            boundary,
208            batch_ids: queued_work_batch_ids(claim),
209            causes,
210        },
211    )
212    .await;
213}
214
215struct TurnFinishInput {
216    turn_pipeline: TurnBoundary,
217    assembler: TurnAssembler,
218    new_messages: crate::MessageSequence,
219    policy: RuntimeSessionPolicy,
220    turn_index: usize,
221    queued_work_claims: Vec<crate::QueuedWorkClaim>,
222    turn_input_claims: Vec<crate::TurnInputClaim>,
223    trace_turn_id: String,
224}
225
226impl LashRuntime {
227    fn max_context_tokens(&self) -> usize {
228        self.state.effective_policy().context_window_tokens()
229    }
230
231    async fn claim_session_execution_lease(
232        &self,
233        cancel: CancellationToken,
234        busy_is_error: bool,
235    ) -> Result<Option<SessionExecutionLeaseGuard>, RuntimeError> {
236        let Some(store) = self
237            .session
238            .as_ref()
239            .and_then(|session| session.history_store())
240        else {
241            return Ok(None);
242        };
243        match SessionExecutionLeaseGuard::try_acquire(
244            store,
245            &self.state.session_id,
246            &self.runtime_lease_owner,
247            self.host.core.control.lease_timings,
248            Arc::clone(&self.host.core.clock),
249            cancel,
250        )
251        .await
252        .map_err(|err| RuntimeError::new(RuntimeErrorCode::StoreCommitFailed, err.to_string()))?
253        {
254            Some(lease) => Ok(Some(lease)),
255            None if busy_is_error => Err(RuntimeError::new(
256                RuntimeErrorCode::SessionExecutionBusy,
257                format!(
258                    "session `{}` is already executing on another runtime owner",
259                    self.state.session_id
260                ),
261            )),
262            None => Ok(None),
263        }
264    }
265
266    async fn settle_session_execution_lease<T>(
267        &self,
268        guard: Option<&SessionExecutionLeaseGuard>,
269        result: Result<T, RuntimeError>,
270    ) -> Result<T, RuntimeError> {
271        match result {
272            Ok(value) => {
273                if let Some(guard) = guard {
274                    guard.release_if_live().await.map_err(|err| {
275                        RuntimeError::new(RuntimeErrorCode::StoreCommitFailed, err.to_string())
276                    })?;
277                }
278                Ok(value)
279            }
280            Err(err) => {
281                if err.code != RuntimeErrorCode::StoreCommitFailed
282                    && let Some(guard) = guard
283                    && let Err(release_err) = guard.release_if_live().await
284                {
285                    tracing::warn!(
286                        error = %release_err,
287                        "failed to release session execution lease after runtime error"
288                    );
289                }
290                Err(err)
291            }
292        }
293    }
294
295    async fn ensure_session_execution_lease_live(
296        &self,
297        guard: Option<&SessionExecutionLeaseGuard>,
298    ) -> Result<(), RuntimeError> {
299        let Some(guard) = guard else {
300            return Ok(());
301        };
302        guard.refresh_or_mark_lost().await.map_err(|err| {
303            RuntimeError::new(
304                RuntimeErrorCode::SessionExecutionLeaseLost,
305                format!(
306                    "session execution lease for session `{}` was lost before commit: {err}",
307                    self.state.session_id
308                ),
309            )
310        })
311    }
312
313    // Prompt handback on lease loss. This is no longer load-bearing for
314    // correctness: a claim is generation-fenced under the session lease, so once
315    // this owner has lost the lease its claims are already superseded and the
316    // next acquirer reclaims them by generation regardless (ADR 0029). Abandoning
317    // eagerly just lets a peer reclaim the rows without waiting to observe the
318    // generation bump.
319    async fn abandon_queued_work_claims_after_lease_loss(
320        &self,
321        err: &RuntimeError,
322        claims: &[crate::QueuedWorkClaim],
323    ) {
324        if err.code != RuntimeErrorCode::SessionExecutionLeaseLost || claims.is_empty() {
325            return;
326        }
327        let Some(store) = self
328            .session
329            .as_ref()
330            .and_then(|session| session.history_store())
331        else {
332            return;
333        };
334        for claim in claims {
335            if let Err(abandon_err) = store.abandon_queued_work_claim(claim).await {
336                tracing::warn!(
337                    error = %abandon_err,
338                    session_id = %claim.session_id,
339                    claim_id = %claim.claim_id,
340                    "failed to abandon queued work claim after session execution lease loss"
341                );
342            }
343        }
344    }
345
346    async fn abandon_turn_input_claims_after_lease_loss(
347        &self,
348        err: &RuntimeError,
349        claims: &[crate::TurnInputClaim],
350    ) {
351        if err.code != RuntimeErrorCode::SessionExecutionLeaseLost || claims.is_empty() {
352            return;
353        }
354        let Some(store) = self
355            .session
356            .as_ref()
357            .and_then(|session| session.history_store())
358        else {
359            return;
360        };
361        for claim in claims {
362            if let Err(abandon_err) = store.abandon_turn_input_claim(claim).await {
363                tracing::warn!(
364                    error = %abandon_err,
365                    session_id = %claim.session_id,
366                    claim_id = %claim.claim_id,
367                    "failed to abandon turn input claim after session execution lease loss"
368                );
369            }
370        }
371    }
372
373    #[doc(hidden)]
374    pub fn set_turn_phase_probe(&mut self, probe: Arc<dyn RuntimeTurnPhaseProbe>) {
375        self.turn_phase_probe = Some(probe);
376    }
377
378    fn mark_phase_begin(&self, phase: RuntimeTurnPhase) {
379        if let Some(probe) = self.turn_phase_probe.as_ref() {
380            probe.begin(phase);
381        }
382    }
383
384    fn mark_phase_end(&self, phase: RuntimeTurnPhase) {
385        if let Some(probe) = self.turn_phase_probe.as_ref() {
386            probe.end(phase);
387        }
388    }
389
390    #[allow(clippy::too_many_arguments)]
391    async fn finish_turn(
392        &mut self,
393        finish: TurnFinishInput,
394        events: &dyn EventSink,
395        scoped_effect_controller: &ScopedEffectController<'_>,
396        cancel_state: &CancellationToken,
397        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
398        session_execution_lease_release_policy: SessionExecutionLeaseReleasePolicy,
399        turn_control: &ActiveTurnControl,
400    ) -> Result<PhysicalTurnExecution, RuntimeError> {
401        let TurnFinishInput {
402            mut turn_pipeline,
403            assembler,
404            new_messages,
405            policy,
406            turn_index,
407            queued_work_claims,
408            turn_input_claims,
409            trace_turn_id,
410        } = finish;
411        self.policy = self.state.effective_policy().clone();
412        turn_pipeline.state_mut().policy = self.policy.clone();
413        turn_pipeline.state_mut().turn_index = turn_index;
414
415        let mut turn_usage_delta = {
416            let mut ledger = self.shared_token_ledger.lock().expect("token ledger lock");
417            std::mem::take(&mut *ledger)
418        };
419        if assembler.token_usage.total() > 0 {
420            turn_usage_delta.push(TokenLedgerEntry {
421                source: "turn".to_string(),
422                model: policy.model.id.clone(),
423                usage: assembler.token_usage.clone(),
424            });
425        }
426        let turn_usage_delta = merge_usage_delta_entries(turn_usage_delta);
427
428        if self.session.is_some() {
429            self.ensure_session_execution_lease_live(session_execution_lease)
430                .await?;
431        }
432        let assembled_cancelled = matches!(
433            assembler.outcome,
434            Some(TurnOutcome::Stopped(TurnStop::Cancelled))
435        );
436        let lease_was_lost = session_execution_lease.is_some_and(|lease| lease.is_lost());
437        let cancellation = turn_control
438            .settle_before_commit(
439                scoped_effect_controller.controller(),
440                assembled_cancelled || (cancel_state.is_cancelled() && !lease_was_lost),
441            )
442            .await?;
443        if cancellation.is_some() {
444            cancel_state.cancel();
445        }
446        let interrupted = cancel_state.is_cancelled();
447
448        turn_pipeline.finalize_turn_read_state(new_messages, interrupted);
449        if assembler.token_usage.total() > 0 {
450            turn_pipeline.state_mut().token_usage = assembler.token_usage.clone();
451        }
452
453        let last_prompt_usage = assembler.last_llm_usage().and_then(normalize_prompt_usage);
454        turn_pipeline.state_mut().last_prompt_usage = last_prompt_usage;
455        let assembled_state = turn_pipeline.export_state_for_assembly();
456        let mut assembled = assembler.finish(
457            assembled_state,
458            interrupted,
459            None,
460            &self.host.core.control.termination,
461        );
462        assembled.cancellation = cancellation;
463
464        let Some(session) = self.session.as_ref() else {
465            self.state.apply_snapshot(&assembled.state);
466            self.emit_completed_turn_trace(&assembled.state, &assembled.outcome, &trace_turn_id);
467            publish_terminal_after_commit(
468                turn_control,
469                scoped_effect_controller.controller(),
470                &TurnTerminal::Committed {
471                    outcome: assembled.outcome.clone(),
472                    cancellation: assembled.cancellation.clone(),
473                    session_revision: None,
474                },
475                &self.state.session_id,
476                &trace_turn_id,
477            )
478            .await;
479            return Ok(PhysicalTurnExecution {
480                turn: assembled,
481                enqueued_queue_batches: Vec::new(),
482            });
483        };
484
485        let plugins = Arc::clone(session.plugins());
486        let manager = match self.runtime_session_services_for_turn(None) {
487            Ok(manager) => manager,
488            Err(err) => {
489                return Err(RuntimeError::new(
490                    RuntimeErrorCode::PluginSessionManager,
491                    err.to_string(),
492                ));
493            }
494        };
495
496        self.mark_phase_begin(RuntimeTurnPhase::FinalizeTurn);
497        let finalized = match plugins
498            .finalize_turn_with_phase_probe(
499                assembled,
500                manager.state_service(),
501                manager.lifecycle_service(),
502                manager.graph_service(),
503                self.turn_phase_probe.clone(),
504            )
505            .await
506        {
507            Ok(finalized) => finalized,
508            Err(err) => {
509                self.mark_phase_end(RuntimeTurnPhase::FinalizeTurn);
510                return Err(RuntimeError::new(
511                    RuntimeErrorCode::PluginFinalizeTurn,
512                    err.to_string(),
513                ));
514            }
515        };
516        self.mark_phase_end(RuntimeTurnPhase::FinalizeTurn);
517        self.ensure_session_execution_lease_live(session_execution_lease)
518            .await?;
519
520        let mut returned_turn = finalized.turn;
521        if returned_turn.cancellation.is_some()
522            && !matches!(
523                returned_turn.outcome,
524                TurnOutcome::Stopped(TurnStop::Cancelled)
525            )
526        {
527            returned_turn.outcome = TurnOutcome::Stopped(TurnStop::Cancelled);
528        }
529        if matches!(
530            returned_turn.outcome,
531            TurnOutcome::Stopped(TurnStop::Cancelled)
532        ) && returned_turn.cancellation.is_none()
533        {
534            return Err(RuntimeError::new(
535                "turn_cancellation_evidence_missing",
536                "cancelled turns must carry cancellation evidence",
537            ));
538        }
539        let release_session_execution_lease =
540            session_execution_lease_release_policy.should_release(&returned_turn.outcome);
541        let commit_effects = LogicalTurnClaims::new(queued_work_claims, turn_input_claims)
542            .into_commit_effects(
543                &returned_turn.outcome,
544                &self.state.session_id,
545                &trace_turn_id,
546                Some(self.state.effective_protocol_turn_options().clone()),
547            );
548        self.mark_phase_begin(RuntimeTurnPhase::PersistTurn);
549        self.mark_phase_begin(RuntimeTurnPhase::FinalCommit);
550        let queued_work_completion_trace = commit_effects.completed_queue_claims.clone();
551        let turn_input_completion_trace = commit_effects.completed_turn_input_claims.clone();
552        let pending_attachment_ids = self
553            .host
554            .core
555            .durability
556            .attachment_store
557            .pending_manifest_commit_ids();
558        let enqueued_queue_batches = match turn_pipeline
559            .final_commit(
560                &mut returned_turn,
561                self.session.as_mut(),
562                &turn_usage_delta,
563                Some(&trace_turn_id),
564                commit_effects.originating_queue_claims,
565                commit_effects.originating_turn_input_claims,
566                commit_effects.completed_queue_claims,
567                commit_effects.completed_turn_input_claims,
568                commit_effects.enqueued_queue_batches,
569                cancel_state.is_cancelled().then(|| trace_turn_id.clone()),
570                pending_attachment_ids.clone(),
571                release_session_execution_lease
572                    .then(|| session_execution_lease.map(SessionExecutionLeaseGuard::completion))
573                    .flatten(),
574            )
575            .await
576        {
577            Ok(batches) => batches,
578            Err(err) => {
579                self.mark_phase_end(RuntimeTurnPhase::FinalCommit);
580                self.mark_phase_end(RuntimeTurnPhase::PersistTurn);
581                return Err(err);
582            }
583        };
584        if release_session_execution_lease && let Some(lease) = session_execution_lease {
585            lease.mark_released();
586        }
587        self.host
588            .core
589            .durability
590            .attachment_store
591            .mark_manifest_committed(&pending_attachment_ids);
592        self.mark_phase_end(RuntimeTurnPhase::FinalCommit);
593
594        emit_session_events_to_sink(events, finalized.events).await;
595        self.state = turn_pipeline.into_final_state();
596        publish_terminal_after_commit(
597            turn_control,
598            scoped_effect_controller.controller(),
599            &TurnTerminal::Committed {
600                outcome: returned_turn.outcome.clone(),
601                cancellation: returned_turn.cancellation.clone(),
602                session_revision: None,
603            },
604            &self.state.session_id,
605            &trace_turn_id,
606        )
607        .await;
608        if matches!(returned_turn.outcome, TurnOutcome::AgentFrameSwitch { .. })
609            && let Some(session) = self.session.as_mut()
610        {
611            let protocol_session = Arc::clone(session.plugins().protocol_session());
612            let session_id = self.state.session_id.clone();
613            protocol_session
614                .restore_session(
615                    crate::plugin::ProtocolSessionContext::new(session, &session_id),
616                    &self.state,
617                )
618                .await
619                .map_err(|err| {
620                    RuntimeError::new(
621                        RuntimeErrorCode::Other("protocol_restore_session".to_string()),
622                        err.to_string(),
623                    )
624                })?;
625        }
626        if !queued_work_completion_trace.is_empty() {
627            crate::trace::emit_trace(
628                &self.host.core.tracing.trace_sink,
629                &self.host.core.tracing.trace_context,
630                lash_trace::TraceContext::default()
631                    .for_session(returned_turn.state.session_id.clone())
632                    .for_turn_index(returned_turn.state.turn_index)
633                    .for_turn(trace_turn_id.clone()),
634                lash_trace::TraceEvent::Custom {
635                    name: "queued_work.completed".to_string(),
636                    payload: queued_work_completion_trace_payload(&queued_work_completion_trace),
637                },
638                self.host.core.clock.as_ref(),
639            );
640        }
641        if !turn_input_completion_trace.is_empty() {
642            crate::trace::emit_trace(
643                &self.host.core.tracing.trace_sink,
644                &self.host.core.tracing.trace_context,
645                lash_trace::TraceContext::default()
646                    .for_session(returned_turn.state.session_id.clone())
647                    .for_turn_index(returned_turn.state.turn_index)
648                    .for_turn(trace_turn_id.clone()),
649                lash_trace::TraceEvent::Custom {
650                    name: "turn_input.completed".to_string(),
651                    payload: turn_input_completion_trace_payload(&turn_input_completion_trace),
652                },
653                self.host.core.clock.as_ref(),
654            );
655        }
656        self.mark_phase_begin(RuntimeTurnPhase::PostPersistHooks);
657        self.emit_turn_persisted_event(&returned_turn, scoped_effect_controller, &trace_turn_id)
658            .await?;
659        self.mark_phase_end(RuntimeTurnPhase::PostPersistHooks);
660        self.mark_phase_end(RuntimeTurnPhase::PersistTurn);
661
662        self.emit_completed_turn_trace(
663            &returned_turn.state,
664            &returned_turn.outcome,
665            &trace_turn_id,
666        );
667        Ok(PhysicalTurnExecution {
668            turn: returned_turn,
669            enqueued_queue_batches,
670        })
671    }
672
673    #[allow(clippy::too_many_arguments)]
674    async fn finish_cancelled_turn_after_effect_abort(
675        &mut self,
676        driver: RuntimeTurnDriver<'_>,
677        mut assembler: TurnAssembler,
678        cancellation_messages: crate::MessageSequence,
679        events: &dyn EventSink,
680        finish_scoped_effect_controller: &ScopedEffectController<'_>,
681        cancel: &CancellationToken,
682        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
683        session_execution_lease_release_policy: SessionExecutionLeaseReleasePolicy,
684        turn_control: &ActiveTurnControl,
685        turn_index: usize,
686        trace_turn_id: String,
687    ) -> Result<PhysicalTurnExecution, RuntimeError> {
688        let RuntimeTurnDriver {
689            session,
690            policy,
691            turn_pipeline,
692            pending_queue_claims,
693            pending_turn_input_claims,
694            ..
695        } = driver;
696        self.session = Some(session);
697        let outcome_event = SessionStreamEvent::TurnOutcome {
698            outcome: TurnOutcome::Stopped(TurnStop::Cancelled),
699        };
700        assembler.push(&outcome_event);
701        emit_session_event_to_sink(events, outcome_event).await;
702        assembler.push(&SessionStreamEvent::Done);
703        emit_session_event_to_sink(events, SessionStreamEvent::Done).await;
704        self.finish_turn(
705            TurnFinishInput {
706                turn_pipeline,
707                assembler,
708                new_messages: cancellation_messages,
709                policy,
710                turn_index,
711                queued_work_claims: pending_queue_claims,
712                turn_input_claims: pending_turn_input_claims,
713                trace_turn_id,
714            },
715            events,
716            finish_scoped_effect_controller,
717            cancel,
718            session_execution_lease,
719            session_execution_lease_release_policy,
720            turn_control,
721        )
722        .await
723    }
724
725    fn emit_completed_turn_trace(
726        &self,
727        state: &SessionSnapshot,
728        outcome: &TurnOutcome,
729        trace_turn_id: &str,
730    ) {
731        if self.host.core.tracing.trace_sink.is_none() {
732            return;
733        }
734
735        let (status, done_reason, agent_frame_switch) = trace_fields_from_outcome(outcome);
736        crate::trace::emit_trace(
737            &self.host.core.tracing.trace_sink,
738            &self.host.core.tracing.trace_context,
739            lash_trace::TraceContext::default()
740                .for_session(state.session_id.clone())
741                .for_turn_index(state.turn_index)
742                .for_turn(trace_turn_id.to_string()),
743            lash_trace::TraceEvent::TurnCompleted {
744                status: status.to_string(),
745                done_reason: done_reason.to_string(),
746                agent_frame_switch,
747            },
748            self.host.core.clock.as_ref(),
749        );
750    }
751
752    #[allow(clippy::too_many_arguments)]
753    pub(super) async fn finish_logical_turn_error(
754        &mut self,
755        message: String,
756        trace_turn_id: String,
757        events: &dyn EventSink,
758        turn_events: &dyn TurnActivitySink,
759        scoped_effect_controller: ScopedEffectController<'_>,
760        cancel: CancellationToken,
761        claims: LogicalTurnClaims,
762        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
763    ) -> Result<PhysicalTurnExecution, RuntimeError> {
764        let turn_control = Arc::new(
765            ActiveTurnControl::new(
766                scoped_effect_controller.controller(),
767                TurnAddress::new(&self.state.session_id, &trace_turn_id),
768            )
769            .await?,
770        );
771        let mut assembler = TurnAssembler::default();
772        let error_event = SessionStreamEvent::Error {
773            message: message.clone(),
774            envelope: Some(crate::session_model::ErrorEnvelope {
775                kind: "runtime".to_string(),
776                code: Some("agent_frame_switch_limit".to_string()),
777                terminal_reason: None,
778                user_message: message.clone(),
779                raw: None,
780                retryable: Some(false),
781                provider_failure_kind: None,
782            }),
783        };
784        assembler.push(&error_event);
785        emit_turn_activity_to_sink(
786            turn_events,
787            TurnActivity::independent(TurnEvent::Error {
788                message: message.clone(),
789            }),
790        )
791        .await;
792        emit_session_event_to_sink(events, error_event).await;
793        let outcome_event = SessionStreamEvent::TurnOutcome {
794            outcome: TurnOutcome::Stopped(TurnStop::RuntimeError),
795        };
796        assembler.push(&outcome_event);
797        emit_session_event_to_sink(events, outcome_event).await;
798        assembler.push(&SessionStreamEvent::Done);
799        emit_session_event_to_sink(events, SessionStreamEvent::Done).await;
800
801        let messages = crate::MessageSequence::from_base(self.state.read_model().messages);
802        let mut turn_pipeline = TurnBoundary::from_state_with_clock(
803            self.state.clone(),
804            Arc::clone(&self.host.core.clock),
805        )
806        .with_session_execution_lease(
807            session_execution_lease.map(SessionExecutionLeaseGuard::fence),
808        );
809        turn_pipeline.apply_prepared_messages(&messages);
810        self.finish_turn(
811            TurnFinishInput {
812                turn_pipeline,
813                assembler,
814                new_messages: messages,
815                policy: RuntimeSessionPolicy::new(
816                    self.state.effective_policy().clone(),
817                    Default::default(),
818                ),
819                turn_index: self.state.turn_index + 1,
820                queued_work_claims: claims.queued,
821                turn_input_claims: claims.turn_inputs,
822                trace_turn_id,
823            },
824            events,
825            &scoped_effect_controller,
826            &cancel,
827            session_execution_lease,
828            SessionExecutionLeaseReleasePolicy::KeepOnAgentFrameSwitch,
829            &turn_control,
830        )
831        .await
832    }
833
834    async fn emit_turn_persisted_event(
835        &self,
836        returned_turn: &AssembledTurn,
837        scoped_effect_controller: &ScopedEffectController<'_>,
838        trace_turn_id: &str,
839    ) -> Result<(), RuntimeError> {
840        let Some(session) = self.session.as_ref() else {
841            return Ok(());
842        };
843        let Ok(manager) = self.runtime_session_services() else {
844            return Ok(());
845        };
846        let phase_turn_id = turn_phase_id(trace_turn_id, "turn-persisted");
847        let phase_controller = scoped_child_turn_controller(
848            scoped_effect_controller,
849            &self.state.session_id,
850            &phase_turn_id,
851        )?;
852        let direct_completions = manager.direct_completion_client(
853            RuntimeEffectControllerHandle::borrowed(phase_controller),
854            Some(phase_turn_id),
855        );
856
857        session
858            .plugins()
859            .emit_runtime_event_with_phase_probe(
860                crate::PluginLifecycleEvent::TurnPersisted(Box::new(
861                    crate::SessionStateChangedContext {
862                        session_id: self.state.session_id.clone(),
863                        state: crate::SessionReadView::from_snapshot(&returned_turn.state),
864                        sessions: manager.state_service(),
865                        session_graph: manager.graph_service(),
866                        direct_completions,
867                    },
868                )),
869                self.turn_phase_probe.clone(),
870            )
871            .await;
872        Ok(())
873    }
874
875    /// Run one logical turn and stream every physical frame to the host sink.
876    pub async fn stream_turn(
877        &mut self,
878        mut input: TurnInput,
879        opts: TurnOptions<'_>,
880    ) -> Result<AssembledTurn, RuntimeError> {
881        if let Some(hint) = opts.local_cancel_origin_hint() {
882            input.turn_context.set_local_cancel_origin_hint(hint);
883        }
884        let stopwatch = TurnStopwatch::start(self.host.core.clock.as_ref());
885        let cancel = opts.cancel.clone();
886        let session_execution_lease = self
887            .claim_session_execution_lease(cancel.clone(), true)
888            .await?;
889        let scoped_effect_controller = opts.scoped_effect_controller();
890        let result = Box::pin(self.drive_logical_turn(
891            LogicalTurnStart::Input(input),
892            opts.events_or_noop(),
893            opts.turn_events_or_noop(),
894            scoped_effect_controller,
895            cancel,
896            LogicalTurnClaims::new(Vec::new(), Vec::new()),
897            session_execution_lease.as_ref(),
898            stopwatch,
899        ))
900        .await
901        .map(|run| {
902            run.into_final_turn()
903                .expect("logical turn always contains a terminal physical turn")
904        });
905        self.settle_session_execution_lease(session_execution_lease.as_ref(), result)
906            .await
907    }
908
909    pub async fn stream_next_queued_work(
910        &mut self,
911        opts: TurnOptions<'_>,
912    ) -> Result<Option<AssembledTurn>, RuntimeError> {
913        self.stream_queued_work(opts, None).await
914    }
915
916    pub async fn stream_selected_queued_work(
917        &mut self,
918        opts: TurnOptions<'_>,
919        batch_ids: &[String],
920    ) -> Result<Option<AssembledTurn>, RuntimeError> {
921        self.stream_queued_work(opts, Some(batch_ids)).await
922    }
923
924    async fn stream_queued_work(
925        &mut self,
926        opts: TurnOptions<'_>,
927        selected_batch_ids: Option<&[String]>,
928    ) -> Result<Option<AssembledTurn>, RuntimeError> {
929        let stopwatch = TurnStopwatch::start(self.host.core.clock.as_ref());
930        let cancel = opts.cancel.clone();
931        let Some(session_execution_lease) = self
932            .claim_session_execution_lease(cancel.clone(), false)
933            .await?
934        else {
935            return Ok(None);
936        };
937        let session_execution_fence = session_execution_lease.fence();
938        let Some(store) = self
939            .session
940            .as_ref()
941            .and_then(|session| session.history_store())
942        else {
943            session_execution_lease
944                .release_if_live()
945                .await
946                .map_err(|err| {
947                    RuntimeError::new(RuntimeErrorCode::StoreCommitFailed, err.to_string())
948                })?;
949            return Ok(None);
950        };
951        let drain_commands_before_turn_input = if selected_batch_ids.is_some() {
952            true
953        } else {
954            self.session_commands_precede_pending_turn_input(store.as_ref())
955                .await?
956        };
957        if drain_commands_before_turn_input {
958            loop {
959                match self
960                    .drain_next_session_command(&session_execution_fence)
961                    .await
962                {
963                    Ok(Some(_)) => {}
964                    Ok(None) => break,
965                    Err(err) => {
966                        let _ = session_execution_lease.release_if_live().await;
967                        return Err(err);
968                    }
969                }
970            }
971        }
972        if selected_batch_ids.is_none() {
973            let input_claim = store
974                .claim_next_turn_inputs(
975                    &self.state.session_id,
976                    &session_execution_fence,
977                    &self.runtime_lease_owner,
978                    64,
979                )
980                .await
981                .map_err(super::runtime_error_from_store_commit)?;
982            if let Some(input_claim) = input_claim {
983                let mut input = input_claim.materialize_for_turn();
984                if let Some(hint) = opts.local_cancel_origin_hint() {
985                    input.turn_context.set_local_cancel_origin_hint(hint);
986                }
987                let turn_id = input
988                    .trace_turn_id
989                    .clone()
990                    .or_else(|| Some(opts.execution_scope_id().to_owned()))
991                    .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
992                input.trace_turn_id = Some(turn_id.clone());
993                crate::trace::emit_trace(
994                    &self.host.core.tracing.trace_sink,
995                    &self.host.core.tracing.trace_context,
996                    lash_trace::TraceContext::default()
997                        .for_session(self.state.session_id.clone())
998                        .for_turn_index(self.state.turn_index + 1)
999                        .for_turn(turn_id.clone()),
1000                    lash_trace::TraceEvent::Custom {
1001                        name: "turn_input.claimed".to_string(),
1002                        payload: serde_json::json!({
1003                            "claim_id": &input_claim.claim_id,
1004                            "input_ids": input_claim.inputs.iter().map(|input| input.input_id.clone()).collect::<Vec<_>>(),
1005                        }),
1006                    },
1007                    self.host.core.clock.as_ref(),
1008                );
1009                let claim_for_abandon = input_claim.clone();
1010                let scoped_effect_controller = opts.scoped_effect_controller();
1011                let result = Box::pin(self.drive_logical_turn(
1012                    LogicalTurnStart::Input(input),
1013                    opts.events_or_noop(),
1014                    opts.turn_events_or_noop(),
1015                    scoped_effect_controller,
1016                    cancel,
1017                    LogicalTurnClaims::new(Vec::new(), vec![input_claim]),
1018                    Some(&session_execution_lease),
1019                    stopwatch,
1020                ))
1021                .await
1022                .map(AgentFrameRun::into_final_turn);
1023                if let Err(err) = &result {
1024                    self.abandon_turn_input_claims_after_lease_loss(
1025                        err,
1026                        std::slice::from_ref(&claim_for_abandon),
1027                    )
1028                    .await;
1029                }
1030                return self
1031                    .settle_session_execution_lease(Some(&session_execution_lease), result)
1032                    .await;
1033            }
1034        }
1035        let claim = if let Some(batch_ids) = selected_batch_ids {
1036            store
1037                .claim_ready_queued_work_by_batch_ids(
1038                    &self.state.session_id,
1039                    &session_execution_fence,
1040                    &self.runtime_lease_owner,
1041                    crate::QueuedWorkClaimBoundary::Idle,
1042                    batch_ids,
1043                )
1044                .await
1045        } else {
1046            store
1047                .claim_ready_queued_work(
1048                    &self.state.session_id,
1049                    &session_execution_fence,
1050                    &self.runtime_lease_owner,
1051                    crate::QueuedWorkClaimBoundary::Idle,
1052                    64,
1053                )
1054                .await
1055        }
1056        .map_err(super::runtime_error_from_store_commit)?;
1057        let Some(claim) = claim else {
1058            session_execution_lease
1059                .release_if_live()
1060                .await
1061                .map_err(|err| {
1062                    RuntimeError::new(RuntimeErrorCode::StoreCommitFailed, err.to_string())
1063                })?;
1064            return Ok(None);
1065        };
1066        let mut work = claim.materialize_for_turn();
1067        if let Some(hint) = opts.local_cancel_origin_hint() {
1068            work.input.turn_context.set_local_cancel_origin_hint(hint);
1069        }
1070        let turn_id = work
1071            .input
1072            .trace_turn_id
1073            .clone()
1074            .or_else(|| Some(opts.execution_scope_id().to_owned()))
1075            .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1076        work.input.trace_turn_id = Some(turn_id.clone());
1077        let causes = work.turn_causes.clone();
1078        emit_queued_work_started_to_sink(
1079            opts.turn_events_or_noop(),
1080            crate::QueuedWorkClaimBoundary::Idle,
1081            &claim,
1082            causes.clone(),
1083        )
1084        .await;
1085        crate::trace::emit_trace(
1086            &self.host.core.tracing.trace_sink,
1087            &self.host.core.tracing.trace_context,
1088            lash_trace::TraceContext::default()
1089                .for_session(self.state.session_id.clone())
1090                .for_turn_index(self.state.turn_index + 1)
1091                .for_turn(turn_id.clone()),
1092            lash_trace::TraceEvent::Custom {
1093                name: "queued_work.claimed".to_string(),
1094                payload: queued_work_trace_payload(
1095                    crate::QueuedWorkClaimBoundary::Idle,
1096                    &claim,
1097                    &causes,
1098                ),
1099            },
1100            self.host.core.clock.as_ref(),
1101        );
1102        let claim_for_abandon = claim.clone();
1103        let scoped_effect_controller = opts.scoped_effect_controller();
1104        let result = Box::pin(self.drive_logical_turn(
1105            LogicalTurnStart::Input(work.input),
1106            opts.events_or_noop(),
1107            opts.turn_events_or_noop(),
1108            scoped_effect_controller,
1109            cancel,
1110            LogicalTurnClaims::new(vec![claim], Vec::new()),
1111            Some(&session_execution_lease),
1112            stopwatch,
1113        ))
1114        .await
1115        .map(AgentFrameRun::into_final_turn);
1116        if let Err(err) = &result {
1117            self.abandon_queued_work_claims_after_lease_loss(
1118                err,
1119                std::slice::from_ref(&claim_for_abandon),
1120            )
1121            .await;
1122        }
1123        self.settle_session_execution_lease(Some(&session_execution_lease), result)
1124            .await
1125    }
1126
1127    async fn session_commands_precede_pending_turn_input(
1128        &self,
1129        store: &dyn crate::RuntimePersistence,
1130    ) -> Result<bool, RuntimeError> {
1131        let pending_inputs = store
1132            .list_pending_turn_inputs(&self.state.session_id)
1133            .await
1134            .map_err(super::runtime_error_from_store_commit)?;
1135        let earliest_input = pending_inputs
1136            .iter()
1137            .filter(|input| input.state.is_next_turn_pending())
1138            .min_by_key(|input| (input.enqueued_at_ms, input.enqueue_seq));
1139        let queued_work = store
1140            .list_pending_queued_work(&self.state.session_id)
1141            .await
1142            .map_err(super::runtime_error_from_store_commit)?;
1143        let earliest_command = queued_work
1144            .iter()
1145            .filter(|batch| batch.is_session_command_work())
1146            .min_by_key(|batch| (batch.enqueued_at_ms, batch.enqueue_seq));
1147        Ok(match (earliest_command, earliest_input) {
1148            (Some(command), Some(input)) => command.enqueued_at_ms < input.enqueued_at_ms,
1149            (Some(_), None) => true,
1150            _ => false,
1151        })
1152    }
1153
1154    /// Enforce the durable-first wiring invariant at a turn-scope boundary: when
1155    /// the host wired a durable effect host, every store reachable from this
1156    /// scope must also be durable. A durable host running against any ephemeral
1157    /// store fails loudly here rather than silently degrading.
1158    ///
1159    /// Inline controllers (the default tier) impose no requirement, so
1160    /// inline/in-memory hosts pass unchanged.
1161    fn ensure_durable_store_facets_for_scope(
1162        &self,
1163        scoped_effect_controller: &ScopedEffectController<'_>,
1164    ) -> Result<(), RuntimeError> {
1165        if scoped_effect_controller.controller().durability_tier() != crate::DurabilityTier::Durable
1166        {
1167            return Ok(());
1168        }
1169        if self
1170            .host
1171            .core
1172            .durability
1173            .attachment_store
1174            .persistence()
1175            .durability_tier()
1176            != crate::DurabilityTier::Durable
1177        {
1178            return Err(RuntimeError::durable_store_required(
1179                crate::DurableStoreFacet::AttachmentStore,
1180            ));
1181        }
1182        if self
1183            .host
1184            .core
1185            .durability
1186            .process_env_store
1187            .durability_tier()
1188            != crate::DurabilityTier::Durable
1189        {
1190            return Err(RuntimeError::durable_store_required(
1191                crate::DurableStoreFacet::ProcessEnvStore,
1192            ));
1193        }
1194        if let Some(store) = self
1195            .session
1196            .as_ref()
1197            .and_then(|session| session.history_store())
1198            && store.durability_tier() != crate::DurabilityTier::Durable
1199        {
1200            return Err(RuntimeError::durable_store_required(
1201                crate::DurableStoreFacet::SessionStore,
1202            ));
1203        }
1204        if let Some(process_registry) = self.host.process_registry.as_ref()
1205            && process_registry.durability_tier() != crate::DurabilityTier::Durable
1206        {
1207            return Err(RuntimeError::durable_store_required(
1208                crate::DurableStoreFacet::ProcessRegistry,
1209            ));
1210        }
1211        Ok(())
1212    }
1213
1214    #[allow(clippy::too_many_arguments)]
1215    pub(super) async fn stream_turn_with_scoped_effect_controller_inner(
1216        &mut self,
1217        mut input: TurnInput,
1218        events: &dyn EventSink,
1219        turn_events: &dyn TurnActivitySink,
1220        scoped_effect_controller: ScopedEffectController<'_>,
1221        cancel: CancellationToken,
1222        queued_claims: Vec<crate::QueuedWorkClaim>,
1223        turn_input_claims: Vec<crate::TurnInputClaim>,
1224        materialize_initial_claims: bool,
1225        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
1226        session_execution_lease_release_policy: SessionExecutionLeaseReleasePolicy,
1227    ) -> Result<PhysicalTurnExecution, RuntimeError> {
1228        if queued_claims.is_empty() && turn_input_claims.is_empty() {
1229            if let Some(lease) = session_execution_lease {
1230                while self
1231                    .drain_next_session_command(&lease.fence())
1232                    .await?
1233                    .is_some()
1234                {}
1235            } else if self
1236                .session
1237                .as_ref()
1238                .and_then(|session| session.history_store())
1239                .is_some()
1240            {
1241                return Err(RuntimeError::new(
1242                    RuntimeErrorCode::StoreCommitFailed,
1243                    "session command drain requires a session execution lease",
1244                ));
1245            }
1246        }
1247        if let Some(input_turn_id) = input.trace_turn_id.as_deref()
1248            && scoped_effect_controller
1249                .execution_scope()
1250                .validates_turn_trace_id()
1251            && input_turn_id != scoped_effect_controller.scope_id()
1252        {
1253            return Err(RuntimeError::new(
1254                RuntimeErrorCode::ExecutionScopeTurnIdMismatch,
1255                format!(
1256                    "input trace_turn_id `{input_turn_id}` does not match execution scope id `{}`",
1257                    scoped_effect_controller.scope_id()
1258                ),
1259            ));
1260        }
1261        self.ensure_durable_store_facets_for_scope(&scoped_effect_controller)?;
1262        input
1263            .trace_turn_id
1264            .get_or_insert_with(|| scoped_effect_controller.scope_id().to_string());
1265        self.stream_turn_inner(
1266            input.clone(),
1267            events,
1268            turn_events,
1269            scoped_effect_controller,
1270            cancel.clone(),
1271            queued_claims,
1272            turn_input_claims,
1273            materialize_initial_claims,
1274            session_execution_lease,
1275            session_execution_lease_release_policy,
1276        )
1277        .await
1278    }
1279
1280    /// Stream one logical host turn, following foreground AgentFrame switches
1281    /// until a terminal outcome is reached.
1282    ///
1283    /// A protocol continuation creates a new frame in the same session. Hosts
1284    /// that only care about the benchmark/app answer should not need to
1285    /// special-case that intermediate outcome; this helper keeps driving the
1286    /// same session through each frame's task with the normal runtime turn
1287    /// guards.
1288    pub async fn stream_turn_with_agent_frames(
1289        &mut self,
1290        mut input: TurnInput,
1291        opts: TurnOptions<'_>,
1292    ) -> Result<AgentFrameRun, RuntimeError> {
1293        if let Some(hint) = opts.local_cancel_origin_hint() {
1294            input.turn_context.set_local_cancel_origin_hint(hint);
1295        }
1296        let stopwatch = TurnStopwatch::start(self.host.core.clock.as_ref());
1297        let cancel = opts.cancel.clone();
1298        let session_execution_lease = self
1299            .claim_session_execution_lease(cancel.clone(), true)
1300            .await?;
1301        let scoped_effect_controller = opts.scoped_effect_controller();
1302        let result = Box::pin(self.drive_logical_turn(
1303            LogicalTurnStart::Input(input),
1304            opts.events_or_noop(),
1305            opts.turn_events_or_noop(),
1306            scoped_effect_controller,
1307            cancel,
1308            LogicalTurnClaims::new(Vec::new(), Vec::new()),
1309            session_execution_lease.as_ref(),
1310            stopwatch,
1311        ))
1312        .await;
1313        self.settle_session_execution_lease(session_execution_lease.as_ref(), result)
1314            .await
1315    }
1316
1317    #[allow(clippy::too_many_arguments)]
1318    async fn stream_turn_inner(
1319        &mut self,
1320        mut input: TurnInput,
1321        events: &dyn EventSink,
1322        turn_events: &dyn TurnActivitySink,
1323        scoped_effect_controller: ScopedEffectController<'_>,
1324        cancel: CancellationToken,
1325        queued_claims: Vec<crate::QueuedWorkClaim>,
1326        turn_input_claims: Vec<crate::TurnInputClaim>,
1327        materialize_initial_claims: bool,
1328        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
1329        session_execution_lease_release_policy: SessionExecutionLeaseReleasePolicy,
1330    ) -> Result<PhysicalTurnExecution, RuntimeError> {
1331        self.refresh_session_graph_from_store()
1332            .await
1333            .map_err(session_head_refresh_error)?;
1334        let input_trace_turn_id = input.trace_turn_id.clone();
1335        let queued_turn_work = materialize_initial_claims
1336            .then(|| queued_claims.first())
1337            .flatten()
1338            .map(crate::QueuedWorkClaim::materialize_for_turn);
1339        let pending_turn_input = materialize_initial_claims
1340            .then(|| turn_input_claims.first())
1341            .flatten()
1342            .map(crate::TurnInputClaim::materialize_for_turn);
1343        if let Some(work) = pending_turn_input.as_ref()
1344            && input.items.is_empty()
1345            && input.image_blobs.is_empty()
1346        {
1347            input = work.clone();
1348            if input.trace_turn_id.is_none() {
1349                input.trace_turn_id = input_trace_turn_id.clone();
1350            }
1351        }
1352        if let Some(work) = queued_turn_work.as_ref()
1353            && input.items.is_empty()
1354            && input.image_blobs.is_empty()
1355        {
1356            input = work.input.clone();
1357            if input.trace_turn_id.is_none() {
1358                input.trace_turn_id = input_trace_turn_id;
1359            }
1360        }
1361        if self
1362            .session
1363            .as_ref()
1364            .and_then(|session| session.history_store())
1365            .is_some()
1366        {
1367            ensure_durable_effect_input(&input)?;
1368        }
1369        if let Some(extension) = &input.protocol_extension
1370            && let Some(session) = self.session.as_ref()
1371        {
1372            let protocol_session = std::sync::Arc::clone(session.plugins().protocol_session());
1373            protocol_session
1374                .validate_turn_extension(extension)
1375                .await
1376                .map_err(|err| {
1377                    RuntimeError::new(RuntimeErrorCode::ProtocolTurnExtension, err.to_string())
1378                })?;
1379        }
1380        let previous_prompt_usage = self.state.last_prompt_usage.clone();
1381        let normalized = match self
1382            .normalize_input_items(&input.items, &input.image_blobs)
1383            .await
1384        {
1385            Ok(items) => items,
1386            Err(e) => {
1387                self.state.last_prompt_usage = None;
1388                let mut assembler = TurnAssembler::default();
1389                let error_event = SessionStreamEvent::Error {
1390                    message: e.clone(),
1391                    envelope: Some(crate::session_model::ErrorEnvelope {
1392                        kind: "input_validation".to_string(),
1393                        code: Some("invalid_turn_input".to_string()),
1394                        terminal_reason: None,
1395                        user_message: e.clone(),
1396                        raw: None,
1397                        retryable: Some(false),
1398                        provider_failure_kind: None,
1399                    }),
1400                };
1401                assembler.push(&error_event);
1402                emit_turn_activity_to_sink(
1403                    turn_events,
1404                    TurnActivity::independent(TurnEvent::Error { message: e }),
1405                )
1406                .await;
1407                emit_session_event_to_sink(events, error_event).await;
1408                let outcome_event = SessionStreamEvent::TurnOutcome {
1409                    outcome: TurnOutcome::Stopped(TurnStop::InvalidInput),
1410                };
1411                assembler.push(&outcome_event);
1412                emit_session_event_to_sink(events, outcome_event).await;
1413                assembler.push(&SessionStreamEvent::Done);
1414                emit_session_event_to_sink(events, SessionStreamEvent::Done).await;
1415                let turn_index = self.state.turn_index + 1;
1416                let trace_turn_id = input
1417                    .trace_turn_id
1418                    .clone()
1419                    .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1420                let turn_control = ActiveTurnControl::new(
1421                    scoped_effect_controller.controller(),
1422                    TurnAddress::new(&self.state.session_id, &trace_turn_id),
1423                )
1424                .await?
1425                .with_local_cancel_origin(input.turn_context.local_cancel_origin_hint());
1426                let messages = crate::MessageSequence::from_base(self.state.read_model().messages);
1427                let mut turn_pipeline = TurnBoundary::from_state_with_clock(
1428                    self.state.clone(),
1429                    Arc::clone(&self.host.core.clock),
1430                )
1431                .with_session_execution_lease(
1432                    session_execution_lease.map(SessionExecutionLeaseGuard::fence),
1433                );
1434                turn_pipeline.apply_prepared_messages(&messages);
1435                return self
1436                    .finish_turn(
1437                        TurnFinishInput {
1438                            turn_pipeline,
1439                            assembler,
1440                            new_messages: messages,
1441                            policy: RuntimeSessionPolicy::new(
1442                                self.state.effective_policy().clone(),
1443                                Default::default(),
1444                            ),
1445                            turn_index,
1446                            queued_work_claims: queued_claims,
1447                            turn_input_claims,
1448                            trace_turn_id,
1449                        },
1450                        events,
1451                        &scoped_effect_controller,
1452                        &cancel,
1453                        session_execution_lease,
1454                        session_execution_lease_release_policy,
1455                        &turn_control,
1456                    )
1457                    .await;
1458            }
1459        };
1460        let turn_index = self.state.turn_index + 1;
1461        let trace_turn_id = input
1462            .trace_turn_id
1463            .clone()
1464            .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
1465        if self.host.core.tracing.trace_sink.is_some() {
1466            let mut trace_metadata = std::collections::BTreeMap::new();
1467            trace_metadata.insert(
1468                "input_item_count".to_string(),
1469                serde_json::json!(normalized.len()),
1470            );
1471            crate::trace::emit_trace(
1472                &self.host.core.tracing.trace_sink,
1473                &self.host.core.tracing.trace_context,
1474                lash_trace::TraceContext::default()
1475                    .for_session(self.state.session_id.clone())
1476                    .for_turn_index(turn_index)
1477                    .for_turn(trace_turn_id.clone()),
1478                lash_trace::TraceEvent::TurnStarted {
1479                    metadata: trace_metadata,
1480                },
1481                self.host.core.clock.as_ref(),
1482            );
1483        }
1484
1485        let base_read_model = self.state.read_model();
1486        let base_messages = base_read_model.messages;
1487        let base_render_cache = base_read_model.prompt_render_cache;
1488        let mut turn_delta = Vec::new();
1489        let initial_turn_causes = queued_turn_work
1490            .as_ref()
1491            .map(|work| work.turn_causes.clone())
1492            .unwrap_or_default();
1493        turn_delta.extend(
1494            initial_turn_causes
1495                .iter()
1496                .map(crate::TurnCause::to_event_message),
1497        );
1498
1499        let user_id = fresh_message_id();
1500        let mut user_parts: Vec<Part> = Vec::new();
1501        for item in normalized {
1502            match item {
1503                NormalizedItem::Text(text) => {
1504                    if text.is_empty() {
1505                        continue;
1506                    }
1507                    user_parts.push(Part {
1508                        id: format!("{}.p{}", user_id, user_parts.len()),
1509                        kind: PartKind::Text,
1510                        content: text,
1511                        attachment: None,
1512                        tool_call_id: None,
1513                        tool_name: None,
1514                        tool_replay: None,
1515                        prune_state: PruneState::Intact,
1516                        reasoning_meta: None,
1517                        response_meta: None,
1518                    });
1519                }
1520                NormalizedItem::Image(reference) => {
1521                    user_parts.push(Part {
1522                        id: format!("{}.p{}", user_id, user_parts.len()),
1523                        kind: PartKind::Image,
1524                        content: String::new(),
1525                        attachment: Some(crate::session_model::message::PartAttachment {
1526                            reference,
1527                        }),
1528                        tool_call_id: None,
1529                        tool_name: None,
1530                        tool_replay: None,
1531                        prune_state: PruneState::Intact,
1532                        reasoning_meta: None,
1533                        response_meta: None,
1534                    });
1535                }
1536            }
1537        }
1538        if user_parts.is_empty() && initial_turn_causes.is_empty() {
1539            user_parts.push(Part {
1540                id: format!("{}.p0", user_id),
1541                kind: PartKind::Text,
1542                content: String::new(),
1543                attachment: None,
1544                tool_call_id: None,
1545                tool_name: None,
1546                tool_replay: None,
1547                prune_state: PruneState::Intact,
1548                reasoning_meta: None,
1549                response_meta: None,
1550            });
1551        }
1552        if !user_parts.is_empty() {
1553            reassign_part_ids(&user_id, &mut user_parts);
1554            turn_delta.push(Message {
1555                id: user_id.clone(),
1556                role: MessageRole::User,
1557                parts: shared_parts(user_parts),
1558                origin: None,
1559            });
1560        }
1561
1562        let manager = self
1563            .runtime_session_services_for_turn(None)
1564            .map_err(|err| {
1565                RuntimeError::new(RuntimeErrorCode::PluginSessionManager, err.to_string())
1566            })?;
1567        let plugin_session = self
1568            .session
1569            .as_ref()
1570            .map(|s| Arc::clone(s.plugins()))
1571            .ok_or_else(|| {
1572                RuntimeError::new(
1573                    RuntimeErrorCode::ContextPrepareTurn,
1574                    "runtime session not available",
1575                )
1576            })?;
1577        let prepare_phase_turn_id = turn_phase_id(&trace_turn_id, "prepare-turn");
1578        let prepare_phase_controller = scoped_child_turn_controller(
1579            &scoped_effect_controller,
1580            &self.state.session_id,
1581            &prepare_phase_turn_id,
1582        )?;
1583        let turn_ctx = crate::TurnTransformContext {
1584            session_id: self.state.session_id.clone(),
1585            state: self.read_view(),
1586            prompt_usage: previous_prompt_usage.clone(),
1587            max_context_tokens: Some(LashRuntime::max_context_tokens(self)),
1588            sessions: manager.state_service(),
1589            session_lifecycle: manager.lifecycle_service(),
1590            session_graph: manager.graph_service(),
1591            scoped_effect_controller: scoped_effect_controller.clone(),
1592            direct_completions: manager.direct_completion_client(
1593                RuntimeEffectControllerHandle::borrowed(prepare_phase_controller),
1594                Some(prepare_phase_turn_id),
1595            ),
1596        };
1597        self.mark_phase_begin(RuntimeTurnPhase::ContextTransform);
1598        let prepared_context = plugin_session
1599            .prepare_turn_context(
1600                &turn_ctx,
1601                crate::session_model::context::PreparedContext {
1602                    messages: crate::MessageSequence::from_base_and_delta(
1603                        base_messages,
1604                        turn_delta,
1605                    )
1606                    .with_base_render_cache(base_render_cache),
1607                    ..Default::default()
1608                },
1609                self.turn_phase_probe.clone(),
1610            )
1611            .await
1612            .map_err(|err| {
1613                RuntimeError::new(RuntimeErrorCode::ContextPrepareTurn, err.to_string())
1614            })?;
1615        self.mark_phase_end(RuntimeTurnPhase::ContextTransform);
1616        // Release the read-view's graph clone before the rest of the turn
1617        // runs. Keeping it alive into `stream_prepared_turn` forces the
1618        // post-turn `append_active_read_delta` to deep-clone the session
1619        // graph (Arc::make_mut with refcount > 1).
1620        drop(turn_ctx);
1621        let messages = prepared_context.messages;
1622        if let Some(session) = self.session.as_mut() {
1623            session
1624                .set_context_overlay(
1625                    prepared_context.tool_providers,
1626                    prepared_context.prompt_contributions,
1627                    prepared_context.include_base_tools,
1628                )
1629                .map_err(|err| {
1630                    RuntimeError::new(
1631                        RuntimeErrorCode::Other("session_tool_registry".to_string()),
1632                        err.to_string(),
1633                    )
1634                })?;
1635        }
1636
1637        self.state.last_prompt_usage = None;
1638        Box::pin(self.stream_prepared_turn_inner(
1639            messages,
1640            previous_prompt_usage,
1641            input.protocol_turn_options.clone(),
1642            input.protocol_extension.clone(),
1643            input.turn_context.clone(),
1644            initial_turn_causes,
1645            trace_turn_id,
1646            turn_index,
1647            events,
1648            turn_events,
1649            scoped_effect_controller,
1650            cancel,
1651            queued_claims,
1652            turn_input_claims,
1653            session_execution_lease,
1654            session_execution_lease_release_policy,
1655        ))
1656        .await
1657    }
1658
1659    /// Run one logical turn and return only its assembled terminal result.
1660    pub async fn run_turn_assembled(
1661        &mut self,
1662        input: TurnInput,
1663        cancel: CancellationToken,
1664        scoped_effect_controller: ScopedEffectController<'_>,
1665    ) -> Result<AssembledTurn, RuntimeError> {
1666        self.stream_turn(input, TurnOptions::new(cancel, scoped_effect_controller))
1667            .await
1668    }
1669
1670    /// Run one logical turn using host-prepared message history.
1671    #[allow(clippy::too_many_arguments)]
1672    pub async fn stream_prepared_turn(
1673        &mut self,
1674        messages: crate::MessageSequence,
1675        previous_prompt_usage: Option<PromptUsage>,
1676        protocol_turn_options: Option<crate::ProtocolTurnOptions>,
1677        protocol_extension: Option<crate::ProtocolTurnExtensionHandle>,
1678        turn_context: crate::TurnContext,
1679        initial_turn_causes: Vec<crate::TurnCause>,
1680        trace_turn_id: String,
1681        turn_index: usize,
1682        events: &dyn EventSink,
1683        turn_events: &dyn TurnActivitySink,
1684        scoped_effect_controller: ScopedEffectController<'_>,
1685        cancel: CancellationToken,
1686        initial_queue_claim: Option<crate::QueuedWorkClaim>,
1687        initial_turn_input_claim: Option<crate::TurnInputClaim>,
1688    ) -> Result<AssembledTurn, RuntimeError> {
1689        let stopwatch = TurnStopwatch::start(self.host.core.clock.as_ref());
1690        let session_execution_lease = self
1691            .claim_session_execution_lease(cancel.clone(), true)
1692            .await?;
1693        let result = Box::pin(self.drive_logical_turn(
1694            LogicalTurnStart::Prepared(PreparedLogicalTurn {
1695                messages,
1696                previous_prompt_usage,
1697                protocol_turn_options,
1698                protocol_extension,
1699                turn_context,
1700                initial_turn_causes,
1701                trace_turn_id,
1702                turn_index,
1703            }),
1704            events,
1705            turn_events,
1706            scoped_effect_controller,
1707            cancel,
1708            LogicalTurnClaims::new(
1709                initial_queue_claim.into_iter().collect(),
1710                initial_turn_input_claim.into_iter().collect(),
1711            ),
1712            session_execution_lease.as_ref(),
1713            stopwatch,
1714        ))
1715        .await
1716        .map(|run| {
1717            run.into_final_turn()
1718                .expect("logical turn always contains a terminal physical turn")
1719        });
1720        self.settle_session_execution_lease(session_execution_lease.as_ref(), result)
1721            .await
1722    }
1723
1724    #[allow(clippy::too_many_arguments)]
1725    pub(super) async fn stream_prepared_turn_inner(
1726        &mut self,
1727        messages: crate::MessageSequence,
1728        _previous_prompt_usage: Option<PromptUsage>,
1729        protocol_turn_options: Option<crate::ProtocolTurnOptions>,
1730        protocol_extension: Option<crate::ProtocolTurnExtensionHandle>,
1731        turn_context: crate::TurnContext,
1732        initial_turn_causes: Vec<crate::TurnCause>,
1733        trace_turn_id: String,
1734        turn_index: usize,
1735        events: &dyn EventSink,
1736        turn_events: &dyn TurnActivitySink,
1737        scoped_effect_controller: ScopedEffectController<'_>,
1738        cancel: CancellationToken,
1739        initial_queue_claims: Vec<crate::QueuedWorkClaim>,
1740        initial_turn_input_claims: Vec<crate::TurnInputClaim>,
1741        session_execution_lease: Option<&SessionExecutionLeaseGuard>,
1742        session_execution_lease_release_policy: SessionExecutionLeaseReleasePolicy,
1743    ) -> Result<PhysicalTurnExecution, RuntimeError> {
1744        let turn_control = Arc::new(
1745            ActiveTurnControl::new(
1746                scoped_effect_controller.controller(),
1747                TurnAddress::new(&self.state.session_id, &trace_turn_id),
1748            )
1749            .await?
1750            .with_local_cancel_origin(turn_context.local_cancel_origin_hint()),
1751        );
1752        if session_execution_lease.is_none()
1753            && self
1754                .session
1755                .as_ref()
1756                .and_then(|session| session.history_store())
1757                .is_some()
1758        {
1759            return Err(RuntimeError::new(
1760                RuntimeErrorCode::StoreCommitFailed,
1761                "prepared turn requires a session execution lease",
1762            ));
1763        }
1764        let session_execution_fence =
1765            session_execution_lease.map(SessionExecutionLeaseGuard::fence);
1766        let (event_tx, mut event_rx) = mpsc::channel::<RuntimeStreamEvent>(100);
1767        let child_usage_event_relay = ChildUsageEventRelay::new(event_tx.clone());
1768        let mut turn_policy = self.state.effective_policy().clone();
1769        let turn_provider_override = turn_context.provider().cloned();
1770        if let Some(provider) = turn_provider_override.as_ref() {
1771            turn_policy.provider_id = provider.kind().to_string();
1772        }
1773        let session_protocol_turn_options = self.state.effective_protocol_turn_options().clone();
1774        let effective_protocol_turn_options = protocol_turn_options
1775            .clone()
1776            .map(|options| session_protocol_turn_options.merged_with_override(&options))
1777            .unwrap_or(session_protocol_turn_options);
1778        let manager = self
1779            .runtime_session_services_for_turn(Some(child_usage_event_relay.clone()))
1780            .map_err(|err| {
1781                RuntimeError::new(RuntimeErrorCode::PluginSessionManager, err.to_string())
1782            })?;
1783        let plugins = {
1784            let session = self
1785                .session
1786                .as_ref()
1787                .expect("lash runtime session must be available");
1788            Arc::clone(session.plugins())
1789        };
1790        let mut assembler = TurnAssembler::new();
1791        self.mark_phase_begin(RuntimeTurnPhase::BeforeTurnHooks);
1792        // Block-scope the pinned future so it (and its captured
1793        // `SessionReadView` clone of the session graph) drops before the
1794        // post-turn `append_active_read_delta` mutation. Keeping it alive
1795        // across the turn forces `Arc::make_mut` to deep-clone
1796        // `SessionGraphData`.
1797        let prepared = {
1798            let prepare_turn = plugins.prepare_turn_with_phase_probe(
1799                PrepareTurnRequest {
1800                    session_id: self.state.session_id.clone(),
1801                    state: crate::SessionReadView::from_runtime_state(
1802                        &self.state,
1803                        turn_policy.clone(),
1804                        effective_protocol_turn_options.clone(),
1805                    ),
1806                    messages,
1807                    sessions: manager.state_service(),
1808                    session_lifecycle: manager.lifecycle_service(),
1809                    session_graph: manager.graph_service(),
1810                    turn_context: turn_context.clone(),
1811                },
1812                self.turn_phase_probe.clone(),
1813            );
1814            let mut prepare_turn = Box::pin(prepare_turn);
1815
1816            loop {
1817                tokio::select! {
1818                    prepared = prepare_turn.as_mut() => {
1819                        let prepared = prepared.map_err(|err| {
1820                            RuntimeError::new(RuntimeErrorCode::PluginPrepareTurn, err.to_string())
1821                        })?;
1822                        self.mark_phase_end(RuntimeTurnPhase::BeforeTurnHooks);
1823                        break prepared;
1824                    }
1825                    maybe_event = event_rx.recv() => {
1826                        if let Some(event) = maybe_event {
1827                            emit_runtime_stream_event_to_sinks(
1828                                events,
1829                                turn_events,
1830                                event,
1831                                &mut assembler,
1832                            )
1833                            .await;
1834                        }
1835                    }
1836                }
1837            }
1838        };
1839        for event in &prepared.events {
1840            assembler.push(event);
1841        }
1842        emit_session_events_to_sink(events, prepared.events).await;
1843        if let Some(abort) = prepared.abort {
1844            drop(event_tx);
1845
1846            let mut turn_pipeline = TurnBoundary::from_state_with_clock(
1847                self.state.clone(),
1848                Arc::clone(&self.host.core.clock),
1849            )
1850            .with_session_execution_lease(session_execution_fence.clone());
1851            turn_pipeline.apply_prepared_messages(&prepared.messages);
1852            let issue = TurnIssue {
1853                kind: "plugin".to_string(),
1854                code: Some(abort.code),
1855                terminal_reason: None,
1856                message: abort.message.clone(),
1857                raw: None,
1858                retryable: None,
1859                provider_failure_kind: None,
1860            };
1861            let error_event = SessionStreamEvent::Error {
1862                message: abort.message,
1863                envelope: Some(crate::session_model::ErrorEnvelope {
1864                    kind: "plugin".to_string(),
1865                    code: issue.code.clone(),
1866                    terminal_reason: None,
1867                    user_message: issue.message.clone(),
1868                    raw: None,
1869                    retryable: None,
1870                    provider_failure_kind: None,
1871                }),
1872            };
1873            assembler.push(&error_event);
1874            emit_turn_activity_to_sink(
1875                turn_events,
1876                TurnActivity::independent(TurnEvent::Error {
1877                    message: issue.message.clone(),
1878                }),
1879            )
1880            .await;
1881            emit_session_event_to_sink(events, error_event).await;
1882            let outcome_event = SessionStreamEvent::TurnOutcome {
1883                outcome: TurnOutcome::Stopped(TurnStop::PluginAbort),
1884            };
1885            assembler.push(&outcome_event);
1886            emit_session_event_to_sink(events, outcome_event).await;
1887            assembler.push(&SessionStreamEvent::Done);
1888            emit_session_event_to_sink(events, SessionStreamEvent::Done).await;
1889            return self
1890                .finish_turn(
1891                    TurnFinishInput {
1892                        turn_pipeline,
1893                        assembler,
1894                        new_messages: prepared.messages,
1895                        policy: RuntimeSessionPolicy::new(
1896                            self.state.effective_policy().clone(),
1897                            Default::default(),
1898                        ),
1899                        turn_index,
1900                        queued_work_claims: initial_queue_claims,
1901                        turn_input_claims: initial_turn_input_claims,
1902                        trace_turn_id,
1903                    },
1904                    events,
1905                    &scoped_effect_controller,
1906                    &cancel,
1907                    session_execution_lease,
1908                    session_execution_lease_release_policy,
1909                    turn_control.as_ref(),
1910                )
1911                .await;
1912        }
1913        let mut turn_pipeline = TurnBoundary::from_state_with_clock(
1914            self.state.clone(),
1915            Arc::clone(&self.host.core.clock),
1916        )
1917        .with_session_execution_lease(session_execution_fence.clone());
1918        let store = self
1919            .session
1920            .as_ref()
1921            .and_then(|session| session.history_store());
1922        // Durable controllers, like Restate, own in-flight replay. Writing
1923        // progress checkpoints directly to the shared store would make handler
1924        // replay observe a newer partial turn and change effect replay keys.
1925        let progress_store = if scoped_effect_controller.controller().durability_tier()
1926            == crate::DurabilityTier::Durable
1927        {
1928            None
1929        } else {
1930            store.as_ref().map(|store| store.as_ref())
1931        };
1932        turn_pipeline
1933            .prepared_checkpoint(
1934                progress_store,
1935                turn_policy.clone(),
1936                turn_index,
1937                &prepared.messages,
1938                self.session.as_mut(),
1939            )
1940            .await
1941            .map_err(super::runtime_error_from_store_commit)?;
1942        let resolved_turn_policy = if let Some(provider) = turn_provider_override {
1943            RuntimeSessionPolicy::from_provider(
1944                turn_policy.clone(),
1945                provider.with_clock(Arc::clone(&self.host.core.clock)),
1946            )
1947            .map_err(|err| RuntimeError::new("llm_provider", err.to_string()))?
1948        } else {
1949            self.host
1950                .resolve_session_policy(&self.state.session_id, turn_policy.clone())
1951                .map_err(|err| RuntimeError::new("llm_provider", err.to_string()))?
1952        };
1953        let manager = self
1954            .runtime_session_services_for_turn(Some(child_usage_event_relay.clone()))
1955            .map_err(|err| {
1956                RuntimeError::new(RuntimeErrorCode::PluginSessionManager, err.to_string())
1957            })?;
1958        let cancel_state = cancel.clone();
1959        let finish_scoped_effect_controller = scoped_effect_controller.clone();
1960        let shared_cancel_controller = match scoped_effect_controller.shared_controller() {
1961            Some(controller) => Some(controller),
1962            None if scoped_effect_controller.controller().durability_tier()
1963                == crate::DurabilityTier::Durable
1964                && self.host.core.control.effect_host.durability_tier()
1965                    == crate::DurabilityTier::Durable =>
1966            {
1967                self.host
1968                    .core
1969                    .control
1970                    .effect_host
1971                    .scoped_static(scoped_effect_controller.execution_scope().clone())?
1972                    .and_then(|scoped| scoped.shared_controller())
1973            }
1974            None => None,
1975        };
1976        let session = self
1977            .session
1978            .take()
1979            .expect("lash runtime session must be available");
1980        let mut driver = Box::new(RuntimeTurnDriver {
1981            session,
1982            policy: resolved_turn_policy,
1983            host: self.host.clone(),
1984            turn_id: scoped_effect_controller.scope_id().to_string(),
1985            scoped_effect_controller,
1986            session_id: self.state.session_id.clone(),
1987            turn_index,
1988            turn_pipeline,
1989            llm_stream_summaries: HashMap::new(),
1990            llm_calls: Vec::new(),
1991            next_llm_ordinal: 0,
1992            session_services: manager,
1993            protocol_turn_options: effective_protocol_turn_options,
1994            protocol_extension,
1995            turn_context,
1996            turn_causes: initial_turn_causes,
1997            pending_queue_claims: initial_queue_claims,
1998            pending_turn_input_claims: initial_turn_input_claims,
1999            checkpoint_messages: crate::tool_dispatch::CheckpointMessageBuffer::default(),
2000            session_execution_lease: session_execution_fence,
2001            runtime_lease_owner: self.runtime_lease_owner.clone(),
2002            turn_phase_probe: self.turn_phase_probe.clone(),
2003        });
2004        let protocol_run_offset = 0;
2005        let cancellation_messages = prepared.messages.clone();
2006        self.mark_phase_begin(RuntimeTurnPhase::EffectLoop);
2007        let run_result = Box::pin(run_turn_effect_loop(
2008            &mut driver,
2009            prepared.messages,
2010            event_tx,
2011            cancel.clone(),
2012            protocol_run_offset,
2013            Arc::clone(&turn_control),
2014            shared_cancel_controller,
2015            finish_scoped_effect_controller.controller(),
2016            &mut event_rx,
2017            &mut assembler,
2018            &child_usage_event_relay,
2019            events,
2020            turn_events,
2021        ))
2022        .await;
2023        let (new_messages, _new_protocol_iteration) = match run_result {
2024            Ok(result) => result,
2025            Err(_err) if cancel.is_cancelled() && turn_control.evidence().is_some() => {
2026                self.mark_phase_end(RuntimeTurnPhase::EffectLoop);
2027                return Box::pin(self.finish_cancelled_turn_after_effect_abort(
2028                    *driver,
2029                    assembler,
2030                    cancellation_messages,
2031                    events,
2032                    &finish_scoped_effect_controller,
2033                    &cancel,
2034                    session_execution_lease,
2035                    session_execution_lease_release_policy,
2036                    turn_control.as_ref(),
2037                    turn_index,
2038                    trace_turn_id,
2039                ))
2040                .await;
2041            }
2042            Err(err) => {
2043                self.mark_phase_end(RuntimeTurnPhase::EffectLoop);
2044                let RuntimeTurnDriver {
2045                    session,
2046                    pending_queue_claims,
2047                    pending_turn_input_claims,
2048                    ..
2049                } = *driver;
2050                self.session = Some(session);
2051                self.abandon_queued_work_claims_after_lease_loss(&err, &pending_queue_claims)
2052                    .await;
2053                self.abandon_turn_input_claims_after_lease_loss(&err, &pending_turn_input_claims)
2054                    .await;
2055                return Err(err);
2056            }
2057        };
2058        self.mark_phase_end(RuntimeTurnPhase::EffectLoop);
2059        tracing::debug!(
2060            new_message_count = new_messages.len(),
2061            tool_call_count = assembler.tool_calls.len(),
2062            "runtime post-run_task"
2063        );
2064
2065        let RuntimeTurnDriver {
2066            session,
2067            policy,
2068            turn_pipeline,
2069            llm_calls,
2070            pending_queue_claims,
2071            pending_turn_input_claims,
2072            ..
2073        } = *driver;
2074        self.session = Some(session);
2075        let pending_queue_claims_for_abandon = pending_queue_claims.clone();
2076        let pending_turn_input_claims_for_abandon = pending_turn_input_claims.clone();
2077        let finish_result = Box::pin(self.finish_turn(
2078            TurnFinishInput {
2079                turn_pipeline,
2080                assembler: assembler.with_llm_calls(llm_calls),
2081                new_messages,
2082                policy,
2083                turn_index,
2084                queued_work_claims: pending_queue_claims,
2085                turn_input_claims: pending_turn_input_claims,
2086                trace_turn_id,
2087            },
2088            events,
2089            &finish_scoped_effect_controller,
2090            &cancel_state,
2091            session_execution_lease,
2092            session_execution_lease_release_policy,
2093            turn_control.as_ref(),
2094        ))
2095        .await;
2096        if let Err(err) = &finish_result {
2097            self.abandon_queued_work_claims_after_lease_loss(
2098                err,
2099                &pending_queue_claims_for_abandon,
2100            )
2101            .await;
2102            self.abandon_turn_input_claims_after_lease_loss(
2103                err,
2104                &pending_turn_input_claims_for_abandon,
2105            )
2106            .await;
2107        }
2108        finish_result
2109    }
2110    async fn normalize_input_items(
2111        &self,
2112        items: &[InputItem],
2113        image_blobs: &HashMap<String, Vec<u8>>,
2114    ) -> Result<Vec<NormalizedItem>, String> {
2115        normalize_input_items(
2116            items,
2117            image_blobs,
2118            self.host.core.durability.attachment_store.as_ref(),
2119        )
2120        .await
2121    }
2122}
2123
2124pub fn ensure_durable_effect_input(input: &TurnInput) -> Result<(), RuntimeError> {
2125    if input.protocol_extension.is_some() {
2126        return Err(RuntimeError::new(
2127            RuntimeErrorCode::DurableEffectLiveProtocolExtension,
2128            "durable effect hosts do not support live protocol_extension inputs; encode replayable data in protocol_turn_options or persisted plugin state",
2129        ));
2130    }
2131    input
2132        .turn_context
2133        .live_plugin_inputs()
2134        .durable_effect_rejection()?;
2135    Ok(())
2136}
2137
2138async fn emit_turn_activity_to_sink(events: &dyn TurnActivitySink, activity: TurnActivity) {
2139    if !events.is_noop() {
2140        events.emit(activity).await;
2141    }
2142}
2143
2144async fn publish_terminal_after_commit(
2145    turn_control: &ActiveTurnControl,
2146    resolver: &dyn AwaitEventResolver,
2147    terminal: &TurnTerminal,
2148    session_id: &str,
2149    turn_id: &str,
2150) {
2151    if let Err(err) = turn_control.publish_terminal(resolver, terminal).await {
2152        tracing::warn!(
2153            error = %err,
2154            session_id,
2155            turn_id,
2156            "turn committed but terminal publication failed"
2157        );
2158    }
2159}
2160
2161#[allow(clippy::too_many_arguments)]
2162async fn run_turn_effect_loop(
2163    driver: &mut RuntimeTurnDriver<'_>,
2164    messages: crate::MessageSequence,
2165    event_tx: mpsc::Sender<RuntimeStreamEvent>,
2166    cancellation: CancellationToken,
2167    protocol_run_offset: usize,
2168    turn_control: Arc<ActiveTurnControl>,
2169    shared_cancel_controller: Option<Arc<dyn RuntimeEffectController>>,
2170    cancel_controller: &dyn RuntimeEffectController,
2171    event_rx: &mut mpsc::Receiver<RuntimeStreamEvent>,
2172    assembler: &mut TurnAssembler,
2173    child_usage_event_relay: &ChildUsageEventRelay,
2174    events: &dyn EventSink,
2175    turn_events: &dyn TurnActivitySink,
2176) -> Result<(crate::MessageSequence, usize), RuntimeError> {
2177    // The start gate can change the handler's control flow before its first
2178    // effect, so durable runtimes must observe it through the handler-scoped
2179    // controller. That controller journals the observation and replays the
2180    // same answer after an owner crash. The shared controller is intentionally
2181    // reserved for the concurrent live watcher below: an out-of-band peek here
2182    // could observe a cancel that arrived after the original attempt and make
2183    // a replay take a different command path.
2184    if await_turn_cancellation_with_retry(|| turn_control.observe_pending_cancel(cancel_controller))
2185        .await
2186        .is_some()
2187    {
2188        cancellation.cancel();
2189    }
2190    let cancel_watcher = shared_cancel_controller.map(|controller| {
2191        let turn_control = Arc::clone(&turn_control);
2192        let cancellation = cancellation.clone();
2193        tokio::spawn(async move {
2194            if await_turn_cancellation_with_retry(|| {
2195                turn_control.await_cancel(controller.as_ref(), CancellationToken::new())
2196            })
2197            .await
2198            .is_some()
2199            {
2200                cancellation.cancel();
2201            }
2202        })
2203    });
2204    let run_future = Box::pin(driver.run(
2205        messages,
2206        event_tx,
2207        cancellation.clone(),
2208        protocol_run_offset,
2209    ));
2210    let result = if cancel_watcher.is_some() {
2211        drive_turn_to_completion(
2212            run_future,
2213            event_rx,
2214            assembler,
2215            child_usage_event_relay,
2216            events,
2217            turn_events,
2218        )
2219        .await
2220    } else {
2221        drive_turn_to_completion_with_cancel(
2222            run_future,
2223            await_turn_cancellation_with_retry(|| {
2224                turn_control.await_cancel(cancel_controller, CancellationToken::new())
2225            }),
2226            cancellation,
2227            event_rx,
2228            assembler,
2229            child_usage_event_relay,
2230            events,
2231            turn_events,
2232        )
2233        .await
2234    };
2235    if let Some(watcher) = cancel_watcher {
2236        watcher.abort();
2237    }
2238    result
2239}
2240
2241const TURN_CANCEL_WATCH_RETRY_INITIAL: std::time::Duration = std::time::Duration::from_millis(25);
2242const TURN_CANCEL_WATCH_RETRY_MAX: std::time::Duration = std::time::Duration::from_secs(1);
2243
2244async fn await_turn_cancellation_with_retry<F, C>(mut watch: F) -> Option<TurnCancellationEvidence>
2245where
2246    F: FnMut() -> C,
2247    C: std::future::Future<Output = Result<Option<TurnCancellationEvidence>, RuntimeError>>,
2248{
2249    let mut backoff = TURN_CANCEL_WATCH_RETRY_INITIAL;
2250    loop {
2251        match watch().await {
2252            Ok(observation) => return observation,
2253            Err(err) => {
2254                tracing::warn!(
2255                    error = %err,
2256                    retry_after_ms = backoff.as_millis(),
2257                    "turn cancellation watcher failed; retrying while the turn remains active"
2258                );
2259                tokio::time::sleep(backoff).await;
2260                backoff = backoff.saturating_mul(2).min(TURN_CANCEL_WATCH_RETRY_MAX);
2261            }
2262        }
2263    }
2264}
2265
2266/// Pump the turn driver's event channel into the host sinks while the run
2267/// future executes, then drain any events emitted between completion and the
2268/// sender dropping.
2269///
2270/// Both the fresh and resumed turn entry points construct a
2271/// `RuntimeTurnDriver`, kick off its run future, and need identical
2272/// event-pump/drain behavior before tearing the driver down. Only the driver
2273/// construction and post-run teardown differ, so each caller owns those and
2274/// shares this loop.
2275async fn drive_turn_to_completion<F>(
2276    run_future: F,
2277    event_rx: &mut mpsc::Receiver<RuntimeStreamEvent>,
2278    assembler: &mut TurnAssembler,
2279    child_usage_event_relay: &ChildUsageEventRelay,
2280    events: &dyn EventSink,
2281    turn_events: &dyn TurnActivitySink,
2282) -> Result<(crate::MessageSequence, usize), RuntimeError>
2283where
2284    F: std::future::Future<Output = Result<(crate::MessageSequence, usize), RuntimeError>>,
2285{
2286    let run_result = {
2287        let mut run_future = Box::pin(run_future);
2288        loop {
2289            tokio::select! {
2290                // Some durable adapter futures are not fused. Once turn
2291                // completion is ready, select it before another ready branch
2292                // so the loop never polls the completed future again.
2293                biased;
2294
2295                completed = run_future.as_mut() => {
2296                    child_usage_event_relay.clear();
2297                    break completed;
2298                }
2299                maybe_event = event_rx.recv() => {
2300                    if let Some(event) = maybe_event {
2301                        emit_runtime_stream_event_to_sinks(
2302                            events,
2303                            turn_events,
2304                            event,
2305                            assembler,
2306                        )
2307                        .await;
2308                    }
2309                }
2310            }
2311        }
2312    };
2313    while let Some(event) = event_rx.recv().await {
2314        emit_runtime_stream_event_to_sinks(events, turn_events, event, assembler).await;
2315    }
2316    run_result
2317}
2318
2319#[allow(clippy::too_many_arguments)]
2320async fn drive_turn_to_completion_with_cancel<F, C>(
2321    run_future: F,
2322    cancel_future: C,
2323    cancellation: CancellationToken,
2324    event_rx: &mut mpsc::Receiver<RuntimeStreamEvent>,
2325    assembler: &mut TurnAssembler,
2326    child_usage_event_relay: &ChildUsageEventRelay,
2327    events: &dyn EventSink,
2328    turn_events: &dyn TurnActivitySink,
2329) -> Result<(crate::MessageSequence, usize), RuntimeError>
2330where
2331    F: std::future::Future<Output = Result<(crate::MessageSequence, usize), RuntimeError>>,
2332    C: std::future::Future<Output = Option<TurnCancellationEvidence>>,
2333{
2334    let run_result = {
2335        let mut run_future = Box::pin(run_future);
2336        let mut cancel_future = Box::pin(cancel_future);
2337        let mut cancellation_observed = false;
2338        loop {
2339            tokio::select! {
2340                // Keep the non-fused turn future from being re-polled when
2341                // cancellation or a stream event becomes ready alongside it.
2342                biased;
2343
2344                completed = run_future.as_mut() => {
2345                    child_usage_event_relay.clear();
2346                    break completed;
2347                }
2348                maybe_event = event_rx.recv() => {
2349                    if let Some(event) = maybe_event {
2350                        emit_runtime_stream_event_to_sinks(
2351                            events,
2352                            turn_events,
2353                            event,
2354                            assembler,
2355                        )
2356                        .await;
2357                    }
2358                }
2359                observation = cancel_future.as_mut(), if !cancellation_observed => {
2360                    cancellation_observed = true;
2361                    if observation.is_some() {
2362                        cancellation.cancel();
2363                    }
2364                }
2365            }
2366        }
2367    };
2368    while let Some(event) = event_rx.recv().await {
2369        emit_runtime_stream_event_to_sinks(events, turn_events, event, assembler).await;
2370    }
2371    run_result
2372}
2373
2374async fn emit_runtime_stream_event_to_sinks(
2375    events: &dyn EventSink,
2376    turn_events: &dyn TurnActivitySink,
2377    event: RuntimeStreamEvent,
2378    assembler: &mut TurnAssembler,
2379) {
2380    match event {
2381        RuntimeStreamEvent::Session(event) => {
2382            assembler.push(&event);
2383            emit_session_event_to_sink(events, event).await;
2384        }
2385        RuntimeStreamEvent::Turn(activity) => {
2386            emit_turn_activity_to_sink(turn_events, activity).await;
2387        }
2388    }
2389}
2390
2391#[cfg(test)]
2392mod tests {
2393    use std::sync::Arc;
2394    use std::sync::atomic::{AtomicUsize, Ordering};
2395
2396    use super::{
2397        ActiveTurnControl, agent_frame_follow_turn_id, await_turn_cancellation_with_retry,
2398        publish_terminal_after_commit,
2399    };
2400    use crate::{
2401        AwaitEventKey, AwaitEventResolver, Resolution, ResolveOutcome, RuntimeError, TurnAddress,
2402        TurnCancellationEvidence, TurnFinish, TurnOutcome, TurnTerminal,
2403    };
2404
2405    #[derive(Default)]
2406    struct RejectTerminalPublication {
2407        attempts: AtomicUsize,
2408    }
2409
2410    #[async_trait::async_trait]
2411    impl AwaitEventResolver for RejectTerminalPublication {
2412        async fn resolve_await_event(
2413            &self,
2414            _key: &AwaitEventKey,
2415            _resolution: Resolution,
2416        ) -> Result<ResolveOutcome, RuntimeError> {
2417            self.attempts.fetch_add(1, Ordering::SeqCst);
2418            Err(RuntimeError::new(
2419                "transient_terminal_publication",
2420                "terminal backend unavailable",
2421            ))
2422        }
2423    }
2424
2425    #[test]
2426    fn agent_frame_follow_turn_ids_are_distinct_and_deterministic() {
2427        assert_eq!(agent_frame_follow_turn_id("root-turn", 0), "root-turn");
2428        assert_eq!(
2429            agent_frame_follow_turn_id("root-turn", 1),
2430            "root-turn:agent-frame:1"
2431        );
2432        assert_eq!(
2433            agent_frame_follow_turn_id("root-turn", 2),
2434            "root-turn:agent-frame:2"
2435        );
2436    }
2437
2438    #[tokio::test]
2439    async fn cancellation_watch_retries_transient_errors_until_evidence_arrives() {
2440        let attempts = Arc::new(AtomicUsize::new(0));
2441        let observed_attempts = Arc::clone(&attempts);
2442        let evidence = await_turn_cancellation_with_retry(move || {
2443            let attempt = observed_attempts.fetch_add(1, Ordering::SeqCst);
2444            async move {
2445                if attempt < 2 {
2446                    Err(RuntimeError::new(
2447                        "transient_cancel_watch",
2448                        "temporary ingress failure",
2449                    ))
2450                } else {
2451                    Ok(Some(TurnCancellationEvidence {
2452                        request_id: "retry-request".to_string(),
2453                        origin: Some("test-user".to_string()),
2454                        reason: None,
2455                    }))
2456                }
2457            }
2458        })
2459        .await
2460        .expect("cancellation evidence after retries");
2461
2462        assert_eq!(attempts.load(Ordering::SeqCst), 3);
2463        assert_eq!(evidence.request_id, "retry-request");
2464    }
2465
2466    #[tokio::test]
2467    async fn terminal_publication_failure_is_non_fatal_after_commit() {
2468        let resolver = RejectTerminalPublication::default();
2469        let control = ActiveTurnControl::new(
2470            &resolver,
2471            TurnAddress::new("committed-session", "committed-turn"),
2472        )
2473        .await
2474        .expect("active turn control");
2475        publish_terminal_after_commit(
2476            &control,
2477            &resolver,
2478            &TurnTerminal::Committed {
2479                outcome: TurnOutcome::Finished(TurnFinish::AssistantMessage {
2480                    text: "committed".to_string(),
2481                }),
2482                cancellation: None,
2483                session_revision: Some(1),
2484            },
2485            "committed-session",
2486            "committed-turn",
2487        )
2488        .await;
2489        assert_eq!(resolver.attempts.load(Ordering::SeqCst), 1);
2490    }
2491}