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#[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 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 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 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 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 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 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 #[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 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 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 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
2266async 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 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 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}