fabro(01KQT1TWWJYWZGDT8F05E29H9D): implement (succeeded)

Fabro-Run: 01KQT1TWWJYWZGDT8F05E29H9D
Fabro-Completed: 5
Fabro-Checkpoint: 05ae1ef055

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-05-04 19:06:12 +00:00
parent b78ff2ff01
commit 11cd91eef1
38 changed files with 1923 additions and 161 deletions

View file

@ -0,0 +1,157 @@
import { useEffect, useRef, useState } from "react";
import { ApiError } from "../lib/api-client";
import { useSteerRun } from "../lib/mutations";
interface SteerComposerProps {
runId: string;
open: boolean;
onClose: () => void;
}
/**
* Modal composer for sending a mid-run steering message to a running run.
*
* Renders a textarea + two action buttons:
* - Send: appends to the steering queue (default).
* - Interrupt: cancels the in-flight LLM stream / tool calls and delivers
* the message as the next user turn.
*
* Surfaces 409 errors inline so users see the rejection reason (e.g., the
* `cli_agent_not_steerable` code).
*/
export function SteerComposer({ runId, open, onClose }: SteerComposerProps) {
const [text, setText] = useState("");
const [errorMessage, setErrorMessage] = useState<string | null>(null);
const textareaRef = useRef<HTMLTextAreaElement | null>(null);
const { trigger, isMutating } = useSteerRun(runId);
// Autofocus when opening; reset state when closing.
useEffect(() => {
if (open) {
requestAnimationFrame(() => textareaRef.current?.focus());
} else {
setText("");
setErrorMessage(null);
}
}, [open]);
// Close on Escape.
useEffect(() => {
if (!open) return;
const onKey = (e: KeyboardEvent) => {
if (e.key === "Escape") {
e.preventDefault();
onClose();
}
};
window.addEventListener("keydown", onKey);
return () => window.removeEventListener("keydown", onKey);
}, [open, onClose]);
if (!open) return null;
const trimmed = text.trim();
const canSubmit = trimmed.length > 0 && !isMutating;
async function send(interrupt: boolean) {
if (!canSubmit) return;
setErrorMessage(null);
try {
await trigger({ text: trimmed, interrupt });
onClose();
} catch (err) {
if (err instanceof ApiError) {
// Try to surface the well-known 409 codes inline.
const body = err.body as { code?: string; detail?: string } | null;
if (body?.code === "cli_agent_not_steerable") {
setErrorMessage(
"All running agent stages are CLI-mode and can't be steered.",
);
} else if (body?.code === "use_answer_endpoint") {
setErrorMessage(
"Run is blocked on a question; answer the question first.",
);
} else {
setErrorMessage(body?.detail ?? err.message ?? "Steer failed.");
}
} else {
setErrorMessage("Steer failed; try again.");
}
}
}
function handleKeyDown(e: React.KeyboardEvent<HTMLTextAreaElement>) {
if (e.key === "Enter" && !e.shiftKey) {
e.preventDefault();
void send(false);
}
}
return (
<div
className="fixed inset-0 z-50 flex items-center justify-center bg-black/40"
onClick={(e) => {
if (e.target === e.currentTarget) onClose();
}}
>
<div
role="dialog"
aria-modal="true"
aria-label="Steer running agent"
className="w-full max-w-md rounded-lg border border-line bg-bg-elevated p-4 shadow-lg"
>
<div className="mb-2 text-sm font-semibold text-fg">Steer agent</div>
<textarea
ref={textareaRef}
rows={4}
value={text}
onChange={(e) => setText(e.target.value)}
onKeyDown={handleKeyDown}
placeholder="Type a steering message…"
className="w-full resize-none rounded-md border border-line bg-bg p-2 text-sm text-fg outline-none focus:border-teal-500/50"
maxLength={8192}
/>
{errorMessage && (
<p
role="alert"
className="mt-2 text-xs text-amber"
data-testid="steer-error"
>
{errorMessage}
</p>
)}
<div className="mt-3 flex items-center justify-between gap-2">
<span className="text-[11px] text-fg-muted">
Enter to send · Shift+Enter for newline
</span>
<div className="flex gap-2">
<button
type="button"
onClick={onClose}
className="rounded-md border border-line px-3 py-1 text-xs text-fg-2 hover:border-line-strong"
>
Cancel
</button>
<button
type="button"
onClick={() => void send(true)}
disabled={!canSubmit}
className="rounded-md border border-amber/30 px-3 py-1 text-xs text-amber hover:border-amber/60 disabled:cursor-not-allowed disabled:opacity-50"
>
Interrupt
</button>
<button
type="button"
onClick={() => void send(false)}
disabled={!canSubmit}
className="rounded-md bg-teal-500 px-3 py-1 text-xs font-medium text-white hover:bg-teal-600 disabled:cursor-not-allowed disabled:opacity-50"
>
Send
</button>
</div>
</div>
</div>
</div>
);
}

View file

@ -3,6 +3,7 @@ import { useSWRConfig } from "swr";
import type {
PreviewUrlResponse,
RunStatusResponse,
SteerRunRequest,
SubmitAnswerRequest,
} from "@qltysh/fabro-api-client";
@ -116,6 +117,26 @@ export function useSubmitInterviewAnswer(runId: string | undefined) {
);
}
export type SteerRunArg = SteerRunRequest;
export function useSteerRun(runId: string | undefined) {
const { mutate } = useSWRConfig();
return useSWRMutation(
runId ? `steer-run:${runId}` : null,
async (_key: string, { arg }: { arg: SteerRunArg }) => {
if (!runId) throw new Error("runId is required");
const path = `/api/v1/runs/${encodeURIComponent(runId)}/steer`;
await apiJsonMutation<void, SteerRunRequest>(path, { arg });
},
{
onSuccess: () => {
if (!runId) return;
void mutate(queryKeys.runs.detail(runId));
},
},
);
}
export function useToggleDemoMode() {
const { mutate } = useSWRConfig();
return useSWRMutation(
@ -147,4 +168,4 @@ export function useLoginDevToken() {
return response.json() as Promise<{ ok: boolean }>;
},
);
}
}

View file

@ -40,6 +40,13 @@ const INTERVIEW_EVENTS = new Set([
"interview.timeout",
"interview.interrupted",
]);
const STEERING_EVENTS = new Set([
"agent.steering.injected",
"agent.steering.attached",
"agent.steering.detached",
"agent.steer.buffered",
"agent.steer.dropped",
]);
export function queryKeysForRunEvent(
runId: string,
@ -97,6 +104,17 @@ export function queryKeysForRunEvent(
return keys;
}
if (STEERING_EVENTS.has(event)) {
const keys = [
queryKeys.runs.events(runId, 1000),
queryKeys.runs.detail(runId),
];
if (stageId) {
keys.push(queryKeys.runs.stageTurns(runId, stageId));
}
return keys;
}
return [];
}
@ -142,4 +160,4 @@ export function useRunEvents(runId: string | undefined) {
if (!runId) return;
return subscribeToRunEvents(runId, mutate as MutateFn);
}, [mutate, runId]);
}
}

View file

