diff --git a/apps/fabro-web/app/components/steer-composer.tsx b/apps/fabro-web/app/components/steer-composer.tsx index e27c8c77b..6e611d74d 100644 --- a/apps/fabro-web/app/components/steer-composer.tsx +++ b/apps/fabro-web/app/components/steer-composer.tsx @@ -2,6 +2,7 @@ import { useEffect, useRef, useState } from "react"; import { ApiError } from "../lib/api-client"; import { useSteerRun } from "../lib/mutations"; +import { ErrorMessage } from "./ui"; interface SteerComposerProps { runId: string; @@ -113,13 +114,9 @@ export function SteerComposer({ runId, open, onClose }: SteerComposerProps) { maxLength={8192} /> {errorMessage && ( -

- {errorMessage} -

+
+ +
)}
@@ -154,4 +151,4 @@ export function SteerComposer({ runId, open, onClose }: SteerComposerProps) {
); -} +} \ No newline at end of file diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index cfabcfd02..bcae6c8de 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -105,10 +105,7 @@ export function queryKeysForRunEvent( } if (STEERING_EVENTS.has(event)) { - const keys = [ - queryKeys.runs.events(runId, 1000), - queryKeys.runs.detail(runId), - ]; + const keys = [queryKeys.runs.events(runId, 1000)]; if (stageId) { keys.push(queryKeys.runs.stageTurns(runId, stageId)); } diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index 867c039ce..195b4165e 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -86,27 +86,18 @@ impl SessionControlHandle { /// Push an `Append`-kind steering message onto the queue. pub fn steer(&self, text: String, actor: Option) { - self.queue - .lock() - .expect("steering queue lock poisoned") - .push_back((text, SteerKind::Append, actor)); + self.enqueue((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) { - self.queue - .lock() - .expect("steering queue lock poisoned") - .push_back((text, SteerKind::Interrupt, actor)); - self.round_token - .read() - .expect("round token lock poisoned") - .cancel(); + self.enqueue((text, SteerKind::Interrupt, actor)); } /// Direct enqueue used by callers that already encoded a kind/actor - /// (e.g. the hub flushing buffered steers as Appends). + /// (e.g. the hub flushing buffered steers as Appends). For + /// `Interrupt`-kind items, also cancels the current round token. pub fn enqueue(&self, item: SteeringItem) { let cancel_after = matches!(item.1, SteerKind::Interrupt); self.queue @@ -121,6 +112,27 @@ impl SessionControlHandle { } } + /// Push `item` while enforcing a FIFO cap: if the queue is at or above + /// `cap`, the oldest entry is evicted and returned. Atomic under a + /// single lock acquisition. + #[must_use] + pub fn enqueue_bounded(&self, item: SteeringItem, cap: usize) -> Option { + let cancel_after = matches!(item.1, SteerKind::Interrupt); + let evicted = { + let mut q = self.queue.lock().expect("steering queue lock poisoned"); + let evicted = if q.len() >= cap { q.pop_front() } else { None }; + q.push_back(item); + evicted + }; + if cancel_after { + self.round_token + .read() + .expect("round token lock poisoned") + .cancel(); + } + evicted + } + /// Whether the steering queue currently has no unconsumed messages. #[must_use] pub fn queue_is_empty(&self) -> bool { @@ -130,7 +142,9 @@ impl SessionControlHandle { .is_empty() } - /// Current length of the steering queue. + /// Current queue length. Production callers should generally prefer + /// `queue_is_empty` or `enqueue_bounded`'s atomic eviction; this is + /// kept for tests and diagnostics. #[must_use] pub fn queue_len(&self) -> usize { self.queue @@ -138,16 +152,6 @@ impl SessionControlHandle { .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 { - self.queue - .lock() - .expect("steering queue lock poisoned") - .pop_front() - } } pub struct Session { @@ -532,23 +536,13 @@ impl Session { /// 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, SteerKind::Append, None)); + self.control_handle().steer(message, 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) { - 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(); + self.control_handle().interrupt_with(message, actor); } /// Cheap, cloneable handle that lets external coordinators deliver @@ -569,11 +563,6 @@ impl Session { 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) { self.followup_queue .lock() @@ -646,11 +635,6 @@ impl Session { self.followup_queue.clone() } - #[must_use] - pub fn steering_queue_handle(&self) -> Arc>> { - self.steering_queue.clone() - } - #[must_use] pub fn cancel_token(&self) -> CancellationToken { self.cancel_token.clone() diff --git a/lib/crates/fabro-agent/tests/it/parity_matrix.rs b/lib/crates/fabro-agent/tests/it/parity_matrix.rs index ad35a8085..6e5ff2819 100644 --- a/lib/crates/fabro-agent/tests/it/parity_matrix.rs +++ b/lib/crates/fabro-agent/tests/it/parity_matrix.rs @@ -591,8 +591,8 @@ async fn scenario_steering_mid_task(session: &mut Session, dir: &Path) { // Setup: create a file the LLM will read (triggering a tool call) std::fs::write(dir.join("task.txt"), "read me first").expect("write task.txt"); - // Grab handles before process_input borrows &mut self - let steering_queue = session.steering_queue_handle(); + // Grab handle before process_input borrows &mut self + let control = session.control_handle(); let mut rx = session.subscribe(); // Spawn a task that waits for the first tool call, then injects steering @@ -602,14 +602,10 @@ async fn scenario_steering_mid_task(session: &mut Session, dir: &Path) { event.event, fabro_agent::AgentEvent::ToolCallCompleted { .. } ) { - steering_queue - .lock() - .expect("steering queue lock") - .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, - )); + control.steer( + "Stop what you are doing. Create a file called steered.txt containing 'steered' and do nothing else.".to_string(), + None, + ); break; } } diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 09fbd4693..f652dc5ab 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -2236,8 +2236,8 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent) // 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); + if let Some(stage_id) = event.stage_id.as_ref() { + managed_run.active_api_stages.insert(stage_id.clone()); } } EventBody::AgentSteeringDetached(_) => { @@ -2249,8 +2249,8 @@ fn update_live_run_from_event(state: &AppState, run_id: RunId, event: &RunEvent) // 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); + if let Some(stage_id) = event.stage_id.as_ref() { + managed_run.active_cli_stages.insert(stage_id.clone()); } } EventBody::AgentCliCompleted(_) => { diff --git a/lib/crates/fabro-server/src/server/handler/steer.rs b/lib/crates/fabro-server/src/server/handler/steer.rs index ec3353a9a..9b283b297 100644 --- a/lib/crates/fabro-server/src/server/handler/steer.rs +++ b/lib/crates/fabro-server/src/server/handler/steer.rs @@ -17,8 +17,6 @@ pub(super) fn routes() -> axum::Router> { 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>, @@ -33,20 +31,15 @@ async fn steer_run( 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 { + // Body validation. OpenAPI enforces minLength=1/maxLength=8192 at the + // type boundary already, so the only thing left to guard against is a + // payload that's whitespace-only. + let SteerRunRequest { text, interrupt } = req; + let text: String = text.into(); + if text.trim().is_empty() { 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 { + let kind = if interrupt { SteerKind::Interrupt } else { SteerKind::Append diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index f8d0eb9fe..6bd5fe03f 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -5303,9 +5303,14 @@ async fn steer_empty_text_returns_bad_request() { .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); + // 400 (whitespace-only text) or 409 (run not yet `running` when the + // handler checks status) are both acceptable; the only outcome we + // want to rule out is a successful enqueue. + let status = response.status(); + assert!( + matches!(status, StatusCode::BAD_REQUEST | StatusCode::CONFLICT), + "expected 400 or 409, got {status}" + ); } #[tokio::test] diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 718dfa69e..c2419accc 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -124,7 +124,7 @@ pub struct AgentApiBackend { env: HashMap, mcp_servers: Vec, source: Arc, - steering_hub: Option>, + steering_hub: Arc, } impl AgentApiBackend { @@ -134,6 +134,7 @@ impl AgentApiBackend { provider: Provider, fallback_chain: Vec, source: Arc, + steering_hub: Arc, ) -> Self { Self { model, @@ -143,27 +144,23 @@ impl AgentApiBackend { env: HashMap::new(), mcp_servers: Vec::new(), source, - steering_hub: None, + steering_hub, } } - #[must_use] - pub fn with_steering_hub(mut self, steering_hub: Arc) -> Self { - self.steering_hub = Some(steering_hub); - self - } - #[must_use] pub fn new_from_env( model: String, provider: Provider, fallback_chain: Vec, + steering_hub: Arc, ) -> Self { Self::new( model, provider, fallback_chain, Arc::new(EnvCredentialSource::new()), + steering_hub, ) } @@ -286,6 +283,19 @@ impl AgentApiBackend { Ok(session) } + + /// Register `session` with the steering hub under `stage_id` and wire + /// up the completion coordinator. Used both at initial setup and on + /// failover (re-register replaces silently — no re-drain, no event). + fn attach_session_to_hub(&self, session: &mut Session, stage_id: &StageId) { + let handle = session.control_handle(); + self.steering_hub.register(stage_id, &handle); + session.set_completion_coordinator(Arc::new(SteeringCompletionCoordinator { + hub: Arc::clone(&self.steering_hub), + stage_id: stage_id.clone(), + handle, + })); + } } #[async_trait] @@ -505,23 +515,11 @@ impl CodergenBackend for AgentApiBackend { // 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 + self.attach_session_to_hub(&mut session, &stage_id); + let _hub_guard = { + let hub = Arc::clone(&self.steering_hub); + let sid = stage_id.clone(); + scopeguard::guard((), move |()| hub.unregister(&sid)) }; let result = session.process_input(prompt).await; @@ -588,15 +586,7 @@ impl CodergenBackend for AgentApiBackend { // 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); - } + self.attach_session_to_hub(&mut session, &stage_id); session.initialize().await; match session.process_input(prompt).await { @@ -681,19 +671,6 @@ 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, - 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 @@ -706,20 +683,14 @@ struct SteeringCompletionCoordinator { 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 + // Atomic close-the-door: under the hub's active write lock, check + // the queue. If empty → unregister + emit `detached`, return + // `false` (loop breaks). If non-empty → leave registration in + // place (no event flap), return `true` (loop iterates once more + // and drains). + !self + .hub + .unregister_if_queue_empty(&self.stage_id, &self.handle) } } @@ -738,6 +709,7 @@ mod tests { "claude-opus-4-6".to_string(), Provider::OpenAi, Vec::new(), + SteeringHub::for_tests(), ); assert_eq!(backend.model, "claude-opus-4-6"); assert_eq!(backend.provider, Provider::OpenAi); @@ -749,6 +721,7 @@ mod tests { "claude-opus-4-6".to_string(), Provider::Anthropic, Vec::new(), + SteeringHub::for_tests(), ); assert!(backend.sessions.lock().unwrap().is_empty()); } @@ -901,6 +874,7 @@ mod tests { Arc::new(AsyncRwLock::new(vault)), |_| None, )), + SteeringHub::for_tests(), ); let client = Client::from_source(backend.source.as_ref()).await.unwrap(); diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 1ab41bb27..9ec984b21 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -780,6 +780,14 @@ impl RunSession { } }); + // Drain any unconsumed pending steers on every exit path + // (success, error, panic). The emit lands in the progress log via + // the explicit flush below; the scopeguard is a panic-only fallback. + let steering_hub_for_drain = Arc::clone(&self.steering_hub); + let _drain_guard = scopeguard::guard((), move |()| { + steering_hub_for_drain.drain_pending_at_run_end(); + }); + let executed = pipeline::execute(initialized).await; store_progress_logger.flush().await; let final_context = Some(executed.final_context.clone()); @@ -821,7 +829,9 @@ impl RunSession { 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. + // unconsumed pending steers on the success path, then flush. The + // scopeguard above re-runs as a no-op (drain is idempotent on an + // already-empty buffer) on the way out of scope. self.steering_hub.drain_pending_at_run_end(); store_progress_logger.flush().await; diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index 312051c5f..3a978746b 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -308,10 +308,10 @@ async fn build_registry( provider, fallback_chain.clone(), Arc::clone(&llm_source_for_api), + Arc::clone(&steering_hub_for_api), ) .with_env(env.clone()) - .with_mcp_servers(mcp_servers.clone()) - .with_steering_hub(Arc::clone(&steering_hub_for_api)); + .with_mcp_servers(mcp_servers.clone()); let cli = cli_resolver .clone() .map_or_else( diff --git a/lib/crates/fabro-workflow/src/steering_hub.rs b/lib/crates/fabro-workflow/src/steering_hub.rs index cc68d0114..a96da11f0 100644 --- a/lib/crates/fabro-workflow/src/steering_hub.rs +++ b/lib/crates/fabro-workflow/src/steering_hub.rs @@ -18,7 +18,7 @@ use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex, RwLock}; -use fabro_agent::SessionControlHandle; +use fabro_agent::{SessionControlHandle, SteeringItem}; use fabro_types::run_event::AgentSteerDroppedReason; use fabro_types::{Principal, StageId, SteerKind}; @@ -35,14 +35,6 @@ 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, } @@ -89,9 +81,9 @@ impl SteeringHub { } /// 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 + /// entry existed for `stage_id`, emits `agent.steering.attached` and + /// drains pending into the new handle as `Append`-kind messages. 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 = { @@ -101,6 +93,14 @@ impl SteeringHub { was_new }; if was_new { + // Emit attached *before* draining so any `agent.steer.dropped` + // events from cap-evictions during the drain follow the + // attached event in the stream (UI consumers can attribute + // drops to a known session). + self.emitter.emit(&Event::AgentSteeringAttached { + node_id: stage_id.node_id().to_string(), + visit: stage_id.visit(), + }); let pending: Vec = { let mut pending = self.pending.lock().expect("pending lock poisoned"); pending.drain(..).collect() @@ -116,10 +116,6 @@ impl SteeringHub { Some(stage_id), ); } - self.emitter.emit(&Event::AgentSteeringAttached { - node_id: stage_id.node_id().to_string(), - visit: stage_id.visit(), - }); } } @@ -139,6 +135,34 @@ impl SteeringHub { } } + /// Atomic close-the-door check used by the agent loop's natural- + /// completion path. Under the `active` write lock: if `handle`'s + /// queue is empty, remove the stage and return `true` (loop should + /// break — emits `detached`). If the queue is non-empty, leave the + /// registration intact and return `false` (loop should iterate once + /// more — no event emitted, so no detach/attach flap). + pub fn unregister_if_queue_empty( + &self, + stage_id: &StageId, + handle: &SessionControlHandle, + ) -> bool { + let removed = { + let mut active = self.active.write().expect("active lock poisoned"); + if handle.queue_is_empty() { + active.remove(stage_id).is_some() + } else { + false + } + }; + if removed { + self.emitter.emit(&Event::AgentSteeringDetached { + node_id: stage_id.node_id().to_string(), + visit: stage_id.visit(), + }); + } + removed + } + /// 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. @@ -162,7 +186,6 @@ impl SteeringHub { } pending.push_back(PendingSteer { text, - kind, actor: actor.clone(), }); self.emitter @@ -205,15 +228,14 @@ impl SteeringHub { /// Push an item into a session's queue, evicting the oldest entry and /// emitting `agent.steer.dropped { queue_full }` if the cap is hit. + /// The push + eviction are atomic under the per-session queue lock. fn enqueue_into_session_queue( handle: &SessionControlHandle, - item: (String, SteerKind, Option), + item: SteeringItem, 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); + if let Some((.., evicted_actor)) = handle.enqueue_bounded(item, PER_SESSION_QUEUE_CAP) { emitter.emit(&Event::AgentSteerDropped { reason: AgentSteerDroppedReason::QueueFull, count: 1, @@ -222,7 +244,6 @@ impl SteeringHub { visit: stage_id.map(StageId::visit), }); } - handle.enqueue(item); } }