Steering a Running Agent
Adding to an agent run in flight, and stopping one.
An ai-agent run is otherwise a closed loop: a request goes in, an answer
comes out, and nothing in between can change its mind or end it early. In a
conversation a person changes what they wanted, or stops wanting it, while the
agent is still working.
A run claims its conversation, the memoryThreadId it already has, for as long
as it is in flight. There is no second identifier to configure, since two runs
on one thread would also overwrite each other's memory. The same block handles
every case, with no branching in the flow.
One block, three outcomes
- type: ai-agent
name: assistant
connector: claude
memoryThreadId: vars.sub + ":" + body.threadId
stopWhen: '"X-Agent-Stop" in vars && vars["X-Agent-Stop"] != ""'
input: body.message
stream: true
# ...| The request | What happens |
|---|---|
| A first message | Nothing is running, so this invocation claims the id and runs the agent. |
| A follow-up | A run is going, so the message is handed to it and this flow stops with an empty body. No second run, no second stream, no second answer to a person already reading the first. |
| A stop | stopWhen holds, so the run in flight is ended and this flow stops with an empty body. |
| An answer | authorizeId resolves to an id, so the request is a person's decision about a tool call the run is holding. It is delivered and this flow stops with an empty body. |
The lookup and the handover happen under one lock, so "somebody took this" is a
fact rather than a timeout. The injected text is the agent's own input
expression evaluated against the new message.
Where a handed-over message lands
It arrives as an ordinary user turn, at the top of the run's next iteration, because every provider requires a tool turn to follow immediately after the assistant turn that requested it. If it arrives while the agent is producing its final answer, the run takes another turn rather than dropping it, since the invocation that handed it over already stopped its own flow.
Stopping a run
stopWhen is a boolean CEL condition over the message: a header, a body field,
whatever the caller can send. When it holds, the model call in flight is
cancelled rather than paid for in full, the agent saves its memory, and
both flows stop. Memory is saved by pruning rather than summarizing,
whatever the agent configured, since summarizing costs a model call. A stop
that finds nothing running is a no-op that answers, so a client can send
one blind.
When the caller hangs up
A closed connection ends a run through the same machinery. An
sse-event with ifClosed: stop in the
agent's events path returns a stop, and that stop cancels the model call in
flight instead of waiting for it.
How quickly a disconnect is noticed depends on stream. With stream: true
the events path fires per token, so a closed connection is caught within a
token. Without it, nothing is emitted between turn_start and turn_end, so
the disconnect is only noticed at the next turn boundary, and that turn is paid
for. Stream any agent whose caller can walk away.
Choosing a thread id
memoryThreadId decides two things at once: what the agent remembers, and
what a later message can join. A thread per user gives that user's conversation
a memory and a run; a thread per channel gives the channel one, shared by
everyone in it. A JWT subject arrives as an ordinary message variable, so
vars.sub + ":" + body.threadId says "this person, in this conversation".
The id is scoped to this block. Two agents whose expressions happen to agree
(body.threadId, say) do not hand each other messages, and a run can only be
joined or stopped through the block that started it, so give a stop route the
same agent block (as the sample does, with a header) rather than a second one.
An agent with no memoryThreadId claims nothing and cannot be joined or
stopped.
Whatever the thread id is derived from is exactly who can reach the
conversation: read it, steer it, or stop it. Build it from an identity the
caller cannot choose: vars.sub, the verified subject of a token a
jwt-validate block checked, rather than a field
of the request body.
Across replicas
A conversation is claimed for the whole deployment. A message that lands on another replica does not start a second run: the claim says who owns the conversation, and the message is delivered there.
The claim is a lease: exclusive, expiring, and refused immediately rather than waited on. The delivery is an internal queue, whose reply is what makes "somebody took this" a fact. On Kubernetes that is a coordination Lease and core NATS; there is nothing to enable and no persistence to configure.
Before the claim, a stop landing on the wrong replica stopped nothing and reported success.
Three consequences follow. A run that loses its claim stops: another replica now
owns the conversation and will save it. A run may hold its conversation for
fifteen minutes at most, so a wedged tool branch cannot keep it claimed; at
the limit the run is stopped exactly as stopWhen stops it, the transcript is
pruned and saved, and the next message picks up cleanly (an agent with no
memoryThreadId is not timed). Delivery is at most once: a message
published between a replica's death and its lease expiring is lost.
Known limits
A client that reconnects cannot rejoin its run. The claim does not move a
stream: a browser that reconnects to another replica is told the conversation is
owned and receives nothing, because the run's events go to the replica it left.
Session affinity on the ingress, keyed on the same thread memoryThreadId
is, is the answer.
A message accepted at the iteration cap is not answered. It is reported as
an unanswered signal event carrying its text. Raise maxIterations if it
happens often.
For the design work behind this, see
docs/distributed-agent-runs.md.
Watching it happen
signal is an agent event like any other,
carrying signal and the text it concerns. Put it in emit and a progress UI
shows the follow-up arriving in the same stream as the answer.
signal | Means |
|---|---|
context | A message was handed to the run and folded into the conversation. |
stop | The run was asked to end. |
authorize | A person answered a tool authorization. Carries authorizationId and allowed rather than text. |
unanswered | A message the run accepted and never got a turn to answer. |
Try it
samples/ai-agent-steering.yaml
is the whole pattern, runnable, in two terminals.
Next steps
- Tool authorization: the third thing a request can be, a person allowing a call.
- Agent Memory: what a stopped run keeps.
- AI Blocks reference: full field tables.