@ -21,8 +21,8 @@ import { CSS } from "@dnd-kit/utilities";
import { ciConfig, columnStatusDisplay, deriveCiStatus, mapRunListItem } from "../data/runs";
import type { CiStatus, CheckRun, CheckStatus, RunItem, RunWithStatus, ColumnStatus } from "../data/runs";
import { EmptyState } from "../components/state";
import { SteerComposer } from "../components/steer-composer";
import { shouldRefreshBoardForEvent, useBoardEvents } from "../lib/board-events";
import { useDemoMode } from "../lib/demo-mode";
import { useAuthConfig, useBoardsRuns, useSystemInfo } from "../lib/queries";
import type { PaginatedBoardRunList } from "@qltysh/fabro-api-client";
@ -305,8 +305,15 @@ function PrCard({
actions?: string[];
}) {
const lifecycleLabel = boardLifecycleStatusLabel(pr);
const [steerOpen, setSteerOpen] = useState(false);
return (
<>
<SteerComposer
runId={pr.id}
open={steerOpen}
onClose={() => setSteerOpen(false)}
/>
<Link to={`/runs/${pr.id}`} className="group block rounded-md border border-line bg-panel p-4 transition-all duration-200 hover:border-line-strong hover:shadow-lg hover:shadow-black/20">
<div className="mb-2 flex items-center gap-1.5">
<Icon className={`size-3.5 shrink-0 ${iconColor}`} />
@ -369,6 +376,13 @@ function PrCard({
key={label}
type="button"
disabled={pr.actionDisabled}
onClick={(e) => {
if (label === "Steer") {
e.preventDefault();
e.stopPropagation();
setSteerOpen(true);
}
}}
className={`inline-flex items-center gap-1.5 rounded-md border px-2.5 py-1 text-[11px] font-medium transition-colors disabled:cursor-not-allowed disabled:text-fg-muted disabled:border-line ${
label === "Merge"
? "border-mint/20 text-mint hover:border-mint/50 hover:text-fg"
@ -410,6 +424,7 @@ function PrCard({
</div>
)}
</Link>
</>
);
}
@ -455,10 +470,9 @@ function SortablePrCard({
function BoardColumn({ column }: { column: Column }) {
const Icon = iconMap[column.iconType];
const demoMode = useDemoMode();
const actions = demoMode
? column.actions
: column.actions.filter((label) => label !== "Steer");
// Steer button is shown in all modes (no longer demo-gated). The modal/
// composer logic lives below in `PrCard`; clicking the button opens it.
const actions = column.actions;
return (
<div className="flex min-w-0 flex-col">
<div className="mb-3 flex items-center gap-3">
@ -894,4 +908,4 @@ export default function Runs() {
</div>
</DndContext>
);
}
}

View file

@ -878,6 +878,68 @@ paths:
schema:
$ref: "#/components/schemas/ErrorResponse"
/api/v1/runs/{id}/steer:
post:
operationId: steerRun
tags: [Human-in-the-Loop]
summary: Steer Run
description: |
Send a mid-run steering message to the live agent session(s) of a
running run. Set `interrupt=true` to cancel the in-flight LLM stream
and tool calls in the current round and deliver the message as the
next user turn; otherwise the message is appended to the steering
queue and picked up at the next turn boundary.
parameters:
- $ref: "#/components/parameters/RunId"
requestBody:
required: true
content:
application/json:
schema:
$ref: "#/components/schemas/SteerRunRequest"
responses:
"202":
description: Steer accepted and forwarded to the worker
"400":
description: Invalid request body
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"
content:
application/json:
schema:
$ref: "#/components/schemas/ErrorResponse"
"404":
description: Run not found
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"
content:
application/json:
schema:
$ref: "#/components/schemas/ErrorResponse"
"409":
description: |
Run is not currently steerable. Returned when the run is in a
terminal state, blocked (use the answer endpoint instead), or
all currently running agent stages are CLI-mode.
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"
content:
application/json:
schema:
$ref: "#/components/schemas/ErrorResponse"
"503":
description: Worker control channel unavailable
headers:
x-request-id:
$ref: "#/components/headers/XRequestId"
content:
application/json:
schema:
$ref: "#/components/schemas/ErrorResponse"
/api/v1/runs/{id}/start:
post:
operationId: startRun
@ -4546,6 +4608,27 @@ components:
warn:
type: boolean
SteerRunRequest:
description: Request body for steering a running run mid-execution.
type: object
required:
- text
properties:
text:
type: string
description: The steering message text to deliver as a user turn.
minLength: 1
maxLength: 8192
example: Try a different approach
interrupt:
type: boolean
description: |
When true, cancel the in-flight LLM stream and tool calls in the
current round before delivering. When false (default), append to
the steering queue and let the agent pick it up at the next
turn boundary.
default: false
StartRunRequest:
description: Request body for starting or resuming a run.
type: object
@ -8322,4 +8405,4 @@ components:
login:
type: string
description: User's login identifier (e.g. GitHub username).
example: octocat
example: octocat

View file

@ -44,7 +44,7 @@ pub use sandbox::{
SandboxEvent, SandboxEventCallback, WorktreeEvent, WorktreeEventCallback, WorktreeOptions,
WorktreeSandbox, format_lines_numbered, shell_quote,
};
pub use session::Session;
pub use session::{CompletionCoordinator, Session, SessionControlHandle, SteeringItem};
pub use skills::Skill;
pub use subagent::{
SubAgent, SubAgentEventCallback, SubAgentManager, SubAgentResult, SubAgentStatus,
@ -55,7 +55,7 @@ pub use tools::{
make_shell_tool, make_shell_tool_with_config, make_write_file_tool, register_core_tools,
};
pub use truncation::{TruncationMode, truncate_lines, truncate_output, truncate_tool_output};
pub use types::{AgentEvent, SessionEvent, SessionState, Turn};
pub use types::{AgentEvent, SessionEvent, SessionState, SteerKind, Turn};
#[cfg(test)]
#[allow(

View file

@ -1,5 +1,5 @@
use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex};
use std::sync::{Arc, Mutex, RwLock};
use std::time::SystemTime;
use fabro_auth::CredentialSource;
@ -13,6 +13,7 @@ use fabro_llm::types::{
use fabro_llm::{Error as LlmError, retry};
use fabro_mcp::config::{McpServerSettings, McpTransport};
use fabro_mcp::connection_manager::McpConnectionManager;
use fabro_types::{Principal, SteerKind};
use futures::StreamExt;
use tokio::sync::{Mutex as AsyncMutex, broadcast};
use tokio::time;
@ -38,26 +39,139 @@ use crate::subagent::{SubAgentCallbackEvent, SubAgentEventCallback, SubAgentMana
use crate::tool_execution::execute_tool_calls;
use crate::types::{AgentEvent, SessionEvent, SessionState, Turn};
/// One queued steering message: text + delivery kind + the principal that
/// authored it (None for direct internal callers like loop-detection).
pub type SteeringItem = (String, SteerKind, Option<Principal>);
/// Trait that lets the workflow layer keep an agent in `process_input` when a
/// natural completion (no tool calls) coincides with an unconsumed steering
/// message. The implementation must coordinate with the steering source so
/// that, once it returns `false`, no further steers can race into the queue
/// for this session.
pub trait CompletionCoordinator: Send + Sync {
/// Called inside the agent loop when the assistant finishes a turn with
/// no tool calls. Return `true` to continue (the session will iterate
/// once more and drain pending steering messages); `false` to break out
/// of the loop normally.
fn on_natural_completion(&self) -> bool;
}
/// Cheap clone of the parts of a `Session` that an external coordinator
/// (e.g. the workflow `SteeringHub`) needs to deliver steering messages and
/// interrupt the current round without holding the session itself.
#[derive(Clone)]
pub struct SessionControlHandle {
queue: Arc<Mutex<VecDeque<SteeringItem>>>,
round_token: Arc<RwLock<CancellationToken>>,
}
impl Default for SessionControlHandle {
fn default() -> Self {
Self::new()
}
}
impl SessionControlHandle {
/// Build an unattached handle for testing or direct construction by
/// callers that want to wire a queue into something other than a live
/// `Session`. Both pieces are independent `Arc` values; cloning the
/// handle clones the `Arc`s.
#[must_use]
pub fn new() -> Self {
Self {
queue: Arc::new(Mutex::new(VecDeque::new())),
round_token: Arc::new(RwLock::new(CancellationToken::new())),
}
}
/// Push an `Append`-kind steering message onto the queue.
pub fn steer(&self, text: String, actor: Option<Principal>) {
self.queue
.lock()
.expect("steering queue lock poisoned")
.push_back((text, SteerKind::Append, actor));
}
/// Push an `Interrupt`-kind steering message onto the queue *and* cancel
/// the current round so the agent loop reacts immediately.
pub fn interrupt_with(&self, text: String, actor: Option<Principal>) {
self.queue
.lock()
.expect("steering queue lock poisoned")
.push_back((text, SteerKind::Interrupt, actor));
self.round_token
.read()
.expect("round token lock poisoned")
.cancel();
}
/// Direct enqueue used by callers that already encoded a kind/actor
/// (e.g. the hub flushing buffered steers as Appends).
pub fn enqueue(&self, item: SteeringItem) {
let cancel_after = matches!(item.1, SteerKind::Interrupt);
self.queue
.lock()
.expect("steering queue lock poisoned")
.push_back(item);
if cancel_after {
self.round_token
.read()
.expect("round token lock poisoned")
.cancel();
}
}
/// Whether the steering queue currently has no unconsumed messages.
#[must_use]
pub fn queue_is_empty(&self) -> bool {
self.queue
.lock()
.expect("steering queue lock poisoned")
.is_empty()
}
/// Current length of the steering queue.
#[must_use]
pub fn queue_len(&self) -> usize {
self.queue
.lock()
.expect("steering queue lock poisoned")
.len()
}
/// Pop the oldest queued steer (used by the hub to enforce its per-session
/// cap with FIFO eviction).
#[must_use]
pub fn pop_oldest(&self) -> Option<SteeringItem> {
self.queue
.lock()
.expect("steering queue lock poisoned")
.pop_front()
}
}
pub struct Session {
id: String,
config: SessionOptions,
history: History,
event_emitter: Emitter,
state: SessionState,
llm_client: Client,
provider_profile: Arc<dyn AgentProfile>,
sandbox: Arc<dyn Sandbox>,
steering_queue: Arc<Mutex<VecDeque<String>>>,
followup_queue: Arc<Mutex<VecDeque<String>>>,
cancel_token: CancellationToken,
interrupt_reason: Arc<Mutex<Option<InterruptReason>>>,
memory: Vec<String>,
env_context: EnvContext,
skills: Vec<Skill>,
system_prompt: String,
file_tracker: FileTracker,
tool_env: Option<HashMap<String, String>>,
subagent_manager: Option<Arc<AsyncMutex<SubAgentManager>>>,
id: String,
config: SessionOptions,
history: History,
event_emitter: Emitter,
state: SessionState,
llm_client: Client,
provider_profile: Arc<dyn AgentProfile>,
sandbox: Arc<dyn Sandbox>,
steering_queue: Arc<Mutex<VecDeque<SteeringItem>>>,
followup_queue: Arc<Mutex<VecDeque<String>>>,
cancel_token: CancellationToken,
round_token: Arc<RwLock<CancellationToken>>,
interrupt_reason: Arc<Mutex<Option<InterruptReason>>>,
memory: Vec<String>,
env_context: EnvContext,
skills: Vec<Skill>,
system_prompt: String,
file_tracker: FileTracker,
tool_env: Option<HashMap<String, String>>,
subagent_manager: Option<Arc<AsyncMutex<SubAgentManager>>>,
completion_coordinator: Option<Arc<dyn CompletionCoordinator>>,
}
impl Session {
@ -81,6 +195,7 @@ impl Session {
steering_queue: Arc::new(Mutex::new(VecDeque::new())),
followup_queue: Arc::new(Mutex::new(VecDeque::new())),
cancel_token: CancellationToken::new(),
round_token: Arc::new(RwLock::new(CancellationToken::new())),
interrupt_reason: Arc::new(Mutex::new(None)),
memory: Vec::new(),
env_context: EnvContext::default(),
@ -89,6 +204,7 @@ impl Session {
file_tracker: FileTracker::default(),
tool_env: None,
subagent_manager,
completion_coordinator: None,
}
}
@ -413,11 +529,49 @@ impl Session {
self.event_emitter.subscribe()
}
/// Push an `Append`-kind steer onto the queue (no actor — internal
/// callers like loop-detection use this).
pub fn steer(&self, message: String) {
self.steering_queue
.lock()
.expect("steering queue lock poisoned")
.push_back(message);
.push_back((message, SteerKind::Append, None));
}
/// Push an `Interrupt`-kind steer onto the queue and cancel the current
/// round token so the agent loop reacts mid-round.
pub fn interrupt_with(&self, message: String, actor: Option<Principal>) {
self.steering_queue
.lock()
.expect("steering queue lock poisoned")
.push_back((message, SteerKind::Interrupt, actor));
self.round_token
.read()
.expect("round token lock poisoned")
.cancel();
}
/// Cheap, cloneable handle that lets external coordinators deliver
/// steers and trigger interrupts without owning the `Session` itself.
#[must_use]
pub fn control_handle(&self) -> SessionControlHandle {
SessionControlHandle {
queue: self.steering_queue.clone(),
round_token: self.round_token.clone(),
}
}
/// Install a coordinator that decides whether `process_input` should
/// keep iterating after a no-tool turn. Used by the workflow layer to
/// race-safely include any steers that arrived during the final
/// response.
pub fn set_completion_coordinator(&mut self, coordinator: Arc<dyn CompletionCoordinator>) {
self.completion_coordinator = Some(coordinator);
}
/// Remove any installed completion coordinator.
pub fn clear_completion_coordinator(&mut self) {
self.completion_coordinator = None;
}
pub fn follow_up(&self, message: String) {
@ -493,7 +647,7 @@ impl Session {
}
#[must_use]
pub fn steering_queue_handle(&self) -> Arc<Mutex<VecDeque<String>>> {
pub fn steering_queue_handle(&self) -> Arc<Mutex<VecDeque<SteeringItem>>> {
self.steering_queue.clone()
}
@ -690,12 +844,30 @@ impl Session {
text: expanded_input.clone(),
});
// Drain steering queue before first LLM call
self.drain_steering();
let mut round_count: usize = 0;
loop {
// Top-of-loop: if the previous round's interrupt token fired,
// swap in a fresh one before draining and rebuilding state.
// (Terminal cancel via `cancel_token` is handled by the explicit
// check below and by `interrupted_error()`.)
{
let needs_refresh = self
.round_token
.read()
.expect("round token lock poisoned")
.is_cancelled();
if needs_refresh {
*self.round_token.write().expect("round token lock poisoned") =
CancellationToken::new();
}
}
// Drain pending steering messages at the top of every iteration
// so an Interrupt-kind steer pushed mid-round is delivered as the
// first turn of the next round.
self.drain_steering();
// Check max_tool_rounds_per_input
if self.config.max_tool_rounds_per_input > 0
&& round_count >= self.config.max_tool_rounds_per_input
@ -722,6 +894,13 @@ impl Session {
return Err(self.interrupted_error());
}
// Snapshot the per-round token; it stays stable for this iteration.
let round_token = self
.round_token
.read()
.expect("round token lock poisoned")
.clone();
// Pre-turn compaction: trim context before building the request
self.compact_if_needed().await;
@ -751,21 +930,54 @@ impl Session {
..Default::default()
};
let client = self.llm_client.clone();
let mut event_stream = self
.open_stream_with_retry(&client, &request, &retry_policy)
.await?;
let cancel_token_for_select = self.cancel_token.clone();
let stream_outcome: Option<Result<StreamEventStream, Error>> = tokio::select! {
biased;
() = round_token.cancelled() => None,
() = cancel_token_for_select.cancelled() => None,
stream = self.open_stream_with_retry(&client, &request, &retry_policy) => Some(stream),
};
let mut event_stream = if let Some(stream) = stream_outcome {
stream?
} else {
if self.cancel_token.is_cancelled() {
self.close();
return Err(self.interrupted_error());
}
// Round-only cancel before stream opened — re-iterate to
// pick up the steer.
continue;
};
// Consume the stream, retrying up to 3 times if the provider
// closes the stream without sending a Finish event. If visible
// output was already emitted, clear it before replaying the turn.
let mut response = None;
// Set true if a steer-interrupt cancelled the round mid-stream so
// we can clear partial output and `continue` after the loop.
let mut steer_interrupted = false;
let mut emitted_anything = false;
for stream_attempt in 0..=STREAM_CONSUME_RETRIES {
'streamattempts: for stream_attempt in 0..=STREAM_CONSUME_RETRIES {
let mut accumulator = StreamAccumulator::new();
let mut emitted_text = String::new();
let mut emitted_reasoning = String::new();
while let Some(event_result) = event_stream.next().await {
loop {
let chunk = tokio::select! {
biased;
() = round_token.cancelled() => None,
() = self.cancel_token.cancelled() => None,
next = event_stream.next() => Some(next),
};
let Some(event_opt) = chunk else {
// One of the cancellation tokens fired.
break;
};
let Some(event_result) = event_opt else {
// Stream ended normally.
break;
};
match event_result {
Ok(event) => {
match &event {
@ -795,21 +1007,28 @@ impl Session {
return Err(self.emit_llm_error(err));
}
}
// Check cancellation between chunks
if self.cancel_token.is_cancelled() {
break;
}
}
// If interrupted during streaming, drop the stream to cancel the HTTP
// connection, then close the session before returning.
// Track whether anything was rendered this attempt.
if !emitted_text.is_empty() || !emitted_reasoning.is_empty() {
emitted_anything = true;
}
// If terminal cancel fired, drop the stream and bail out.
if self.cancel_token.is_cancelled() {
drop(event_stream);
self.close();
return Err(self.interrupted_error());
}
// If only the round token fired (steer interrupt), drop the
// stream now; we'll clear partial output and continue below.
if round_token.is_cancelled() {
drop(event_stream);
steer_interrupted = true;
break 'streamattempts;
}
if let Some(resp) = accumulator.response().cloned() {
response = Some(resp);
break;
@ -831,12 +1050,37 @@ impl Session {
},
);
}
event_stream = self
.open_stream_with_retry(&client, &request, &retry_policy)
.await?;
let cancel_token_for_select = self.cancel_token.clone();
let retry_outcome: Option<Result<StreamEventStream, Error>> = tokio::select! {
biased;
() = round_token.cancelled() => None,
() = cancel_token_for_select.cancelled() => None,
stream = self.open_stream_with_retry(&client, &request, &retry_policy) => Some(stream),
};
event_stream = if let Some(stream) = retry_outcome {
stream?
} else {
steer_interrupted =
round_token.is_cancelled() && !self.cancel_token.is_cancelled();
break 'streamattempts;
};
}
}
// Mid-LLM steer interrupt: drop the unrecorded turn, clear any
// partial visible output, and re-iterate. The next turn's
// top-of-loop drain delivers the steer as the next user message.
if steer_interrupted {
if emitted_anything {
self.event_emitter
.emit(self.id.clone(), AgentEvent::AssistantOutputReplace {
text: String::new(),
reasoning: None,
});
}
continue;
}
let Some(response) = response else {
return Err(self.emit_llm_error(LlmError::Stream {
message: "Stream ended without a Finish event (after retries)".into(),
@ -877,13 +1121,38 @@ impl Session {
// Post-response compaction: trim context after appending assistant turn
self.compact_if_needed().await;
// If no tool calls, natural completion
// If no tool calls, natural completion. Consult the optional
// completion coordinator: it can return `true` to force one more
// iteration when a steer arrived during the final response.
if tool_calls.is_empty() {
let should_continue = self
.completion_coordinator
.as_ref()
.is_some_and(|c| c.on_natural_completion());
if should_continue {
continue;
}
break;
}
round_count += 1;
// Build a composite cancellation token covering both terminal
// cancel and round (steer) interrupt. Tools observe it
// cooperatively — they synthesize "Cancelled" results rather
// than being dropped mid-flight, which preserves the
// tool_use ↔ tool_result invariant.
let composite_token = CancellationToken::new();
let composite_for_cancel = composite_token.clone();
let cancel_token_clone = self.cancel_token.clone();
let round_token_clone = round_token.clone();
let composite_watcher = tokio::spawn(async move {
tokio::select! {
() = cancel_token_clone.cancelled() => composite_for_cancel.cancel(),
() = round_token_clone.cancelled() => composite_for_cancel.cancel(),
}
});
// Execute tool calls (parallel or sequential based on provider)
self.transition(SessionState::Executing);
let results = execute_tool_calls(
@ -892,36 +1161,39 @@ impl Session {
self.provider_profile.tool_registry(),
self.sandbox.clone(),
self.config.tool_hooks.as_ref(),
&self.cancel_token,
&composite_token,
&self.config,
&self.event_emitter,
&self.id,
self.tool_env.as_ref(),
)
.await;
composite_watcher.abort();
// Track file operations from tool calls
self.file_tracker
.record_from_tool_calls(&tool_calls, &results);
// Check cancellation after tool execution
if self.cancel_token.is_cancelled() {
self.history.push(Turn::ToolResults {
results,
timestamp: SystemTime::now(),
});
self.close();
return Err(self.interrupted_error());
}
// Record tool results turn
// Always append tool_results so the tool_use ↔ tool_result
// invariant holds, regardless of which token fired.
self.history.push(Turn::ToolResults {
results,
timestamp: SystemTime::now(),
});
// Drain steering after tool execution
self.drain_steering();
// Terminal cancel takes precedence: close and return.
if self.cancel_token.is_cancelled() {
self.close();
return Err(self.interrupted_error());
}
// Round-only cancel (steer interrupt mid-tool): re-iterate;
// the next top-of-loop drain delivers the steer.
if round_token.is_cancelled() {
self.transition(SessionState::Thinking);
continue;
}
self.transition(SessionState::Thinking);
// Loop detection
@ -970,20 +1242,23 @@ impl Session {
}
fn drain_steering(&mut self) {
let messages: Vec<String> = self
let messages: Vec<SteeringItem> = self
.steering_queue
.lock()
.expect("steering queue lock poisoned")
.drain(..)
.collect();
for msg in messages {
let text = msg.clone();
for (text, kind, actor) in messages {
self.history.push(Turn::Steering {
content: msg,
content: text.clone(),
timestamp: SystemTime::now(),
});
self.event_emitter
.emit(self.id.clone(), AgentEvent::SteeringInjected { text });
.emit(self.id.clone(), AgentEvent::SteeringInjected {
text,
kind,
actor,
});
}
}
@ -1266,6 +1541,98 @@ mod tests {
assert!(matches!(&turns[2], Turn::Assistant { .. }));
}
#[tokio::test]
async fn steer_event_carries_append_kind() {
let mut session = make_session(vec![text_response("OK")]).await;
let mut rx = session.subscribe();
session.steer("hi there".to_string());
session.process_input("Do something").await.unwrap();
// Drain events; find the SteeringInjected one and assert kind=Append
let mut found_kind = None;
while let Ok(ev) = rx.try_recv() {
if let AgentEvent::SteeringInjected { kind, .. } = ev.event {
found_kind = Some(kind);
break;
}
}
assert_eq!(found_kind, Some(SteerKind::Append));
}
#[tokio::test]
async fn interrupt_with_pushes_interrupt_kind_event() {
let mut session = make_session(vec![text_response("OK")]).await;
let mut rx = session.subscribe();
// Get a control handle, then push an interrupt steer before any
// round runs. Since no LLM call is in flight, the round_token
// cancel just makes the first iteration loop once before
// proceeding — drain_steering at top-of-loop will pick up the
// queued (text, Interrupt) item.
let handle = session.control_handle();
handle.interrupt_with("stop now".to_string(), None);
session.process_input("start").await.unwrap();
let mut found_kind = None;
while let Ok(ev) = rx.try_recv() {
if let AgentEvent::SteeringInjected { kind, .. } = ev.event {
found_kind = Some(kind);
break;
}
}
assert_eq!(found_kind, Some(SteerKind::Interrupt));
}
#[tokio::test]
async fn append_during_final_response_triggers_extra_round_when_coordinator_returns_true() {
use std::sync::atomic::{AtomicUsize, Ordering};
struct OnceCoordinator {
calls: AtomicUsize,
handle: SessionControlHandle,
}
impl CompletionCoordinator for OnceCoordinator {
fn on_natural_completion(&self) -> bool {
let n = self.calls.fetch_add(1, Ordering::SeqCst);
if n == 0 {
// Simulate a steer that arrived during the first
// completion: enqueue and report "keep going".
self.handle
.steer("after-completion steer".to_string(), None);
true
} else {
false
}
}
}
// First scripted response is a no-tool natural completion; second
// also natural completion. The completion coordinator forces the
// loop to iterate once more — that iteration must drain the queued
// steer and produce a second Assistant turn.
let responses = vec![
text_response("First reply"),
text_response("Second reply, after steer"),
];
let mut session = make_session(responses).await;
let handle = session.control_handle();
session.set_completion_coordinator(Arc::new(OnceCoordinator {
calls: AtomicUsize::new(0),
handle,
}));
session.process_input("hi").await.unwrap();
let turns = session.history().turns();
// User + Assistant + Steering + Assistant = 4
assert_eq!(turns.len(), 4);
assert!(matches!(&turns[0], Turn::User { .. }));
assert!(matches!(&turns[1], Turn::Assistant { content, .. } if content == "First reply"));
assert!(matches!(&turns[2], Turn::Steering { content, .. }
if content == "after-completion steer"));
assert!(matches!(&turns[3], Turn::Assistant { content, .. }
if content == "Second reply, after steer"));
}
#[tokio::test]
async fn follow_up_triggers_new_cycle() {
let responses = vec![

View file

@ -2,6 +2,7 @@ use std::time::SystemTime;
use fabro_llm::Error as LlmError;
use fabro_llm::types::{ContentPart, ThinkingData, TokenCounts, ToolCall, ToolResult};
pub use fabro_types::SteerKind;
use serde::{Deserialize, Serialize};
use crate::error::Error;
@ -157,7 +158,13 @@ pub enum AgentEvent {
skill_name: String,
},
SteeringInjected {
text: String,
text: String,
kind: SteerKind,
/// Principal that authored the steer. Lifted to top-level
/// `RunEvent.actor` by the workflow event-conversion layer; never
/// serialized into event props.
#[serde(default, skip_serializing_if = "Option::is_none")]
actor: Option<fabro_types::Principal>,
},
CompactionStarted {
estimated_tokens: usize,
@ -309,8 +316,13 @@ impl AgentEvent {
Self::SkillExpanded { skill_name } => {
debug!(session_id, skill = skill_name.as_str(), "Skill expanded");
}
Self::SteeringInjected { text } => {
debug!(session_id, text_len = text.len(), "Steering injected");
Self::SteeringInjected { text, kind, .. } => {
debug!(
session_id,
text_len = text.len(),
kind = kind.as_str(),
"Steering injected"
);
}
Self::CompactionStarted {
estimated_tokens,

View file

@ -605,9 +605,11 @@ async fn scenario_steering_mid_task(session: &mut Session, dir: &Path) {
steering_queue
.lock()
.expect("steering queue lock")
.push_back(
.push_back((
"Stop what you are doing. Create a file called steered.txt containing 'steered' and do nothing else.".to_string(),
);
fabro_agent::SteerKind::Append,
None,
));
break;
}
}

View file

@ -670,6 +670,27 @@ pub(crate) struct WaitArgs {
pub(crate) interval: u64,
}
#[derive(Args)]
pub(crate) struct SteerArgs {
#[command(flatten)]
pub(crate) server: ServerTargetArgs,
/// Run ID prefix to steer
pub(crate) run: String,
/// Steer message text (omit when --text-stdin is used)
pub(crate) text: Option<String>,
/// Read steer text from stdin instead of a positional arg
#[arg(long, conflicts_with = "text")]
pub(crate) text_stdin: bool,
/// Cancel the in-flight LLM stream / tool calls and deliver the message
/// as the next user turn (default: append to the steering queue).
#[arg(long)]
pub(crate) interrupt: bool,
}
#[derive(Args)]
pub(crate) struct WorkflowListArgs;
@ -944,6 +965,8 @@ pub(crate) enum RunCommands {
Fork(ForkArgs),
/// Block until a workflow run completes
Wait(WaitArgs),
/// Steer a running agent mid-execution
Steer(SteerArgs),
}
impl RunCommands {
@ -958,6 +981,7 @@ impl RunCommands {
Self::Logs(_) => "logs",
Self::Resume(_) => "resume",
Self::Rewind(_) => "rewind",
Self::Steer(_) => "steer",
Self::Fork(_) => "fork",
Self::Wait(_) => "wait",
}

View file

@ -25,6 +25,7 @@ pub(crate) mod run_progress;
pub(crate) mod runner;
pub(crate) mod ssh;
pub(crate) mod start;
pub(crate) mod steer;
pub(crate) mod wait;
pub(crate) async fn dispatch(
@ -123,5 +124,6 @@ pub(crate) async fn dispatch(
let styles = Styles::detect_stderr();
wait::run(&args, &styles, base_ctx).await
}
RunCommands::Steer(args) => steer::run(args, base_ctx).await,
}
}

View file

@ -87,7 +87,13 @@ pub(crate) async fn execute(
)));
let interviewer = Arc::new(ControlInterviewer::new());
let cancel_token = Arc::new(AtomicBool::new(false));
spawn_worker_control_stream(Arc::clone(&interviewer), Arc::clone(&cancel_token))?;
let emitter = Arc::new(Emitter::new(run_id));
let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::clone(&emitter)));
spawn_worker_control_stream(
Arc::clone(&interviewer),
Arc::clone(&cancel_token),
Arc::clone(&steering_hub),
)?;
let run_control = RunControlState::new();
install_signal_handlers(Arc::clone(&run_control), Arc::clone(&cancel_token))?;
let vault = load_worker_vault(storage_dir.as_deref())?;
@ -101,8 +107,9 @@ pub(crate) async fn execute(
let services = StartServices {
run_id,
cancel_token: Some(Arc::clone(&cancel_token)),
emitter: Arc::new(Emitter::new(run_id)),
emitter,
interviewer,
steering_hub,
run_store: run_store.clone(),
event_sink: RunEventSink::map(
stamp_system_worker,
@ -163,11 +170,13 @@ enum WorkerControlStreamEvent {
fn spawn_worker_control_stream(
interviewer: Arc<ControlInterviewer>,
cancel_token: Arc<AtomicBool>,
steering_hub: Arc<fabro_workflow::SteeringHub>,
) -> Result<()> {
let (event_tx, event_rx) = mpsc::unbounded_channel();
tokio::spawn(handle_worker_control_stream_events(
interviewer,
cancel_token,
steering_hub,
event_rx,
));
std::thread::Builder::new()
@ -206,12 +215,13 @@ fn read_worker_control_stream_blocking<R>(
async fn handle_worker_control_stream_events(
interviewer: Arc<ControlInterviewer>,
cancel_token: Arc<AtomicBool>,
steering_hub: Arc<fabro_workflow::SteeringHub>,
mut event_rx: mpsc::UnboundedReceiver<WorkerControlStreamEvent>,
) {
while let Some(event) = event_rx.recv().await {
match event {
WorkerControlStreamEvent::Line(line) => {
apply_worker_control_line(&interviewer, &cancel_token, &line).await;
apply_worker_control_line(&interviewer, &cancel_token, &steering_hub, &line).await;
}
WorkerControlStreamEvent::Eof => {
interviewer.interrupt_all().await;
@ -226,6 +236,7 @@ async fn handle_worker_control_stream_events(
async fn apply_worker_control_line(
interviewer: &ControlInterviewer,
cancel_token: &AtomicBool,
steering_hub: &fabro_workflow::SteeringHub,
line: &str,
) {
if line.trim().is_empty() {
@ -246,6 +257,9 @@ async fn apply_worker_control_line(
cancel_token.store(true, Ordering::SeqCst);
interviewer.interrupt_all().await;
}
WorkerControlMessage::Steer { text, kind, actor } => {
steering_hub.deliver(text, kind, Some(actor));
}
}
}
@ -632,6 +646,11 @@ mod tests {
};
use crate::args::RunWorkerMode;
fn test_steering_hub() -> Arc<fabro_workflow::SteeringHub> {
let emitter = Arc::new(fabro_workflow::event::Emitter::new(fixtures::RUN_1));
Arc::new(fabro_workflow::SteeringHub::new(emitter))
}
#[test]
fn clone_sandbox_credentials_are_required_for_clone_based_providers() {
assert!(super::clone_sandbox_requires_github_credentials("docker"));
@ -829,9 +848,11 @@ mod tests {
let ask_interviewer = Arc::clone(&interviewer);
let answer_task = tokio::spawn(async move { ask_interviewer.ask(question).await });
let hub = test_steering_hub();
apply_worker_control_line(
&interviewer,
&cancel_token,
&hub,
r#"{"v":1,"type":"interview.answer","qid":"q-1","answer":{"kind":"yes"},"actor":{"kind":"system","system_kind":"engine"}}"#,
)
.await;
@ -851,9 +872,11 @@ mod tests {
let answer_task = tokio::spawn(async move { ask_interviewer.ask(question).await });
tokio::task::yield_now().await;
let hub = test_steering_hub();
apply_worker_control_line(
&interviewer,
&cancel_token,
&hub,
r#"{"v":1,"type":"run.cancel"}"#,
)
.await;
@ -903,9 +926,11 @@ mod tests {
event_tx.send(WorkerControlStreamEvent::Eof).unwrap();
drop(event_tx);
let hub = test_steering_hub();
handle_worker_control_stream_events(
Arc::clone(&interviewer),
Arc::clone(&cancel_token),
hub,
event_rx,
)
.await;

View file

@ -0,0 +1,32 @@
use anyhow::{Result, bail};
use tokio::io::{AsyncReadExt as _, stdin};
use tracing::info;
use crate::args::SteerArgs;
use crate::command_context::CommandContext;
pub(crate) async fn run(args: SteerArgs, base_ctx: &CommandContext) -> Result<()> {
let ctx = base_ctx.with_target(&args.server)?;
let client = ctx.server().await?;
let run_id = client.resolve_run(&args.run).await?.run_id;
let text = match (args.text_stdin, args.text.clone()) {
(true, _) => {
let mut buf = String::new();
stdin().read_to_string(&mut buf).await?;
buf
}
(false, Some(text)) => text,
(false, None) => {
bail!("missing steer text — pass it as a positional argument or use --text-stdin")
}
};
let text = text.trim().to_string();
if text.is_empty() {
bail!("steer text must not be empty");
}
info!(run_id = %run_id, interrupt = args.interrupt, "Sending steer");
client.steer_run(&run_id, text, args.interrupt).await?;
Ok(())
}

View file

@ -21,6 +21,7 @@ fn help() {
rewind Rewind a workflow run to an earlier checkpoint
fork Fork a workflow run from an earlier checkpoint into a new run
wait Block until a workflow run completes
steer Steer a running agent mid-execution
preflight Validate run configuration without executing
validate Validate a workflow
graph Render a workflow graph as SVG

View file

@ -805,6 +805,27 @@ impl Client {
Ok(())
}
pub async fn steer_run(&self, run_id: &RunId, text: String, interrupt: bool) -> Result<()> {
let body: types::SteerRunRequest = types::SteerRunRequest::builder()
.text(text)
.interrupt(interrupt)
.try_into()
.map_err(|e| anyhow!("failed to build SteerRunRequest: {e}"))?;
self.send_api(|client| {
let body = body.clone();
async move {
client
.steer_run()
.id(run_id.to_string())
.body(body)
.send()
.await
}
})
.await?;
Ok(())
}
pub async fn archive_run(&self, run_id: &RunId) -> Result<()> {
self.send_api(
|client| async move { client.archive_run().id(run_id.to_string()).send().await },

View file

@ -1,4 +1,5 @@
use fabro_types::Principal;
pub use fabro_types::SteerKind;
use serde::{Deserialize, Serialize};
use crate::{Answer, AnswerSubmission, AnswerValue};
@ -32,6 +33,18 @@ impl WorkerControlEnvelope {
message: WorkerControlMessage::RunCancel,
}
}
#[must_use]
pub fn steer(text: impl Into<String>, kind: SteerKind, actor: Principal) -> Self {
Self {
v: WORKER_CONTROL_PROTOCOL_VERSION,
message: WorkerControlMessage::Steer {
text: text.into(),
kind,
actor,
},
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@ -45,6 +58,12 @@ pub enum WorkerControlMessage {
},
#[serde(rename = "run.cancel")]
RunCancel,
#[serde(rename = "run.steer")]
Steer {
text: String,
kind: SteerKind,
actor: Principal,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@ -129,4 +148,37 @@ mod tests {
let parsed: WorkerControlEnvelope = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, envelope);
}
#[test]
fn steer_append_round_trips_through_json() {
let envelope = WorkerControlEnvelope::steer(
"try again",
SteerKind::Append,
fabro_types::Principal::System {
system_kind: fabro_types::SystemActorKind::Engine,
},
);
let json = serde_json::to_string(&envelope).unwrap();
assert_eq!(
json,
r#"{"v":1,"type":"run.steer","text":"try again","kind":"append","actor":{"kind":"system","system_kind":"engine"}}"#
);
let parsed: WorkerControlEnvelope = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, envelope);
}
#[test]
fn steer_interrupt_round_trips_through_json() {
let envelope = WorkerControlEnvelope::steer(
"stop, do X instead",
SteerKind::Interrupt,
fabro_types::Principal::System {
system_kind: fabro_types::SystemActorKind::Engine,
},
);
let json = serde_json::to_string(&envelope).unwrap();
assert!(json.contains(r#""kind":"interrupt""#));
let parsed: WorkerControlEnvelope = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, envelope);
}
}

View file

@ -222,7 +222,7 @@ pub use callback::CallbackInterviewer;
pub use console::ConsoleInterviewer;
pub use control::{ControlInterviewer, SubmitError};
pub use control_protocol::{
WORKER_CONTROL_PROTOCOL_VERSION, WorkerControlAnswer, WorkerControlEnvelope,
SteerKind, WORKER_CONTROL_PROTOCOL_VERSION, WorkerControlAnswer, WorkerControlEnvelope,
WorkerControlMessage,
};
pub use queue::QueueInterviewer;

View file

@ -194,6 +194,14 @@ struct ManagedRun {
// Populated when running:
answer_transport: Option<RunAnswerTransport>,
accepted_questions: HashSet<String>,
/// Stage IDs of currently running API-mode (SDK) agent sessions, as
/// observed from the worker's `agent.steering.attached/detached`
/// events. Used by the steerability predicate.
active_api_stages: HashSet<StageId>,
/// Stage IDs of currently running CLI-mode agent sessions, observed
/// from `agent.cli.started/completed` plus `stage.completed`/
/// `stage.failed` backstops.
active_cli_stages: HashSet<StageId>,
event_tx: Option<broadcast::Sender<RunEvent>>,
checkpoint: Option<Checkpoint>,
cancel_tx: Option<oneshot::Sender<()>>,
@ -243,7 +251,8 @@ enum RunAnswerTransport {
control_tx: mpsc::Sender<WorkerControlEnvelope>,
},
InProcess {
interviewer: Arc<ControlInterviewer>,
interviewer: Arc<ControlInterviewer>,
steering_hub: Arc<fabro_workflow::SteeringHub>,
},
}
@ -267,7 +276,7 @@ impl RunAnswerTransport {
.map_err(|_| AnswerTransportError::Timeout)?
.map_err(|_| AnswerTransportError::Closed)
}
Self::InProcess { interviewer } => interviewer
Self::InProcess { interviewer, .. } => interviewer
.submit(qid, submission)
.await
.map_err(|_| AnswerTransportError::Closed),
@ -283,12 +292,35 @@ impl RunAnswerTransport {
.map_err(|_| AnswerTransportError::Timeout)?
.map_err(|_| AnswerTransportError::Closed)
}
Self::InProcess { interviewer } => {
Self::InProcess { interviewer, .. } => {
interviewer.cancel_all().await;
Ok(())
}
}
}
/// Forward a steer to the worker (subprocess) or directly into the
/// in-process steering hub.
async fn steer(
&self,
text: String,
kind: fabro_types::SteerKind,
actor: Principal,
) -> Result<(), AnswerTransportError> {
match self {
Self::Subprocess { control_tx } => {
let message = WorkerControlEnvelope::steer(text, kind, actor);
timeout(WORKER_CONTROL_ENQUEUE_TIMEOUT, control_tx.send(message))
.await
.map_err(|_| AnswerTransportError::Timeout)?
.map_err(|_| AnswerTransportError::Closed)
}
Self::InProcess { steering_hub, .. } => {
steering_hub.deliver(text, kind, Some(actor));
Ok(())
}
}
}
}
#[derive(Debug, Clone)]
@ -1752,6 +1784,8 @@ fn octet_stream_response(bytes: Bytes) -> Response {
fn clear_live_run_state(run: &mut ManagedRun) {
run.answer_transport = None;
run.accepted_questions.clear();
run.active_api_stages.clear();
run.active_cli_stages.clear();
run.event_tx = None;
run.cancel_tx = None;
run.cancel_token = None;
@ -2079,6 +2113,8 @@ fn managed_run(
enqueued_at: Instant::now(),
answer_transport: None,
accepted_questions: HashSet::new(),
active_api_stages: HashSet::new(),
active_cli_stages: HashSet::new(),
event_tx: None,
checkpoint: None,
cancel_tx: None,
@ -2174,12 +2210,16 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent)
reason: props.reason,
};
managed_run.error = None;
managed_run.active_api_stages.clear();
managed_run.active_cli_stages.clear();
}
EventBody::RunFailed(props) => {
managed_run.status = RunStatus::Failed {
reason: props.reason,
};
managed_run.error = Some(props.error.clone());
managed_run.active_api_stages.clear();
managed_run.active_cli_stages.clear();
}
EventBody::RunArchived(_) => {
if let Some(prior) = managed_run.status.terminal_status() {
@ -2191,6 +2231,41 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent)
managed_run.status = prior.into();
}
}
// Track API-mode steerable sessions. attached/detached fire
// deterministically inside `AgentApiBackend::run` (register/
// unregister), so they're the authoritative window in which a
// steer can be delivered to a live session.
EventBody::AgentSteeringAttached(_) => {
if let Some(stage_id) = event.stage_id.clone() {
managed_run.active_api_stages.insert(stage_id);
}
}
EventBody::AgentSteeringDetached(_) => {
if let Some(stage_id) = &event.stage_id {
managed_run.active_api_stages.remove(stage_id);
}
}
// Track CLI-mode agent stages. CLI started/completed are coarser
// and sometimes fail to emit `completed` on error paths — the
// stage.completed/stage.failed handler below is the backstop.
EventBody::AgentCliStarted(_) => {
if let Some(stage_id) = event.stage_id.clone() {
managed_run.active_cli_stages.insert(stage_id);
}
}
EventBody::AgentCliCompleted(_) => {
if let Some(stage_id) = &event.stage_id {
managed_run.active_cli_stages.remove(stage_id);
}
}
// Stage lifecycle backstop: cover both completion and failure
// paths so a failing CLI stage doesn't strand its entry.
EventBody::StageCompleted(_) | EventBody::StageFailed(_) => {
if let Some(stage_id) = &event.stage_id {
managed_run.active_api_stages.remove(stage_id);
managed_run.active_cli_stages.remove(stage_id);
}
}
_ => {}
}
}
@ -2614,6 +2689,7 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
.as_ref()
.map(|factory| Arc::new(factory(Arc::clone(&interview_runtime))));
let emitter = Arc::new(emitter);
let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::clone(&emitter)));
// Transition to Running, populate interviewer
let cancelled_during_setup = {
@ -2622,7 +2698,8 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
if managed_run.status == RunStatus::Starting {
managed_run.status = RunStatus::Running;
managed_run.answer_transport = Some(RunAnswerTransport::InProcess {
interviewer: Arc::clone(&interviewer),
interviewer: Arc::clone(&interviewer),
steering_hub: Arc::clone(&steering_hub),
});
false
} else {
@ -2744,6 +2821,7 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
cancel_token: Some(Arc::clone(&cancel_token)),
emitter: Arc::clone(&emitter),
interviewer: Arc::clone(&interview_runtime),
steering_hub: Arc::clone(&steering_hub),
run_store: run_store.clone().into(),
event_sink: workflow_event::RunEventSink::store(run_store.clone()),
artifact_sink: Some(ArtifactSink::Store(state.artifact_store.clone())),

View file

@ -16,6 +16,7 @@ mod pull_requests;
mod runs;
mod sandbox;
mod secrets;
mod steer;
pub(in crate::server) mod system;
pub(super) use system::{health, openapi_spec};
@ -114,7 +115,6 @@ pub(super) fn demo_routes() -> Router<Arc<AppState>> {
pub(super) fn real_routes() -> Router<Arc<AppState>> {
Router::new()
.route("/runs/{id}/stages/{stageId}/turns", get(not_implemented))
.route("/runs/{id}/steer", post(not_implemented))
.route("/workflows", get(not_implemented))
.route("/workflows/{name}", get(not_implemented))
.route("/workflows/{name}/runs", get(not_implemented))
@ -137,6 +137,7 @@ pub(super) fn real_routes() -> Router<Arc<AppState>> {
.merge(artifacts::routes())
.merge(sandbox::routes())
.merge(lifecycle::routes())
.merge(steer::routes())
.merge(graph::manifest_routes())
.merge(graph::run_routes())
.merge(models::routes())

View file

@ -0,0 +1,126 @@
use std::sync::Arc;
use axum::Json;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use axum::routing::post;
use fabro_api::types::SteerRunRequest;
use fabro_types::{Principal, SteerKind};
use fabro_workflow::run_status::RunStatus;
use super::super::{AnswerTransportError, AppState, parse_run_id_path, reject_if_archived};
use crate::error::ApiError;
use crate::principal_middleware::RequiredUser;
pub(super) fn routes() -> axum::Router<Arc<AppState>> {
axum::Router::new().route("/runs/{id}/steer", post(steer_run))
}
const MAX_STEER_TEXT_LEN: usize = 8192;
async fn steer_run(
auth: RequiredUser,
State(state): State<Arc<AppState>>,
Path(id): Path<String>,
Json(req): Json<SteerRunRequest>,
) -> Response {
let id = match parse_run_id_path(&id) {
Ok(id) => id,
Err(response) => return response,
};
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
return response;
}
// Body validation. (OpenAPI enforces minLength=1/maxLength=8192 at the
// type boundary already; we re-check defensively for trims/whitespace.)
let text = req.text.to_string();
let trimmed_len = text.trim().len();
if trimmed_len == 0 {
return ApiError::bad_request("Steer text must not be empty.").into_response();
}
if text.len() > MAX_STEER_TEXT_LEN {
return ApiError::bad_request(format!(
"Steer text must be ≤ {MAX_STEER_TEXT_LEN} characters."
))
.into_response();
}
let kind = if req.interrupt {
SteerKind::Interrupt
} else {
SteerKind::Append
};
// Status + steerability gate. Take the answer_transport snapshot under
// the same lock so we can hand it off without further state races.
let answer_transport = {
let runs = state.runs.lock().expect("runs lock poisoned");
let Some(managed_run) = runs.get(&id) else {
return ApiError::not_found("Run not found.").into_response();
};
match managed_run.status {
RunStatus::Blocked { .. } => {
return ApiError::with_code(
StatusCode::CONFLICT,
"Run is blocked on a question; use the interview-answer endpoint instead.",
"use_answer_endpoint",
)
.into_response();
}
RunStatus::Submitted
| RunStatus::Queued
| RunStatus::Starting
| RunStatus::Paused { .. } => {
return ApiError::new(StatusCode::CONFLICT, "Run is not currently running.")
.into_response();
}
RunStatus::Failed { .. }
| RunStatus::Succeeded { .. }
| RunStatus::Removing
| RunStatus::Dead
| RunStatus::Archived { .. } => {
return ApiError::new(StatusCode::CONFLICT, "Run is no longer steerable.")
.into_response();
}
RunStatus::Running => {}
}
// Steerability predicate. Best-effort, target-oriented:
// - If at least one API-mode session is active → forward.
// - Else if no agent stages are active at all → forward (worker hub buffers
// for the next session).
// - Else (active agents exist but all are CLI-mode) → 409.
if managed_run.active_api_stages.is_empty() && !managed_run.active_cli_stages.is_empty() {
return ApiError::with_code(
StatusCode::CONFLICT,
"All currently running agent stages are CLI-mode and cannot be steered.",
"cli_agent_not_steerable",
)
.into_response();
}
managed_run.answer_transport.clone()
};
let Some(answer_transport) = answer_transport else {
return ApiError::new(
StatusCode::SERVICE_UNAVAILABLE,
"Run has no live worker control channel.",
)
.into_response();
};
let actor = Principal::User(auth.0);
match answer_transport.steer(text, kind, actor).await {
Ok(()) => StatusCode::ACCEPTED.into_response(),
Err(AnswerTransportError::Timeout) => ApiError::new(
StatusCode::SERVICE_UNAVAILABLE,
"Worker control channel timed out.",
)
.into_response(),
Err(AnswerTransportError::Closed) => ApiError::new(
StatusCode::SERVICE_UNAVAILABLE,
"Worker control channel is closed.",
)
.into_response(),
}
}

View file

@ -1878,8 +1878,13 @@ async fn subprocess_answer_transport_cancel_run_enqueues_cancel_message() {
#[tokio::test]
async fn in_process_answer_transport_cancel_run_cancels_pending_interviews() {
let interviewer = Arc::new(ControlInterviewer::new());
let emitter = Arc::new(fabro_workflow::event::Emitter::new(
fabro_types::RunId::new(),
));
let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(emitter));
let transport = RunAnswerTransport::InProcess {
interviewer: Arc::clone(&interviewer),
interviewer: Arc::clone(&interviewer),
steering_hub: Arc::clone(&steering_hub),
};
let mut question = Question::new("Approve?", QuestionType::YesNo);
question.id = "q-1".to_string();
@ -5264,6 +5269,45 @@ async fn cancel_nonexistent_run_returns_not_found() {
assert_status!(response, StatusCode::NOT_FOUND).await;
}
#[tokio::test]
async fn steer_nonexistent_run_returns_not_found() {
let app = test_app_with();
let missing_run_id = fixtures::RUN_64;
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{missing_run_id}/steer")))
.header("content-type", "application/json")
.body(Body::from(r#"{"text":"try again"}"#))
.unwrap();
let response = app.oneshot(req).await.unwrap();
assert_status!(response, StatusCode::NOT_FOUND).await;
}
#[tokio::test]
async fn steer_empty_text_returns_bad_request() {
let state = test_app_state();
let app = crate::test_support::build_test_router(Arc::clone(&state));
let run_id = create_and_start_run(&app, MINIMAL_DOT)
.await
.parse::<RunId>()
.unwrap();
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/steer")))
.header("content-type", "application/json")
.body(Body::from(r#"{"text":" "}"#))
.unwrap();
let response = app.oneshot(req).await.unwrap();
// Could be 400 (empty text) or 409 (run not running yet) depending on
// timing — either way it must NOT be 202.
assert_ne!(response.status(), StatusCode::ACCEPTED);
}
#[tokio::test]
async fn get_graph_returns_svg() {
let state = test_app_state();

View file

@ -31,6 +31,7 @@ pub mod stage_completion;
pub mod stage_id;
pub mod start;
pub mod status;
pub mod steering;
pub use artifact::ArtifactUpload;
pub use auth::{IdpIdentity, IdpIdentityError};
@ -86,3 +87,4 @@ pub use status::{
BlockedReason, FailureReason, InvalidTransition, ParseFailureReasonError,
ParseSuccessReasonError, RunControlAction, RunStatus, SuccessReason, TerminalStatus,
};
pub use steering::SteerKind;

View file

@ -2,6 +2,7 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;
use super::BilledTokenCounts;
use crate::SteerKind;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentSessionStartedProps {
@ -82,9 +83,34 @@ pub struct AgentTurnLimitReachedProps {
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentSteeringInjectedProps {
pub text: String,
pub kind: SteerKind,
pub visit: u32,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct AgentSteeringAttachedProps;
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct AgentSteeringDetachedProps;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentSteerBufferedProps {
pub kind: SteerKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AgentSteerDroppedReason {
QueueFull,
RunEnded,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentSteerDroppedProps {
pub reason: AgentSteerDroppedReason,
pub count: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentCompactionStartedProps {
pub estimated_tokens: usize,

View file

@ -170,6 +170,14 @@ pub enum EventBody {
AgentTurnLimitReached(AgentTurnLimitReachedProps),
#[serde(rename = "agent.steering.injected")]
AgentSteeringInjected(AgentSteeringInjectedProps),
#[serde(rename = "agent.steering.attached")]
AgentSteeringAttached(AgentSteeringAttachedProps),
#[serde(rename = "agent.steering.detached")]
AgentSteeringDetached(AgentSteeringDetachedProps),
#[serde(rename = "agent.steer.buffered")]
AgentSteerBuffered(AgentSteerBufferedProps),
#[serde(rename = "agent.steer.dropped")]
AgentSteerDropped(AgentSteerDroppedProps),
#[serde(rename = "agent.compaction.started")]
AgentCompactionStarted(AgentCompactionStartedProps),
#[serde(rename = "agent.compaction.completed")]
@ -392,6 +400,10 @@ impl EventBody {
Self::AgentLoopDetected(_) => "agent.loop.detected",
Self::AgentTurnLimitReached(_) => "agent.turn.limit",
Self::AgentSteeringInjected(_) => "agent.steering.injected",
Self::AgentSteeringAttached(_) => "agent.steering.attached",
Self::AgentSteeringDetached(_) => "agent.steering.detached",
Self::AgentSteerBuffered(_) => "agent.steer.buffered",
Self::AgentSteerDropped(_) => "agent.steer.dropped",
Self::AgentCompactionStarted(_) => "agent.compaction.started",
Self::AgentCompactionCompleted(_) => "agent.compaction.completed",
Self::AgentLlmRetry(_) => "agent.llm.retry",
@ -524,6 +536,10 @@ fn is_known_event_name(event: &str) -> bool {
| "agent.loop.detected"
| "agent.turn.limit"
| "agent.steering.injected"
| "agent.steering.attached"
| "agent.steering.detached"
| "agent.steer.buffered"
| "agent.steer.dropped"
| "agent.compaction.started"
| "agent.compaction.completed"
| "agent.llm.retry"

View file

@ -0,0 +1,52 @@
use serde::{Deserialize, Serialize};
/// Two flavors of mid-run steering messages delivered to a live agent
/// session.
///
/// - `Append` — push to the steering queue; the agent picks it up at the next
/// turn boundary.
/// - `Interrupt` — cancel the in-flight LLM stream / tool call in the current
/// round, then deliver the message as the next user turn.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum SteerKind {
Append,
Interrupt,
}
impl SteerKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Append => "append",
Self::Interrupt => "interrupt",
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn append_round_trips_through_json() {
let json = serde_json::to_string(&SteerKind::Append).unwrap();
assert_eq!(json, "\"append\"");
let parsed: SteerKind = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, SteerKind::Append);
}
#[test]
fn interrupt_round_trips_through_json() {
let json = serde_json::to_string(&SteerKind::Interrupt).unwrap();
assert_eq!(json, "\"interrupt\"");
let parsed: SteerKind = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, SteerKind::Interrupt);
}
#[test]
fn unknown_value_fails_to_deserialize() {
let result: Result<SteerKind, _> = serde_json::from_str("\"unknown\"");
assert!(result.is_err());
}
}

View file

@ -605,9 +605,10 @@ fn event_body_from_event(event: &Event) -> EventBody {
visit: *visit,
})
}
AgentEvent::SteeringInjected { text } => {
AgentEvent::SteeringInjected { text, kind, .. } => {
EventBody::AgentSteeringInjected(fabro_types::AgentSteeringInjectedProps {
text: text.clone(),
kind: *kind,
visit: *visit,
})
}
@ -1006,6 +1007,21 @@ fn event_body_from_event(event: &Event) -> EventBody {
exit_code: *exit_code,
duration_ms: *duration_ms,
}),
Event::AgentSteeringAttached { .. } => {
EventBody::AgentSteeringAttached(fabro_types::AgentSteeringAttachedProps {})
}
Event::AgentSteeringDetached { .. } => {
EventBody::AgentSteeringDetached(fabro_types::AgentSteeringDetachedProps {})
}
Event::AgentSteerBuffered { kind, .. } => {
EventBody::AgentSteerBuffered(fabro_types::AgentSteerBufferedProps { kind: *kind })
}
Event::AgentSteerDropped { reason, count, .. } => {
EventBody::AgentSteerDropped(fabro_types::AgentSteerDroppedProps {
reason: *reason,
count: *count,
})
}
Event::PullRequestCreated {
pr_url,
pr_number,

View file

@ -3,7 +3,7 @@ use std::collections::BTreeMap;
use ::fabro_types::{
BilledTokenCounts, BlockedReason, CommandTermination, FailureReason, ForkSourceRef, GitContext,
ParallelBranchId, Principal, PullRequestRecord, RunBlobId, RunId, RunNoticeLevel,
RunProvenance, StageId, SuccessReason, run_event as fabro_types,
RunProvenance, StageId, SteerKind, SuccessReason, run_event as fabro_types,
};
use fabro_agent::{AgentEvent, SandboxEvent};
use serde::{Deserialize, Serialize};
@ -523,6 +523,36 @@ pub enum Event {
model: String,
command: String,
},
/// A `SteeringHub` registered an active API-mode session for a stage.
/// Emitted once per `register` insert (not on replace).
AgentSteeringAttached {
node_id: String,
visit: u32,
},
/// The corresponding session was unregistered from the hub.
AgentSteeringDetached {
node_id: String,
visit: u32,
},
/// A steer arrived with no active session and was parked in the run-wide
/// pending buffer. The actor (steer author) is lifted to top-level.
AgentSteerBuffered {
kind: SteerKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
actor: Option<Principal>,
},
/// One or more buffered/queued steers were dropped because a cap was
/// reached or the run ended before they could be delivered.
AgentSteerDropped {
reason: fabro_types::AgentSteerDroppedReason,
count: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
actor: Option<Principal>,
#[serde(default, skip_serializing_if = "Option::is_none")]
node_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
visit: Option<u32>,
},
AgentCliCompleted {
node_id: String,
stdout: String,
@ -1247,6 +1277,18 @@ impl Event {
} => {
debug!(node_id, exit_code, duration_ms, "Agent CLI completed");
}
Self::AgentSteeringAttached { node_id, visit } => {
debug!(node_id, visit, "Steering hub attached to session");
}
Self::AgentSteeringDetached { node_id, visit } => {
debug!(node_id, visit, "Steering hub detached from session");
}
Self::AgentSteerBuffered { kind, .. } => {
debug!(kind = kind.as_str(), "Steer buffered (no active session)");
}
Self::AgentSteerDropped { reason, count, .. } => {
warn!(?reason, count, "Steer dropped");
}
Self::PullRequestCreated {
pr_url,
pr_number,

View file

@ -116,6 +116,10 @@ pub fn event_name(event: &Event) -> &'static str {
Event::CommandCompleted { .. } => "command.completed",
Event::AgentCliStarted { .. } => "agent.cli.started",
Event::AgentCliCompleted { .. } => "agent.cli.completed",
Event::AgentSteeringAttached { .. } => "agent.steering.attached",
Event::AgentSteeringDetached { .. } => "agent.steering.detached",
Event::AgentSteerBuffered { .. } => "agent.steer.buffered",
Event::AgentSteerDropped { .. } => "agent.steer.dropped",
Event::PullRequestCreated { .. } => "pull_request.created",
Event::PullRequestFailed { .. } => "pull_request.failed",
Event::DevcontainerResolved { .. } => "devcontainer.resolved",

View file

@ -65,7 +65,8 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields {
| Event::RunUnpauseRequested { actor }
| Event::RunArchived { actor }
| Event::RunUnarchived { actor, .. }
| Event::InterviewCompleted { actor, .. } => StoredEventFields {
| Event::InterviewCompleted { actor, .. }
| Event::AgentSteerBuffered { actor, .. } => StoredEventFields {
actor: actor.clone(),
..StoredEventFields::default()
},
@ -117,6 +118,39 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields {
| Event::CommandCompleted { node_id, .. }
| Event::AgentCliStarted { node_id, .. }
| Event::AgentCliCompleted { node_id, .. } => node_stored_fields(Some(node_id.clone())),
Event::AgentSteeringAttached { node_id, visit }
| Event::AgentSteeringDetached { node_id, visit } => {
let node_id_str = node_id.clone();
let node_label = default_node_label(Some(&node_id_str), None);
StoredEventFields {
node_id: Some(node_id_str.clone()),
node_label,
stage_id: Some(StageId::new(node_id_str, *visit)),
..StoredEventFields::default()
}
}
Event::AgentSteerDropped {
actor,
node_id,
visit,
..
} => {
let node_id_str = node_id.clone();
let node_label = node_id_str
.as_ref()
.and_then(|n| default_node_label(Some(n), None));
let stage_id = match (node_id.clone(), visit) {
(Some(n), Some(v)) => Some(StageId::new(n, *v)),
_ => None,
};
StoredEventFields {
node_id: node_id_str,
node_label,
stage_id,
actor: actor.clone(),
..StoredEventFields::default()
}
}
Event::Agent {
stage,
visit,
@ -213,6 +247,7 @@ fn agent_actor_for_event(
parent_session_id: parent_session_id.map(str::to_string),
model: None,
}),
AgentEvent::SteeringInjected { actor, .. } => actor.clone(),
_ => None,
}
}

View file

@ -4,8 +4,8 @@ use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use fabro_agent::subagent::{SessionFactory, SubAgentManager};
use fabro_agent::{
AgentEvent, AgentProfile, AnthropicProfile, GeminiProfile, OpenAiProfile, Sandbox, Session,
SessionOptions, Turn,
AgentEvent, AgentProfile, AnthropicProfile, CompletionCoordinator, GeminiProfile,
OpenAiProfile, Sandbox, Session, SessionControlHandle, SessionOptions, Turn,
};
use fabro_auth::{CredentialSource, EnvCredentialSource};
use fabro_graphviz::graph::Node;
@ -13,6 +13,7 @@ use fabro_llm::client::Client;
use fabro_llm::types::{Message, Request, TokenCounts};
use fabro_mcp::config::McpServerSettings;
use fabro_model::{FallbackTarget, Provider};
use fabro_types::StageId;
use tokio::sync::Mutex as TokioMutex;
use super::super::agent::{CodergenBackend, CodergenResult};
@ -21,6 +22,7 @@ use crate::context::{Context, WorkflowContext};
use crate::error::Error;
use crate::event::{Emitter, Event, StageScope};
use crate::outcome::billed_model_usage_from_llm;
use crate::steering_hub::SteeringHub;
fn build_profile(model: &str, provider: Provider) -> Box<dyn AgentProfile> {
match provider {
@ -122,6 +124,7 @@ pub struct AgentApiBackend {
env: HashMap<String, String>,
mcp_servers: Vec<McpServerSettings>,
source: Arc<dyn CredentialSource>,
steering_hub: Option<Arc<SteeringHub>>,
}
impl AgentApiBackend {
@ -140,9 +143,16 @@ impl AgentApiBackend {
env: HashMap::new(),
mcp_servers: Vec::new(),
source,
steering_hub: None,
}
}
#[must_use]
pub fn with_steering_hub(mut self, steering_hub: Arc<SteeringHub>) -> Self {
self.steering_hub = Some(steering_hub);
self
}
#[must_use]
pub fn new_from_env(
model: String,
@ -491,6 +501,29 @@ impl CodergenBackend for AgentApiBackend {
session.initialize().await;
}
// Register with the steering hub so HTTP `POST /runs/{id}/steer`
// calls reach this session. The RAII guard below unregisters on
// every exit path (success, error, failover replace).
let stage_id = stage_scope.stage_id();
let _hub_guard = if let Some(ref hub) = self.steering_hub {
hub.register(&stage_id, &session.control_handle());
// Wire the completion coordinator so a steer that arrives
// during a final no-tool response triggers an extra round
// rather than being silently dropped.
let coordinator = Arc::new(SteeringCompletionCoordinator {
hub: Arc::clone(hub),
stage_id: stage_id.clone(),
handle: session.control_handle(),
});
session.set_completion_coordinator(coordinator);
Some(SteeringHubGuard {
hub: Arc::clone(hub),
stage_id: stage_id.clone(),
})
} else {
None
};
let result = session.process_input(prompt).await;
// On failover-eligible error, try fallback providers.
@ -553,6 +586,18 @@ impl CodergenBackend for AgentApiBackend {
Arc::clone(&file_tracking),
);
// Re-register the new session's handle under the same
// stage_id (replace, no re-drain, no attached event).
if let Some(ref hub) = self.steering_hub {
hub.register(&stage_id, &session.control_handle());
let coordinator = Arc::new(SteeringCompletionCoordinator {
hub: Arc::clone(hub),
stage_id: stage_id.clone(),
handle: session.control_handle(),
});
session.set_completion_coordinator(coordinator);
}
session.initialize().await;
match session.process_input(prompt).await {
Ok(()) => {
@ -636,6 +681,48 @@ impl CodergenBackend for AgentApiBackend {
}
}
/// RAII guard that unregisters a session from the steering hub when the
/// stage completes (success, error, or panic).
struct SteeringHubGuard {
hub: Arc<SteeringHub>,
stage_id: StageId,
}
impl Drop for SteeringHubGuard {
fn drop(&mut self) {
self.hub.unregister(&self.stage_id);
}
}
/// Coordinator that lets the agent loop ask the workflow layer whether to
/// keep iterating after a no-tool natural completion. Implements the
/// "close-the-door" pattern: unregister, check the queue, then either
/// break or re-register and report `true` so the loop drains.
struct SteeringCompletionCoordinator {
hub: Arc<SteeringHub>,
stage_id: StageId,
handle: SessionControlHandle,
}
impl CompletionCoordinator for SteeringCompletionCoordinator {
fn on_natural_completion(&self) -> bool {
// Close the door: take this session out of the active set so no
// further deliver() races with the queue check. (Idempotent if
// already unregistered.)
self.hub.unregister(&self.stage_id);
if self.handle.queue_is_empty() {
return false;
}
// A steer landed in the queue (or arrived between unregister and
// this check — impossible under the lock discipline since
// deliver/unregister serialize on `active`'s RwLock). Re-register
// so future deliveries continue to land here, and tell the loop
// to drain.
self.hub.register(&self.stage_id, &self.handle);
true
}
}
#[cfg(test)]
mod tests {
use fabro_agent::subagent::SessionFactory;

View file

@ -144,6 +144,7 @@ pub mod run_lookup;
pub use error::{Error, FailureCategory, FailureSignature, FailureSignatureExt, Result};
pub use manifest_path::ManifestPath;
pub use steering_hub::SteeringHub;
pub mod run_materialization;
pub(crate) mod run_metadata;
pub mod run_options;
@ -153,6 +154,7 @@ pub mod sandbox_git;
pub(crate) mod sandbox_git_runtime;
pub mod services;
mod stage_scope;
pub mod steering_hub;
#[doc(hidden)]
pub mod test_support;
#[doc(hidden)]

View file

@ -51,6 +51,7 @@ use crate::run_metadata::metadata_branch_name;
use crate::run_options::{GitCheckpointOptions, LifecycleOptions, RunOptions};
use crate::run_status::{FailureReason, RunStatus};
use crate::runtime_store::RunStoreHandle;
use crate::steering_hub::SteeringHub;
use crate::workflow_bundle::{RunDefinition, WorkflowBundle};
struct RunSession {
@ -59,6 +60,7 @@ struct RunSession {
sandbox: SandboxSpec,
llm: LlmSpec,
interviewer: Arc<dyn Interviewer>,
steering_hub: Arc<SteeringHub>,
on_node: crate::OnNodeCallback,
lifecycle: LifecycleOptions,
hooks: fabro_hooks::HookSettings,
@ -89,6 +91,7 @@ pub struct StartServices {
pub cancel_token: Option<Arc<AtomicBool>>,
pub emitter: Arc<Emitter>,
pub interviewer: Arc<dyn Interviewer>,
pub steering_hub: Arc<SteeringHub>,
pub run_store: RunStoreHandle,
pub event_sink: RunEventSink,
pub artifact_sink: Option<ArtifactSink>,
@ -426,6 +429,7 @@ impl RunSession {
dry_run: resolved.execution.mode == RunMode::DryRun,
},
interviewer,
steering_hub: services.steering_hub,
on_node: services.on_node,
lifecycle: LifecycleOptions {
setup_commands: resolved.prepare.commands.clone(),
@ -744,6 +748,7 @@ impl RunSession {
sandbox: self.sandbox,
llm: self.llm,
interviewer: self.interviewer,
steering_hub: Arc::clone(&self.steering_hub),
lifecycle: self.lifecycle,
run_options,
workflow_path: self.workflow_path,
@ -815,6 +820,9 @@ impl RunSession {
let retro = retroed.retro.clone();
let concluded = Box::pin(pipeline::finalize(retroed, &finalize_opts)).await?;
let finalized = Box::pin(pipeline::pull_request(concluded, &pr_opts)).await;
// Emit `agent.steer.dropped { reason: run_ended }` for any
// unconsumed pending steers, then flush so it lands in the store.
self.steering_hub.drain_pending_at_run_end();
store_progress_logger.flush().await;
scopeguard::ScopeGuard::into_inner(cleanup_guard);
@ -1106,11 +1114,13 @@ mod tests {
emitter: Arc<Emitter>,
registry: Arc<HandlerRegistry>,
) -> StartServices {
let steering_hub = Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone()));
StartServices {
run_id: fixtures::RUN_1,
cancel_token: None,
emitter,
interviewer: Arc::new(fabro_interview::AutoApproveInterviewer::engine()),
steering_hub,
run_store: store.open_run(&fixtures::RUN_1).await.unwrap().into(),
event_sink: RunEventSink::store(store.open_run(&fixtures::RUN_1).await.unwrap()),
artifact_sink: None,

View file

@ -201,7 +201,7 @@ async fn execute_test_run_with_options(
run_id: run_id_value,
run_store: run_store.into(),
dry_run: false,
emitter,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
@ -213,6 +213,7 @@ async fn execute_test_run_with_options(
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
@ -271,6 +272,9 @@ async fn execute_runs_start_to_exit_and_returns_final_context() {
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(test_emitter_arc(
"run-test",
))),
lifecycle: LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
@ -331,7 +335,7 @@ async fn run_with_lifecycle(
run_id,
run_store: test_run_store(&run_id).await.into(),
dry_run: false,
emitter,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: PathBuf::from(sandbox.working_directory()),
},
@ -343,6 +347,7 @@ async fn run_with_lifecycle(
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle,
run_options,
workflow_path: None,

View file

@ -38,6 +38,7 @@ use crate::run_options::{GitCheckpointOptions, RunOptions};
use crate::sandbox_git::GIT_REMOTE;
use crate::sandbox_git_runtime::SandboxGitRuntime;
use crate::services::{EngineServices, RunServices};
use crate::steering_hub::SteeringHub;
struct WorktreePlan {
branch_name: String,
@ -255,6 +256,7 @@ async fn build_sandbox_env(
async fn build_registry(
spec: &LlmSpec,
interviewer: Arc<dyn fabro_interview::Interviewer>,
steering_hub: Arc<SteeringHub>,
sandbox_env: &HashMap<String, String>,
graph: &graph::Graph,
llm_source: Arc<dyn CredentialSource>,
@ -299,6 +301,7 @@ async fn build_registry(
let fallback_chain = spec.fallback_chain.clone();
let mcp_servers = spec.mcp_servers.clone();
let llm_source_for_api = Arc::clone(&llm_source);
let steering_hub_for_api = Arc::clone(&steering_hub);
let registry = Arc::new(default_registry(interviewer, move || {
let api = AgentApiBackend::new(
model.clone(),
@ -307,7 +310,8 @@ async fn build_registry(
Arc::clone(&llm_source_for_api),
)
.with_env(env.clone())
.with_mcp_servers(mcp_servers.clone());
.with_mcp_servers(mcp_servers.clone())
.with_steering_hub(Arc::clone(&steering_hub_for_api));
let cli = cli_resolver
.clone()
.map_or_else(
@ -571,6 +575,7 @@ pub async fn initialize(
build_registry(
&options.llm,
Arc::clone(&options.interviewer),
Arc::clone(&options.steering_hub),
&env,
&graph,
Arc::clone(&llm_source),
@ -904,49 +909,50 @@ mod tests {
});
let result = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: {
run_id: test_run_id(),
run_store: {
let store = memory_store();
let inner = store.create_run(&test_run_id()).await.unwrap();
inner.into()
},
dry_run: false,
emitter,
sandbox: SandboxSpec::Local {
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
llm: LlmSpec {
llm: LlmSpec {
model: "test-model".to_string(),
provider: fabro_llm::Provider::Anthropic,
fallback_chain: Vec::new(),
mcp_servers: Vec::new(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
lifecycle: crate::run_options::LifecycleOptions {
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![command.to_string()],
setup_command_timeout_ms: 1_000,
devcontainer_phases: vec![],
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
devcontainer_env: HashMap::new(),
toml_env: HashMap::new(),
github_permissions: None,
origin_url: None,
},
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
registry_override: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
})
.await;
let events = seen.lock().unwrap().clone();
@ -959,6 +965,7 @@ mod tests {
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let store = memory_store();
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let mut options = InitOptions {
run_id: test_run_id(),
run_store: {
@ -966,7 +973,7 @@ mod tests {
inner.into()
},
dry_run: false,
emitter: Arc::new(crate::event::Emitter::new(test_run_id())),
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
@ -978,6 +985,7 @@ mod tests {
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
@ -1021,49 +1029,50 @@ mod tests {
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let initialized = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: {
run_id: test_run_id(),
run_store: {
let store = memory_store();
let inner = store.create_run(&test_run_id()).await.unwrap();
inner.into()
},
dry_run: false,
emitter,
sandbox: SandboxSpec::Local {
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
llm: LlmSpec {
llm: LlmSpec {
model: "test-model".to_string(),
provider: fabro_llm::Provider::Anthropic,
fallback_chain: Vec::new(),
mcp_servers: Vec::new(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
lifecycle: crate::run_options::LifecycleOptions {
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 1_000,
devcontainer_phases: vec![],
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
devcontainer_env: HashMap::new(),
toml_env: HashMap::from([("TEST_KEY".to_string(), "value".to_string())]),
github_permissions: None,
origin_url: None,
},
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
registry_override: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
})
.await
.unwrap();
@ -1115,6 +1124,7 @@ mod tests {
let (graph, _) = llm_graph();
let vault = Arc::new(AsyncRwLock::new(vault));
let test_emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let (_registry, effective_dry_run) = build_registry(
&LlmSpec {
model: "claude-opus-4-6".to_string(),
@ -1124,6 +1134,7 @@ mod tests {
dry_run: false,
},
Arc::new(AutoApproveInterviewer::engine()),
Arc::new(crate::steering_hub::SteeringHub::new(test_emitter)),
&HashMap::new(),
&graph,
Arc::new(VaultCredentialSource::new(Arc::clone(&vault))),
@ -1154,45 +1165,46 @@ mod tests {
store_logger.register(&emitter);
let initialized = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: run_store.into(),
dry_run: false,
emitter,
sandbox: SandboxSpec::Local {
run_id: test_run_id(),
run_store: run_store.into(),
dry_run: false,
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
llm: LlmSpec {
llm: LlmSpec {
model: "test-model".to_string(),
provider: fabro_llm::Provider::Anthropic,
fallback_chain: Vec::new(),
mcp_servers: Vec::new(),
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
lifecycle: crate::run_options::LifecycleOptions {
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec!["true".to_string()],
setup_command_timeout_ms: 1_000,
devcontainer_phases: vec![],
},
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
run_options: test_settings(&run_dir),
workflow_path: None,
workflow_bundle: None,
hooks: fabro_hooks::HookSettings { hooks: vec![] },
sandbox_env: SandboxEnvSpec {
devcontainer_env: HashMap::new(),
toml_env: HashMap::new(),
github_permissions: None,
origin_url: None,
},
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
vault: None,
devcontainer: None,
git: None,
worktree_mode: None,
run_control: None,
registry_override: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
artifact_sink: None,
checkpoint: None,
seed_context: None,
})
.await
.unwrap();
@ -1261,6 +1273,7 @@ mod tests {
let mut run_options = test_settings(&run_dir);
run_options.cancel_token = Some(cancel_token);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let result = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: {
@ -1269,7 +1282,7 @@ mod tests {
inner.into()
},
dry_run: false,
emitter: Arc::new(crate::event::Emitter::new(test_run_id())),
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
@ -1281,6 +1294,7 @@ mod tests {
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec!["sleep 5".to_string()],
setup_command_timeout_ms: 5_000,
@ -1322,6 +1336,7 @@ mod tests {
let mut run_options = test_settings(&run_dir);
run_options.cancel_token = Some(cancel_token);
let emitter = Arc::new(crate::event::Emitter::new(test_run_id()));
let result = initialize(persisted, InitOptions {
run_id: test_run_id(),
run_store: {
@ -1330,7 +1345,7 @@ mod tests {
inner.into()
},
dry_run: false,
emitter: Arc::new(crate::event::Emitter::new(test_run_id())),
emitter: emitter.clone(),
sandbox: SandboxSpec::Local {
working_directory: std::env::current_dir().unwrap(),
},
@ -1342,6 +1357,7 @@ mod tests {
dry_run: true,
},
interviewer: Arc::new(AutoApproveInterviewer::engine()),
steering_hub: Arc::new(crate::steering_hub::SteeringHub::new(emitter.clone())),
lifecycle: crate::run_options::LifecycleOptions {
setup_commands: vec![],
setup_command_timeout_ms: 5_000,

View file

@ -29,6 +29,7 @@ use crate::run_control::RunControlState;
use crate::run_options::{GitCheckpointOptions, LifecycleOptions, RunOptions};
use crate::runtime_store::RunStoreHandle;
use crate::services::{EngineServices, RunServices};
use crate::steering_hub::SteeringHub;
use crate::transforms::Transform;
use crate::workflow_bundle::WorkflowBundle;
@ -240,6 +241,7 @@ pub struct InitOptions {
pub sandbox: SandboxSpec,
pub llm: LlmSpec,
pub interviewer: Arc<dyn Interviewer>,
pub steering_hub: Arc<SteeringHub>,
pub lifecycle: LifecycleOptions,
pub run_options: RunOptions,
pub workflow_path: Option<ManifestPath>,

View file

@ -0,0 +1,337 @@
//! Bridge between the worker's HTTP control plane and live agent
//! `Session`s. The hub owns:
//!
//! - A map of currently steerable API-mode sessions, keyed by `StageId` →
//! `SessionControlHandle`.
//! - A bounded run-wide pending buffer for steers that arrive when no session
//! is registered (between stages, before the first agent stage, or after a
//! session ends but before the next registers).
//!
//! Lock discipline (race safety):
//! - `active` is `std::sync::RwLock`; deliver takes the read lock for the
//! entire decide-and-push step.
//! - `pending` is `std::sync::Mutex` taken under the active read lock.
//! - All methods are sync — no `.await` while holding any lock — so the
//! `CompletionCoordinator::on_natural_completion` close-the-door dance can
//! call `unregister(...)` synchronously from the agent loop.
use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex, RwLock};
use fabro_agent::SessionControlHandle;
use fabro_types::run_event::AgentSteerDroppedReason;
use fabro_types::{Principal, StageId, SteerKind};
use crate::event::{Emitter, Event};
/// Cap on the steering queue length kept per active session. Overflow
/// evicts the oldest entry (FIFO) and emits `agent.steer.dropped`.
pub const PER_SESSION_QUEUE_CAP: usize = 32;
/// Cap on the run-wide pending buffer used when no session is registered.
/// Overflow evicts the oldest entry (FIFO) and emits `agent.steer.dropped`.
pub const PER_RUN_PENDING_CAP: usize = 32;
#[derive(Debug, Clone)]
struct PendingSteer {
text: String,
/// Original kind. Buffered steers always flush as `Append` when a
/// session registers (see `register`), but we keep the original so a
/// future per-stage targeting feature can preserve it.
#[allow(
dead_code,
reason = "captured for future per-stage targeting; see TODO"
)]
kind: SteerKind,
actor: Option<Principal>,
}
#[allow(
clippy::module_name_repetitions,
reason = "external callers refer to it as SteeringHub"
)]
pub struct SteeringHub {
active: RwLock<HashMap<StageId, SessionControlHandle>>,
pending: Mutex<VecDeque<PendingSteer>>,
emitter: Arc<Emitter>,
}
impl SteeringHub {
#[must_use]
pub fn new(emitter: Arc<Emitter>) -> Self {
Self {
active: RwLock::new(HashMap::new()),
pending: Mutex::new(VecDeque::new()),
emitter,
}
}
/// Test-only constructor with an isolated emitter.
#[cfg(test)]
#[must_use]
pub fn for_tests() -> Arc<Self> {
use fabro_types::RunId;
Arc::new(Self::new(Arc::new(Emitter::new(RunId::new()))))
}
/// Test-only: snapshot of pending buffer length.
#[cfg(test)]
#[must_use]
pub fn pending_len(&self) -> usize {
self.pending.lock().expect("pending lock poisoned").len()
}
/// Test-only: snapshot of registered stage count.
#[cfg(test)]
#[must_use]
pub fn active_count(&self) -> usize {
self.active.read().expect("active lock poisoned").len()
}
/// Register an API-mode session as steerable for this stage. If no
/// entry existed for `stage_id`, drains pending into the new handle as
/// `Append`-kind messages and emits `agent.steering.attached`. If an
/// entry already existed (e.g. failover replaced the underlying
/// session), the handle is overwritten silently — no drain, no event.
pub fn register(&self, stage_id: &StageId, handle: &SessionControlHandle) {
let was_new = {
let mut active = self.active.write().expect("active lock poisoned");
let was_new = !active.contains_key(stage_id);
active.insert(stage_id.clone(), handle.clone());
was_new
};
if was_new {
let pending: Vec<PendingSteer> = {
let mut pending = self.pending.lock().expect("pending lock poisoned");
pending.drain(..).collect()
};
for item in pending {
// Buffered steers always flush as Append — the original
// Interrupt semantics no longer make sense once the round
// has rolled over.
Self::enqueue_into_session_queue(
handle,
(item.text, SteerKind::Append, item.actor),
&self.emitter,
Some(stage_id),
);
}
self.emitter.emit(&Event::AgentSteeringAttached {
node_id: stage_id.node_id().to_string(),
visit: stage_id.visit(),
});
}
}
/// Unregister the session previously registered for this stage. Emits
/// `agent.steering.detached` only when an entry was actually removed
/// (idempotent — safe to call multiple times from RAII guards).
pub fn unregister(&self, stage_id: &StageId) {
let removed = {
let mut active = self.active.write().expect("active lock poisoned");
active.remove(stage_id).is_some()
};
if removed {
self.emitter.emit(&Event::AgentSteeringDetached {
node_id: stage_id.node_id().to_string(),
visit: stage_id.visit(),
});
}
}
/// Deliver a steer from the HTTP control plane. Broadcasts to every
/// active session if any are registered, otherwise parks the message
/// in the run-wide pending buffer.
pub fn deliver(&self, text: String, kind: SteerKind, actor: Option<Principal>) {
// Hold the active read lock for the entire decide-and-dispatch
// step so register/unregister cannot race with this push.
let active = self.active.read().expect("active lock poisoned");
if active.is_empty() {
drop(active);
let mut pending = self.pending.lock().expect("pending lock poisoned");
if pending.len() >= PER_RUN_PENDING_CAP {
let dropped = pending.pop_front();
let dropped_actor = dropped.and_then(|d| d.actor);
self.emitter.emit(&Event::AgentSteerDropped {
reason: AgentSteerDroppedReason::QueueFull,
count: 1,
actor: dropped_actor,
node_id: None,
visit: None,
});
}
pending.push_back(PendingSteer {
text,
kind,
actor: actor.clone(),
});
self.emitter
.emit(&Event::AgentSteerBuffered { kind, actor });
return;
}
// Broadcast to every active session.
for (stage_id, handle) in active.iter() {
Self::enqueue_into_session_queue(
handle,
(text.clone(), kind, actor.clone()),
&self.emitter,
Some(stage_id),
);
}
}
/// Drain any unconsumed pending steers and emit a single
/// `agent.steer.dropped` event with `reason: run_ended`. Called from
/// `operations::start` after the pipeline finishes (success or
/// failure) but before the emitter is flushed.
pub fn drain_pending_at_run_end(&self) {
let count: u32 = {
let mut pending = self.pending.lock().expect("pending lock poisoned");
let n = u32::try_from(pending.len()).unwrap_or(u32::MAX);
pending.clear();
n
};
if count > 0 {
self.emitter.emit(&Event::AgentSteerDropped {
reason: AgentSteerDroppedReason::RunEnded,
count,
actor: None,
node_id: None,
visit: None,
});
}
}
/// Push an item into a session's queue, evicting the oldest entry and
/// emitting `agent.steer.dropped { queue_full }` if the cap is hit.
fn enqueue_into_session_queue(
handle: &SessionControlHandle,
item: (String, SteerKind, Option<Principal>),
emitter: &Emitter,
stage_id: Option<&StageId>,
) {
if handle.queue_len() >= PER_SESSION_QUEUE_CAP {
let evicted = handle.pop_oldest();
let evicted_actor = evicted.and_then(|(.., a)| a);
emitter.emit(&Event::AgentSteerDropped {
reason: AgentSteerDroppedReason::QueueFull,
count: 1,
actor: evicted_actor,
node_id: stage_id.map(|s| s.node_id().to_string()),
visit: stage_id.map(StageId::visit),
});
}
handle.enqueue(item);
}
}
#[cfg(test)]
mod tests {
use fabro_agent::SessionControlHandle;
use fabro_types::{Principal, StageId, SteerKind, SystemActorKind};
use super::SteeringHub;
#[test]
fn deliver_with_no_active_buffers_message() {
let hub = SteeringHub::for_tests();
hub.deliver(
"hi".into(),
SteerKind::Append,
Some(Principal::System {
system_kind: SystemActorKind::Engine,
}),
);
assert_eq!(hub.pending_len(), 1);
}
#[test]
fn drain_pending_at_run_end_clears_buffer() {
let hub = SteeringHub::for_tests();
hub.deliver("a".into(), SteerKind::Append, None);
hub.deliver("b".into(), SteerKind::Append, None);
assert_eq!(hub.pending_len(), 2);
hub.drain_pending_at_run_end();
assert_eq!(hub.pending_len(), 0);
}
#[test]
fn pending_buffer_evicts_oldest_at_cap() {
let hub = SteeringHub::for_tests();
for i in 0..(super::PER_RUN_PENDING_CAP + 5) {
hub.deliver(format!("msg{i}"), SteerKind::Append, None);
}
assert_eq!(hub.pending_len(), super::PER_RUN_PENDING_CAP);
}
#[test]
fn unregister_is_idempotent() {
let hub = SteeringHub::for_tests();
let stage = StageId::new("agent-node", 1);
hub.unregister(&stage);
hub.unregister(&stage);
}
#[test]
fn register_drains_pending_into_first_session() {
let hub = SteeringHub::for_tests();
hub.deliver("queued1".into(), SteerKind::Append, None);
hub.deliver("queued2".into(), SteerKind::Interrupt, None);
assert_eq!(hub.pending_len(), 2);
let stage = StageId::new("agent-node", 1);
let handle = SessionControlHandle::new();
hub.register(&stage, &handle);
assert_eq!(handle.queue_len(), 2);
assert_eq!(hub.pending_len(), 0);
assert_eq!(hub.active_count(), 1);
}
#[test]
fn deliver_broadcasts_to_active_sessions() {
let hub = SteeringHub::for_tests();
let stage_a = StageId::new("a", 1);
let stage_b = StageId::new("b", 1);
let handle_a = SessionControlHandle::new();
let handle_b = SessionControlHandle::new();
hub.register(&stage_a, &handle_a);
hub.register(&stage_b, &handle_b);
hub.deliver("hello".into(), SteerKind::Append, None);
assert_eq!(handle_a.queue_len(), 1);
assert_eq!(handle_b.queue_len(), 1);
assert_eq!(hub.pending_len(), 0);
}
#[test]
fn re_register_same_stage_does_not_redrain() {
let hub = SteeringHub::for_tests();
let stage = StageId::new("a", 1);
let handle1 = SessionControlHandle::new();
hub.register(&stage.clone(), &handle1);
hub.deliver("x".into(), SteerKind::Append, None);
assert_eq!(handle1.queue_len(), 1);
// Replace handle (failover) — must not redrain pending or emit
// attached again.
let handle2 = SessionControlHandle::new();
hub.register(&stage, &handle2);
assert_eq!(handle2.queue_len(), 0);
}
#[test]
fn per_session_queue_evicts_oldest_at_cap() {
let hub = SteeringHub::for_tests();
let stage = StageId::new("a", 1);
let handle = SessionControlHandle::new();
hub.register(&stage, &handle);
for i in 0..(super::PER_SESSION_QUEUE_CAP + 5) {
hub.deliver(format!("m{i}"), SteerKind::Append, None);
}
assert_eq!(handle.queue_len(), super::PER_SESSION_QUEUE_CAP);
}
}

View file

@ -281,6 +281,7 @@ export * from './stage-projection';
export * from './stage-state';
export * from './stage-turn';
export * from './start-run-request';
export * from './steer-run-request';
export * from './submit-answer-request';
export * from './success-reason';
export * from './system-actor-kind';
@ -302,4 +303,4 @@ export * from './workflow-namespace';
export * from './workflow-reference';
export * from './workflow-settings';
export * from './worktree-mode';
export * from './write-blob-response';
export * from './write-blob-response';

View file

@ -0,0 +1,29 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.1.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
/**
* Request body for steering a running run mid-execution.
*/
export interface SteerRunRequest {
/**
* The steering message text to deliver as a user turn.
*/
'text': string;
/**
* When true, cancel the in-flight LLM stream and tool calls in the current round before delivering. When false (default), append to the steering queue and let the agent pick it up at the next turn boundary.
*/
'interrupt'?: boolean;
}