mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
fabro(01KQT1TWWJYWZGDT8F05E29H9D): simplify_opus (succeeded)
Fabro-Run: 01KQT1TWWJYWZGDT8F05E29H9D
Fabro-Completed: 6
Fabro-Checkpoint: 4c456fe262
⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
parent
11cd91eef1
commit
dd94986c9b
11 changed files with 152 additions and 175 deletions
|
|
@ -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 && (
|
||||
<p
|
||||
role="alert"
|
||||
className="mt-2 text-xs text-amber"
|
||||
data-testid="steer-error"
|
||||
>
|
||||
{errorMessage}
|
||||
</p>
|
||||
<div className="mt-2">
|
||||
<ErrorMessage message={errorMessage} />
|
||||
</div>
|
||||
)}
|
||||
<div className="mt-3 flex items-center justify-between gap-2">
|
||||
<span className="text-[11px] text-fg-muted">
|
||||
|
|
@ -154,4 +151,4 @@ export function SteerComposer({ runId, open, onClose }: SteerComposerProps) {
|
|||
</div>
|
||||
</div>
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -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));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -86,27 +86,18 @@ impl SessionControlHandle {
|
|||
|
||||
/// 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));
|
||||
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<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();
|
||||
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<SteeringItem> {
|
||||
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<SteeringItem> {
|
||||
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<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();
|
||||
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<Mutex<VecDeque<SteeringItem>>> {
|
||||
self.steering_queue.clone()
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn cancel_token(&self) -> CancellationToken {
|
||||
self.cancel_token.clone()
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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(_) => {
|
||||
|
|
|
|||
|
|
@ -17,8 +17,6 @@ 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>>,
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -124,7 +124,7 @@ pub struct AgentApiBackend {
|
|||
env: HashMap<String, String>,
|
||||
mcp_servers: Vec<McpServerSettings>,
|
||||
source: Arc<dyn CredentialSource>,
|
||||
steering_hub: Option<Arc<SteeringHub>>,
|
||||
steering_hub: Arc<SteeringHub>,
|
||||
}
|
||||
|
||||
impl AgentApiBackend {
|
||||
|
|
@ -134,6 +134,7 @@ impl AgentApiBackend {
|
|||
provider: Provider,
|
||||
fallback_chain: Vec<FallbackTarget>,
|
||||
source: Arc<dyn CredentialSource>,
|
||||
steering_hub: Arc<SteeringHub>,
|
||||
) -> 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<SteeringHub>) -> Self {
|
||||
self.steering_hub = Some(steering_hub);
|
||||
self
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn new_from_env(
|
||||
model: String,
|
||||
provider: Provider,
|
||||
fallback_chain: Vec<FallbackTarget>,
|
||||
steering_hub: Arc<SteeringHub>,
|
||||
) -> 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<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
|
||||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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<Principal>,
|
||||
}
|
||||
|
||||
|
|
@ -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<PendingSteer> = {
|
||||
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<Principal>),
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue