Octov0.11.7
AI

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 requestWhat happens
A first messageNothing is running, so this invocation claims the id and runs the agent.
A follow-upA 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 stopstopWhen holds, so the run in flight is ended and this flow stops with an empty body.
An answerauthorizeId 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.

signalMeans
contextA message was handed to the run and folded into the conversation.
stopThe run was asked to end.
authorizeA person answered a tool authorization. Carries authorizationId and allowed rather than text.
unansweredA 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

On this page