Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 6 additions & 1 deletion crates/agent-core/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,12 @@ builder on `Agent` or `ModelRequest`, and each is exercised by unit tests):
A queued message is a `SteeringMessage` (text plus optional image attachments), not a bare `String` —
the same shape a fresh `prompt` accepts, so `steer`/`follow_up` aren't a lesser channel that silently
drops attachments a `prompt` would have kept. `From<&str>`/`From<String>` build a text-only message,
so every plain-text `push`/`push_steer` call site keeps compiling unchanged. On drain, a follow-up
so every plain-text `push`/`push_steer` call site keeps compiling unchanged. A message may carry an
opaque `tag` (`SteeringMessage::with_tag`); the `Steered` event that reports its injection lists the
tags injected (`AgentEvent::Steered::tags`, not serialized), emitted after the transcript holding
them is checkpointed on both paths — so a host learns exactly which of its queued messages reached
the model, whatever later happens to the transcript (`serve` records MCP Events delivery this way).
On drain, a follow-up
with images becomes a real `Message::user_with_images` turn instead of a plain-text one; a mid-run
steer's images are appended as `ContentBlock::Image` blocks after its text block, onto the same
tool-results turn its text rides on.
Expand Down
124 changes: 115 additions & 9 deletions crates/agent-core/src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,8 +66,15 @@ pub enum AgentEvent {
},
/// One model turn completed.
TurnEnd { stop_reason: StopReason, step: u32 },
/// Queued steering/follow-up messages were injected at a stop boundary; the run continues.
Steered { messages: usize },
/// Queued steering/follow-up messages were injected (folded into a tool-results turn, or at a
/// stop boundary) and the transcript holding them checkpointed; the run continues. `tags` are
/// the [`crate::steering::SteeringMessage::tag`]s of the messages injected, for the host that
/// queued them (not serialized).
Steered {
messages: usize,
#[serde(skip)]
tags: Vec<u64>,
},
/// The run finished normally (the model ended its turn and no steering was queued).
AgentEnd { steps: u32 },
/// A compaction round has begun — [`Agent::compact`] confirmed a worthwhile prefix exists and is
Expand Down Expand Up @@ -1250,6 +1257,7 @@ impl Agent {
return Ok(());
}
let count = injected.len();
let tags: Vec<u64> = injected.iter().filter_map(|m| m.tag).collect();
for msg in injected {
if msg.images.is_empty() {
session.user(msg.text);
Expand All @@ -1265,10 +1273,16 @@ impl Agent {
if let Some(context) = steering.take_turn_context() {
session.attach_request_block(context);
}
sink(AgentEvent::Steered { messages: count });
// A plain user message ends the visible history here — a valid, resumable checkpoint
// (see `CheckpointHook`) before the next model call.
// (see `CheckpointHook`) before the next model call. Checkpointed *before* the
// `Steered` event, as on the tool-results path below: a host that records delivery
// on that event (by its tags) must find the messages already in the persisted
// transcript.
self.checkpoint_guarded(session).await;
sink(AgentEvent::Steered {
messages: count,
tags,
});
continue;
}

Expand Down Expand Up @@ -1811,6 +1825,7 @@ impl Agent {
// lane, injected only at the stop boundary below.
let steered = steering.drain_steer();
let steered_count = steered.len();
let steered_tags: Vec<u64> = steered.iter().filter_map(|m| m.tag).collect();
for msg in steered {
result_blocks.push(ContentBlock::text(msg.text));
for source in msg.images {
Expand All @@ -1828,6 +1843,14 @@ impl Agent {
// both committed now, so this is a valid, resumable checkpoint (see `CheckpointHook`) — the
// one mid-run point a crash between here and the run's eventual end would otherwise lose.
self.checkpoint_guarded(session).await;
// Reported before any way out of the loop below: the steered messages are in the
// checkpointed transcript now, whether or not the run goes on.
if steered_count > 0 {
sink(AgentEvent::Steered {
messages: steered_count,
tags: steered_tags,
});
}
if terminate {
// A tool requested completion (e.g. an `attempt_completion`/`exit` tool) and the whole
// batch agreed. The results are already recorded; end the run as if the model had
Expand All @@ -1838,11 +1861,6 @@ impl Agent {
});
return Ok(());
}
if steered_count > 0 {
sink(AgentEvent::Steered {
messages: steered_count,
});
}
// A graceful-stop request is honored here too, after this turn's tool results (and any
// folded-in steer text) are already committed — the same turn-boundary contract as the
// tool-less branch above. Checked *after* the `Steered` event so a client sees its steer
Expand Down Expand Up @@ -4898,6 +4916,57 @@ mod tests {
));
}

/// A steer folded into the tool-results turn of a batch that ends the run (`terminate`) did
/// reach the model's transcript: its `Steered` event, with its tag, is still reported.
#[tokio::test]
async fn a_steer_folded_into_a_terminating_batch_is_still_reported() {
struct ExitTool;
#[async_trait]
impl Tool for ExitTool {
fn name(&self) -> &str {
"exit"
}
fn description(&self) -> &str {
"End the run."
}
fn input_schema(&self) -> Value {
serde_json::json!({ "type": "object" })
}
async fn run(
&self,
_: Value,
) -> std::result::Result<crate::tool::ToolOutput, crate::error::ToolError> {
Ok(crate::tool::ToolOutput::text("done").with_terminate(true))
}
}
let mut tools = ToolRegistry::new();
tools.register(Arc::new(ExitTool));
let (agent, _mock) = agent_with(vec![turn::tool_call("tu_1", "exit", "{}")], tools);
let mut session = Session::new();
session.user("finish up");
let steering = Steering::new();
steering.push_steer(SteeringMessage::new("one more thing", Vec::new()).with_tag(42));
let mut tags = Vec::new();
agent
.run_events_steered(
&mut session,
|ev| {
if let AgentEvent::Steered { tags: t, .. } = ev {
tags.extend(t);
}
},
CancellationToken::new(),
steering,
)
.await
.unwrap();
assert_eq!(
tags,
vec![42],
"the steer reached the transcript and is reported"
);
}

#[tokio::test]
async fn a_terminate_request_only_wins_when_every_call_in_the_batch_agrees() {
// pi-parity coverage (`packages/agent/test/agent-loop.test.ts`, "should continue after
Expand Down Expand Up @@ -9584,6 +9653,43 @@ mod tests {
);
}

/// A queued message's tag comes back in the `Steered` event that reports it reached the model —
/// mid-run (folded into a tool-results turn) and at a stop boundary alike — so the host that
/// queued it knows exactly what was received, whatever later happens to the transcript.
#[tokio::test]
async fn a_steered_message_reports_its_tag_when_it_reaches_the_model() {
let mut tools = ToolRegistry::new();
tools.register(Arc::new(EchoTool));
let (agent, _mock) = agent_with(
vec![
turn::tool_call("tu_1", "echo", r#"{"text":"pong"}"#),
turn::text("acknowledged"),
turn::text("and again"),
],
tools,
);
let mut session = Session::new();
session.user("start");
let steering = Steering::new();
steering.push_steer(SteeringMessage::new("mid-run", Vec::new()).with_tag(7));
steering.push(SteeringMessage::new("at the stop", Vec::new()).with_tag(8));
let mut tags = Vec::new();
agent
.run_events_steered(
&mut session,
|ev| {
if let AgentEvent::Steered { tags: t, .. } = ev {
tags.extend(t);
}
},
CancellationToken::new(),
steering,
)
.await
.unwrap();
assert_eq!(tags, vec![7, 8]);
}

#[tokio::test]
async fn steering_is_injected_mid_run_between_tool_turns() {
// A steering message queued while the agent is mid-tool-call must be folded into the *same*
Expand Down
13 changes: 13 additions & 0 deletions crates/agent-core/src/steering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,10 @@ use crate::tool::ToolRegistry;
pub struct SteeringMessage {
pub text: String,
pub images: Vec<ImageSource>,
/// An opaque id the queuer chose, reported back in [`crate::AgentEvent::Steered`] at the moment
/// the message reaches the model — so a host can tell exactly which of its queued messages the
/// model received, whatever later happens to the transcript (a compaction, say).
pub tag: Option<u64>,
}

impl SteeringMessage {
Expand All @@ -88,15 +92,23 @@ impl SteeringMessage {
Self {
text: text.into(),
images,
tag: None,
}
}

/// The same message, tagged (see [`Self::tag`]).
pub fn with_tag(mut self, tag: u64) -> Self {
self.tag = Some(tag);
self
}
}

impl From<&str> for SteeringMessage {
fn from(text: &str) -> Self {
Self {
text: text.to_string(),
images: Vec::new(),
tag: None,
}
}
}
Expand All @@ -106,6 +118,7 @@ impl From<String> for SteeringMessage {
Self {
text,
images: Vec::new(),
tag: None,
}
}
}
Expand Down
Loading
Loading