Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
8e3eb51
test(agent_loop): cover provider_started for failed summarizers
senamakel Oct 9, 2026
8fd9c62
fix(harness): mark provider_started when a summarizer dispatched with…
senamakel Oct 9, 2026
bea2bca
test(integration): assert graph driver emits turn lifecycle events
senamakel Oct 9, 2026
3e4dafc
refactor(agent_loop): split compile and runtime into separate modules
senamakel Oct 9, 2026
09a2823
fix(agent_loop): emit turn lifecycle events on graph loop exit
senamakel Oct 9, 2026
e5d31e0
feat(graph): emit turn and message lifecycle events from the graph dr…
senamakel Oct 9, 2026
d0a96d3
fix(orchestration): record cancelled and failed subagent children wit…
senamakel Oct 9, 2026
3965260
docs: document subagent completion recording and graph loop event parity
senamakel Oct 9, 2026
35b3f5b
docs(agent_loop): document graph lifecycle events and summarizer disp…
senamakel Oct 9, 2026
1a8167e
docs: describe graph lifecycle events, summarizer dispatch tracking a…
senamakel Oct 9, 2026
b340c66
chore: drop accidentally tracked embedded worktree
senamakel Oct 9, 2026
b61c350
feat(tinyagents-integration-tests): add graph lifecycle event tests
senamakel Oct 9, 2026
21b9dcb
feat(agent_loop): retract announced appends when a node interrupts
senamakel Oct 9, 2026
a23134b
feat(tinyagents-integration-tests): add graph lifecycle event tests
senamakel Oct 9, 2026
5a5cd26
feat(graph): emit lifecycle events for agent loop runs
senamakel Oct 9, 2026
ef0aaa6
fix(graph): keep partial tool results on error, close tool turn after…
senamakel Oct 9, 2026
8b8e7b9
test: tighten retraction check and model terminal conflicts in the co…
senamakel Oct 9, 2026
465fd99
feat(agent-loop): emit lifecycle events for graph runs
senamakel Oct 9, 2026
4e4d13a
refactor(agent_loop): extract runtime helpers into dedicated module
senamakel Oct 9, 2026
74b5545
fix: continue turn numbering and re-open tool turns on resume; empty …
senamakel Oct 9, 2026
26e518d
fix(graph): close the turn and announce partial results when a node f…
senamakel Oct 9, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion crates/tinyagents-graph/src/agent_loop/compile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
})
Expand Down
16 changes: 12 additions & 4 deletions crates/tinyagents-graph/src/agent_loop/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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}`"
Expand Down Expand Up @@ -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);
Comment thread
senamakel marked this conversation as resolved.
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
Expand Down
4 changes: 3 additions & 1 deletion crates/tinyagents-graph/src/agent_loop/iter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}`"
Expand Down
11 changes: 11 additions & 0 deletions crates/tinyagents-graph/src/agent_loop/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
147 changes: 141 additions & 6 deletions crates/tinyagents-graph/src/agent_loop/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 => {
Expand Down Expand Up @@ -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<State, Ctx>(
async fn model_node_inner<State, Ctx>(
harness: &AgentHarness<State, Ctx>,
app_state: &State,
ctx: &mut RunContext<Ctx>,
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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);
Comment thread
senamakel marked this conversation as resolved.
Comment thread
senamakel marked this conversation as resolved.
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(),
Expand All @@ -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(),
};
Expand Down Expand Up @@ -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();
Expand All @@ -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
Expand All @@ -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<State, Ctx>(
harness: &AgentHarness<State, Ctx>,
app_state: &State,
ctx: &mut RunContext<Ctx>,
run: &mut AgentRun,
status: &mut HarnessRunStatus,
loop_state: LoopState,
) -> Result<NodeResult<LoopState>>
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<State, Ctx>(
harness: &AgentHarness<State, Ctx>,
app_state: &State,
ctx: &mut RunContext<Ctx>,
run: &mut AgentRun,
status: &mut HarnessRunStatus,
loop_state: LoopState,
) -> Result<NodeResult<LoopState>>
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<State, Ctx>(
async fn tools_node_inner<State, Ctx>(
harness: &AgentHarness<State, Ctx>,
app_state: &State,
ctx: &mut RunContext<Ctx>,
Expand All @@ -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,
Expand Down Expand Up @@ -570,7 +645,14 @@ where
}
return Ok(goto(loop_state, node::SETTLE));
}
Err(error) => return Err(error),
Err(error) => {
Comment thread
senamakel marked this conversation as resolved.
// 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();
Comment thread
senamakel marked this conversation as resolved.
return Err(error);
Comment thread
senamakel marked this conversation as resolved.
}
};
loop_state.tool_calls = run.tool_calls;
loop_state.executed_tools = run.executed_tools.clone();
Expand All @@ -581,24 +663,77 @@ 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))
}

/// 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<State: Send + Sync, Ctx: Send + Sync>(
harness: &AgentHarness<State, Ctx>,
ctx: &mut RunContext<Ctx>,
result: &Result<NodeResult<LoopState>>,
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<State, Ctx>(
harness: &AgentHarness<State, Ctx>,
ctx: &mut RunContext<Ctx>,
run: &mut AgentRun,
mut loop_state: LoopState,
) -> Result<NodeResult<LoopState>>
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(),
Expand Down
5 changes: 3 additions & 2 deletions crates/tinyagents-harness/src/agent_loop/entry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -504,8 +504,9 @@ impl<State: Send + Sync, Ctx: Send + Sync> AgentHarness<State, Ctx> {
}
// `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());
Expand Down
Loading
Loading