diff --git a/crates/tinyagents-graph/src/agent_loop/compile.rs b/crates/tinyagents-graph/src/agent_loop/compile.rs index 7ad4b9ac5..4c0669880 100644 --- a/crates/tinyagents-graph/src/agent_loop/compile.rs +++ b/crates/tinyagents-graph/src/agent_loop/compile.rs @@ -137,8 +137,9 @@ 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/driver.rs b/crates/tinyagents-graph/src/agent_loop/driver.rs index 374983571..ab31865ff 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}`" @@ -256,16 +256,24 @@ 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) } }; + // 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}`" 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-graph/src/agent_loop/runtime.rs b/crates/tinyagents-graph/src/agent_loop/runtime.rs index 30543ff4a..eae609349 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 => { @@ -279,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, @@ -316,6 +320,9 @@ where return Err(TinyAgentsError::LimitExceeded(error.to_string())); } + phases::lifecycle_seed(ctx, loop_state.messages.len()); + phases::lifecycle_resume(ctx, loop_state.turn, None); + let entry_len = loop_state.messages.len(); let request = loop_state .pending_request .take() @@ -424,6 +431,16 @@ where request.reasoning = Some(mapped.clone()); } + // Same point as the direct loop: pending appends (steering) are announced, + // 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", + run_id = %ctx.run_id(), + turn, + "[graph_loop] turn started" + ); + let started_record = ctx.emit(AgentEvent::ModelStarted { call_id: call_id.clone(), model: model_name.clone(), @@ -432,7 +449,6 @@ where 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(), }; @@ -488,6 +504,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(); @@ -504,7 +521,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(harness, ctx, &result, &loop_state.messages, entry_len, true); + return result; } // Stash the response for `settle` to extract structured output from. // Reusing `pending_request`'s sibling field would need a new field; keep @@ -521,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, @@ -538,6 +608,11 @@ where State: Send + Sync, Ctx: Send + Sync, { + 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, Some(entry_len.saturating_sub(1))); let calls = std::mem::take(&mut loop_state.pending_tool_calls); let outcome = phases::execute_tool_batch( harness, @@ -570,7 +645,14 @@ 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); + } }; loop_state.tool_calls = run.tool_calls; loop_state.executed_tools = run.executed_tools.clone(); @@ -581,8 +663,24 @@ 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); + // 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); + } + return result; } + phases::lifecycle_close_turn(harness, ctx, &loop_state.messages); Ok(goto(loop_state, node::PLAN)) } @@ -590,8 +688,40 @@ 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( + harness: &AgentHarness, + ctx: &mut RunContext, + result: &Result>, + messages: &[tinyinference_llm::message::Message], + entry_len: usize, + close_turn: bool, +) -> 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); + // 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 +} + pub(crate) async fn settle_node( harness: &AgentHarness, + ctx: &mut RunContext, run: &mut AgentRun, mut loop_state: LoopState, ) -> Result> @@ -599,6 +729,11 @@ 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). + 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/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/lifecycle.rs b/crates/tinyagents-harness/src/agent_loop/lifecycle.rs index 0d1df4db3..5a72db8ea 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,39 @@ impl TurnTracker { seed_len, turn: 0, open: None, + seeded: true, + } + } + + /// 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) { + 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 +244,15 @@ impl RunContext { self.turns.rebase(&self.events, new_len, reason); } } + +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 af49461d9..9cb7e53c2 100644 --- a/crates/tinyagents-harness/src/agent_loop/phases.rs +++ b/crates/tinyagents-harness/src/agent_loop/phases.rs @@ -213,3 +213,68 @@ 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); +} + +/// 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); +} + +/// 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-harness/src/agent_loop/terminal_outcome_tests.rs b/crates/tinyagents-harness/src/agent_loop/terminal_outcome_tests.rs index 98bc9559c..a4f1ff37d 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,84 @@ 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); +} 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/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-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(( 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 new file mode 100644 index 000000000..e4cf79e08 --- /dev/null +++ b/crates/tinyagents-integration-tests/tests/graph_lifecycle_events.rs @@ -0,0 +1,604 @@ +//! 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:?}" + ); +} + +/// 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 = |armed: bool| { + 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( + 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 rt = Arc::new(LoopRuntime::for_run(harness, Arc::new(()), ctx)); + compile_loop(rt) + .expect("compiles") + .with_checkpointer(checkpointer.clone()) + }; + + let first = graph_for(build(true)) + .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. + let second = build(false); + 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); + } + _ => {} + } + } +} + +/// 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" + } + async fn after_tool( + &self, + ctx: &mut RunContext<()>, + _state: &(), + _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 { + 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 checkpointer = Arc::new(InMemoryCheckpointer::::default()); + let graph = compile_loop(rt) + .expect("compiles") + .with_checkpointer(checkpointer.clone()); + 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:?}" + ); + // 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!( + !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:?}"); +} + +/// `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/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:?}"); - } -} diff --git a/crates/tinyagents-orchestration/src/subagent/README.md b/crates/tinyagents-orchestration/src/subagent/README.md index dc690f56d..9d554ba02 100644 --- a/crates/tinyagents-orchestration/src/subagent/README.md +++ b/crates/tinyagents-orchestration/src/subagent/README.md @@ -151,12 +151,25 @@ 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 -driver is unchanged. Detached children tracked by a status channel use +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. 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. 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/crates/tinyagents-orchestration/src/subagent/completion.rs b/crates/tinyagents-orchestration/src/subagent/completion.rs index fcab99b7d..0632c8bba 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,20 +82,21 @@ 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). + let status = CompletionStatus::try_from(&outcome.status).ok()?; if matches!(&outcome.status, SubagentOutcomeKind::Cancelled) { - return None; + // 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 status = CompletionStatus::try_from(&outcome.status).ok()?; let text = match &outcome.status { SubagentOutcomeKind::Incomplete(incomplete) if outcome.output.is_empty() => { incomplete.reason.clone() @@ -108,8 +110,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..f0c59f550 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,12 @@ impl SubagentPersistence for Memory { outcome: &SubagentOutcome, _: Option<&SubagentResume>, ) -> Result { - self.0.lock().unwrap().insert(key.clone(), outcome.clone()); + let mut terminals = self.terminals.lock().unwrap(); + if terminals.contains_key(key) { + return Ok(SubagentTerminalPersistenceDisposition::Existing); + } + self.pauses.lock().unwrap().remove(key); + terminals.insert(key.clone(), outcome.clone()); Ok(SubagentTerminalPersistenceDisposition::Inserted) } } @@ -227,9 +244,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 +288,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_eq!(router.pending_for("p")[0].status, CompletionStatus::Success); + 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(); + 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 +435,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 +452,94 @@ 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); + // 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] +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); +} + +/// 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(); + 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; diff --git a/docs/modules/harness/terminal-outcome.md b/docs/modules/harness/terminal-outcome.md index a25c66b7c..07816a5da 100644 --- a/docs/modules/harness/terminal-outcome.md +++ b/docs/modules/harness/terminal-outcome.md @@ -34,8 +34,40 @@ 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 (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. 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 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 + +`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`): 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 +model call is still classified by that call's own dispatch.