From a70750f3f6b2ef04ac7cce5e09ee405d30899d55 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 19:59:23 -0400 Subject: [PATCH] Steer a Petri run by stage `SteerRunRequest` takes an optional `stage`: the label the projection shows (`node@visit`, or `node/e@visit` when two executions share one) or the node's name. The server passes it on the worker control message; the worker's `RunControls` resolves a label to the live agent firing and steers that firing, and a node name through Petri's own live-stage index. Unnamed, the one-live-agent rule stays, and the refusal now names the live stages by their labels. `fabro steer --stage` sets it. A controls scenario runs two agent stages side by side, sees the unnamed steer refused with both named, and steers each apart, one over the API and one through the flag. Co-Authored-By: Claude Fable 5.1 --- docs/public/api-reference/fabro-api.yaml | 10 + docs/public/reference/cli.mdx | 1 + lib/apps/fabro-cli/src/args.rs | 5 + .../src/commands/run/petri_worker.rs | 13 +- lib/apps/fabro-cli/src/commands/run/steer.rs | 12 +- .../tests/it/scenario/petri_controls.rs | 218 +++++++++++++++++- lib/apps/fabro-server/src/server.rs | 14 +- .../fabro-server/src/server/handler/steer.rs | 24 +- lib/apps/fabro-server/src/server/tests.rs | 48 +++- .../fabro-interview/src/control_protocol.rs | 14 +- lib/components/fabro-petri/src/controls.rs | 176 +++++++++++--- lib/components/fabro-tool/src/fabro_client.rs | 2 +- lib/foundation/fabro-client/src/client.rs | 17 +- .../src/models/steer-run-request.ts | 4 + 14 files changed, 491 insertions(+), 67 deletions(-) diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 52c6a58eb..e7e3dbc69 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -9946,6 +9946,16 @@ components: (default), append to the steering queue and let the agent pick it up at the next turn boundary. default: false + stage: + type: string + description: | + The agent stage to steer: its stage identifier (`node_id@visit`) + or its node name. Omit it to steer the run's one live agent + stage; a run with several live agent stages then refuses the + steer with a `run.notice` record. + minLength: 1 + maxLength: 200 + example: code@2 StartRunRequest: description: Request body for starting or resuming a run. diff --git a/docs/public/reference/cli.mdx b/docs/public/reference/cli.mdx index dfca14157..6b3f4ec76 100644 --- a/docs/public/reference/cli.mdx +++ b/docs/public/reference/cli.mdx @@ -1360,6 +1360,7 @@ fabro steer [OPTIONS] [TEXT] | --- | --- | | `--interrupt` | Cancel the in-flight LLM stream / tool calls and deliver the message as the next user turn (default: append to the steering queue) | | `--server ` | Fabro server target: http(s) URL or absolute Unix socket path | +| `--stage ` | Agent stage to steer, as its stage id (node@visit) or node name (default: the run's one live agent stage) | | `--text-stdin` | Read steer text from stdin instead of a positional arg | ### `fabro system` diff --git a/lib/apps/fabro-cli/src/args.rs b/lib/apps/fabro-cli/src/args.rs index 747768bdc..5cca0f072 100644 --- a/lib/apps/fabro-cli/src/args.rs +++ b/lib/apps/fabro-cli/src/args.rs @@ -798,6 +798,11 @@ pub(crate) struct SteerArgs { /// as the next user turn (default: append to the steering queue). #[arg(long)] pub(crate) interrupt: bool, + + /// Agent stage to steer, as its stage id (node@visit) or node name + /// (default: the run's one live agent stage) + #[arg(long, value_name = "STAGE")] + pub(crate) stage: Option, } #[derive(Args)] diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index f98002bf6..5bfe1fc52 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -23,8 +23,9 @@ //! questions wait on (`fabro_petri::interview`), so a human gate answered //! through the API continues; pause and unpause (and `SIGUSR1`/`SIGUSR2`) //! hold and release admission through the run's [`RunControls`]; a steer -//! goes to the run's one live -//! agent stage, or is refused with a `run.notice` record saying why. The +//! goes to the agent stage it names (`node@visit`, or the node name) or, +//! unnamed, to the run's one live agent stage, and is refused with a +//! `run.notice` record saying why when neither resolves. The //! paused state is mirrored to Fabro's lifecycle: a `paused` lifecycle //! record when admission is held and `unpaused` when it is released, so //! the server's live status and the projection agree with Petri's own @@ -331,10 +332,10 @@ impl PetriControls { self.controls.unpause().await; info!(run_id = %self.run_id, "unpause recorded: admission is released"); } - WorkerControlMessage::Steer { text, actor } => { - match self.controls.steer(None, &text).await { - Ok(node) => { - info!(run_id = %self.run_id, node, actor = ?actor, "steer delivered"); + WorkerControlMessage::Steer { text, stage, actor } => { + match self.controls.steer(stage.as_deref(), &text).await { + Ok(stage) => { + info!(run_id = %self.run_id, stage, actor = ?actor, "steer delivered"); } Err(error) => { warn!(run_id = %self.run_id, error = %error, "steer refused"); diff --git a/lib/apps/fabro-cli/src/commands/run/steer.rs b/lib/apps/fabro-cli/src/commands/run/steer.rs index fbc908440..09b0d8df8 100644 --- a/lib/apps/fabro-cli/src/commands/run/steer.rs +++ b/lib/apps/fabro-cli/src/commands/run/steer.rs @@ -26,7 +26,15 @@ pub(crate) async fn run(args: SteerArgs, base_ctx: &CommandContext) -> Result<() 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?; + let stage = args + .stage + .as_deref() + .map(str::trim) + .filter(|stage| !stage.is_empty()) + .map(str::to_owned); + info!(run_id = %run_id, interrupt = args.interrupt, stage = ?stage, "Sending steer"); + client + .steer_run(&run_id, text, args.interrupt, stage) + .await?; Ok(()) } diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs b/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs index 2de57fde8..f20211af9 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs @@ -3,8 +3,9 @@ //! `paused` in between; `SIGUSR1` and `SIGUSR2` on the worker do the same //! without the API; a steer reaches the agent stage on the twin, which //! sees it in its next request, and the stream carries the control record; -//! a run paused when its server and worker die resumes paused and goes on -//! once unpaused. +//! two live agent stages are steered apart by their stage labels, and an +//! unnamed steer between them is refused; a run paused when its server and +//! worker die resumes paused and goes on once unpaused. //! //! The harness is `petri.rs`'s: a foreground server on disk storage, the //! run started with `fabro run --detach`, and the host scope through the @@ -43,6 +44,12 @@ const HOLD: Duration = Duration::from_secs(1); const MODEL: &str = "gpt-5.4"; const PROMPT: &str = "Wait for the gate, then report."; const STEER: &str = "Steer: mention the word lighthouse in your report."; +/// Two agent stages side by side: each waits on its own gate, each is +/// steered apart. +const PROMPT_A: &str = "Alpha: wait for the gate, then report."; +const PROMPT_B: &str = "Bravo: wait for the gate, then report."; +const STEER_A: &str = "Steer alpha: mention the word lighthouse."; +const STEER_B: &str = "Steer bravo: mention the word windmill."; /// Two command stages: `a` waits on `gate`, `b` leaves `marker`. fn two_stage_workspace(context: &fabro_test::TestContext, gate: &Path, marker: &Path) -> PathBuf { @@ -93,16 +100,57 @@ async fn unpause(server: &RunningServer, run_id: &str) { } async fn steer(server: &RunningServer, run_id: &str, text: &str) { - let (status, body) = control( - server, - run_id, - "steer", - Some(json!({ "text": text, "interrupt": false })), - ) - .await; + steer_stage(server, run_id, text, None).await; +} + +/// `POST /runs/{id}/steer` naming `stage`, or no stage. +async fn steer_stage(server: &RunningServer, run_id: &str, text: &str, stage: Option<&str>) { + let mut body = json!({ "text": text, "interrupt": false }); + if let Some(stage) = stage { + body["stage"] = json!(stage); + } + let (status, body) = control(server, run_id, "steer", Some(body)).await; assert_eq!(status, 202, "steer: {body}"); } +/// `fabro steer --stage ` against the server. +fn steer_by_cli( + context: &fabro_test::TestContext, + server: &RunningServer, + run_id: &str, + stage: &str, + text: &str, +) { + let output = context + .command() + .args(["steer", "--server", &server.target(), run_id]) + .args(["--stage", stage, text]) + .output() + .expect("the steer command executes"); + assert!( + output.status.success(), + "fabro steer failed\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); +} + +/// The twin's request inputs that carry `prompt`, in order. +fn inputs_with(logs: &Value, prompt: &str) -> Vec { + logs["requests"] + .as_array() + .expect("the twin request log is an array") + .iter() + .map(|request| { + request["input_text"] + .as_str() + .unwrap_or_default() + .to_string() + }) + .filter(|input| input.contains(prompt)) + .collect() +} + /// The run's pending control, as the API shows it. async fn pending_control(server: &RunningServer, run_id: &str) -> Value { run_json(server, &format!("runs/{run_id}")).await["lifecycle"]["pending_control"].clone() @@ -422,6 +470,158 @@ async fn a_steer_reaches_the_agent_stage_on_the_twin() { server.shutdown(); } +/// Two agent stages live at once, as the branches of a parallel node: a +/// steer that names no stage is refused with a notice naming both, and a +/// steer to each label (`a@1` over the API, `b@1` through the CLI flag) +/// reaches that stage's session and no other. +#[tokio::test(flavor = "multi_thread")] +async fn two_live_agent_stages_are_steered_apart_by_their_labels() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let twin = twin_openai().await; + let namespace = format!("{}::{}", module_path!(), line!()); + let server = RunningServer::start_with( + &format!( + "\n[llm.providers.openai]\nbase_url = \"{}\"\n", + twin.base_url + ), + &[(EnvVars::OPENAI_API_KEY, &namespace)], + ) + .await; + let gate_a = context.temp_dir.join("a.gate"); + let gate_b = context.temp_dir.join("b.gate"); + let wait_on = |gate: &Path| { + TwinToolCall::new( + "shell", + json!({ "command": format!("while [ ! -f {} ]; do sleep 0.05; done", gate.display()) }), + ) + }; + TwinScenarios::new(namespace.clone()) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(PROMPT_A) + .tool_call(wait_on(&gate_a)), + ) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(PROMPT_A) + .text("Alpha's gate opened."), + ) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(STEER_A) + .text("Lighthouse noted."), + ) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(PROMPT_B) + .tool_call(wait_on(&gate_b)), + ) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(PROMPT_B) + .text("Bravo's gate opened."), + ) + .scenario( + TwinScenario::responses(MODEL) + .input_contains(STEER_B) + .text("Windmill noted."), + ) + .load(twin) + .await; + let workspace = write_petri_workflow( + &context, + &format!( + "digraph Pair {{\n graph [goal=\"Two agents wait then report\", \ + default_max_retries=0]\n start [shape=Mdiamond]\n exit [shape=Msquare]\n fan \ + [shape=component]\n a [shape=box, prompt=\"{PROMPT_A}\", max_retries=0]\n b \ + [shape=box, prompt=\"{PROMPT_B}\", max_retries=0]\n join \ + [shape=tripleoctagon]\n start -> fan\n fan -> a\n fan -> b\n a -> join\n b -> \ + join\n join -> exit\n}}\n" + ), + ); + let run_id = run_detached_with(&context, &server, &workspace, &[ + "--auto-approve", + "--provider", + "openai", + "--model", + MODEL, + ]); + + wait_for_status(&server, &run_id, &["running"]).await; + wait_until_gate_is_polled(&gate_a); + wait_until_gate_is_polled(&gate_b); + eprintln!("run {run_id}: both agents' tools are waiting on their gates"); + + // Unnamed, the steer has two candidates and is refused with both named. + steer(&server, &run_id, "Steer nobody.").await; + let names = wait_for_stream_count(&server, &run_id, "run.notice", 1).await; + assert_eq!(count_of(&names, "control.requested"), 0, "{names:?}"); + let notice = run_stream(&server, &run_id) + .await + .into_iter() + .find(|item| item["item"]["record"]["kind"] == "run.notice") + .expect("the refusal is recorded"); + let message = notice["item"]["record"]["message"] + .as_str() + .unwrap_or_default() + .to_string(); + assert_eq!( + notice["item"]["record"]["code"], "steer_refused", + "{notice}" + ); + assert!( + message.contains("a@1") && message.contains("b@1"), + "the notice names both live stages: {message}" + ); + + steer_stage(&server, &run_id, STEER_A, Some("a@1")).await; + steer_by_cli(&context, &server, &run_id, "b@1", STEER_B); + wait_for_stream_count(&server, &run_id, "control.requested", 2).await; + eprintln!("run {run_id}: both steers are recorded"); + std::fs::write(&gate_a, "go").expect("gate a opens"); + std::fs::write(&gate_b, "go").expect("gate b opens"); + + let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await; + let items = settled_stream(&server, &run_id).await; + let names = stream_names(&items); + assert_eq!( + status, + "succeeded", + "stream: {names:?}\nserver stderr:\n{}", + server.stderr_text() + ); + assert_petri_succeeded(&server, &run_id).await; + + let logs = twin.request_logs(&namespace).await; + for (prompt, steer, other) in [(PROMPT_A, STEER_A, STEER_B), (PROMPT_B, STEER_B, STEER_A)] { + let inputs = inputs_with(&logs, prompt); + assert_eq!( + inputs.len(), + 3, + "{prompt}: the tool call, its answer, the steer: {inputs:?}" + ); + assert!( + !inputs[1].contains(steer), + "{prompt}: the answer's request came before the steer's turn: {}", + inputs[1] + ); + assert!( + inputs[2].contains(steer), + "{prompt}: the follow-up request carries its own steer: {}", + inputs[2] + ); + assert!( + !inputs[2].contains(other), + "{prompt}: the other stage's steer stayed away: {}", + inputs[2] + ); + } + server.shutdown(); +} + /// A run paused with its next stage held at admission, whose server and /// worker then die, resumes paused: the resumed worker reports the pause /// again, admits nothing until the unpause, then finishes the run. (A diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 074f9593b..01808ace5 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -393,12 +393,18 @@ impl RunAnswerTransport { } } - /// Forward a steer to the worker. The in-process test path drives no - /// steer: its run has no live agent session to steer. - async fn steer(&self, text: String, actor: Principal) -> Result<(), AnswerTransportError> { + /// Forward a steer to the worker, for the stage it names or the run's + /// one live agent stage. The in-process test path drives no steer: its + /// run has no live agent session to steer. + async fn steer( + &self, + text: String, + stage: Option, + actor: Principal, + ) -> Result<(), AnswerTransportError> { match self { Self::Worker { run_id, bus } => { - let message = WorkerControlEnvelope::steer(text, actor); + let message = WorkerControlEnvelope::steer(text, stage, actor); Self::publish_worker_control(*run_id, bus, message) .await .map_err(|err| Self::answer_error_from_bus(&err)) diff --git a/lib/apps/fabro-server/src/server/handler/steer.rs b/lib/apps/fabro-server/src/server/handler/steer.rs index dc2bb0df9..d44de3882 100644 --- a/lib/apps/fabro-server/src/server/handler/steer.rs +++ b/lib/apps/fabro-server/src/server/handler/steer.rs @@ -20,7 +20,10 @@ pub(super) fn routes() -> axum::Router> { } enum RunControlRequest { - Steer { text: String }, + Steer { + text: String, + stage: Option, + }, } async fn steer_run( @@ -30,15 +33,26 @@ async fn steer_run( ) -> Response { // OpenAPI enforces minLength=1/maxLength=8192 already; only whitespace-only // payloads can slip through. - let SteerRunRequest { text, interrupt } = req; + let SteerRunRequest { + text, + interrupt, + stage, + } = req; let text: String = text.into(); if text.trim().is_empty() { return ApiError::bad_request("Steer text must not be empty.").into_response(); } + let stage = stage.map(String::from); + if stage + .as_deref() + .is_some_and(|stage| stage.trim().is_empty()) + { + return ApiError::bad_request("Steer stage must not be empty.").into_response(); + } if interrupt { return interrupt_unsupported(); } - control_run(actor, state, id, RunControlRequest::Steer { text }).await + control_run(actor, state, id, RunControlRequest::Steer { text, stage }).await } /// Interrupting a live agent turn has no adapter over Petri's control @@ -139,8 +153,8 @@ async fn control_run( .into_response(); }; - let RunControlRequest::Steer { text } = control; - let result = answer_transport.steer(text, actor).await; + let RunControlRequest::Steer { text, stage } = control; + let result = answer_transport.steer(text, stage, actor).await; match result { Ok(()) => StatusCode::ACCEPTED.into_response(), diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 57de156fd..f834d88a7 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -2788,13 +2788,13 @@ async fn worker_answer_transport_steer_publishes_plain_steer_message() { }; transport - .steer("try again".to_string(), actor.clone()) + .steer("try again".to_string(), None, actor.clone()) .await .unwrap(); assert_eq!( recv_worker_control_envelope(&mut control_rx).await, - WorkerControlEnvelope::steer("try again", actor) + WorkerControlEnvelope::steer("try again", None, actor) ); } @@ -8318,6 +8318,50 @@ async fn steer_without_active_steerable_session_forwards_plain_steer_for_bufferi )); } +#[tokio::test] +async fn steer_with_a_stage_forwards_the_stage_to_the_worker() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = fixtures::RUN_1; + let (transport, mut control_rx) = worker_transport_with_receiver(run_id).await; + let _temp_dir = insert_running_control_run(&state, run_id, Some(transport)); + + let req = Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/steer"))) + .header("content-type", "application/json") + .body(Body::from(r#"{"text":"try again","stage":"code@2"}"#)) + .unwrap(); + + let response = app.oneshot(req).await.unwrap(); + assert_status!(response, StatusCode::ACCEPTED).await; + let envelope = recv_worker_control_envelope(&mut control_rx).await; + assert!(matches!( + envelope.message, + WorkerControlMessage::Steer { ref text, ref stage, .. } + if text == "try again" && stage.as_deref() == Some("code@2") + )); +} + +#[tokio::test] +async fn steer_with_a_blank_stage_returns_bad_request() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = fixtures::RUN_1; + let (transport, _control_rx) = worker_transport_with_receiver(run_id).await; + let _temp_dir = insert_running_control_run(&state, run_id, Some(transport)); + + let req = Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/steer"))) + .header("content-type", "application/json") + .body(Body::from(r#"{"text":"try again","stage":" "}"#)) + .unwrap(); + + let response = app.oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); +} + #[tokio::test] async fn steer_with_active_non_steerable_session_returns_conflict() { let state = test_app_state(); diff --git a/lib/components/fabro-interview/src/control_protocol.rs b/lib/components/fabro-interview/src/control_protocol.rs index f76b61121..83fb273f9 100644 --- a/lib/components/fabro-interview/src/control_protocol.rs +++ b/lib/components/fabro-interview/src/control_protocol.rs @@ -71,11 +71,12 @@ impl WorkerControlEnvelope { } #[must_use] - pub fn steer(text: impl Into, actor: Principal) -> Self { + pub fn steer(text: impl Into, stage: Option, actor: Principal) -> Self { Self { v: WORKER_CONTROL_PROTOCOL_VERSION, message: WorkerControlMessage::Steer { text: text.into(), + stage, actor, }, } @@ -163,7 +164,14 @@ pub enum WorkerControlMessage { #[serde(rename = "run.unpause")] RunUnpause, #[serde(rename = "run.steer")] - Steer { text: String, actor: Principal }, + Steer { + text: String, + /// The stage to steer (`node@visit`, or the node name); `None` + /// steers the run's one live agent stage. + #[serde(default, skip_serializing_if = "Option::is_none")] + stage: Option, + actor: Principal, + }, #[serde(rename = "run.interrupt")] Interrupt { actor: Principal }, #[serde(rename = "run.interrupt_then_steer")] @@ -298,7 +306,7 @@ mod tests { #[test] fn steer_append_round_trips_through_json() { - let envelope = WorkerControlEnvelope::steer("try again", Principal::System { + let envelope = WorkerControlEnvelope::steer("try again", None, Principal::System { system_kind: SystemActorKind::Engine, }); let json = serde_json::to_string(&envelope).unwrap(); diff --git a/lib/components/fabro-petri/src/controls.rs b/lib/components/fabro-petri/src/controls.rs index d056eb6d1..ae32e1a3e 100644 --- a/lib/components/fabro-petri/src/controls.rs +++ b/lib/components/fabro-petri/src/controls.rs @@ -21,9 +21,11 @@ //! - steer delivers a text to a live agent stage as guidance for its session: //! the stage's firing records `control.requested` with the `{"$steer": …}` //! value, and the agent runs the text as a follow-up turn once its current -//! answer is reached. Fabro's steer names no stage, so the steer goes to the -//! one live agent stage; with none, or several, it is refused with the -//! reason, and nothing is recorded. +//! answer is reached. A steer names its stage by the label the projection +//! shows (`node@visit`, or `node/e@visit` when two executions +//! share one) or by the node's name; unnamed, it goes to the one live agent +//! stage. With no live agent, several unnamed, or a name that is not running, +//! it is refused with the reason, and nothing is recorded. //! - cancel is the caller's cancellation token ([`RunRequest::cancel`]); the //! service's own cancel is here for a host that holds only this. //! @@ -56,20 +58,70 @@ pub enum SteerError { /// control that needs a live agent session. #[error("Run has no active steerable agent session.")] NoLiveAgent, - /// More than one agent stage is running and the steer names none. - #[error("Run has several active agent stages ({}); the steer names none.", .0.join(", "))] + /// More than one agent stage is running and the steer names none, or + /// names a label several live firings answer to. + #[error("Run has several active agent stages ({}); the steer names none of them.", .0.join(", "))] SeveralLiveAgents(Vec), /// The named stage is not running, or the run has ended. #[error(transparent)] Control(#[from] ControlError), } -/// The live agent firings, by node name: what a steer that names no stage -/// is routed by. +/// One live agent firing: the node's name and which firing of the node it +/// is within its execution, which is the visit its stage label carries. +#[derive(Clone, Debug, PartialEq, Eq)] +struct LiveAgent { + node: String, + visit: u32, +} + +/// The live agent firings: what a steer is routed by. #[derive(Default)] struct LiveAgents { - stages: BTreeMap, - firings: BTreeMap<(ExecutionId, FiringId), String>, + firings: BTreeMap<(ExecutionId, FiringId), LiveAgent>, +} + +impl LiveAgents { + /// Every live agent firing with its label: `node@visit`, or + /// `node/e@visit` when another execution's firing has the + /// same node and visit, as the projection labels them. + fn labelled(&self) -> Vec<((ExecutionId, FiringId), String)> { + let mut counts: BTreeMap<(&str, u32), usize> = BTreeMap::new(); + for agent in self.firings.values() { + *counts + .entry((agent.node.as_str(), agent.visit)) + .or_default() += 1; + } + self.firings + .iter() + .map(|(key, agent)| { + let label = if counts[&(agent.node.as_str(), agent.visit)] > 1 { + format!("{}/e{}@{}", agent.node, key.0.raw(), agent.visit) + } else { + format!("{}@{}", agent.node, agent.visit) + }; + (*key, label) + }) + .collect() + } +} + +/// A stage label taken apart: the node name, the execution when the label +/// names one, and the visit. `None` when `stage` is not a label. +fn parse_label(stage: &str) -> Option<(&str, Option, u32)> { + let (node, visit) = stage.rsplit_once('@')?; + let visit = visit.parse().ok()?; + let (node, execution) = match node.rsplit_once("/e") { + Some((name, execution)) => match execution.parse::() { + Ok(execution) => (name, Some(execution)), + Err(_) => (node, None), + }, + None => (node, None), + }; + if node.is_empty() { + return None; + } + Some((node, execution, visit)) } /// One run's controls. Clone freely: every clone drives the same service. @@ -115,27 +167,66 @@ impl RunControls { self.service.paused_changes() } - /// The names of the agent stages running now. + /// The labels of the agent stages running now. #[must_use] pub fn live_agents(&self) -> Vec { - self.agents().stages.keys().cloned().collect() + self.agents() + .labelled() + .into_iter() + .map(|(_, label)| label) + .collect() } - /// Deliver `text` to the named agent stage, or to the one live agent - /// stage when `node` is `None`. The name of the stage steered. - pub async fn steer(&self, node: Option<&str>, text: &str) -> Result { - let node = if let Some(node) = node { - node.to_owned() - } else { - let mut live = self.live_agents(); - match live.len() { + /// Deliver `text` to the stage `stage` names (a label, `node@visit`, or + /// a node name), or to the one live agent stage when `stage` is `None`. + /// The label of the stage steered. + pub async fn steer(&self, stage: Option<&str>, text: &str) -> Result { + let live = self.agents().labelled(); + let ((execution, firing), label) = match stage { + None => match live.len() { 0 => return Err(SteerError::NoLiveAgent), - 1 => live.remove(0), - _ => return Err(SteerError::SeveralLiveAgents(live)), - } + 1 => live.into_iter().next().expect("one live agent"), + _ => { + return Err(SteerError::SeveralLiveAgents( + live.into_iter().map(|(_, label)| label).collect(), + )); + } + }, + Some(stage) => match parse_label(stage) { + Some((node, execution, visit)) => { + let agents = self.agents(); + let mut matches: Vec<_> = live + .into_iter() + .filter(|(key, _)| { + let agent = &agents.firings[key]; + agent.node == node + && agent.visit == visit + && execution.is_none_or(|execution| key.0.raw() == execution) + }) + .collect(); + match matches.len() { + 0 => { + return Err(SteerError::Control(ControlError::NoSuchStage( + stage.to_owned(), + ))); + } + 1 => matches.remove(0), + _ => { + return Err(SteerError::SeveralLiveAgents( + matches.into_iter().map(|(_, label)| label).collect(), + )); + } + } + } + None => { + // A node name: the service's own live-stage index. + self.service.steer(stage, text).await?; + return Ok(stage.to_owned()); + } + }, }; - self.service.steer(&node, text).await?; - Ok(node) + self.service.steer_firing(execution, firing, text).await?; + Ok(label) } /// Cancel the whole run politely; a second call reaches the kill tier. @@ -186,18 +277,18 @@ impl ExecutionObserver for RunControls { if node.step.kind != AGENT_KIND { return; } - let name = node.name.to_string(); - let mut agents = self.agents(); - agents.stages.insert(name.clone(), (execution, *firing)); - agents.firings.insert((execution, *firing), name); + // The visit is the firing's ordinal among the node's firings + // in this execution: what the projection labels the stage by. + let visit = state.firing_count(node.id).max(1); + self.agents() + .firings + .insert((execution, *firing), LiveAgent { + node: node.name.to_string(), + visit, + }); } Event::StepFinished { firing, .. } => { - let mut agents = self.agents(); - if let Some(name) = agents.firings.remove(&(execution, *firing)) { - if agents.stages.get(&name) == Some(&(execution, *firing)) { - agents.stages.remove(&name); - } - } + self.agents().firings.remove(&(execution, *firing)); } _ => {} } @@ -238,6 +329,23 @@ mod tests { "work".to_string() ))) ); + assert_eq!( + controls.steer(Some("work@1"), "hurry up").await, + Err(SteerError::Control(ControlError::NoSuchStage( + "work@1".to_string() + ))) + ); + } + + #[test] + fn a_stage_label_names_its_node_visit_and_execution() { + assert_eq!(parse_label("work@1"), Some(("work", None, 1))); + assert_eq!(parse_label("work/e2@3"), Some(("work", Some(2), 3))); + assert_eq!(parse_label("a/b@1"), Some(("a/b", None, 1))); + assert_eq!(parse_label("a/ex@1"), Some(("a/ex", None, 1))); + assert_eq!(parse_label("work"), None); + assert_eq!(parse_label("work@one"), None); + assert_eq!(parse_label("@1"), None); } #[test] diff --git a/lib/components/fabro-tool/src/fabro_client.rs b/lib/components/fabro-tool/src/fabro_client.rs index 7103f7d31..84948afd0 100644 --- a/lib/components/fabro-tool/src/fabro_client.rs +++ b/lib/components/fabro-tool/src/fabro_client.rs @@ -134,7 +134,7 @@ impl FabroToolBackend for ClientBackend { async fn steer_run(&self, run_id: &RunId, text: String, interrupt: bool) -> anyhow::Result<()> { self.ensure_run_scope(run_id)?; - self.client.steer_run(run_id, text, interrupt).await + self.client.steer_run(run_id, text, interrupt, None).await } async fn archive_run(&self, run_id: &RunId) -> anyhow::Result { diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 241b09431..4b0e6a2e4 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1136,10 +1136,25 @@ impl Client { Ok(()) } - pub async fn steer_run(&self, run_id: &RunId, text: String, interrupt: bool) -> Result<()> { + /// Steer a run: the named stage (`node@visit`, or the node name), or + /// the run's one live agent stage when `stage` is `None`. + pub async fn steer_run( + &self, + run_id: &RunId, + text: String, + interrupt: bool, + stage: Option, + ) -> Result<()> { + let stage = stage + .map(|stage| { + types::SteerRunRequestStage::try_from(stage) + .map_err(|e| anyhow!("invalid steer stage: {e}")) + }) + .transpose()?; let body: types::SteerRunRequest = types::SteerRunRequest::builder() .text(text) .interrupt(interrupt) + .stage(stage) .try_into() .map_err(|e| anyhow!("failed to build SteerRunRequest: {e}"))?; self.send_api(|client| { diff --git a/lib/packages/fabro-api-client/src/models/steer-run-request.ts b/lib/packages/fabro-api-client/src/models/steer-run-request.ts index daa65e248..4d8f4be5e 100644 --- a/lib/packages/fabro-api-client/src/models/steer-run-request.ts +++ b/lib/packages/fabro-api-client/src/models/steer-run-request.ts @@ -26,4 +26,8 @@ export interface SteerRunRequest { * When true, apply a worker-control interrupt first, then deliver this text as steering in the same control operation. When false (default), append to the steering queue and let the agent pick it up at the next turn boundary. */ 'interrupt'?: boolean; + /** + * The agent stage to steer: its stage identifier (`node_id@visit`) or its node name. Omit it to steer the run\'s one live agent stage; a run with several live agent stages then refuses the steer with a `run.notice` record. + */ + 'stage'?: string; }