From 17c629136a6df759501ad1aa077df59796a23cc3 Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 26 Aug 2026 13:47:35 +0900 Subject: [PATCH] fix: scope max turns to logical runs --- crates/agen/src/engine.rs | 85 +++++++-- crates/agen/tests/engine_state_test.rs | 187 +++++++++++++++++++- crates/session-store/src/segment.rs | 2 + crates/session-store/src/segment_log.rs | 134 +++++++++++++- crates/session-store/tests/fs_store_test.rs | 1 + crates/session-store/tests/session_test.rs | 2 + crates/worker/src/worker.rs | 20 +++ 7 files changed, 412 insertions(+), 19 deletions(-) diff --git a/crates/agen/src/engine.rs b/crates/agen/src/engine.rs index 15c22d82..6619997a 100644 --- a/crates/agen/src/engine.rs +++ b/crates/agen/src/engine.rs @@ -179,14 +179,20 @@ pub struct Engine { history: Vec, /// History length at lock time (only meaningful in Locked state) locked_prefix_len: usize, - /// AgentTurn count. + /// AgentTurn count across the lifetime of this Engine. /// /// Once retry (`agen-stream-continuation`) is implemented, an /// AgentTurn collapses N retried `LlmCall`s with identical input; /// today retry is not implemented so AgentTurn and LlmCall fire 1:1 /// and the increment site (the LLM-call loop) is shared. - /// `max_turns` is interpreted as a per-`run()` AgentTurn cap. turn_count: usize, + /// AgentTurns consumed by the currently active logical run. + /// + /// A fresh [`run`](Self::run) starts at zero. Pause and Yield retain the + /// count for [`resume`](Self::resume), while terminal outcomes clear it. + /// `max_turns` is enforced against this run-scoped count rather than the + /// cumulative `turn_count` above. + active_run_turn_count: Option, /// LlmCall count (per-Engine running counter, monotonic). Unlike /// `turn_count` this never collapses retries. llm_call_count: usize, @@ -268,6 +274,20 @@ impl Engine { self.last_run_interrupted = false; } + fn start_logical_run(&mut self) { + self.active_run_turn_count = Some(0); + } + + fn ensure_logical_run(&mut self) { + self.active_run_turn_count.get_or_insert(0); + } + + fn finish_logical_run(&mut self, result: &Result) { + if !matches!(result, Ok(EngineResult::Paused) | Ok(EngineResult::Yielded)) { + self.active_run_turn_count = None; + } + } + fn drain_cancel_queue(&mut self) { while self.cancel_rx.try_recv().is_ok() {} } @@ -650,6 +670,23 @@ impl Engine { self.turn_count } + /// Get the AgentTurns consumed by an interrupted logical run. + /// + /// `Some` is retained only while Pause or Yield permits a later + /// [`resume`](Self::resume). Terminal outcomes return this to `None`. + pub fn active_run_turn_count(&self) -> Option { + self.active_run_turn_count + } + + /// Restore the persisted turn budget of an interrupted logical run. + /// + /// Session owners restore this together with the cumulative turn count and + /// history. `None` means there is no resumable logical run and the next + /// [`resume`](Self::resume) starts a fresh budget. + pub fn set_active_run_turn_count(&mut self, turn_count: Option) { + self.active_run_turn_count = turn_count; + } + /// Get the current LlmCall count (per-Engine running counter, never /// collapsed by retry). pub fn llm_call_count(&self) -> usize { @@ -1123,6 +1160,19 @@ impl Engine { return Err(EngineError::Cancelled); } + if let Some(max) = self.max_turns + && self.active_run_turn_count.unwrap_or(0) >= max as usize + { + info!( + active_run_turn_count = self.active_run_turn_count.unwrap_or(0), + total_turn_count = self.turn_count, + max_turns = max, + "Logical run turn limit reached" + ); + self.last_run_interrupted = false; + return Ok(EngineResult::LimitReached); + } + let current_turn = self.turn_count; if !continuing_stream { debug!(turn = current_turn, "Turn start"); @@ -1314,6 +1364,7 @@ impl Engine { cb(current_turn); } self.turn_count += 1; + *self.active_run_turn_count.get_or_insert(0) += 1; // Collect and commit assistant items. Routed through // `append_history_items` so observers see each item as it lands. @@ -1344,18 +1395,6 @@ impl Engine { if let Some(result) = self.execute_and_commit_tools(tool_calls).await? { return Ok(result); } - - if let Some(max) = self.max_turns { - if self.turn_count >= max as usize { - info!( - turn_count = self.turn_count, - max_turns = max, - "Turn limit reached" - ); - self.last_run_interrupted = false; - return Ok(EngineResult::LimitReached); - } - } } } @@ -1664,6 +1703,7 @@ impl Engine { history: Vec::new(), locked_prefix_len: 0, turn_count: 0, + active_run_turn_count: None, llm_call_count: 0, tool_execution_batch_count: 0, max_turns: None, @@ -1866,6 +1906,9 @@ impl Engine { /// Set the last_run_interrupted flag (for session restoration) pub fn set_last_run_interrupted(&mut self, interrupted: bool) { self.last_run_interrupted = interrupted; + if !interrupted { + self.active_run_turn_count = None; + } } /// Apply configuration (reserved for future extensions) @@ -1934,6 +1977,7 @@ impl Engine { history: self.history, locked_prefix_len, turn_count: self.turn_count, + active_run_turn_count: self.active_run_turn_count, llm_call_count: self.llm_call_count, tool_execution_batch_count: self.tool_execution_batch_count, max_turns: self.max_turns, @@ -1974,6 +2018,8 @@ impl Engine { &mut self, user_input: impl Into, ) -> Result { + // Supplying new user input abandons any paused/yielded logical run. + self.active_run_turn_count = None; self.reset_interruption_state(); // Interceptor: on_prompt_submit let mut user_item = Item::user_message(user_input); @@ -1991,8 +2037,11 @@ impl Engine { if !extras.is_empty() { self.append_history_items(extras)?; } + self.start_logical_run(); let result = self.run_turn_loop().await; - self.finalize_interruption(result).await + let result = self.finalize_interruption(result).await; + self.finish_logical_run(&result); + result } /// Resume execution (from Paused state) @@ -2000,8 +2049,11 @@ impl Engine { /// Resumes turn processing from current state without adding a new user message. pub async fn resume(&mut self) -> Result { self.reset_interruption_state(); + self.ensure_logical_run(); let result = self.run_turn_loop().await; - self.finalize_interruption(result).await + let result = self.finalize_interruption(result).await; + self.finish_logical_run(&result); + result } /// Get the prefix length at lock time @@ -2027,6 +2079,7 @@ impl Engine { history: self.history, locked_prefix_len: 0, turn_count: self.turn_count, + active_run_turn_count: self.active_run_turn_count, llm_call_count: self.llm_call_count, tool_execution_batch_count: self.tool_execution_batch_count, max_turns: self.max_turns, diff --git a/crates/agen/tests/engine_state_test.rs b/crates/agen/tests/engine_state_test.rs index 794fa88d..5bffd7e4 100644 --- a/crates/agen/tests/engine_state_test.rs +++ b/crates/agen/tests/engine_state_test.rs @@ -9,9 +9,12 @@ use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use agen::Item; +use agen::interceptor::{ + Interceptor, PreRequestAction, PreToolAction, ToolCallInfo, TurnEndAction, +}; use agen::llm_client::event::{Event, ResponseStatus, StatusEvent}; use agen::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput}; -use agen::{Engine, EngineError}; +use agen::{Engine, EngineError, EngineResult}; use async_trait::async_trait; use common::MockLlmClient; @@ -561,3 +564,185 @@ fn test_system_prompt_change_after_unlock() { let relocked = unlocked.lock(); assert_eq!(relocked.get_system_prompt(), Some("New prompt")); } + +fn completed_text_events() -> Vec { + vec![ + Event::text_block_start(0), + Event::text_delta(0, "done"), + Event::text_block_stop(0, None), + Event::Status(StatusEvent { + status: ResponseStatus::Completed, + }), + ] +} + +struct YieldOnce { + calls: AtomicUsize, +} + +#[async_trait] +impl Interceptor for YieldOnce { + async fn pre_llm_request(&self, _context: &mut Vec) -> PreRequestAction { + if self.calls.fetch_add(1, Ordering::SeqCst) == 0 { + PreRequestAction::Yield + } else { + PreRequestAction::Continue + } + } +} + +struct PauseToolOnce { + calls: AtomicUsize, +} + +#[async_trait] +impl Interceptor for PauseToolOnce { + async fn pre_tool_call(&self, _info: &mut ToolCallInfo) -> PreToolAction { + if self.calls.fetch_add(1, Ordering::SeqCst) == 0 { + PreToolAction::Pause + } else { + PreToolAction::Continue + } + } +} + +struct ContinueTurnOnce { + calls: AtomicUsize, +} + +#[async_trait] +impl Interceptor for ContinueTurnOnce { + async fn on_turn_end(&self, _history: &[Item]) -> TurnEndAction { + if self.calls.fetch_add(1, Ordering::SeqCst) == 0 { + TurnEndAction::ContinueWithMessages(vec![Item::system_message("continue")]) + } else { + TurnEndAction::Finish + } + } +} + +#[tokio::test] +async fn max_turns_is_scoped_to_each_fresh_run() { + let responses = vec![completed_text_events(), completed_text_events()]; + let mut engine = Engine::new(MockLlmClient::with_responses(responses)); + engine.set_max_turns(Some(1)); + let mut engine = engine.lock(); + + assert_eq!(engine.run("first").await.unwrap(), EngineResult::Finished); + assert_eq!(engine.turn_count(), 1); + assert_eq!(engine.active_run_turn_count(), None); + + assert_eq!(engine.run("second").await.unwrap(), EngineResult::Finished); + assert_eq!(engine.turn_count(), 2); + assert_eq!(engine.active_run_turn_count(), None); +} + +#[tokio::test] +async fn yielded_resume_keeps_the_same_unspent_turn_budget() { + let mut engine = Engine::new(MockLlmClient::new(completed_text_events())); + engine.set_max_turns(Some(1)); + engine.set_interceptor(YieldOnce { + calls: AtomicUsize::new(0), + }); + let mut engine = engine.lock(); + + assert_eq!(engine.run("start").await.unwrap(), EngineResult::Yielded); + assert_eq!(engine.turn_count(), 0); + assert_eq!(engine.active_run_turn_count(), Some(0)); + + assert_eq!(engine.resume().await.unwrap(), EngineResult::Finished); + assert_eq!(engine.turn_count(), 1); + assert_eq!(engine.active_run_turn_count(), None); +} + +#[tokio::test] +async fn paused_tool_resume_does_not_reset_the_consumed_turn_budget() { + let events = vec![ + Event::tool_use_start(0, "call_1", "count_tool"), + Event::tool_input_delta(0, "{}"), + Event::tool_use_stop(0), + Event::Status(StatusEvent { + status: ResponseStatus::Completed, + }), + ]; + let tool = CountingTool::new("count_tool"); + let mut engine = Engine::new(MockLlmClient::new(events)); + engine.set_max_turns(Some(1)); + engine.register_tool(tool.definition()); + engine.set_interceptor(PauseToolOnce { + calls: AtomicUsize::new(0), + }); + let mut engine = engine.lock(); + + assert_eq!(engine.run("call it").await.unwrap(), EngineResult::Paused); + assert_eq!(engine.turn_count(), 1); + assert_eq!(engine.active_run_turn_count(), Some(1)); + assert_eq!(tool.call_count(), 0); + + assert_eq!(engine.resume().await.unwrap(), EngineResult::LimitReached); + assert_eq!(engine.turn_count(), 1); + assert_eq!(engine.active_run_turn_count(), None); + assert_eq!(tool.call_count(), 1, "the consumed turn's tool still runs"); +} + +#[tokio::test] +async fn fresh_input_abandons_a_paused_run_and_starts_a_new_budget() { + let tool_events = vec![ + Event::tool_use_start(0, "call_1", "count_tool"), + Event::tool_input_delta(0, "{}"), + Event::tool_use_stop(0), + Event::Status(StatusEvent { + status: ResponseStatus::Completed, + }), + ]; + let client = MockLlmClient::with_responses(vec![tool_events, completed_text_events()]); + let tool = CountingTool::new("count_tool"); + let mut engine = Engine::new(client); + engine.set_max_turns(Some(1)); + engine.register_tool(tool.definition()); + engine.set_interceptor(PauseToolOnce { + calls: AtomicUsize::new(0), + }); + let mut engine = engine.lock(); + + assert_eq!(engine.run("pause").await.unwrap(), EngineResult::Paused); + assert_eq!(engine.active_run_turn_count(), Some(1)); + + assert_eq!(engine.run("replace").await.unwrap(), EngineResult::Finished); + assert_eq!(engine.turn_count(), 2); + assert_eq!(engine.active_run_turn_count(), None); + assert_eq!(tool.call_count(), 1, "pending-tool semantics are unchanged"); +} + +#[tokio::test] +async fn interceptor_continuation_consumes_the_logical_run_budget() { + let mut engine = Engine::new(MockLlmClient::new(completed_text_events())); + engine.set_max_turns(Some(1)); + engine.set_interceptor(ContinueTurnOnce { + calls: AtomicUsize::new(0), + }); + let mut engine = engine.lock(); + + assert_eq!( + engine.run("start").await.unwrap(), + EngineResult::LimitReached + ); + assert_eq!(engine.turn_count(), 1); + assert_eq!(engine.llm_call_count(), 1); + assert_eq!(engine.active_run_turn_count(), None); +} + +#[tokio::test] +async fn restored_active_run_budget_is_enforced_before_another_llm_call() { + let mut engine = Engine::new(MockLlmClient::new(completed_text_events())); + engine.set_max_turns(Some(1)); + engine.set_turn_count(7); + engine.set_last_run_interrupted(true); + engine.set_active_run_turn_count(Some(1)); + let mut engine = engine.lock(); + + assert_eq!(engine.resume().await.unwrap(), EngineResult::LimitReached); + assert_eq!(engine.turn_count(), 7); + assert_eq!(engine.llm_call_count(), 0); + assert_eq!(engine.active_run_turn_count(), None); +} diff --git a/crates/session-store/src/segment.rs b/crates/session-store/src/segment.rs index a0250ad6..5be5b320 100644 --- a/crates/session-store/src/segment.rs +++ b/crates/session-store/src/segment.rs @@ -307,6 +307,7 @@ pub fn save_run_completed( segment_id: SegmentId, result: EngineResult, interrupted: bool, + active_run_turn_count: Option, ) -> Result<(), StoreError> { append_entry( store, @@ -316,6 +317,7 @@ pub fn save_run_completed( ts: segment_log::now_millis(), interrupted, result, + active_run_turn_count, }, ) } diff --git a/crates/session-store/src/segment_log.rs b/crates/session-store/src/segment_log.rs index a97a1de6..5bf47454 100644 --- a/crates/session-store/src/segment_log.rs +++ b/crates/session-store/src/segment_log.rs @@ -125,11 +125,16 @@ pub enum LogEntry { TurnEnd { ts: u64, turn_count: usize }, /// `run()` / `resume()` が `EngineResult` で正常終了した。 - /// Audit-only metadata: replay は `interrupted` のみ反映する。 + /// Replay restores both interruption state and any resumable logical-run + /// turn budget. RunCompleted { ts: u64, interrupted: bool, result: EngineResult, + /// AgentTurns consumed by a paused/yielded logical run. Terminal + /// outcomes persist `None`. + #[serde(default, skip_serializing_if = "Option::is_none")] + active_run_turn_count: Option, }, /// `run()` / `resume()` が `EngineError` で終了した。 @@ -141,6 +146,15 @@ pub enum LogEntry { message: String, }, + /// Restores an active logical-run budget at a segment boundary, notably + /// after compaction replaced the segment that held the original Invoke and + /// RunCompleted entries. + ActiveRunCheckpoint { + ts: u64, + active_turn_count: usize, + total_turn_count: usize, + }, + /// A paused interrupted turn was explicitly abandoned without calling /// `run()` or `resume()` again. Replay clears the interrupted marker so /// the restored Worker is idle and future user input starts a normal new turn. @@ -209,6 +223,8 @@ pub struct RestoredState { pub config: RequestConfig, pub history: Vec, pub turn_count: usize, + /// AgentTurns consumed by the active paused/yielded logical run. + pub active_run_turn_count: Option, pub last_run_interrupted: bool, /// Number of entries replayed. `0` means the segment log was empty. /// Writers track their own append count via the same counter so @@ -238,6 +254,7 @@ pub fn collect_state(entries: &[LogEntry]) -> RestoredState { config: RequestConfig::default(), history: Vec::new(), turn_count: 0, + active_run_turn_count: None, last_run_interrupted: false, entries_count: 0, usage_history: Vec::new(), @@ -265,6 +282,7 @@ pub fn collect_state(entries: &[LogEntry]) -> RestoredState { // A terminal run record below clears or refines this. If the // log ends first, restore must treat the turn as interrupted. state.last_run_interrupted = true; + state.active_run_turn_count = Some(0); } LogEntry::UserInput { segments, @@ -290,16 +308,44 @@ pub fn collect_state(entries: &[LogEntry]) -> RestoredState { state.history.push(item.to_history_item()); } LogEntry::TurnEnd { turn_count, .. } => { + if let Some(active_turn_count) = &mut state.active_run_turn_count { + *active_turn_count += turn_count.saturating_sub(state.turn_count); + } state.turn_count = *turn_count; } - LogEntry::RunCompleted { interrupted, .. } => { + LogEntry::RunCompleted { + interrupted, + result, + active_run_turn_count, + .. + } => { state.last_run_interrupted = *interrupted; + if *interrupted && matches!(result, EngineResult::Paused | EngineResult::Yielded) { + // Legacy entries omit the explicit field; retain the + // Invoke/TurnEnd-derived count in that case. + if let Some(turn_count) = active_run_turn_count { + state.active_run_turn_count = Some(*turn_count); + } + } else { + state.active_run_turn_count = None; + } } LogEntry::RunErrored { interrupted, .. } => { state.last_run_interrupted = *interrupted; + state.active_run_turn_count = None; + } + LogEntry::ActiveRunCheckpoint { + active_turn_count, + total_turn_count, + .. + } => { + state.active_run_turn_count = Some(*active_turn_count); + state.turn_count = *total_turn_count; + state.last_run_interrupted = true; } LogEntry::PausedTurnAbandoned { .. } => { state.last_run_interrupted = false; + state.active_run_turn_count = None; } LogEntry::ConfigChanged { config, .. } => { state.config = config.clone(); @@ -397,6 +443,7 @@ mod tests { ts: 3200, interrupted: false, result: EngineResult::Finished, + active_run_turn_count: None, }, ]); assert_eq!(state.history.len(), 2); @@ -695,10 +742,93 @@ mod tests { ts: 100, interrupted: true, result: EngineResult::Paused, + active_run_turn_count: Some(1), }, LogEntry::PausedTurnAbandoned { ts: 200 }, ]); assert!(!state.last_run_interrupted); + assert_eq!(state.active_run_turn_count, None); + } + + #[test] + fn replay_restores_active_run_budget_across_compaction_checkpoint() { + let state = collect_state(&[ + LogEntry::SegmentStart { + ts: 0, + session_id: uuid::Uuid::nil(), + system_prompt: None, + config: RequestConfig::default(), + history: vec![], + forked_from: None, + compacted_from: None, + }, + LogEntry::ActiveRunCheckpoint { + ts: 100, + active_turn_count: 3, + total_turn_count: 9, + }, + ]); + + assert_eq!(state.turn_count, 9); + assert_eq!(state.active_run_turn_count, Some(3)); + assert!(state.last_run_interrupted); + } + + #[test] + fn legacy_interrupted_run_derives_budget_from_invoke_and_turn_end() { + let entry: LogEntry = serde_json::from_value(serde_json::json!({ + "kind": "run_completed", + "ts": 300, + "interrupted": true, + "result": "paused" + })) + .expect("legacy run-completed entry"); + let state = collect_state(&[ + LogEntry::SegmentStart { + ts: 0, + session_id: uuid::Uuid::nil(), + system_prompt: None, + config: RequestConfig::default(), + history: vec![], + forked_from: None, + compacted_from: None, + }, + LogEntry::Invoke { + ts: 100, + trigger: InvokeKind::UserSend, + }, + LogEntry::TurnEnd { + ts: 200, + turn_count: 2, + }, + entry, + ]); + + assert_eq!(state.active_run_turn_count, Some(2)); + assert!(state.last_run_interrupted); + } + + #[test] + fn non_resumable_interruption_clears_the_active_run_budget() { + let state = collect_state(&[ + LogEntry::Invoke { + ts: 100, + trigger: InvokeKind::UserSend, + }, + LogEntry::TurnEnd { + ts: 200, + turn_count: 2, + }, + LogEntry::RunCompleted { + ts: 300, + interrupted: true, + result: EngineResult::LimitReached, + active_run_turn_count: None, + }, + ]); + + assert!(state.last_run_interrupted); + assert_eq!(state.active_run_turn_count, None); } #[test] diff --git a/crates/session-store/tests/fs_store_test.rs b/crates/session-store/tests/fs_store_test.rs index 7838ad4d..63028d29 100644 --- a/crates/session-store/tests/fs_store_test.rs +++ b/crates/session-store/tests/fs_store_test.rs @@ -51,6 +51,7 @@ fn round_trip_write_and_read() { ts: 3200, interrupted: false, result: EngineResult::Finished, + active_run_turn_count: None, }, ]; diff --git a/crates/session-store/tests/session_test.rs b/crates/session-store/tests/session_test.rs index 557f4272..dda7205d 100644 --- a/crates/session-store/tests/session_test.rs +++ b/crates/session-store/tests/session_test.rs @@ -132,6 +132,7 @@ async fn run_and_persist( segment_id, r.clone(), worker.last_run_interrupted(), + worker.active_run_turn_count(), ) .unwrap(); } @@ -309,6 +310,7 @@ async fn session_resume_after_pause() { // Restore state and verify let state = session_store::restore(&store, sid, segid).unwrap(); assert!(state.last_run_interrupted); + assert_eq!(state.active_run_turn_count, Some(2)); } #[tokio::test] diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 03430c10..b85480ce 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -781,6 +781,7 @@ struct EmptyTurnRollbackSnapshot { usage_history_len: usize, ai_activity_count: usize, last_run_interrupted: bool, + active_run_turn_count: Option, flow_runtime_state: Option, } @@ -1777,6 +1778,8 @@ impl Worker { self.engine_mut().set_turn_count(state.turn_count); self.engine_mut() .set_last_run_interrupted(state.last_run_interrupted); + self.engine_mut() + .set_active_run_turn_count(state.active_run_turn_count); self.user_segments = state.user_segments; *self.usage_history.lock().expect("usage_history poisoned") = state.usage_history; *self @@ -2350,6 +2353,7 @@ impl Worker { usage_history_len, ai_activity_count: self.ai_activity_counter.load(Ordering::SeqCst), last_run_interrupted: self.engine().last_run_interrupted(), + active_run_turn_count: self.engine().active_run_turn_count(), flow_runtime_state: self .flow_runtime_state .lock() @@ -2381,6 +2385,8 @@ impl Worker { self.engine_mut().truncate_history(snapshot.history_len); self.engine_mut() .set_last_run_interrupted(snapshot.last_run_interrupted); + self.engine_mut() + .set_active_run_turn_count(snapshot.active_run_turn_count); *self .flow_runtime_state .lock() @@ -2774,6 +2780,10 @@ impl Worker { // can enter `Engine::resume` and execute them again after a crash. if self.engine.as_ref().unwrap().last_run_interrupted() { self.apply_interrupt_prep()?; + // Notification delivery begins a new logical run over the + // prepared history rather than inheriting the abandoned run's + // max-turn budget. + self.engine_mut().set_last_run_interrupted(false); } self.prepare_for_run().await?; @@ -3218,12 +3228,14 @@ impl Worker { } let interrupted = self.engine.as_ref().unwrap().last_run_interrupted(); + let active_run_turn_count = self.engine.as_ref().unwrap().active_run_turn_count(); match result { Ok(r) => { self.commit_entry(LogEntry::RunCompleted { ts: segment_log::now_millis(), interrupted, result: r.clone(), + active_run_turn_count, })?; } Err(e) => { @@ -3754,6 +3766,13 @@ impl Worker { }), }; let mut initial_entries = vec![entry.clone()]; + if let Some(active_turn_count) = w.active_run_turn_count() { + initial_entries.push(LogEntry::ActiveRunCheckpoint { + ts: segment_log::now_millis(), + active_turn_count, + total_turn_count: source_turn_count, + }); + } if let Some(flow_state) = self .flow_runtime_state .lock() @@ -5179,6 +5198,7 @@ where worker.set_request_config(state.config.clone()); worker.set_turn_count(state.turn_count); worker.set_last_run_interrupted(state.last_run_interrupted); + worker.set_active_run_turn_count(state.active_run_turn_count); if anchored_on_summary { worker.set_cache_anchor(Some(0)); }