From 8e3eb51813e0cb50eaee80c5fcec3b502f5387db Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:42:13 +0300 Subject: [PATCH 01/21] test(agent_loop): cover provider_started for failed summarizers Add tests asserting that a summarizer which reached its provider before failing without usage still marks provider_started, while one rejected before dispatch leaves it false. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/terminal_outcome_tests.rs | 77 +++++++++++++++++++ 1 file changed, 77 insertions(+) diff --git a/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs b/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs index 98bc9559c..4f0260ba0 100644 --- a/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs @@ -528,3 +528,80 @@ async fn a_cache_hit_followed_by_an_after_model_error_does_not_claim_the_provide assert!(partial.error.is_some()); assert!(!partial.run.terminal.expect("outcome").provider_started); } + +// --- provider_started for a summarizer that failed without usage --------- + +struct RejectingSummarizer; + +#[async_trait] +impl crate::summarization::Summarizer for RejectingSummarizer { + async fn summarize( + &self, + _: &[Message], + ) -> crate::error::Result { + Err(TinyAgentsError::Validation("rejected before dispatch".into())) + } +} + +fn long_input() -> Vec { + let mut input = vec![Message::system("sys")]; + for i in 0..12 { + input.push(Message::user(format!("question {i} {}", "x".repeat(400)))); + input.push(Message::Assistant(AssistantMessage { + id: None, + content: vec![ContentBlock::Text(format!("answer {i} {}", "y".repeat(400)))], + tool_calls: Vec::new(), + usage: None, + origin: None, + })); + } + input.push(Message::user("now")); + input +} + +fn aborting_harness( + summarizer: Box, +) -> AgentHarness<()> { + use crate::middleware::{CompressionFailurePolicy, ContextCompressionMiddleware}; + let policy = crate::summarization::SummarizationPolicy { + keep_last: 2, + ..Default::default() + } + .with_trigger_override(100); + let mw = Arc::new( + ContextCompressionMiddleware::with_summarizer(policy, summarizer) + .with_failure_policy(CompressionFailurePolicy::Abort), + ); + let mut harness = harness_with(Arc::new(ScriptedModel::new(vec![response(vec![], "x")]))); + harness.push_middleware(mw.clone()); + harness.push_model_middleware(mw); + harness +} + +#[tokio::test] +async fn a_summarizer_that_dispatched_and_failed_without_usage_marks_provider_started() { + let summarizer = crate::summarization::ModelSummarizer::new(Arc::new(FailingModel("boom")), "m"); + let harness = aborting_harness(Box::new(summarizer)); + let ctx = RunContext::new(RunConfig::new("sum-dispatched"), ()); + let partial = harness + .invoke_in_context_collecting_partial(&(), ctx, long_input()) + .await; + assert!(partial.error.is_some()); + let outcome = partial.run.terminal.expect("outcome"); + assert!( + outcome.provider_started, + "the summarizer reached its provider" + ); +} + +#[tokio::test] +async fn a_summarizer_rejected_before_dispatch_leaves_provider_started_false() { + let harness = aborting_harness(Box::new(RejectingSummarizer)); + let ctx = RunContext::new(RunConfig::new("sum-rejected"), ()); + let partial = harness + .invoke_in_context_collecting_partial(&(), ctx, long_input()) + .await; + assert!(partial.error.is_some()); + let outcome = partial.run.terminal.expect("outcome"); + assert!(!outcome.provider_started); +} From 8fd9c624c11662a68bc0356605e5e345618c1b80 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:42:53 +0300 Subject: [PATCH 02/21] fix(harness): mark provider_started when a summarizer dispatched without usage Co-authored-by: Medulla --- crates/tinyagents-harness/src/context/mod.rs | 8 +++++ .../src/middleware/library/context.rs | 2 +- .../middleware/library/context/overflow.rs | 2 +- .../src/middleware/library/context/summary.rs | 19 ++++++++++- .../src/summarization/dispatch.rs | 34 +++++++++++++++++++ .../src/summarization/mod.rs | 1 + .../src/summarization/model_summarizer.rs | 1 + 7 files changed, 64 insertions(+), 3 deletions(-) create mode 100644 crates/tinyagents-harness/src/summarization/dispatch.rs diff --git a/crates/tinyagents-harness/src/context/mod.rs b/crates/tinyagents-harness/src/context/mod.rs index 3b02c122b..f0e797201 100644 --- a/crates/tinyagents-harness/src/context/mod.rs +++ b/crates/tinyagents-harness/src/context/mod.rs @@ -880,6 +880,14 @@ impl RunContext { self.call_provider_started = true; } + /// Records that a summarizer (a compaction, outside the model-call layer) + /// dispatched a provider call. Sets only the run-wide flag: the per-call + /// flag still describes the main model call. + pub(crate) fn mark_summarizer_dispatched(&mut self) { + tracing::debug!("[tinyagents::run] summarizer dispatched a provider call"); + self.provider_started = true; + } + #[doc(hidden)] pub fn mark_model_call_failed(&mut self) { self.model_call_failed = true; diff --git a/crates/tinyagents-harness/src/middleware/library/context.rs b/crates/tinyagents-harness/src/middleware/library/context.rs index 03cc711d3..24647a3f7 100644 --- a/crates/tinyagents-harness/src/middleware/library/context.rs +++ b/crates/tinyagents-harness/src/middleware/library/context.rs @@ -619,7 +619,7 @@ impl ContextCompressionMiddleware { ); let started = std::time::Instant::now(); let record = match self - .summarize_batch(&to_summarize, &to_keep, previous_summary) + .summarize_batch(ctx, &to_summarize, &to_keep, previous_summary) .await { Ok(record) => record, diff --git a/crates/tinyagents-harness/src/middleware/library/context/overflow.rs b/crates/tinyagents-harness/src/middleware/library/context/overflow.rs index fee192a20..ee659886e 100644 --- a/crates/tinyagents-harness/src/middleware/library/context/overflow.rs +++ b/crates/tinyagents-harness/src/middleware/library/context/overflow.rs @@ -225,7 +225,7 @@ impl ContextCompressionMiddleware { ); let started = std::time::Instant::now(); let record = match self - .summarize_batch(&to_summarize, &to_keep, previous_summary) + .summarize_batch(ctx, &to_summarize, &to_keep, previous_summary) .await { Ok(record) => record, diff --git a/crates/tinyagents-harness/src/middleware/library/context/summary.rs b/crates/tinyagents-harness/src/middleware/library/context/summary.rs index 3119ddc2d..430bb8b55 100644 --- a/crates/tinyagents-harness/src/middleware/library/context/summary.rs +++ b/crates/tinyagents-harness/src/middleware/library/context/summary.rs @@ -34,7 +34,24 @@ impl ContextCompressionMiddleware { /// with the lists the previous summary carried. The previous summary /// reaches the summarizer without its lists, and any the summarizer /// echoes are dropped, so each list appears exactly once. - pub(in crate::middleware::library) async fn summarize_batch( + pub(in crate::middleware::library) async fn summarize_batch( + &self, + ctx: &mut RunContext, + to_summarize: &[Message], + to_keep: &[Message], + previous_summary: Option, + ) -> Result { + let (result, dispatched) = crate::summarization::dispatch::track_dispatch( + self.summarize_batch_inner(to_summarize, to_keep, previous_summary), + ) + .await; + if dispatched { + ctx.mark_summarizer_dispatched(); + } + result + } + + async fn summarize_batch_inner( &self, to_summarize: &[Message], to_keep: &[Message], diff --git a/crates/tinyagents-harness/src/summarization/dispatch.rs b/crates/tinyagents-harness/src/summarization/dispatch.rs new file mode 100644 index 000000000..22f290318 --- /dev/null +++ b/crates/tinyagents-harness/src/summarization/dispatch.rs @@ -0,0 +1,34 @@ +//! Tracks whether a summarizer reached its provider. +//! +//! Summarizer calls bypass the run context's dispatch marker, so a compaction +//! that fails without usage (a transport error before any metered response) +//! would otherwise leave `TerminalOutcome::provider_started` false although a +//! provider call was made. The middleware scopes each summarization with +//! [`track_dispatch`]; a model-backed summarizer calls [`mark_dispatched`] at +//! the moment it hands the request to its model. A rejection before that point +//! (empty input, validation) never sets the flag. + +use std::future::Future; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; + +tokio::task_local! { + static DISPATCHED: Arc; +} + +/// Records that the summarizer in scope is about to dispatch a provider call. +/// A no-op outside a [`track_dispatch`] scope. +pub(crate) fn mark_dispatched() { + let _ = DISPATCHED.try_with(|flag| { + flag.store(true, Ordering::Relaxed); + }); +} + +/// Runs `fut`, returning its output and whether it dispatched a provider call. +pub(crate) async fn track_dispatch(fut: F) -> (F::Output, bool) { + let flag = Arc::new(AtomicBool::new(false)); + let out = DISPATCHED.scope(Arc::clone(&flag), fut).await; + let dispatched = flag.load(Ordering::Relaxed); + tracing::trace!(dispatched, "[tinyagents::summarize] dispatch tracked"); + (out, dispatched) +} diff --git a/crates/tinyagents-harness/src/summarization/mod.rs b/crates/tinyagents-harness/src/summarization/mod.rs index 250a18d69..63b252e38 100644 --- a/crates/tinyagents-harness/src/summarization/mod.rs +++ b/crates/tinyagents-harness/src/summarization/mod.rs @@ -25,6 +25,7 @@ mod checkpoint; pub mod compaction; +pub(crate) mod dispatch; mod file_ops; mod model_summarizer; pub mod pairing; diff --git a/crates/tinyagents-harness/src/summarization/model_summarizer.rs b/crates/tinyagents-harness/src/summarization/model_summarizer.rs index d69e827ef..2a67e0cc6 100644 --- a/crates/tinyagents-harness/src/summarization/model_summarizer.rs +++ b/crates/tinyagents-harness/src/summarization/model_summarizer.rs @@ -198,6 +198,7 @@ impl ModelSummarizer { let mut last_chars = 0; let mut usage: Option = None; for attempt in 1..=SUMMARY_MARKUP_ATTEMPTS { + super::dispatch::mark_dispatched(); let response = self.model.invoke(&(), request.clone()).await.map_err(|e| { tracing::warn!(error = %e, "[tinyagents::summarize] summarizer model call failed"); Box::new(( From bea2bca7b08bfd432c474dacf4499367370083c2 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:45:07 +0300 Subject: [PATCH 03/21] test(integration): assert graph driver emits turn lifecycle events The graph driver now announces turn and message lifecycle events like the direct loop, so the parity filter that excluded them is gone and the expected kinds are matched in full. The test pinning the old direct-loop-only gap was removed, with the exact event shape covered by the new graph_lifecycle_events test. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../tests/graph_lifecycle_events.rs | 252 ++++++++++++++++++ .../tests/loop_as_graph.rs | 36 +-- 2 files changed, 256 insertions(+), 32 deletions(-) create mode 100644 crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs new file mode 100644 index 000000000..c0dbb2c86 --- /dev/null +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -0,0 +1,252 @@ +//! The graph loop driver announces the same turn/message lifecycle events as +//! the harness's direct loop (`TurnStarted`, `TurnCompleted`, +//! `MessageAppended`), once each, and never for nested tool calls. + +use std::sync::Arc; + +use async_trait::async_trait; +use serde_json::{Value, json}; + +use tinyagents_graph::agent_loop::{AgentLoopGraphExt, GraphLoopDriver}; +use tinyagents_harness::context::{RunConfig, RunContext}; +use tinyagents_harness::events::AgentEvent; +use tinyagents_harness::limits::RunLimits; +use tinyagents_harness::runtime::{AgentHarness, LoopExecution, RunPolicy}; +use tinyagents_harness::steering::{SteeringCommand, SteeringHandle}; +use tinyagents_harness::testkit::{EventRecorder, FakeTool}; +use tinyagents_harness::tool::ToolExecutionContext; +use tinyinference_llm::message::{AssistantMessage, Message}; +use tinyinference_llm::model::ModelResponse; +use tinyinference_llm::providers::MockModel; +use tinyinference_llm::tool::ToolCall; +use tinytools::{Tool, ToolResult}; + +fn tool_call_response(id: &str, name: &str) -> ModelResponse { + ModelResponse { + message: AssistantMessage { + id: Some(format!("msg-{id}")), + content: Vec::new(), + tool_calls: vec![ToolCall::new(id, name, json!({}))], + usage: None, + origin: None, + }, + finish_reason: Some("tool_calls".to_string()), + ..ModelResponse::assistant("") + } +} + +fn harness_for( + execution: LoopExecution, + responses: Vec, + limits: RunLimits, +) -> AgentHarness<()> { + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::new(MockModel::with_responses(responses))) + .set_default_model("mock"); + if matches!(execution, LoopExecution::Graph) { + harness.with_loop_driver(Arc::new(GraphLoopDriver::new())); + } + harness.with_policy(RunPolicy { + execution, + limits, + ..RunPolicy::default() + }); + harness +} + +/// The lifecycle events of a run, flattened to comparable strings. +fn lifecycle(events: &[AgentEvent]) -> Vec { + events + .iter() + .filter_map(|event| match event { + AgentEvent::TurnStarted { turn } => Some(format!("turn.started:{turn}")), + AgentEvent::TurnCompleted { + turn, + tool_call_ids, + .. + } => Some(format!( + "turn.completed:{turn}:{}", + tool_call_ids + .iter() + .map(ToString::to_string) + .collect::>() + .join(",") + )), + AgentEvent::MessageAppended { + role, + index, + call_id, + .. + } => Some(format!( + "append:{index}:{role}:{}", + call_id.as_ref().map(ToString::to_string).unwrap_or_default() + )), + _ => None, + }) + .collect() +} + +async fn run_with( + execution: LoopExecution, + responses: Vec, + input: Vec, + steering: Option, +) -> Vec { + let mut harness = harness_for(execution, responses, RunLimits::default()); + harness.register_tool(Arc::new(FakeTool::returning("lookup", "tool-output"))); + let recorder = EventRecorder::new(); + let mut ctx = RunContext::new(RunConfig::new("lifecycle"), ()).with_events(recorder.sink()); + if let Some(steering) = steering { + ctx = ctx.with_steering(steering); + } + harness + .invoke_in_context(&(), ctx, input) + .await + .expect("run completes"); + lifecycle(&recorder.events()) +} + +#[tokio::test] +async fn a_plain_run_announces_the_same_lifecycle_as_the_direct_loop() { + let script = || vec![ModelResponse::assistant("done")]; + let input = || vec![Message::user("hi")]; + let direct = run_with(LoopExecution::Direct, script(), input(), None).await; + let graph = run_with(LoopExecution::Graph, script(), input(), None).await; + assert!(!direct.is_empty(), "the direct loop emits lifecycle events"); + assert_eq!(graph, direct); +} + +#[tokio::test] +async fn a_tool_run_announces_each_message_once_and_closes_each_turn() { + let script = || { + vec![ + tool_call_response("call-1", "lookup"), + ModelResponse::assistant("done"), + ] + }; + let input = || vec![Message::user("look it up")]; + let direct = run_with(LoopExecution::Direct, script(), input(), None).await; + let graph = run_with(LoopExecution::Graph, script(), input(), None).await; + assert_eq!(graph, direct); + + // Input is never announced; every later index appears exactly once. + let mut indices: Vec<&String> = graph.iter().filter(|e| e.starts_with("append:")).collect(); + let total = indices.len(); + indices.dedup(); + assert_eq!(indices.len(), total, "no duplicate MessageAppended: {graph:?}"); + assert!(!graph.iter().any(|e| e.starts_with("append:0:")), "{graph:?}"); + assert!(graph.iter().any(|e| e == "turn.completed:1:call-1"), "{graph:?}"); +} + +#[tokio::test] +async fn a_steering_injected_message_is_announced_like_the_direct_loop() { + let run = |execution| async move { + let steering = SteeringHandle::allow_all(); + steering.send(SteeringCommand::InjectMessage(Message::user("extra"))); + run_with( + execution, + vec![ModelResponse::assistant("done")], + vec![Message::user("hello")], + Some(steering), + ) + .await + }; + assert_eq!(run(LoopExecution::Graph).await, run(LoopExecution::Direct).await); +} + +/// A tool that calls another tool through the harness. +struct Caller; + +#[async_trait] +impl Tool for Caller { + fn name(&self) -> &str { + "caller" + } + fn description(&self) -> &str { + "calls another tool" + } + fn parameters_schema(&self) -> Value { + json!({"type": "object"}) + } + async fn execute(&self, _arguments: Value) -> anyhow::Result { + unreachable!("dispatched through execute_with_context") + } + async fn execute_with_context( + &self, + _arguments: Value, + _options: tinytools::ToolCallOptions, + context: Option<&dyn tinytools::ToolRunContext>, + ) -> anyhow::Result { + let harness = context + .and_then(tinytools::ToolRunContext::host_extension) + .and_then(|any| any.downcast_ref::()) + .expect("the harness installs its context") + .clone(); + harness.call_tool("lookup", json!({})).await?; + Ok(ToolResult::success("caller-out")) + } +} + +#[tokio::test] +async fn nested_tool_calls_do_not_emit_message_appended_on_the_graph_driver() { + let mut rows = Vec::new(); + for execution in [LoopExecution::Direct, LoopExecution::Graph] { + let mut harness = harness_for( + execution, + vec![ + tool_call_response("p1", "caller"), + ModelResponse::assistant("done"), + ], + RunLimits::default().with_max_nested_depth(3), + ); + harness.register_tool(Arc::new(FakeTool::returning("lookup", "leaf-out"))); + harness.register_tool(Arc::new(Caller)); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("nested"), ()).with_events(recorder.sink()); + harness + .invoke_in_context(&(), ctx, vec![Message::user("go")]) + .await + .expect("run completes"); + let tool_rows: Vec = lifecycle(&recorder.events()) + .into_iter() + .filter(|e| e.contains(":tool:")) + .collect(); + assert_eq!(tool_rows.len(), 1, "{execution:?}: {tool_rows:?}"); + assert!(tool_rows[0].ends_with(":tool:p1"), "{tool_rows:?}"); + rows.push(tool_rows); + } + assert_eq!(rows[0], rows[1]); +} + +#[tokio::test] +async fn iter_stepping_announces_appends_once_and_never_the_input() { + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model( + "mock", + Arc::new(MockModel::with_responses(vec![ + tool_call_response("c1", "lookup"), + ModelResponse::assistant("done"), + ])), + ) + .set_default_model("mock") + .register_tool(Arc::new(FakeTool::returning("lookup", "out"))); + let harness = Arc::new(harness); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("iter-lifecycle"), ()).with_events(recorder.sink()); + let mut iter = harness + .iter(Arc::new(()), ctx, vec![Message::user("go")]) + .expect("iter starts"); + iter.run_to_end().await.expect("finishes"); + + let events = lifecycle(&recorder.events()); + assert!(!events.iter().any(|e| e.starts_with("append:0:")), "{events:?}"); + let appended: Vec<&String> = events.iter().filter(|e| e.starts_with("append:")).collect(); + assert_eq!(appended.len(), 3, "assistant, tool, assistant: {events:?}"); + assert_eq!( + events.iter().filter(|e| e.starts_with("turn.started")).count(), + events.iter().filter(|e| e.starts_with("turn.completed")).count(), + "{events:?}" + ); +} diff --git a/crates/tinyagents-integration-tests/tests/loop_as_graph.rs b/crates/tinyagents-integration-tests/tests/loop_as_graph.rs index 525337612..6ea9c0342 100644 --- a/crates/tinyagents-integration-tests/tests/loop_as_graph.rs +++ b/crates/tinyagents-integration-tests/tests/loop_as_graph.rs @@ -67,14 +67,12 @@ fn harness_for(execution: LoopExecution, model: Arc) -> AgentHarness< /// subsequence of `actual` — the "same kind sequence, extra graph events /// allowed" contract. /// -/// The turn/message lifecycle events (`turn.*`, `message.appended`) are emitted -/// by the direct loop only for now; the graph driver does not announce them, so -/// they are excluded from the expected sequence. +/// The turn/message lifecycle events (`turn.*`, `message.appended`) are part of +/// the contract: the graph driver announces them like the direct loop does, and +/// `graph_lifecycle_events.rs` pins their exact shape. fn assert_kinds_subsequence(expected: &[String], actual: &[String]) { let mut cursor = 0; - let lifecycle = - |kind: &&String| !(kind.starts_with("turn.") || kind.as_str() == "message.appended"); - for kind in expected.iter().filter(lifecycle) { + for kind in expected { let Some(offset) = actual[cursor..].iter().position(|k| k == kind) else { panic!( "expected event kind `{kind}` not found (in order) in graph run's kinds: \ @@ -868,29 +866,3 @@ async fn a_middleware_limit_error_is_not_a_tool_cap_partial_stop() { ); } } - -/// The turn/message lifecycle events are direct-loop only for now (a documented -/// follow-up). Pin that so the gap is explicit and this parity file's filter -/// cannot silently hide a change in either direction. -#[tokio::test] -async fn lifecycle_events_are_direct_loop_only_until_the_graph_driver_emits_them() { - for (execution, expect_lifecycle) in - [(LoopExecution::Direct, true), (LoopExecution::Graph, false)] - { - let model = Arc::new(MockModel::with_responses(vec![ModelResponse::assistant( - "done", - )])); - let harness = harness_for(execution, model); - let recorder = EventRecorder::new(); - let ctx = RunContext::new(RunConfig::new("lc"), ()).with_events(recorder.sink()); - harness - .invoke_in_context(&(), ctx, vec![Message::user("hi")]) - .await - .unwrap(); - let has_lifecycle = recorder - .events() - .iter() - .any(|event| event.kind() == "turn.started" || event.kind() == "message.appended"); - assert_eq!(has_lifecycle, expect_lifecycle, "{execution:?}"); - } -} From 3e4dafc1a5de2ac64bfcfb89e3e21d42b3bb6bbc Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:45:28 +0300 Subject: [PATCH 04/21] refactor(agent_loop): split compile and runtime into separate modules The agent loop's compilation and runtime logic now live in dedicated modules, with lifecycle and phase handling moved into the harness crate. This separates graph construction from execution so each can evolve independently. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/compile.rs | 4 +- .../src/agent_loop/runtime.rs | 22 ++++++++++ .../src/agent_loop/lifecycle.rs | 25 +++++++++++ .../src/agent_loop/phases.rs | 42 +++++++++++++++++++ 4 files changed, 92 insertions(+), 1 deletion(-) diff --git a/crates/tinyagents-graph/src/agent_loop/compile.rs b/crates/tinyagents-graph/src/agent_loop/compile.rs index 7ad4b9ac5..306a355b1 100644 --- a/crates/tinyagents-graph/src/agent_loop/compile.rs +++ b/crates/tinyagents-graph/src/agent_loop/compile.rs @@ -137,8 +137,10 @@ where let rt = rt.clone(); async move { let harness = rt.harness.clone(); + let mut ctx_guard = rt.ctx.lock().await; let mut run_guard = rt.run.lock().await; - runtime::settle_node(&harness, &mut run_guard, loop_state).await + runtime::settle_node(&harness, &mut ctx_guard, &mut run_guard, loop_state) + .await } } }) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 30543ff4a..f9636eabf 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -209,6 +209,10 @@ where if ctx.cancellation.is_cancelled() { return Err(TinyAgentsError::Cancelled); } + // Everything on the transcript when the loop is first entered is input (or + // a resumed run's already-announced history); only later appends are + // announced. A no-op after the first activation. + phases::lifecycle_seed(ctx, loop_state.messages.len()); match apply_pending_steering(ctx, &mut loop_state.messages)? { SteeringOutcome::Cancel => return Err(TinyAgentsError::Cancelled), SteeringOutcome::Pause => { @@ -432,6 +436,15 @@ where status.active_model_call = Some(call_id.clone()); ctx.active_model_call = Some(call_id.clone()); ctx.begin_model_call(); + // Same point as the direct loop: pending appends (steering) are announced, + // the previous turn closed, and this one opened, just before the call. + let turn = phases::lifecycle_start_turn(harness, ctx, &loop_state.messages); + tracing::debug!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + turn, + "[graph_loop] turn started" + ); let base = DirectModelBase { model: binding.model.as_ref(), @@ -488,6 +501,7 @@ where response.message.clone(), )); loop_state.turn += 1; + phases::lifecycle_flush(harness, ctx, &loop_state.messages); let tool_calls = response.tool_calls().to_vec(); loop_state.pending_tool_calls = tool_calls.clone(); @@ -572,6 +586,9 @@ where } Err(error) => return Err(error), }; + // Every tool result of this batch is on the transcript: announce them and + // close the turn, as the direct loop does after its batch. + phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); loop_state.tool_calls = run.tool_calls; loop_state.executed_tools = run.executed_tools.clone(); let _ = outcome; @@ -592,6 +609,7 @@ where /// (`RunPolicy::output_retry`), and finishes the run. pub(crate) async fn settle_node( harness: &AgentHarness, + ctx: &mut RunContext, run: &mut AgentRun, mut loop_state: LoopState, ) -> Result> @@ -599,6 +617,10 @@ where State: Send + Sync, Ctx: Send + Sync, { + // The final assistant message is the last append of its turn; close the + // turn before any output-retry prompt is pushed (that prompt is announced + // with the next turn's start). + phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); if let Some(plan) = loop_state.pending_structured.take() { let extractor = StructuredExtractor::new( plan.strategy.clone(), diff --git a/crates/tinyagents-harness/src/agent_loop/lifecycle.rs b/crates/tinyagents-harness/src/agent_loop/lifecycle.rs index 0d1df4db3..e25702dad 100644 --- a/crates/tinyagents-harness/src/agent_loop/lifecycle.rs +++ b/crates/tinyagents-harness/src/agent_loop/lifecycle.rs @@ -29,6 +29,10 @@ pub(crate) struct TurnTracker { turn: u32, /// The open turn and the transcript index it started at. open: Option<(u32, usize)>, + /// Whether the seed has been fixed. A context that never went through + /// [`TurnTracker::new`] (a graph run entered mid-flight) is seeded on first + /// use by [`TurnTracker::ensure_seeded`]. + seeded: bool, } pub(crate) fn role_of(message: &Message) -> &'static str { @@ -50,6 +54,20 @@ impl TurnTracker { seed_len, turn: 0, open: None, + seeded: true, + } + } + + /// Treats the first `len` messages as seed input if no seed was fixed yet. + /// A no-op once seeded, so repeated node entries never swallow appends. + pub(crate) fn ensure_seeded(&mut self, len: usize) { + if !self.seeded { + tracing::debug!( + target: "tinyagents::agent_loop", + seed_len = len, + "[agent_loop] lifecycle tracker seeded on first use" + ); + *self = Self::new(len); } } @@ -207,3 +225,10 @@ impl RunContext { self.turns.rebase(&self.events, new_len, reason); } } + +impl RunContext { + /// Fixes the lifecycle seed at `len` messages unless one is already set. + pub(crate) fn ensure_turn_tracker_seeded(&mut self, len: usize) { + self.turns.ensure_seeded(len); + } +} diff --git a/crates/tinyagents-harness/src/agent_loop/phases.rs b/crates/tinyagents-harness/src/agent_loop/phases.rs index af49461d9..1275e90b0 100644 --- a/crates/tinyagents-harness/src/agent_loop/phases.rs +++ b/crates/tinyagents-harness/src/agent_loop/phases.rs @@ -213,3 +213,45 @@ pub async fn execute_tool_batch( executed_tools: run.executed_tools[executed_before..].to_vec(), }) } + +// ── Lifecycle events for alternate drivers ───────────────────────────────── + +/// Fixes the lifecycle seed for a driver that enters a run mid-flight (a node +/// activation of the compiled-graph loop, possibly resumed from a checkpoint): +/// the first `transcript_len` messages are input and are never announced. +/// A no-op once a seed exists, so it is safe to call on every node entry. +pub fn lifecycle_seed(ctx: &mut RunContext, transcript_len: usize) { + ctx.ensure_turn_tracker_seeded(transcript_len); +} + +/// Announces transcript appends not yet announced (`MessageAppended`), using +/// the harness's payload-capture policy. Mirrors the direct loop's flush +/// points; tool calls nested inside a tool never reach the transcript, so they +/// are never announced. +pub fn lifecycle_flush( + harness: &AgentHarness, + ctx: &mut RunContext, + messages: &[Message], +) { + ctx.flush_transcript(harness.policy().capture, messages); +} + +/// Opens the next turn (`TurnStarted`), first announcing pending appends and +/// closing any turn still open. Call right before the model call. +pub fn lifecycle_start_turn( + harness: &AgentHarness, + ctx: &mut RunContext, + messages: &[Message], +) -> u32 { + ctx.start_turn(harness.policy().capture, messages) +} + +/// Announces pending appends and closes the open turn (`TurnCompleted`), if +/// any. Idempotent: closing with no open turn only flushes. +pub fn lifecycle_close_turn( + harness: &AgentHarness, + ctx: &mut RunContext, + messages: &[Message], +) { + ctx.close_turn(harness.policy().capture, messages); +} From 09a28231aa7bfd80cdb38fbf0ff1e0f286da5d0b Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:46:00 +0300 Subject: [PATCH 05/21] fix(agent_loop): emit turn lifecycle events on graph loop exit The graph driver now closes the current turn and flushes any messages appended by the after-agent middleware, matching the direct loop so transcript events are announced on every exit path. The settle node receives the run context so it can emit those events. Auto-committed-on: dragonfly Co-authored-by: Medulla --- crates/tinyagents-graph/src/agent_loop/driver.rs | 11 +++++++++-- crates/tinyagents-graph/src/agent_loop/iter.rs | 4 +++- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/driver.rs b/crates/tinyagents-graph/src/agent_loop/driver.rs index 374983571..5c881353d 100644 --- a/crates/tinyagents-graph/src/agent_loop/driver.rs +++ b/crates/tinyagents-graph/src/agent_loop/driver.rs @@ -12,7 +12,7 @@ use async_trait::async_trait; -use tinyagents_harness::agent_loop::phases::LoopDriver; +use tinyagents_harness::agent_loop::phases::{self, LoopDriver}; use tinyagents_harness::context::RunContext; use tinyagents_harness::error::{Result, TinyAgentsError}; use tinyagents_harness::events::{AgentEvent, HarnessRunStatus}; @@ -107,7 +107,7 @@ where node::TOOLS => { runtime::tools_node(harness, state, ctx, run, status, loop_state).await } - node::SETTLE => runtime::settle_node(harness, run, loop_state).await, + node::SETTLE => runtime::settle_node(harness, ctx, run, loop_state).await, other => { break Err(TinyAgentsError::Validation(format!( "GraphLoopDriver: unknown loop node `{other}`" @@ -263,9 +263,16 @@ where Some(outcome) } }; + // Announce whatever the last turn appended and close it, on every exit + // path, as the direct loop does before the transcript moves onto the + // run (`run.messages` is the transcript as of the last node boundary). + phases::lifecycle_close_turn(harness, ctx, &run.messages); run.terminal = terminal.clone(); status.mark_running(HarnessPhase::Middleware); let after_agent = harness.middleware().run_after_agent(ctx, state, run).await; + // `after_agent` may post-process `run.messages`; announce anything it + // appended (even if it then failed), as the direct loop does. + phases::lifecycle_flush(harness, ctx, &run.messages); if let Err(hook_error) = after_agent { if outcome.is_err() { // The originating node failure stays authoritative, as in the diff --git a/crates/tinyagents-graph/src/agent_loop/iter.rs b/crates/tinyagents-graph/src/agent_loop/iter.rs index 9e374d1a2..431aa1b04 100644 --- a/crates/tinyagents-graph/src/agent_loop/iter.rs +++ b/crates/tinyagents-graph/src/agent_loop/iter.rs @@ -206,7 +206,9 @@ where ) .await? } - node::SETTLE => runtime::settle_node(&harness, &mut run_guard, loop_state).await?, + node::SETTLE => { + runtime::settle_node(&harness, &mut ctx_guard, &mut run_guard, loop_state).await? + } other => { return Err(TinyAgentsError::Validation(format!( "LoopIter::next: unknown loop node `{other}`" From e5d31e0b1cd3aeb2bf065367e3060d0cb4da49e1 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:46:38 +0300 Subject: [PATCH 06/21] feat(graph): emit turn and message lifecycle events from the graph driver Co-authored-by: Medulla --- .../tinyagents-graph/src/agent_loop/runtime.rs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index f9636eabf..5106713e3 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -428,16 +428,8 @@ where request.reasoning = Some(mapped.clone()); } - let started_record = ctx.emit(AgentEvent::ModelStarted { - call_id: call_id.clone(), - model: model_name.clone(), - }); - status.set_last_event(started_record.id); - status.active_model_call = Some(call_id.clone()); - ctx.active_model_call = Some(call_id.clone()); - ctx.begin_model_call(); // Same point as the direct loop: pending appends (steering) are announced, - // the previous turn closed, and this one opened, just before the call. + // the previous turn closed, and this one opened, just before `ModelStarted`. let turn = phases::lifecycle_start_turn(harness, ctx, &loop_state.messages); tracing::debug!( target: "tinyagents::agent_loop", @@ -446,6 +438,14 @@ where "[graph_loop] turn started" ); + let started_record = ctx.emit(AgentEvent::ModelStarted { + call_id: call_id.clone(), + model: model_name.clone(), + }); + status.set_last_event(started_record.id); + status.active_model_call = Some(call_id.clone()); + ctx.active_model_call = Some(call_id.clone()); + ctx.begin_model_call(); let base = DirectModelBase { model: binding.model.as_ref(), }; From d0a96d3b7e7a5f01b91b3cfc4ccc2a6bb85f102a Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:48:57 +0300 Subject: [PATCH 07/21] fix(orchestration): record cancelled and failed subagent children with the completion router Co-authored-by: Medulla --- .../src/subagent/completion.rs | 47 +++- .../src/subagent/driver.rs | 70 ++++- .../src/subagent/driver_completion_tests.rs | 263 ++++++++++++++++-- .../src/subagent/outcome_status_map.rs | 3 +- 4 files changed, 339 insertions(+), 44 deletions(-) diff --git a/crates/tinyagents-orchestration/src/subagent/completion.rs b/crates/tinyagents-orchestration/src/subagent/completion.rs index fcab99b7d..d2e2e8f52 100644 --- a/crates/tinyagents-orchestration/src/subagent/completion.rs +++ b/crates/tinyagents-orchestration/src/subagent/completion.rs @@ -1,4 +1,5 @@ -//! Recording a finished subagent with the durable completion router. +//! Recording a finished subagent with the durable completion router: normal +//! finishes, cancellations and executor errors; never a pause. //! //! Opt-in: the driver only touches any of this when a //! [`CompletionRouter`](tinyagents_tasks::CompletionRouter) was configured with @@ -12,8 +13,8 @@ use tinyagents_tasks::{ }; use super::{ - AppliedResult, ArtifactReference, PreparedSubagent, SubagentOutcome, SubagentOutcomeKind, - SubagentTaskKey, + AppliedResult, ArtifactReference, PreparedSubagent, SubagentError, SubagentOutcome, + SubagentOutcomeKind, SubagentTaskKey, }; const LOG_PREFIX: &str = "[subagent-completion]"; @@ -81,19 +82,14 @@ impl CompletionOrigin { } /// The record for a finished lifecycle, or `None` when the outcome is not a - /// completion: a cancellation is the parent's own doing, and a pause is not - /// final (the same task id completes later). + /// completion: a pause is not final (the same task id completes later, and + /// that finish is what gets recorded). A cancellation is recorded as + /// [`CompletionStatus::Cancelled`] with an empty result. pub(crate) fn record_for_outcome( &self, outcome: &SubagentOutcome, omitted_chars: usize, ) -> Option { - // A cancellation is the parent's own doing: a mapping exists, but it is - // routing policy not to record it. A pause is not final (and has no - // completion status). - if matches!(&outcome.status, SubagentOutcomeKind::Cancelled) { - return None; - } let status = CompletionStatus::try_from(&outcome.status).ok()?; let text = match &outcome.status { SubagentOutcomeKind::Incomplete(incomplete) if outcome.output.is_empty() => { @@ -108,8 +104,37 @@ impl CompletionOrigin { // reference is the one holding the full output. artifact: outcome.artifacts.last().map(CompletionArtifact::from), }; + tracing::debug!( + "{LOG_PREFIX} task_id={} status={} building record from outcome", + self.task_id, + status.as_str() + ); Some(self.record(status, result)) } + + /// The record for an executor error, or `None` when the error is not the + /// child failing: a host seam fault (a task id mismatch, a persistence + /// failure, a missing capability) says nothing about the child's own run. + pub(crate) fn record_for_error(&self, error: &SubagentError) -> Option { + if !matches!( + error, + SubagentError::Execution(_) | SubagentError::Transient { .. } + ) { + return None; + } + tracing::debug!( + "{LOG_PREFIX} task_id={} status=failed building record from executor error", + self.task_id + ); + Some(self.record( + CompletionStatus::Failed, + CompletionResult { + text: error.to_string(), + omitted_chars: 0, + artifact: None, + }, + )) + } } /// Hands `record` to the router. A router failure is logged, never raised: the diff --git a/crates/tinyagents-orchestration/src/subagent/driver.rs b/crates/tinyagents-orchestration/src/subagent/driver.rs index afee4e5d1..854fa5005 100644 --- a/crates/tinyagents-orchestration/src/subagent/driver.rs +++ b/crates/tinyagents-orchestration/src/subagent/driver.rs @@ -169,13 +169,25 @@ impl SubagentDriver { /// /// Only the invocation that wins the durable terminal write records, so a /// coalesced follower or a replayed terminal never produces a second - /// completion. Only a spawn that set a notify mode is recorded. A - /// cancellation is not recorded (the parent asked for it), neither is a - /// pause (the same task completes later), nor an executor error (nothing - /// terminal was persisted, so the task may be re-run; the caller of `run` - /// has the error). A failure to record - /// is logged and never fails the run. Without this call the driver behaves - /// exactly as before. + /// completion. Only a spawn that set a notify mode is recorded: + /// + /// - a normal finish is recorded as success or incomplete; + /// - a cancellation is recorded as `Cancelled` (also when it lands after + /// the planner but before the child launches; one before planning has no + /// notify mode yet and is not recorded); + /// - an executor error (`Execution`, or `Transient` once retries are + /// exhausted) is recorded as `Failed` and still returned to the caller. + /// Host seam faults (task id mismatch, persistence, missing capability) + /// are not the child failing and are not recorded. The router keeps the + /// first record per task id, so a task re-run under the same id after a + /// recorded failure does not add a second completion; + /// - a pause (awaiting input) is not recorded: it is not terminal, and the + /// resume that eventually finishes the same task id records then. For + /// the same reason an error while resuming a paused task is not + /// recorded: its pause is still durable. + /// + /// A failure to record is logged and never fails the run. Without this + /// call the driver behaves exactly as before. pub fn with_completion_router(mut self, router: Arc) -> Self { self.completions = Some(router); self @@ -361,14 +373,23 @@ impl SubagentDriver { actual: prepared.task_id, }); } + // Known as soon as the planner has named the spawn's notify mode, so a + // cancel from here on is recorded too. + let completion_origin = self + .completions + .as_ref() + .and_then(|_| CompletionOrigin::new(&task_key, &prepared)); if cancellation.is_cancelled() { - return self + let result = self .persist_cancelled( task_key, SubagentOutcome::cancelled(prepared.task_id), expected_pause, ) + .await?; + self.record_completion(completion_origin.as_ref(), &result, 0) .await; + return Ok(result); } prepared.tools = restrict_tools( &prepared.tools, @@ -391,10 +412,6 @@ impl SubagentDriver { let policy = prepared.policy.clone(); let result_policy = prepared.result_policy.clone(); - let completion_origin = self - .completions - .as_ref() - .and_then(|_| CompletionOrigin::new(&task_key, &prepared)); let mut attempts = AttemptSource::from_prepared(&prepared); let mut attempt = 0usize; let mut execution = SubagentExecution { @@ -478,18 +495,41 @@ impl SubagentDriver { } Err(SubagentError::Cancelled) => SubagentOutcome::cancelled(task_id), Err(_) if cancellation.is_cancelled() => SubagentOutcome::cancelled(task_id), - Err(error) => return Err(error), + Err(error) => { + // A resumed child still has its durable pause: it is not + // finished, and the router keeps the first record per task id, + // so a failed record here would shadow the eventual finish. + if expected_pause.is_none() + && let (Some(router), Some(origin)) = (&self.completions, &completion_origin) + && let Some(record) = origin.record_for_error(&error) + { + deliver(router, record).await; + } + return Err(error); + } }; let result = self .persist(task_key, outcome, expected_pause, &cancellation) .await?; - if let (Some(router), Some(origin)) = (&self.completions, &completion_origin) + self.record_completion(completion_origin.as_ref(), &result, omitted_chars) + .await; + Ok(result) + } + + /// Records `result` when this invocation owns its host effects and the + /// outcome maps to a completion (a pause does not). + async fn record_completion( + &self, + origin: Option<&CompletionOrigin>, + result: &SubagentRunResult, + omitted_chars: usize, + ) { + if let (Some(router), Some(origin)) = (&self.completions, origin) && result.should_emit_host_effects() && let Some(record) = origin.record_for_outcome(&result.outcome, omitted_chars) { deliver(router, record).await; } - Ok(result) } async fn persist_cancelled( diff --git a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs index 865783f13..768231daf 100644 --- a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs +++ b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs @@ -58,7 +58,10 @@ impl SubagentExecutor for Executor { } #[derive(Default)] -struct Memory(Mutex>); +struct Memory { + terminals: Mutex>, + pauses: Mutex>, +} #[async_trait] impl SubagentPersistence for Memory { @@ -66,21 +69,30 @@ impl SubagentPersistence for Memory { &self, key: &SubagentTaskKey, ) -> Result, SubagentError> { - Ok(self.0.lock().unwrap().get(key).cloned()) + Ok(self.terminals.lock().unwrap().get(key).cloned()) } - async fn load(&self, _: &SubagentTaskKey) -> Result, SubagentError> { - Ok(None) + async fn load(&self, key: &SubagentTaskKey) -> Result, SubagentError> { + Ok(self + .pauses + .lock() + .unwrap() + .get(key) + .and_then(|outcome| match &outcome.status { + SubagentOutcomeKind::AwaitingInput(pause) => Some(pause.resume.clone()), + _ => None, + })) } async fn load_pause( &self, - _: &SubagentTaskKey, + key: &SubagentTaskKey, ) -> Result, SubagentError> { - Ok(None) + Ok(self.pauses.lock().unwrap().get(key).cloned()) } async fn save_pause( &self, - _: PersistedSubagentPause, + pause: PersistedSubagentPause, ) -> Result { + self.pauses.lock().unwrap().insert(pause.key, pause.outcome); Ok(SubagentPausePersistenceDisposition::Inserted) } async fn record_terminal( @@ -89,7 +101,11 @@ impl SubagentPersistence for Memory { outcome: &SubagentOutcome, _: Option<&SubagentResume>, ) -> Result { - self.0.lock().unwrap().insert(key.clone(), outcome.clone()); + self.pauses.lock().unwrap().remove(key); + self.terminals + .lock() + .unwrap() + .insert(key.clone(), outcome.clone()); Ok(SubagentTerminalPersistenceDisposition::Inserted) } } @@ -227,9 +243,42 @@ async fn an_incomplete_child_is_recorded_with_its_reason() { } #[tokio::test] -async fn an_executor_failure_is_returned_and_not_recorded() { +async fn an_executor_failure_is_returned_and_recorded_as_failed() { let router = router(); let behaviour: Behaviour = Arc::new(|_| Err(SubagentError::Execution("boom".into()))); + let failing = + driver(Some(NotifyMode::Off), Some("p"), behaviour).with_completion_router(router.clone()); + // The caller of `run` still gets the error, unchanged. + let error = failing + .run(request("t1"), CancellationToken::new()) + .await + .unwrap_err(); + assert!(matches!(error, SubagentError::Execution(_)), "{error:?}"); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].task_id, "t1"); + assert_eq!(pending[0].status, CompletionStatus::Failed); + assert!(pending[0].result.text.contains("boom"), "{:?}", pending[0]); + // The router keeps the first record per task id: a re-run under the same + // id does not add a second completion. + let retry = + driver(Some(NotifyMode::Off), Some("p"), ok()).with_completion_router(router.clone()); + retry + .run(request("t1"), CancellationToken::new()) + .await + .unwrap(); + assert_eq!(router.pending_for("p").len(), 1); +} + +#[tokio::test] +async fn a_transient_failure_that_exhausts_its_retries_is_recorded_as_failed() { + let router = router(); + let behaviour: Behaviour = Arc::new(|_| { + Err(SubagentError::Transient { + message: "provider unavailable".into(), + tools_ran: false, + }) + }); let failing = driver(Some(NotifyMode::Off), Some("p"), behaviour).with_completion_router(router.clone()); assert!( @@ -238,15 +287,136 @@ async fn an_executor_failure_is_returned_and_not_recorded() { .await .is_err() ); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].status, CompletionStatus::Failed); +} + +#[tokio::test] +async fn a_host_seam_error_is_not_a_child_failure() { + let router = router(); + let behaviour: Behaviour = Arc::new(|_| Ok(SubagentOutcome::completed("other-task", "x"))); + let mismatched = + driver(Some(NotifyMode::Off), Some("p"), behaviour).with_completion_router(router.clone()); + assert!(matches!( + mismatched + .run(request("t1"), CancellationToken::new()) + .await, + Err(SubagentError::TaskIdMismatch { .. }) + )); assert!(router.pending_for("p").is_empty()); - // The same task can be re-run, and its success is not shadowed. - let retry = - driver(Some(NotifyMode::Off), Some("p"), ok()).with_completion_router(router.clone()); - retry +} + +#[tokio::test] +async fn an_executor_failure_without_a_notify_mode_is_not_recorded() { + let router = router(); + let behaviour: Behaviour = Arc::new(|_| Err(SubagentError::Execution("boom".into()))); + let failing = driver(None, Some("p"), behaviour).with_completion_router(router.clone()); + assert!( + failing + .run(request("t1"), CancellationToken::new()) + .await + .is_err() + ); + assert!(router.pending_for("p").is_empty()); +} + +#[tokio::test] +async fn an_executor_failure_without_a_router_is_unchanged() { + let behaviour: Behaviour = Arc::new(|_| Err(SubagentError::Execution("boom".into()))); + let plain = driver(Some(NotifyMode::Followup), Some("p"), behaviour); + let error = plain + .run(request("t1"), CancellationToken::new()) + .await + .unwrap_err(); + assert!(matches!(error, SubagentError::Execution(m) if m == "boom")); +} + +/// A behaviour that pauses on its first call and completes on the next. +fn pause_then_complete() -> Behaviour { + let calls = Arc::new(Mutex::new(0usize)); + Arc::new(move |task| { + let mut calls = calls.lock().unwrap(); + *calls += 1; + if *calls == 1 { + let mut outcome = SubagentOutcome::completed(task, "waiting"); + outcome.status = SubagentOutcomeKind::AwaitingInput(SubagentPause { + reason: "need approval".into(), + resume: SubagentResume::default(), + }); + Ok(outcome) + } else { + Ok(SubagentOutcome::completed(task, "approved and done")) + } + }) +} + +#[tokio::test] +async fn a_paused_child_is_not_recorded_until_it_finishes() { + let router = router(); + let driver = driver(Some(NotifyMode::Off), Some("p"), pause_then_complete()) + .with_completion_router(router.clone()); + let first = driver + .run(request("t1"), CancellationToken::new()) + .await + .unwrap(); + assert!(matches!( + first.outcome.status, + SubagentOutcomeKind::AwaitingInput(_) + )); + assert!( + router.pending_for("p").is_empty(), + "a pause is not terminal" + ); + // The resume finishes the same task id and records exactly once. + driver .run(request("t1"), CancellationToken::new()) .await .unwrap(); - assert_eq!(router.pending_for("p")[0].status, CompletionStatus::Success); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].status, CompletionStatus::Success); + assert_eq!(pending[0].result.text, "approved and done"); +} + +#[tokio::test] +async fn a_resume_that_errors_leaves_the_paused_task_unrecorded() { + let router = router(); + let calls = Arc::new(Mutex::new(0usize)); + let behaviour: Behaviour = { + let pause = pause_then_complete(); + Arc::new(move |task| { + let mut calls = calls.lock().unwrap(); + *calls += 1; + match *calls { + 1 => pause(task), + 2 => Err(SubagentError::Execution("resume crashed".into())), + _ => Ok(SubagentOutcome::completed(task, "finally done")), + } + }) + }; + let driver = + driver(Some(NotifyMode::Off), Some("p"), behaviour).with_completion_router(router.clone()); + driver + .run(request("t1"), CancellationToken::new()) + .await + .unwrap(); + assert!( + driver + .run(request("t1"), CancellationToken::new()) + .await + .is_err() + ); + // The durable pause is still there, so the task is not finished: no + // failed record may shadow its eventual success. + assert!(router.pending_for("p").is_empty()); + driver + .run(request("t1"), CancellationToken::new()) + .await + .unwrap(); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].status, CompletionStatus::Success); } #[tokio::test] @@ -264,7 +434,7 @@ async fn a_spawn_without_a_notify_mode_is_not_recorded() { } #[tokio::test] -async fn a_cancelled_child_is_not_recorded() { +async fn a_cancelled_child_is_recorded_as_cancelled() { let router = router(); let token = CancellationToken::new(); let behaviour: Behaviour = { @@ -281,7 +451,68 @@ async fn a_cancelled_child_is_not_recorded() { result.outcome.status, SubagentOutcomeKind::Cancelled )); - assert!(router.pending_for("p").is_empty()); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].task_id, "t1"); + assert_eq!(pending[0].status, CompletionStatus::Cancelled); +} + +#[tokio::test] +async fn a_cancel_before_the_child_launches_is_recorded_once_it_has_a_notify_mode() { + struct CancellingPlanner(CancellationToken); + + #[async_trait] + impl SubagentPlanner for CancellingPlanner { + async fn prepare( + &self, + request: SubagentRequest, + ) -> Result, SubagentError> { + let prepared = Planner { + mode: Some(NotifyMode::Off), + parent: Some("p"), + } + .prepare(request) + .await?; + self.0.cancel(); + Ok(prepared) + } + } + + let router = router(); + let token = CancellationToken::new(); + let driver = SubagentDriver::new(SubagentCapabilities { + planner: Some(Arc::new(CancellingPlanner(token.clone()))), + executor: Some(Arc::new(Executor(ok()))), + persistence: Some(Arc::new(Memory::default())), + }) + .unwrap() + .with_completion_router(router.clone()); + let result = driver.run(request("t1"), token).await.unwrap(); + assert!(matches!( + result.outcome.status, + SubagentOutcomeKind::Cancelled + )); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].status, CompletionStatus::Cancelled); +} + +#[tokio::test] +async fn a_cancelled_child_without_a_router_is_unchanged() { + let token = CancellationToken::new(); + let behaviour: Behaviour = { + let token = token.clone(); + Arc::new(move |task| { + token.cancel(); + Ok(SubagentOutcome::completed(task, "late")) + }) + }; + let plain = driver(Some(NotifyMode::Off), Some("p"), behaviour); + let result = plain.run(request("t1"), token).await.unwrap(); + assert!(matches!( + result.outcome.status, + SubagentOutcomeKind::Cancelled + )); } #[tokio::test] diff --git a/crates/tinyagents-orchestration/src/subagent/outcome_status_map.rs b/crates/tinyagents-orchestration/src/subagent/outcome_status_map.rs index c17800b66..53bca7dac 100644 --- a/crates/tinyagents-orchestration/src/subagent/outcome_status_map.rs +++ b/crates/tinyagents-orchestration/src/subagent/outcome_status_map.rs @@ -56,8 +56,7 @@ impl TryFrom<&SubagentOutcomeKind> for SubAgentJobStatus { } /// `AwaitingInput` is a pause, not a completion, and fails with -/// [`NoEquivalentStatus`]. (The completion router separately declines to record -/// a `Cancelled` outcome; that is a routing policy, not a mapping gap.) +/// [`NoEquivalentStatus`]; that is also why the driver never records a pause. impl TryFrom<&SubagentOutcomeKind> for CompletionStatus { type Error = NoEquivalentStatus; From 3965260bdd8772bca2e087b8c4c58f23ce8eca92 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:49:22 +0300 Subject: [PATCH 08/21] docs: document subagent completion recording and graph loop event parity The subagent README now spells out which terminal outcomes are recorded by the completion router, covering cancellations, executor errors, and the pause/resume case, and notes that a re-run under the same task id does not add a second record. The harness terminal-outcome doc records that the graph loop driver now publishes the turn and message lifecycle events with the same semantics as the direct loop, and explains how `provider_started` is set for summarizer calls. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/subagent/README.md | 25 +++++++++++++---- docs/modules/harness/terminal-outcome.md | 28 +++++++++++++++---- 2 files changed, 43 insertions(+), 10 deletions(-) diff --git a/crates/tinyagents-orchestration/src/subagent/README.md b/crates/tinyagents-orchestration/src/subagent/README.md index dc690f56d..7bcb2c3b5 100644 --- a/crates/tinyagents-orchestration/src/subagent/README.md +++ b/crates/tinyagents-orchestration/src/subagent/README.md @@ -151,11 +151,26 @@ a foreground one whose result already returns to the parent, is never recorded. The parent key is `PreparedSubagent::with_completion_parent`, else the request's thread id, else the parent run id. Only the invocation that wins the durable terminal write records, so coalesced followers and replayed terminals add -nothing. Recorded: `Completed` (success) and `Incomplete`. Not recorded: a -cancellation (the parent's own doing), a pause (the same task id completes -later), and an executor error (nothing terminal was persisted and the task may be -re-run; `run` returns the error, and a host that wants a failed push records it -itself). A router failure is logged, never raised. With no router configured the +nothing. Recorded: + +- `Completed` as success and `Incomplete` as incomplete; +- a cancellation as `Cancelled` with an empty result, including a cancel that + lands after the planner but before the child launches (a cancel before + planning has no notify mode yet and is not recorded); +- an executor error (`Execution`, or `Transient` once retries are exhausted) as + `Failed` with the error text. `run` still returns the error. A host seam fault + (`TaskIdMismatch`, `Persistence`, `MissingCapability`) is not the child + failing and is not recorded. + +Not recorded: a pause (`AwaitingInput`). It is not terminal, so nothing is +recorded until the resume that finishes the same task id, which records once +(the durable terminal write is what is recorded). For the same reason an error +while resuming a paused task is not recorded: its pause is still durable. The +router keeps the first record per task id, so a task re-run under the same id +after a recorded failure or cancellation does not add a second completion; use a +fresh task id to run it again. Without a router, or without a notify mode, none +of this applies and behaviour is unchanged. A router failure is logged, never +raised. With no router configured the driver is unchanged. Detached children tracked by a status channel use `spawn_status_watcher_with_completions` (see `detached/README.md`), which keeps watching across a pause and skips a child whose ledger shows a cancellation. diff --git a/docs/modules/harness/terminal-outcome.md b/docs/modules/harness/terminal-outcome.md index a25c66b7c..6bbe4dd54 100644 --- a/docs/modules/harness/terminal-outcome.md +++ b/docs/modules/harness/terminal-outcome.md @@ -34,8 +34,26 @@ On `RunCompleted` / `RunFailed` the `outcome` field is `Some` for every run this crate ends; it is `None` only when deserializing journals written before the field existed. -Scope: the typed outcome is published by both the direct and the graph loop -driver, but the turn and message lifecycle events (`TurnStarted`, -`MessageAppended`, `MessageRetracted`, ...) are emitted by the direct loop only. -The graph driver does not emit them yet, so graph-engine hosts should not rely -on them. +Scope: the typed outcome and the turn and message lifecycle events +(`TurnStarted`, `TurnCompleted`, `MessageAppended`) are published by both the +direct loop and the graph loop driver, with the same semantics: input messages +are never announced, each appended message is announced exactly once, and a +nested tool call never produces a `MessageAppended`. The graph driver announces +appends at the same points as the direct loop (before a turn's `ModelStarted`, +after the assistant reply, after a tool batch, and at the end of the run, +including after `after_agent`). `MessageRetracted` and `TranscriptRewritten` +come from direct-loop recovery paths the graph rendition does not implement. The +step-by-step `LoopIter` and `compile_loop` graph announce the same events and +treat the transcript present at their first node activation as the seed. + +## `provider_started` and summarizers + +`provider_started` is true once any provider call was dispatched, including a +context-window summarizer's. Summarizer calls bypass the run context's dispatch +marker, so the compaction middleware scopes each summarization with a dispatch +tracker (`summarization::dispatch`): a model-backed summarizer marks it right +before it calls its model, and the middleware then sets the run-wide flag. A +summarizer that fails with no usage after dispatching (a transport error) still +reports `provider_started: true`; one rejected before dispatch (empty input, +validation) does not. Only the run-wide flag is set, so a timeout during the main +model call is still classified by that call's own dispatch. From 35b3f5b6740d68b14a8138aecd42790edd1b5294 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:50:16 +0300 Subject: [PATCH 09/21] docs(agent_loop): document graph lifecycle events and summarizer dispatch tracking The graph agent loop now documents how it announces TurnStarted, TurnCompleted and MessageAppended at the same points as the direct loop, including that the seed transcript is never announced and nested tool calls never reach the transcript. The summarization README gains a section on dispatch tracking, explaining why a summarizer that fails without usage after dispatching still reports provider_started while a pre-dispatch rejection does not. Remaining edits are rustfmt reflows and README wording cleanups. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/compile.rs | 3 +- crates/tinyagents-graph/src/agent_loop/mod.rs | 11 +++++ .../src/agent_loop/terminal_outcome_tests.rs | 16 +++++--- .../src/summarization/README.md | 7 ++++ .../tests/graph_lifecycle_events.rs | 41 +++++++++++++++---- .../src/subagent/README.md | 6 +-- 6 files changed, 64 insertions(+), 20 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/compile.rs b/crates/tinyagents-graph/src/agent_loop/compile.rs index 306a355b1..4c0669880 100644 --- a/crates/tinyagents-graph/src/agent_loop/compile.rs +++ b/crates/tinyagents-graph/src/agent_loop/compile.rs @@ -139,8 +139,7 @@ where let harness = rt.harness.clone(); let mut ctx_guard = rt.ctx.lock().await; let mut run_guard = rt.run.lock().await; - runtime::settle_node(&harness, &mut ctx_guard, &mut run_guard, loop_state) - .await + runtime::settle_node(&harness, &mut ctx_guard, &mut run_guard, loop_state).await } } }) diff --git a/crates/tinyagents-graph/src/agent_loop/mod.rs b/crates/tinyagents-graph/src/agent_loop/mod.rs index 107f85503..4ff4ec5ac 100644 --- a/crates/tinyagents-graph/src/agent_loop/mod.rs +++ b/crates/tinyagents-graph/src/agent_loop/mod.rs @@ -52,6 +52,17 @@ //! defaults to [`tinyagents_harness::runtime::LoopExecution::Direct`], so //! every existing caller is unaffected unless it opts in. //! +//! # Lifecycle events +//! +//! The node bodies announce `TurnStarted`, `TurnCompleted` and +//! `MessageAppended` at the same points as the direct loop, through +//! [`tinyagents_harness::agent_loop::phases`]'s `lifecycle_*` functions: pending +//! appends and the turn open right before `ModelStarted`, the assistant reply +//! after the model call, the tool results (and the turn close) after a tool +//! batch, and the final close in `settle`. The transcript present at the first +//! `plan` activation is the seed and is never announced. Nested tool calls never +//! reach the transcript, so they never produce `MessageAppended`. +//! //! # Tool batch execution shape //! //! The `tools` node runs a turn's whole tool-call batch in one node diff --git a/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs b/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs index 4f0260ba0..a4f1ff37d 100644 --- a/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs +++ b/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs @@ -539,7 +539,9 @@ impl crate::summarization::Summarizer for RejectingSummarizer { &self, _: &[Message], ) -> crate::error::Result { - Err(TinyAgentsError::Validation("rejected before dispatch".into())) + Err(TinyAgentsError::Validation( + "rejected before dispatch".into(), + )) } } @@ -549,7 +551,10 @@ fn long_input() -> Vec { input.push(Message::user(format!("question {i} {}", "x".repeat(400)))); input.push(Message::Assistant(AssistantMessage { id: None, - content: vec![ContentBlock::Text(format!("answer {i} {}", "y".repeat(400)))], + content: vec![ContentBlock::Text(format!( + "answer {i} {}", + "y".repeat(400) + ))], tool_calls: Vec::new(), usage: None, origin: None, @@ -559,9 +564,7 @@ fn long_input() -> Vec { input } -fn aborting_harness( - summarizer: Box, -) -> AgentHarness<()> { +fn aborting_harness(summarizer: Box) -> AgentHarness<()> { use crate::middleware::{CompressionFailurePolicy, ContextCompressionMiddleware}; let policy = crate::summarization::SummarizationPolicy { keep_last: 2, @@ -580,7 +583,8 @@ fn aborting_harness( #[tokio::test] async fn a_summarizer_that_dispatched_and_failed_without_usage_marks_provider_started() { - let summarizer = crate::summarization::ModelSummarizer::new(Arc::new(FailingModel("boom")), "m"); + let summarizer = + crate::summarization::ModelSummarizer::new(Arc::new(FailingModel("boom")), "m"); let harness = aborting_harness(Box::new(summarizer)); let ctx = RunContext::new(RunConfig::new("sum-dispatched"), ()); let partial = harness diff --git a/crates/tinyagents-harness/src/summarization/README.md b/crates/tinyagents-harness/src/summarization/README.md index 9e312c2c7..d60bb9e47 100644 --- a/crates/tinyagents-harness/src/summarization/README.md +++ b/crates/tinyagents-harness/src/summarization/README.md @@ -23,6 +23,13 @@ loop. `SummarizationPolicy::plan` decides the split between `to_summarize` and `to_keep`; a `Summarizer` then condenses the former into a `SummaryRecord` carrying `CompressionProvenance`. +- **Dispatch tracking** (`dispatch.rs`): summarizer calls bypass the run + context's provider-dispatch marker. `ContextCompressionMiddleware` scopes each + summarization with `track_dispatch`, `ModelSummarizer` calls `mark_dispatched` + right before it invokes its model, and the middleware then sets the run's + `provider_started`, so a summarizer that fails without usage after dispatching + still reports `TerminalOutcome::provider_started`. A rejection before dispatch + does not. - **Tool-call pairing** (`pairing.rs`) is the structural safety net both of the above rely on: a naive length-based cut point routinely separates an assistant tool-call turn from the tool results answering it, producing a diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index c0dbb2c86..cd52f9f17 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -80,7 +80,10 @@ fn lifecycle(events: &[AgentEvent]) -> Vec { .. } => Some(format!( "append:{index}:{role}:{}", - call_id.as_ref().map(ToString::to_string).unwrap_or_default() + call_id + .as_ref() + .map(ToString::to_string) + .unwrap_or_default() )), _ => None, }) @@ -134,9 +137,19 @@ async fn a_tool_run_announces_each_message_once_and_closes_each_turn() { let mut indices: Vec<&String> = graph.iter().filter(|e| e.starts_with("append:")).collect(); let total = indices.len(); indices.dedup(); - assert_eq!(indices.len(), total, "no duplicate MessageAppended: {graph:?}"); - assert!(!graph.iter().any(|e| e.starts_with("append:0:")), "{graph:?}"); - assert!(graph.iter().any(|e| e == "turn.completed:1:call-1"), "{graph:?}"); + assert_eq!( + indices.len(), + total, + "no duplicate MessageAppended: {graph:?}" + ); + assert!( + !graph.iter().any(|e| e.starts_with("append:0:")), + "{graph:?}" + ); + assert!( + graph.iter().any(|e| e == "turn.completed:1:call-1"), + "{graph:?}" + ); } #[tokio::test] @@ -152,7 +165,10 @@ async fn a_steering_injected_message_is_announced_like_the_direct_loop() { ) .await }; - assert_eq!(run(LoopExecution::Graph).await, run(LoopExecution::Direct).await); + assert_eq!( + run(LoopExecution::Graph).await, + run(LoopExecution::Direct).await + ); } /// A tool that calls another tool through the harness. @@ -241,12 +257,21 @@ async fn iter_stepping_announces_appends_once_and_never_the_input() { iter.run_to_end().await.expect("finishes"); let events = lifecycle(&recorder.events()); - assert!(!events.iter().any(|e| e.starts_with("append:0:")), "{events:?}"); + assert!( + !events.iter().any(|e| e.starts_with("append:0:")), + "{events:?}" + ); let appended: Vec<&String> = events.iter().filter(|e| e.starts_with("append:")).collect(); assert_eq!(appended.len(), 3, "assistant, tool, assistant: {events:?}"); assert_eq!( - events.iter().filter(|e| e.starts_with("turn.started")).count(), - events.iter().filter(|e| e.starts_with("turn.completed")).count(), + events + .iter() + .filter(|e| e.starts_with("turn.started")) + .count(), + events + .iter() + .filter(|e| e.starts_with("turn.completed")) + .count(), "{events:?}" ); } diff --git a/crates/tinyagents-orchestration/src/subagent/README.md b/crates/tinyagents-orchestration/src/subagent/README.md index 7bcb2c3b5..9d554ba02 100644 --- a/crates/tinyagents-orchestration/src/subagent/README.md +++ b/crates/tinyagents-orchestration/src/subagent/README.md @@ -163,15 +163,13 @@ nothing. Recorded: failing and is not recorded. Not recorded: a pause (`AwaitingInput`). It is not terminal, so nothing is -recorded until the resume that finishes the same task id, which records once -(the durable terminal write is what is recorded). For the same reason an error +recorded until the resume that finishes the same task id, which records once. For the same reason an error while resuming a paused task is not recorded: its pause is still durable. The router keeps the first record per task id, so a task re-run under the same id after a recorded failure or cancellation does not add a second completion; use a fresh task id to run it again. Without a router, or without a notify mode, none of this applies and behaviour is unchanged. A router failure is logged, never -raised. With no router configured the -driver is unchanged. Detached children tracked by a status channel use +raised. Detached children tracked by a status channel use `spawn_status_watcher_with_completions` (see `detached/README.md`), which keeps watching across a pause and skips a child whose ledger shows a cancellation. From 1a8167e1b2d5d0283e1eb83c1826bd41c4af4881 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:52:14 +0300 Subject: [PATCH 10/21] docs: describe graph lifecycle events, summarizer dispatch tracking and completion recording Co-authored-by: Medulla --- .claude/worktrees/run-1 | 1 + 1 file changed, 1 insertion(+) create mode 160000 .claude/worktrees/run-1 diff --git a/.claude/worktrees/run-1 b/.claude/worktrees/run-1 new file mode 160000 index 000000000..35b3f5b67 --- /dev/null +++ b/.claude/worktrees/run-1 @@ -0,0 +1 @@ +Subproject commit 35b3f5b6740d68b14a8138aecd42790edd1b5294 From b340c66a37731c744b7343d47df798c10612e439 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:52:20 +0300 Subject: [PATCH 11/21] chore: drop accidentally tracked embedded worktree Co-authored-by: Medulla --- .claude/worktrees/run-1 | 1 - 1 file changed, 1 deletion(-) delete mode 160000 .claude/worktrees/run-1 diff --git a/.claude/worktrees/run-1 b/.claude/worktrees/run-1 deleted file mode 160000 index 35b3f5b67..000000000 --- a/.claude/worktrees/run-1 +++ /dev/null @@ -1 +0,0 @@ -Subproject commit 35b3f5b6740d68b14a8138aecd42790edd1b5294 From b61c3505567e9862467fe12abcbe0e83496fae9c Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 08:57:53 +0300 Subject: [PATCH 12/21] feat(tinyagents-integration-tests): add graph lifecycle event tests Add integration tests covering graph lifecycle events, verifying that the expected events are emitted as a graph runs through its lifecycle. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../tests/graph_lifecycle_events.rs | 102 ++++++++++++++++++ 1 file changed, 102 insertions(+) diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index cd52f9f17..fc60cb092 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -275,3 +275,105 @@ async fn iter_stepping_announces_appends_once_and_never_the_input() { "{events:?}" ); } + +/// Requests an approval interrupt on the first `after_model` hook only. +struct PauseOnce(std::sync::atomic::AtomicBool); + +#[async_trait] +impl tinyagents_harness::middleware::Middleware<(), ()> for PauseOnce { + fn name(&self) -> &str { + "pause_once" + } + + async fn after_model( + &self, + ctx: &mut RunContext<()>, + _state: &(), + _response: &mut ModelResponse, + ) -> tinyagents_harness::Result<()> { + if self.0.swap(false, std::sync::atomic::Ordering::SeqCst) { + ctx.request_control(tinyagents_harness::context::MiddlewareControl::Interrupt { + node: "review".into(), + message: "needs approval".into(), + }); + } + Ok(()) + } +} + +/// An interrupted node discards its state and re-runs on resume, in a fresh +/// runtime (a restart). Lifecycle events must stay consistent: the input is +/// never announced and a message index is never announced twice without a +/// retraction in between. +#[tokio::test] +async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_duplicates() { + use tinyagents_graph::InMemoryCheckpointer; + use tinyagents_graph::agent_loop::{LoopRuntime, LoopState, compile_loop}; + + let build = || { + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model( + "mock", + Arc::new(MockModel::with_responses(vec![ + tool_call_response("c1", "lookup"), + ModelResponse::assistant("done"), + ])), + ) + .set_default_model("mock") + .register_tool(Arc::new(FakeTool::returning("lookup", "out"))) + .push_middleware(Arc::new(PauseOnce(std::sync::atomic::AtomicBool::new( + true, + )))); + Arc::new(harness) + }; + let checkpointer = Arc::new(InMemoryCheckpointer::::default()); + let recorder = EventRecorder::new(); + let graph_for = |harness: Arc>| { + let ctx = RunContext::new(RunConfig::new("resume-lifecycle"), ()) + .with_events(recorder.sink()); + let rt = Arc::new(LoopRuntime::for_run(harness, Arc::new(()), ctx)); + compile_loop(rt) + .expect("compiles") + .with_checkpointer(checkpointer.clone()) + }; + + let first = graph_for(build()) + .run_with_thread("t", LoopState::seed(vec![Message::user("go")])) + .await + .expect("first leg reaches the interrupt"); + assert_eq!(first.interrupts.len(), 1); + + // A different harness and runtime resumes from the checkpoint. The pause + // already fired, so the second leg runs to the end. + // (`build` arms a fresh PauseOnce; disarm it by consuming the one shot.) + let second = build(); + let resumed = graph_for(second) + .resume( + "t", + tinyagents_graph::Command { + update: None, + goto: Vec::new(), + resume: Some(json!({ "approved": true })), + resume_by_task: Default::default(), + }, + ) + .await + .expect("resume completes"); + assert!(resumed.state.finished); + + let events = lifecycle(&recorder.events()); + assert!(!events.iter().any(|e| e.starts_with("append:0:")), "{events:?}"); + let mut live = std::collections::BTreeSet::new(); + for event in recorder.events() { + match event { + AgentEvent::MessageAppended { index, .. } => { + assert!(live.insert(index), "index {index} announced twice: {events:?}"); + } + AgentEvent::MessageRetracted { index } => { + live.remove(&index); + } + _ => {} + } + } +} From 21b9dcb7b11c68026b890c70a460878125e9c379 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:00:34 +0300 Subject: [PATCH 13/21] feat(agent_loop): retract announced appends when a node interrupts Interrupted nodes discard their state and re-run from entry on resume, so the appends they had already announced are now retracted with MessageRetracted, keeping event mirrors consistent with the transcript that is actually kept. The built-in model-backed summarizers also mark dispatch before calling their model, so provider_started is reported correctly for host-defined summarizers that cannot call mark_dispatched themselves. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../tinyagents-graph/src/agent_loop/driver.rs | 5 +-- .../src/agent_loop/runtime.rs | 33 +++++++++++++++++-- .../src/agent_loop/entry.rs | 5 +-- .../src/agent_loop/phases.rs | 9 +++++ .../src/summarization/task_state/mod.rs | 1 + .../tests/graph_lifecycle_events.rs | 23 ++++++++----- docs/modules/harness/terminal-outcome.md | 15 +++++++-- 7 files changed, 73 insertions(+), 18 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/driver.rs b/crates/tinyagents-graph/src/agent_loop/driver.rs index 5c881353d..ab31865ff 100644 --- a/crates/tinyagents-graph/src/agent_loop/driver.rs +++ b/crates/tinyagents-graph/src/agent_loop/driver.rs @@ -256,8 +256,9 @@ where .then(|| ctx.peek_last_limit()) .flatten(); let mut outcome = TerminalOutcome::from_error(error, site).with_limit_kind(kind); - // A failed summarizer already received a provider response, though - // summarizer calls bypass the context's dispatch marker. + // A failed summarizer already received a provider response; keep the + // match for summarizers that report usage but never marked dispatch + // (host-defined `Summarizer` impls cannot call `mark_dispatched`). outcome.provider_started = ctx.provider_started() || matches!(error, TinyAgentsError::SummarizationUsage { .. }); Some(outcome) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 5106713e3..e50fea1a8 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -320,6 +320,8 @@ where return Err(TinyAgentsError::LimitExceeded(error.to_string())); } + phases::lifecycle_seed(ctx, loop_state.messages.len()); + let entry_len = loop_state.messages.len(); let request = loop_state .pending_request .take() @@ -518,7 +520,9 @@ where // has the real tool-routing decision to fall through to instead of an // arbitrary default. if let Some(control) = ctx.take_control() { - return apply_control(ctx, &mut loop_state, control, node::MODEL, route); + let result = apply_control(ctx, &mut loop_state, control, node::MODEL, route); + retract_on_interrupt(ctx, &result, entry_len); + return result; } // Stash the response for `settle` to extract structured output from. // Reusing `pending_request`'s sibling field would need a new field; keep @@ -552,6 +556,8 @@ where State: Send + Sync, Ctx: Send + Sync, { + phases::lifecycle_seed(ctx, loop_state.messages.len()); + let entry_len = loop_state.messages.len(); let calls = std::mem::take(&mut loop_state.pending_tool_calls); let outcome = phases::execute_tool_batch( harness, @@ -598,7 +604,9 @@ where } if let Some(control) = ctx.take_control() { - return apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); + let result = apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); + retract_on_interrupt(ctx, &result, entry_len); + return result; } Ok(goto(loop_state, node::PLAN)) @@ -607,6 +615,26 @@ where /// The `settle` node body: extracts/validates structured output when the /// turn planned one, drives the output-validation retry loop /// (`RunPolicy::output_retry`), and finishes the run. +/// An interrupted node's state is discarded and the node re-runs from its entry +/// state on resume, so the appends it already announced are retracted: the +/// re-run announces them again, and events never name a message the kept +/// transcript lacks. +fn retract_on_interrupt( + ctx: &mut RunContext, + result: &Result>, + entry_len: usize, +) { + if matches!(result, Ok(NodeResult::Interrupt(_))) { + tracing::debug!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + entry_len, + "[graph_loop] node interrupted; retracting its announced appends" + ); + phases::lifecycle_retract(ctx, entry_len); + } +} + pub(crate) async fn settle_node( harness: &AgentHarness, ctx: &mut RunContext, @@ -617,6 +645,7 @@ where State: Send + Sync, Ctx: Send + Sync, { + phases::lifecycle_seed(ctx, loop_state.messages.len()); // The final assistant message is the last append of its turn; close the // turn before any output-retry prompt is pushed (that prompt is announced // with the next turn's start). diff --git a/crates/tinyagents-harness/src/agent_loop/entry.rs b/crates/tinyagents-harness/src/agent_loop/entry.rs index 2ca39301e..4327e2dbf 100644 --- a/crates/tinyagents-harness/src/agent_loop/entry.rs +++ b/crates/tinyagents-harness/src/agent_loop/entry.rs @@ -504,8 +504,9 @@ impl AgentHarness { } // `site` describes this failure; the run may still have reached // the provider on an earlier call. - // A failed summarizer already received a provider response, though - // summarizer calls bypass the context's dispatch marker. + // A failed summarizer already received a provider response; keep the + // match for summarizers that report usage but never marked dispatch + // (host-defined `Summarizer` impls cannot call `mark_dispatched`). outcome.provider_started = ctx.provider_started() || matches!(&error, TinyAgentsError::SummarizationUsage { .. }); terminal.run.terminal = Some(outcome.clone()); diff --git a/crates/tinyagents-harness/src/agent_loop/phases.rs b/crates/tinyagents-harness/src/agent_loop/phases.rs index 1275e90b0..3d556c8c7 100644 --- a/crates/tinyagents-harness/src/agent_loop/phases.rs +++ b/crates/tinyagents-harness/src/agent_loop/phases.rs @@ -255,3 +255,12 @@ pub fn lifecycle_close_turn( ) { ctx.close_turn(harness.policy().capture, messages); } + +/// Reports that the transcript was truncated to `new_len` messages, emitting +/// `MessageRetracted` for each announced message removed (highest index first). +/// A driver whose node discards its state (an interrupt re-runs the node from +/// its entry state) calls this so a mirror stays consistent with the transcript +/// that will actually be kept. +pub fn lifecycle_retract(ctx: &mut RunContext, new_len: usize) { + ctx.retract_transcript(new_len); +} diff --git a/crates/tinyagents-harness/src/summarization/task_state/mod.rs b/crates/tinyagents-harness/src/summarization/task_state/mod.rs index 6f2884efe..65076ee21 100644 --- a/crates/tinyagents-harness/src/summarization/task_state/mod.rs +++ b/crates/tinyagents-harness/src/summarization/task_state/mod.rs @@ -260,6 +260,7 @@ impl TaskStateSummarizer { request = request.with_response_format(format.clone()); } for attempt in 1..=STATE_ATTEMPTS { + crate::summarization::dispatch::mark_dispatched(); let response = self.model.invoke(&(), request.clone()).await.map_err(|e| { tracing::warn!(error = %e, "[tinyagents::task_state] state-update call failed"); TinyAgentsError::Model(format!("task-state model call failed: {e}")) diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index fc60cb092..f03766628 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -310,7 +310,7 @@ async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_du use tinyagents_graph::InMemoryCheckpointer; use tinyagents_graph::agent_loop::{LoopRuntime, LoopState, compile_loop}; - let build = || { + let build = |armed: bool| { let mut harness: AgentHarness<()> = AgentHarness::new(); harness .register_model( @@ -323,22 +323,22 @@ async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_du .set_default_model("mock") .register_tool(Arc::new(FakeTool::returning("lookup", "out"))) .push_middleware(Arc::new(PauseOnce(std::sync::atomic::AtomicBool::new( - true, + armed, )))); Arc::new(harness) }; let checkpointer = Arc::new(InMemoryCheckpointer::::default()); let recorder = EventRecorder::new(); let graph_for = |harness: Arc>| { - let ctx = RunContext::new(RunConfig::new("resume-lifecycle"), ()) - .with_events(recorder.sink()); + let ctx = + RunContext::new(RunConfig::new("resume-lifecycle"), ()).with_events(recorder.sink()); let rt = Arc::new(LoopRuntime::for_run(harness, Arc::new(()), ctx)); compile_loop(rt) .expect("compiles") .with_checkpointer(checkpointer.clone()) }; - let first = graph_for(build()) + let first = graph_for(build(true)) .run_with_thread("t", LoopState::seed(vec![Message::user("go")])) .await .expect("first leg reaches the interrupt"); @@ -346,8 +346,7 @@ async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_du // A different harness and runtime resumes from the checkpoint. The pause // already fired, so the second leg runs to the end. - // (`build` arms a fresh PauseOnce; disarm it by consuming the one shot.) - let second = build(); + let second = build(false); let resumed = graph_for(second) .resume( "t", @@ -363,12 +362,18 @@ async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_du assert!(resumed.state.finished); let events = lifecycle(&recorder.events()); - assert!(!events.iter().any(|e| e.starts_with("append:0:")), "{events:?}"); + assert!( + !events.iter().any(|e| e.starts_with("append:0:")), + "{events:?}" + ); let mut live = std::collections::BTreeSet::new(); for event in recorder.events() { match event { AgentEvent::MessageAppended { index, .. } => { - assert!(live.insert(index), "index {index} announced twice: {events:?}"); + assert!( + live.insert(index), + "index {index} announced twice: {events:?}" + ); } AgentEvent::MessageRetracted { index } => { live.remove(&index); diff --git a/docs/modules/harness/terminal-outcome.md b/docs/modules/harness/terminal-outcome.md index 6bbe4dd54..979632862 100644 --- a/docs/modules/harness/terminal-outcome.md +++ b/docs/modules/harness/terminal-outcome.md @@ -44,15 +44,24 @@ after the assistant reply, after a tool batch, and at the end of the run, including after `after_agent`). `MessageRetracted` and `TranscriptRewritten` come from direct-loop recovery paths the graph rendition does not implement. The step-by-step `LoopIter` and `compile_loop` graph announce the same events and -treat the transcript present at their first node activation as the seed. +treat the transcript present at their first node activation (any node) as the seed. A node that +interrupts discards its state and re-runs on resume, so the appends it announced +are retracted with `MessageRetracted` first. + +Known differences from the direct loop: the graph rendition closes a turn before +an output-retry prompt is announced (the direct loop announces the prompt first), +and `LoopIter` / `compile_loop` do not close an open turn when a node errors +(`GraphLoopDriver` does, on every exit). ## `provider_started` and summarizers `provider_started` is true once any provider call was dispatched, including a context-window summarizer's. Summarizer calls bypass the run context's dispatch marker, so the compaction middleware scopes each summarization with a dispatch -tracker (`summarization::dispatch`): a model-backed summarizer marks it right -before it calls its model, and the middleware then sets the run-wide flag. A +tracker (`summarization::dispatch`): the built-in model-backed summarizers (`ModelSummarizer`, `TaskStateSummarizer`) +mark it right before they call their model (host-defined `Summarizer` impls +cannot, since `mark_dispatched` is crate-private; they still count when they +fail with usage), and the middleware then sets the run-wide flag. A summarizer that fails with no usage after dispatching (a transport error) still reports `provider_started: true`; one rejected before dispatch (empty input, validation) does not. Only the run-wide flag is set, so a timeout during the main From a23134bd82e0f88622f6aa2d28c0500266989545 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:06:24 +0300 Subject: [PATCH 14/21] feat(tinyagents-integration-tests): add graph lifecycle event tests Add integration tests covering graph lifecycle events, verifying that the expected events are emitted as a graph runs through its lifecycle. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../tests/graph_lifecycle_events.rs | 106 ++++++++++++++++++ 1 file changed, 106 insertions(+) diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index f03766628..38841b4e9 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -382,3 +382,109 @@ async fn a_checkpoint_resume_in_a_fresh_runtime_neither_reannounces_input_nor_du } } } + +/// A serial batch whose second call trips the tool cap with an error: the +/// first call's result is already on the transcript and must be announced and +/// counted by the closing `TurnCompleted`, as in the direct loop. +#[tokio::test] +async fn a_partially_executed_tool_batch_is_announced_when_the_batch_errors() { + let mut traces = Vec::new(); + for execution in [LoopExecution::Direct, LoopExecution::Graph] { + let mut response = tool_call_response("a", "lookup"); + response + .message + .tool_calls + .push(ToolCall::new("b", "lookup", json!({}))); + let mut harness = harness_for( + execution, + vec![response, ModelResponse::assistant("done")], + RunLimits::default().with_max_tool_calls(1), + ); + harness.register_tool(Arc::new(FakeTool::returning("lookup", "out"))); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("partial"), ()).with_events(recorder.sink()); + let result = harness + .invoke_in_context(&(), ctx, vec![Message::user("go")]) + .await; + assert!(result.is_err(), "{execution:?}: the cap errors the run"); + traces.push(lifecycle(&recorder.events())); + } + assert!( + traces[0].iter().any(|e| e == "append:2:tool:a"), + "{:?}", + traces[0] + ); + assert_eq!(traces[1], traces[0]); +} + +/// Approval interrupt raised after the tool batch: the batch's messages and +/// its turn close are not left announced for a state the graph discards. +#[tokio::test] +async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn() { + use tinyagents_graph::InMemoryCheckpointer; + use tinyagents_graph::agent_loop::{LoopRuntime, LoopState, compile_loop}; + use tinyagents_harness::context::MiddlewareControl; + + struct PauseAfterTools(std::sync::atomic::AtomicBool); + + #[async_trait] + impl tinyagents_harness::middleware::Middleware<(), ()> for PauseAfterTools { + fn name(&self) -> &str { + "pause_after_tools" + } + fn should_stop_after_turn(&self, _: &mut RunContext<()>, _: &tinyagents_harness::middleware::AgentRun) -> bool { + false + } + async fn after_tool( + &self, + ctx: &mut RunContext<()>, + _state: &(), + _call: &ToolCall, + _result: &mut tinytools::ToolResult, + ) -> tinyagents_harness::Result<()> { + if self.0.swap(false, std::sync::atomic::Ordering::SeqCst) { + ctx.request_control(MiddlewareControl::Interrupt { + node: "review".into(), + message: "needs approval".into(), + }); + } + Ok(()) + } + } + + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model( + "mock", + Arc::new(MockModel::with_responses(vec![ + tool_call_response("c1", "lookup"), + ModelResponse::assistant("done"), + ])), + ) + .set_default_model("mock") + .register_tool(Arc::new(FakeTool::returning("lookup", "out"))) + .push_middleware(Arc::new(PauseAfterTools( + std::sync::atomic::AtomicBool::new(true), + ))); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("pause-tools"), ()).with_events(recorder.sink()); + let rt = Arc::new(LoopRuntime::for_run(Arc::new(harness), Arc::new(()), ctx)); + let graph = compile_loop(rt) + .expect("compiles") + .with_checkpointer(Arc::new(InMemoryCheckpointer::::default())); + let first = graph + .run_with_thread("t", LoopState::seed(vec![Message::user("go")])) + .await + .expect("reaches the interrupt"); + assert_eq!(first.interrupts.len(), 1); + + let events = lifecycle(&recorder.events()); + assert!( + !events.iter().any(|e| e.starts_with("turn.completed:1:c1")), + "a turn holding discarded tool results must not be reported complete: {events:?}" + ); + assert!( + !events.iter().any(|e| e.starts_with("append:2:tool")) || recorder.kinds().contains(&"message.retracted".to_string()), + "{events:?}" + ); +} From 5a5cd269041207e013ce84d6586597962c1feb3c Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:06:42 +0300 Subject: [PATCH 15/21] feat(graph): emit lifecycle events for agent loop runs The agent loop runtime now emits graph lifecycle events as a run progresses, so subscribers can observe start, step, and completion transitions. Integration tests cover the new event stream. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/runtime.rs | 47 ++++++++++++------- .../tests/graph_lifecycle_events.rs | 7 +-- 2 files changed, 33 insertions(+), 21 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index e50fea1a8..094c78e9f 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -521,7 +521,7 @@ where // arbitrary default. if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::MODEL, route); - retract_on_interrupt(ctx, &result, entry_len); + retract_on_interrupt(harness, ctx, &result, entry_len); return result; } // Stash the response for `settle` to extract structured output from. @@ -590,11 +590,15 @@ where } return Ok(goto(loop_state, node::SETTLE)); } - Err(error) => return Err(error), + Err(error) => { + // Results of calls that ran before the failure are on the node's + // transcript, which the error path would otherwise drop; keep them + // on the run, as the direct loop does, so the driver's final + // lifecycle close announces and counts them. + run.messages = loop_state.messages.clone(); + return Err(error); + } }; - // Every tool result of this batch is on the transcript: announce them and - // close the turn, as the direct loop does after its batch. - phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); loop_state.tool_calls = run.tool_calls; loop_state.executed_tools = run.executed_tools.clone(); let _ = outcome; @@ -605,9 +609,14 @@ where if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); - retract_on_interrupt(ctx, &result, entry_len); + if !retract_on_interrupt(harness, ctx, &result, entry_len) { + // Every tool result of this batch is on the transcript: announce + // them and close the turn, as the direct loop does after its batch. + phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); + } return result; } + phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); Ok(goto(loop_state, node::PLAN)) } @@ -619,20 +628,26 @@ where /// state on resume, so the appends it already announced are retracted: the /// re-run announces them again, and events never name a message the kept /// transcript lacks. -fn retract_on_interrupt( +fn retract_on_interrupt( + harness: &AgentHarness, ctx: &mut RunContext, result: &Result>, entry_len: usize, -) { - if matches!(result, Ok(NodeResult::Interrupt(_))) { - tracing::debug!( - target: "tinyagents::agent_loop", - run_id = %ctx.run_id(), - entry_len, - "[graph_loop] node interrupted; retracting its announced appends" - ); - phases::lifecycle_retract(ctx, entry_len); +) -> bool { + if !matches!(result, Ok(NodeResult::Interrupt(_))) { + return false; } + tracing::debug!( + target: "tinyagents::agent_loop", + run_id = %ctx.run_id(), + entry_len, + "[graph_loop] node interrupted; retracting its announced appends" + ); + phases::lifecycle_retract(ctx, entry_len); + // Close the turn the node opened: a fresh runtime resuming from the + // checkpoint cannot carry this tracker's open turn over. + phases::lifecycle_close_turn(harness, ctx, &[]); + true } pub(crate) async fn settle_node( diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index 38841b4e9..533b58b1d 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -432,15 +432,12 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( fn name(&self) -> &str { "pause_after_tools" } - fn should_stop_after_turn(&self, _: &mut RunContext<()>, _: &tinyagents_harness::middleware::AgentRun) -> bool { - false - } async fn after_tool( &self, ctx: &mut RunContext<()>, _state: &(), - _call: &ToolCall, - _result: &mut tinytools::ToolResult, + _invocation: &tinyagents_harness::middleware::ToolInvocationIdentity, + _result: &mut ToolResult, ) -> tinyagents_harness::Result<()> { if self.0.swap(false, std::sync::atomic::Ordering::SeqCst) { ctx.request_control(MiddlewareControl::Interrupt { From ef0aaa618715fd768bf8e3f8c3d981c595ab0b07 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:08:17 +0300 Subject: [PATCH 16/21] fix(graph): keep partial tool results on error, close tool turn after interrupt check, close interrupted model turn Co-authored-by: Medulla --- .../src/agent_loop/runtime.rs | 7 +++--- .../tests/graph_lifecycle_events.rs | 3 ++- .../src/subagent/driver_completion_tests.rs | 22 +++++++++++++++++++ 3 files changed, 28 insertions(+), 4 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 094c78e9f..8b4eab29e 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -521,7 +521,7 @@ where // arbitrary default. if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::MODEL, route); - retract_on_interrupt(harness, ctx, &result, entry_len); + retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len); return result; } // Stash the response for `settle` to extract structured output from. @@ -609,7 +609,7 @@ where if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); - if !retract_on_interrupt(harness, ctx, &result, entry_len) { + if !retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len) { // Every tool result of this batch is on the transcript: announce // them and close the turn, as the direct loop does after its batch. phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); @@ -632,6 +632,7 @@ fn retract_on_interrupt( harness: &AgentHarness, ctx: &mut RunContext, result: &Result>, + messages: &[tinyinference_llm::message::Message], entry_len: usize, ) -> bool { if !matches!(result, Ok(NodeResult::Interrupt(_))) { @@ -646,7 +647,7 @@ fn retract_on_interrupt( phases::lifecycle_retract(ctx, entry_len); // Close the turn the node opened: a fresh runtime resuming from the // checkpoint cannot carry this tracker's open turn over. - phases::lifecycle_close_turn(harness, ctx, &[]); + phases::lifecycle_close_turn(harness, ctx, &messages[..entry_len.min(messages.len())]); true } diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index 533b58b1d..5a54e4453 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -481,7 +481,8 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( "a turn holding discarded tool results must not be reported complete: {events:?}" ); assert!( - !events.iter().any(|e| e.starts_with("append:2:tool")) || recorder.kinds().contains(&"message.retracted".to_string()), + !events.iter().any(|e| e.starts_with("append:2:tool")) + || recorder.kinds().contains(&"message.retracted".to_string()), "{events:?}" ); } diff --git a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs index 768231daf..ddd15520d 100644 --- a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs +++ b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs @@ -497,6 +497,28 @@ async fn a_cancel_before_the_child_launches_is_recorded_once_it_has_a_notify_mod assert_eq!(pending[0].status, CompletionStatus::Cancelled); } +/// An executor that returns `SubagentError::Cancelled` is converted to a +/// cancelled outcome before persistence, so it is recorded as `Cancelled` and +/// the call returns `Ok`. +#[tokio::test] +async fn an_executor_cancelled_error_is_recorded_as_cancelled() { + let router = router(); + let behaviour: Behaviour = Arc::new(|_| Err(SubagentError::Cancelled)); + let driver = + driver(Some(NotifyMode::Off), Some("p"), behaviour).with_completion_router(router.clone()); + let result = driver + .run(request("t1"), CancellationToken::new()) + .await + .unwrap(); + assert!(matches!( + result.outcome.status, + SubagentOutcomeKind::Cancelled + )); + let pending = router.pending_for("p"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].status, CompletionStatus::Cancelled); +} + #[tokio::test] async fn a_cancelled_child_without_a_router_is_unchanged() { let token = CancellationToken::new(); From 8b8e7b92520268f54200764ce29eb1fa12b4a448 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:14:10 +0300 Subject: [PATCH 17/21] test: tighten retraction check and model terminal conflicts in the completion test double Co-authored-by: Medulla --- .../tests/graph_lifecycle_events.rs | 18 +++++++++++++++--- .../src/subagent/driver_completion_tests.rs | 9 +++++---- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index 5a54e4453..6058fd17a 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -480,9 +480,21 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( !events.iter().any(|e| e.starts_with("turn.completed:1:c1")), "a turn holding discarded tool results must not be reported complete: {events:?}" ); + // Net of retractions, the discarded tool message (index 2) is not live. + let mut live = std::collections::BTreeSet::new(); + for event in recorder.events() { + match event { + AgentEvent::MessageAppended { index, .. } => { + live.insert(index); + } + AgentEvent::MessageRetracted { index } => { + live.remove(&index); + } + _ => {} + } + } assert!( - !events.iter().any(|e| e.starts_with("append:2:tool")) - || recorder.kinds().contains(&"message.retracted".to_string()), - "{events:?}" + !live.contains(&2), + "discarded tool message left live: {events:?}" ); } diff --git a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs index ddd15520d..a6ad7104c 100644 --- a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs +++ b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs @@ -101,11 +101,12 @@ impl SubagentPersistence for Memory { outcome: &SubagentOutcome, _: Option<&SubagentResume>, ) -> Result { + let mut terminals = self.terminals.lock().unwrap(); + if terminals.contains_key(key) { + return Ok(SubagentTerminalPersistenceDisposition::Existing); + } self.pauses.lock().unwrap().remove(key); - self.terminals - .lock() - .unwrap() - .insert(key.clone(), outcome.clone()); + terminals.insert(key.clone(), outcome.clone()); Ok(SubagentTerminalPersistenceDisposition::Inserted) } } From 465fd99e93e343691b64d3dfd37e2b898ebacc7c Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:16:10 +0300 Subject: [PATCH 18/21] feat(agent-loop): emit lifecycle events for graph runs The graph agent loop now reports lifecycle events through the harness so subagent completion and phase transitions are observable. This lets orchestration track subagent progress and integration tests assert on the event stream. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/runtime.rs | 8 ++++ .../src/agent_loop/lifecycle.rs | 24 ++++++++++ .../src/agent_loop/phases.rs | 14 ++++++ .../tests/graph_lifecycle_events.rs | 47 ++++++++++++++++++- .../src/subagent/completion.rs | 6 +++ .../src/subagent/driver_completion_tests.rs | 4 ++ 6 files changed, 102 insertions(+), 1 deletion(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 8b4eab29e..22acec26a 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -321,6 +321,7 @@ where } phases::lifecycle_seed(ctx, loop_state.messages.len()); + phases::lifecycle_resume(ctx, loop_state.turn as u32, None); let entry_len = loop_state.messages.len(); let request = loop_state .pending_request @@ -558,6 +559,13 @@ where { phases::lifecycle_seed(ctx, loop_state.messages.len()); let entry_len = loop_state.messages.len(); + // The turn that issued these calls began at (or before) the assistant + // message; a fresh runtime resuming an interrupted batch has no open turn. + phases::lifecycle_resume( + ctx, + loop_state.turn as u32, + Some(entry_len.saturating_sub(1)), + ); let calls = std::mem::take(&mut loop_state.pending_tool_calls); let outcome = phases::execute_tool_batch( harness, diff --git a/crates/tinyagents-harness/src/agent_loop/lifecycle.rs b/crates/tinyagents-harness/src/agent_loop/lifecycle.rs index e25702dad..5a72db8ea 100644 --- a/crates/tinyagents-harness/src/agent_loop/lifecycle.rs +++ b/crates/tinyagents-harness/src/agent_loop/lifecycle.rs @@ -58,6 +58,25 @@ impl TurnTracker { } } + /// Continues numbering from `completed_turns` (a resumed run in a fresh + /// runtime starts its tracker at zero) and, when `open_from` is given and no + /// turn is open, re-opens the in-flight turn that began at that transcript + /// index, without announcing it again. + pub(crate) fn adopt(&mut self, completed_turns: u32, open_from: Option) { + self.turn = self.turn.max(completed_turns); + if let Some(start) = open_from + && self.open.is_none() + { + tracing::debug!( + target: "tinyagents::agent_loop", + turn = self.turn, + start, + "[agent_loop] re-opened the in-flight turn after a resume" + ); + self.open = Some((self.turn, start)); + } + } + /// Treats the first `len` messages as seed input if no seed was fixed yet. /// A no-op once seeded, so repeated node entries never swallow appends. pub(crate) fn ensure_seeded(&mut self, len: usize) { @@ -227,6 +246,11 @@ impl RunContext { } impl RunContext { + /// See [`TurnTracker::adopt`]. + pub(crate) fn adopt_turn_state(&mut self, completed_turns: u32, open_from: Option) { + self.turns.adopt(completed_turns, open_from); + } + /// Fixes the lifecycle seed at `len` messages unless one is already set. pub(crate) fn ensure_turn_tracker_seeded(&mut self, len: usize) { self.turns.ensure_seeded(len); diff --git a/crates/tinyagents-harness/src/agent_loop/phases.rs b/crates/tinyagents-harness/src/agent_loop/phases.rs index 3d556c8c7..9cb7e53c2 100644 --- a/crates/tinyagents-harness/src/agent_loop/phases.rs +++ b/crates/tinyagents-harness/src/agent_loop/phases.rs @@ -264,3 +264,17 @@ pub fn lifecycle_close_turn( pub fn lifecycle_retract(ctx: &mut RunContext, new_len: usize) { ctx.retract_transcript(new_len); } + +/// Re-aligns the lifecycle tracker with a run entered mid-flight (resumed from a +/// checkpoint in a fresh runtime, whose tracker starts at zero): turn numbering +/// continues from `completed_turns`, and when `open_turn_start` is `Some(index)` +/// and no turn is open, the in-flight turn that began at that transcript index +/// is re-opened (not announced again) so the next close reports it. Never moves +/// numbering backwards, so it is safe to call on every node entry. +pub fn lifecycle_resume( + ctx: &mut RunContext, + completed_turns: u32, + open_turn_start: Option, +) { + ctx.adopt_turn_state(completed_turns, open_turn_start); +} diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index 6058fd17a..1b980affc 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -466,9 +466,10 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( let recorder = EventRecorder::new(); let ctx = RunContext::new(RunConfig::new("pause-tools"), ()).with_events(recorder.sink()); let rt = Arc::new(LoopRuntime::for_run(Arc::new(harness), Arc::new(()), ctx)); + let checkpointer = Arc::new(InMemoryCheckpointer::::default()); let graph = compile_loop(rt) .expect("compiles") - .with_checkpointer(Arc::new(InMemoryCheckpointer::::default())); + .with_checkpointer(checkpointer.clone()); let first = graph .run_with_thread("t", LoopState::seed(vec![Message::user("go")])) .await @@ -497,4 +498,48 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( !live.contains(&2), "discarded tool message left live: {events:?}" ); + + // A fresh runtime (a restart) resumes from the checkpoint: the re-run tools + // node announces its results and closes the turn it belongs to, with the + // turn numbering continuing from the checkpoint. + let mut second: AgentHarness<()> = AgentHarness::new(); + second + .register_model( + "mock", + Arc::new(MockModel::with_responses(vec![ModelResponse::assistant( + "done", + )])), + ) + .set_default_model("mock") + .register_tool(Arc::new(FakeTool::returning("lookup", "out"))); + let rt2 = Arc::new(LoopRuntime::for_run( + Arc::new(second), + Arc::new(()), + RunContext::new(RunConfig::new("pause-tools"), ()).with_events(recorder.sink()), + )); + let resumed = compile_loop(rt2) + .expect("compiles") + .with_checkpointer(checkpointer) + .resume( + "t", + tinyagents_graph::Command { + update: None, + goto: Vec::new(), + resume: Some(json!({ "approved": true })), + resume_by_task: Default::default(), + }, + ) + .await + .expect("resume completes"); + assert!(resumed.state.finished); + let events = lifecycle(&recorder.events()); + let count = |prefix: &str| events.iter().filter(|e| e.starts_with(prefix)).count(); + assert_eq!(count("turn.completed:1:c1"), 1, "{events:?}"); + assert_eq!(count("turn.started:1"), 1, "{events:?}"); + assert_eq!( + count("turn.started:2"), + 1, + "numbering continues: {events:?}" + ); + assert_eq!(count("turn.started"), count("turn.completed"), "{events:?}"); } diff --git a/crates/tinyagents-orchestration/src/subagent/completion.rs b/crates/tinyagents-orchestration/src/subagent/completion.rs index d2e2e8f52..0632c8bba 100644 --- a/crates/tinyagents-orchestration/src/subagent/completion.rs +++ b/crates/tinyagents-orchestration/src/subagent/completion.rs @@ -91,6 +91,12 @@ impl CompletionOrigin { omitted_chars: usize, ) -> Option { let status = CompletionStatus::try_from(&outcome.status).ok()?; + if matches!(&outcome.status, SubagentOutcomeKind::Cancelled) { + // A cancel that lands after the executor returned keeps the child's + // late output on the outcome (`cancelled_preserving`); it must not + // reach the parent as a usable answer. + return Some(self.record(status, CompletionResult::default())); + } let text = match &outcome.status { SubagentOutcomeKind::Incomplete(incomplete) if outcome.output.is_empty() => { incomplete.reason.clone() diff --git a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs index a6ad7104c..f0c59f550 100644 --- a/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs +++ b/crates/tinyagents-orchestration/src/subagent/driver_completion_tests.rs @@ -456,6 +456,10 @@ async fn a_cancelled_child_is_recorded_as_cancelled() { assert_eq!(pending.len(), 1); assert_eq!(pending[0].task_id, "t1"); assert_eq!(pending[0].status, CompletionStatus::Cancelled); + // The executor's late answer stays on the outcome but must not be pushed + // to the parent as a usable result. + assert!(pending[0].result.text.is_empty(), "{:?}", pending[0].result); + assert!(pending[0].result.artifact.is_none()); } #[tokio::test] From 4e4d13ae37b840e059ecd8c9e842187152ec883e Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:16:17 +0300 Subject: [PATCH 19/21] refactor(agent_loop): extract runtime helpers into dedicated module Moved the agent loop runtime helpers out of the main module into their own file to keep the loop logic easier to follow. No behaviour change. Auto-committed-on: dragonfly Co-authored-by: Medulla --- .../src/agent_loop/runtime.rs | 20 ++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 22acec26a..21be83c1e 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -321,7 +321,7 @@ where } phases::lifecycle_seed(ctx, loop_state.messages.len()); - phases::lifecycle_resume(ctx, loop_state.turn as u32, None); + phases::lifecycle_resume(ctx, loop_state.turn, None); let entry_len = loop_state.messages.len(); let request = loop_state .pending_request @@ -522,7 +522,7 @@ where // arbitrary default. if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::MODEL, route); - retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len); + retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len, true); return result; } // Stash the response for `settle` to extract structured output from. @@ -563,7 +563,7 @@ where // message; a fresh runtime resuming an interrupted batch has no open turn. phases::lifecycle_resume( ctx, - loop_state.turn as u32, + loop_state.turn, Some(entry_len.saturating_sub(1)), ); let calls = std::mem::take(&mut loop_state.pending_tool_calls); @@ -617,7 +617,10 @@ where if let Some(control) = ctx.take_control() { let result = apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); - if !retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len) { + // The tool turn stays open on an interrupt: the re-run closes it with the + // results it produces (a fresh runtime re-opens it via `lifecycle_resume`). + if !retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len, false) + { // Every tool result of this batch is on the transcript: announce // them and close the turn, as the direct loop does after its batch. phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); @@ -642,6 +645,7 @@ fn retract_on_interrupt( result: &Result>, messages: &[tinyinference_llm::message::Message], entry_len: usize, + close_turn: bool, ) -> bool { if !matches!(result, Ok(NodeResult::Interrupt(_))) { return false; @@ -653,9 +657,11 @@ fn retract_on_interrupt( "[graph_loop] node interrupted; retracting its announced appends" ); phases::lifecycle_retract(ctx, entry_len); - // Close the turn the node opened: a fresh runtime resuming from the - // checkpoint cannot carry this tracker's open turn over. - phases::lifecycle_close_turn(harness, ctx, &messages[..entry_len.min(messages.len())]); + // A model node's turn has no results to wait for, and a fresh runtime + // cannot carry this tracker's open turn over, so close it. A tools node + // leaves its turn open for the re-run to close with the real results. + if close_turn { + phases::lifecycle_close_turn(harness, ctx, &messages[..entry_len.min(messages.len())]); true } From 74b554511ebe5c76f24f1019a1a495ca5cc51495 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:18:27 +0300 Subject: [PATCH 20/21] fix: continue turn numbering and re-open tool turns on resume; empty result for cancelled completions Co-authored-by: Medulla --- .../tinyagents-graph/src/agent_loop/runtime.rs | 17 ++++++++++------- docs/modules/harness/terminal-outcome.md | 6 +++++- 2 files changed, 15 insertions(+), 8 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 21be83c1e..4646ab62d 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -561,11 +561,7 @@ where let entry_len = loop_state.messages.len(); // The turn that issued these calls began at (or before) the assistant // message; a fresh runtime resuming an interrupted batch has no open turn. - phases::lifecycle_resume( - ctx, - loop_state.turn, - Some(entry_len.saturating_sub(1)), - ); + phases::lifecycle_resume(ctx, loop_state.turn, Some(entry_len.saturating_sub(1))); let calls = std::mem::take(&mut loop_state.pending_tool_calls); let outcome = phases::execute_tool_batch( harness, @@ -619,8 +615,14 @@ where let result = apply_control(ctx, &mut loop_state, control, node::TOOLS, node::PLAN); // The tool turn stays open on an interrupt: the re-run closes it with the // results it produces (a fresh runtime re-opens it via `lifecycle_resume`). - if !retract_on_interrupt(harness, ctx, &result, &loop_state.messages, entry_len, false) - { + if !retract_on_interrupt( + harness, + ctx, + &result, + &loop_state.messages, + entry_len, + false, + ) { // Every tool result of this batch is on the transcript: announce // them and close the turn, as the direct loop does after its batch. phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); @@ -662,6 +664,7 @@ fn retract_on_interrupt( // leaves its turn open for the re-run to close with the real results. if close_turn { phases::lifecycle_close_turn(harness, ctx, &messages[..entry_len.min(messages.len())]); + } true } diff --git a/docs/modules/harness/terminal-outcome.md b/docs/modules/harness/terminal-outcome.md index 979632862..215068db7 100644 --- a/docs/modules/harness/terminal-outcome.md +++ b/docs/modules/harness/terminal-outcome.md @@ -46,7 +46,11 @@ come from direct-loop recovery paths the graph rendition does not implement. The step-by-step `LoopIter` and `compile_loop` graph announce the same events and treat the transcript present at their first node activation (any node) as the seed. A node that interrupts discards its state and re-runs on resume, so the appends it announced -are retracted with `MessageRetracted` first. +are retracted with `MessageRetracted` first. A model node's turn is closed on an +interrupt; a tools node leaves its turn open, and the re-run closes it with the +real results. A fresh runtime resuming from a checkpoint continues the turn +numbering from the checkpoint and re-opens the in-flight tool turn (without a +second `TurnStarted`). Known differences from the direct loop: the graph rendition closes a turn before an output-retry prompt is announced (the direct loop announces the prompt first), From 26e518d96b14534c6a593566d4234dfc4e7cda0a Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 9 Oct 2026 09:25:49 +0300 Subject: [PATCH 21/21] fix(graph): close the turn and announce partial results when a node fails Co-authored-by: Medulla --- .../src/agent_loop/runtime.rs | 55 ++++++++++++++++- .../tests/graph_lifecycle_events.rs | 59 +++++++++++++++++++ docs/modules/harness/terminal-outcome.md | 9 +-- 3 files changed, 117 insertions(+), 6 deletions(-) diff --git a/crates/tinyagents-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 4646ab62d..eae609349 100644 --- a/crates/tinyagents-graph/src/agent_loop/runtime.rs +++ b/crates/tinyagents-graph/src/agent_loop/runtime.rs @@ -283,7 +283,7 @@ where /// The `model` node body: dispatches the request [`plan_node`] built, /// records usage, appends the assistant message, and routes to `tools` or /// `settle`. -pub(crate) async fn model_node( +async fn model_node_inner( harness: &AgentHarness, app_state: &State, ctx: &mut RunContext, @@ -540,12 +540,63 @@ where /// the call site above without over-cloning it into `LoopState`. struct ModelOutcomeShadow<'a>(#[allow(dead_code)] &'a ModelResponse); +/// The `model` node: [`model_node_inner`], closing the turn it opened if it +/// fails. `GraphLoopDriver` closes on every exit as well, but `LoopIter` and +/// the compiled graph propagate a node error with no epilogue. +pub(crate) async fn model_node( + harness: &AgentHarness, + app_state: &State, + ctx: &mut RunContext, + run: &mut AgentRun, + status: &mut HarnessRunStatus, + loop_state: LoopState, +) -> Result> +where + State: Send + Sync, + Ctx: Send + Sync, +{ + let entry = loop_state.messages.clone(); + let result = model_node_inner(harness, app_state, ctx, run, status, loop_state).await; + if result.is_err() { + phases::lifecycle_close_turn(harness, ctx, &entry); + } + result +} + +/// The `tools` node: [`tools_node_inner`], announcing the results of calls that +/// ran before a failure and closing the turn if it fails. +pub(crate) async fn tools_node( + harness: &AgentHarness, + app_state: &State, + ctx: &mut RunContext, + run: &mut AgentRun, + status: &mut HarnessRunStatus, + loop_state: LoopState, +) -> Result> +where + State: Send + Sync, + Ctx: Send + Sync, +{ + let entry = loop_state.messages.clone(); + let result = tools_node_inner(harness, app_state, ctx, run, status, loop_state).await; + if result.is_err() { + // A batch error leaves the partial transcript on the run. + let messages = if run.messages.len() > entry.len() { + &run.messages + } else { + &entry + }; + phases::lifecycle_close_turn(harness, ctx, messages); + } + result +} + /// The `tools` node body: executes the batch [`model_node`] requested via /// [`phases::execute_tool_batch`] (the exact same admission / /// serial-or-concurrent execution / middleware pipeline the direct loop /// uses — see that function's docs), then routes back to `plan` for the next /// turn. -pub(crate) async fn tools_node( +async fn tools_node_inner( harness: &AgentHarness, app_state: &State, ctx: &mut RunContext, diff --git a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs index 1b980affc..e4cf79e08 100644 --- a/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -543,3 +543,62 @@ async fn an_interrupt_after_the_tool_batch_retracts_instead_of_closing_the_turn( ); assert_eq!(count("turn.started"), count("turn.completed"), "{events:?}"); } + +/// `LoopIter` and `compile_loop` propagate a node error directly, with no +/// driver epilogue: the nodes themselves close the turn they opened and +/// announce what had been appended. +#[tokio::test] +async fn a_node_error_in_the_iterator_closes_its_turn_and_announces_partial_results() { + // Provider error in the model node. + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness + .register_model("mock", Arc::new(FailingModel)) + .set_default_model("mock"); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("model-error"), ()).with_events(recorder.sink()); + let mut iter = Arc::new(harness) + .iter(Arc::new(()), ctx, vec![Message::user("go")]) + .expect("iter starts"); + assert!(iter.run_to_end().await.is_err()); + let events = lifecycle(&recorder.events()); + let count = |prefix: &str| events.iter().filter(|e| e.starts_with(prefix)).count(); + assert_eq!(count("turn.started"), count("turn.completed"), "{events:?}"); + + // A serial batch whose second call trips the tool cap. + let mut response = tool_call_response("a", "lookup"); + response + .message + .tool_calls + .push(ToolCall::new("b", "lookup", json!({}))); + let mut harness = harness_for( + LoopExecution::Direct, + vec![response, ModelResponse::assistant("done")], + RunLimits::default().with_max_tool_calls(1), + ); + harness.register_tool(Arc::new(FakeTool::returning("lookup", "out"))); + let recorder = EventRecorder::new(); + let ctx = RunContext::new(RunConfig::new("batch-error"), ()).with_events(recorder.sink()); + let mut iter = Arc::new(harness) + .iter(Arc::new(()), ctx, vec![Message::user("go")]) + .expect("iter starts"); + assert!(iter.run_to_end().await.is_err()); + let events = lifecycle(&recorder.events()); + assert!(events.iter().any(|e| e == "append:2:tool:a"), "{events:?}"); + assert!( + events.iter().any(|e| e == "turn.completed:1:a"), + "{events:?}" + ); +} + +struct FailingModel; + +#[async_trait] +impl tinyinference_llm::model::ChatModel<()> for FailingModel { + async fn invoke( + &self, + _: &(), + _: tinyinference_llm::model::ModelRequest, + ) -> tinyinference_llm::Result { + Err(tinyinference_llm::Error::Model("boom".to_string())) + } +} diff --git a/docs/modules/harness/terminal-outcome.md b/docs/modules/harness/terminal-outcome.md index 215068db7..07816a5da 100644 --- a/docs/modules/harness/terminal-outcome.md +++ b/docs/modules/harness/terminal-outcome.md @@ -52,10 +52,11 @@ real results. A fresh runtime resuming from a checkpoint continues the turn numbering from the checkpoint and re-opens the in-flight tool turn (without a second `TurnStarted`). -Known differences from the direct loop: the graph rendition closes a turn before -an output-retry prompt is announced (the direct loop announces the prompt first), -and `LoopIter` / `compile_loop` do not close an open turn when a node errors -(`GraphLoopDriver` does, on every exit). +Known difference from the direct loop: the graph rendition closes a turn before +an output-retry prompt is announced (the direct loop announces the prompt first). +A model or tools node that fails closes the turn it opened (and announces the +results of calls that ran before the failure), so `LoopIter` and `compile_loop` +need no driver epilogue. ## `provider_started` and summarizers