mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Answer steer and interrupt with the worker's acknowledgement
The worker control bus was publish-only: the steer and interrupt
endpoints answered 202 once the control was forwarded, and a refusal
showed up only later as a `run.notice` on the run's stream.
A steer or an interrupt now carries a request id. The worker answers it
over the control stream it arrived on with `{request_id, outcome}`,
where the outcome is `delivered` (with the stage's label) or `refused`
(with the code and the reason). The server keeps the outstanding
requests in a registry and waits up to 5 s for the answer: the endpoint
answers 202 `{"outcome":"delivered","stage":…}`, 409 with the refusal's
code (`no_live_turn`, `no_such_stage`, `steer_refused`,
`interrupt_refused`) and message, or 202 `{"outcome":"pending"}` when
the worker gave no answer in time. The `run.notice` record on refusal
stays, under the same code, so a steer to a stage that is not running is
now `no_such_stage` there too. Pause and unpause are unchanged.
The in-process test path answers a steer or an interrupt from the run's
own controls at once. `FABRO_TEST_CONTROL_ACKS_MUTED=1` on the server
mutes the worker's answers, so a test can see the pending fallback.
`fabro steer` prints the worker's answer, and a refusal is its error.
The OpenAPI spec documents the 202 body and the 409 codes; the Rust and
TypeScript clients are regenerated.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
86ac13713b
commit
cb26c5603d
26 changed files with 1619 additions and 221 deletions
|
|
@ -1726,12 +1726,14 @@ paths:
|
|||
`interrupt=true`, the stage's current model turn (the model request
|
||||
and the tool calls it is running) is stopped first, the session is
|
||||
kept, and the text is the stage's next input. The control is
|
||||
forwarded to the run's worker; a control the worker cannot deliver
|
||||
(no live agent stage, several unnamed, a stage that is not running,
|
||||
or, for an interrupt, a stage with no model turn in flight) is
|
||||
refused on the run's event stream as a `run.notice` record whose
|
||||
code says why (`steer_refused`, `no_live_turn`, `no_such_stage`,
|
||||
`interrupt_refused`).
|
||||
forwarded to the run's worker, and the worker's answer is this
|
||||
response: `202` with `outcome: delivered` (and the stage's label)
|
||||
once the worker delivered it, `409` with the refusal's code when
|
||||
the worker or Petri refused it, and `202` with `outcome: pending`
|
||||
when the worker gave no answer within the wait (5 s). A refused
|
||||
control is also a `run.notice` record on the run's event stream
|
||||
under the same code (`steer_refused`, `no_live_turn`,
|
||||
`no_such_stage`, `interrupt_refused`).
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
requestBody:
|
||||
|
|
@ -1742,7 +1744,13 @@ paths:
|
|||
$ref: "#/components/schemas/SteerRunRequest"
|
||||
responses:
|
||||
"202":
|
||||
description: Steer accepted and forwarded to the worker
|
||||
description: |
|
||||
Steer forwarded to the worker: delivered, or pending when the
|
||||
worker gave no answer within the wait.
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/RunControlAcknowledgement"
|
||||
"400":
|
||||
description: Invalid request body
|
||||
headers:
|
||||
|
|
@ -1763,9 +1771,16 @@ paths:
|
|||
$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
|
||||
active agent sessions have no live control channel.
|
||||
Run is not currently steerable, or the worker refused the
|
||||
steer. The error's `code` says which: `run_not_steerable` (a
|
||||
terminal run, or one not running yet), `use_answer_endpoint`
|
||||
(the run is blocked on a question), `agent_not_steerable` (the
|
||||
active agent sessions have no live control channel),
|
||||
`steer_refused` (no live agent stage, or several and none
|
||||
named), `no_such_stage` (the named stage is not running),
|
||||
`no_live_turn` (with `interrupt=true`: the stage has no model
|
||||
turn in flight), `interrupt_refused` (with `interrupt=true`:
|
||||
no live agent stage, or several and none named).
|
||||
headers:
|
||||
x-request-id:
|
||||
$ref: "#/components/headers/XRequestId"
|
||||
|
|
@ -2049,14 +2064,19 @@ paths:
|
|||
names, or the run's one live agent stage) and keep its session. With
|
||||
`text`, the text is the stage's next input; without, the stage waits
|
||||
for the next steer message before starting another model turn. The
|
||||
control is forwarded to the run's worker; an interrupt the worker
|
||||
cannot deliver (a stage with no model turn in flight, such as an
|
||||
agent between turns or a human gate; a stage that is not running; no
|
||||
live agent stage, or several unnamed) is refused on the run's event
|
||||
stream as a `run.notice` record whose code says why (`no_live_turn`,
|
||||
`no_such_stage`, `interrupt_refused`). A delivered interrupt is the
|
||||
stage's `control.requested` record with `$interrupt`, followed by an
|
||||
`attractor.turn.interrupted` progress record.
|
||||
control is forwarded to the run's worker, and the worker's answer is
|
||||
this response: `202` with `outcome: delivered` (and the stage's
|
||||
label) once the worker stopped the turn, `409` with the refusal's
|
||||
code when the worker or Petri refused it (a stage with no model
|
||||
turn in flight, such as an agent between turns or a human gate; a
|
||||
stage that is not running; no live agent stage, or several
|
||||
unnamed), and `202` with `outcome: pending` when the worker gave no
|
||||
answer within the wait (5 s). A refused interrupt is also a
|
||||
`run.notice` record on the run's event stream under the same code
|
||||
(`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered
|
||||
interrupt is the stage's `control.requested` record with
|
||||
`$interrupt`, followed by an `attractor.turn.interrupted` progress
|
||||
record.
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
requestBody:
|
||||
|
|
@ -2067,7 +2087,13 @@ paths:
|
|||
$ref: "#/components/schemas/InterruptRunRequest"
|
||||
responses:
|
||||
"202":
|
||||
description: Interrupt accepted and forwarded to the worker
|
||||
description: |
|
||||
Interrupt forwarded to the worker: delivered, or pending when
|
||||
the worker gave no answer within the wait.
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/RunControlAcknowledgement"
|
||||
"400":
|
||||
description: Invalid request body
|
||||
headers:
|
||||
|
|
@ -2088,8 +2114,14 @@ paths:
|
|||
$ref: "#/components/schemas/ErrorResponse"
|
||||
"409":
|
||||
description: |
|
||||
Run is not currently interruptible. Returned when the run is in a
|
||||
terminal state or is not running yet.
|
||||
Run is not currently interruptible, or the worker refused the
|
||||
interrupt. The error's `code` says which: `run_not_interruptible`
|
||||
(a terminal run, or one not running yet), `agent_not_steerable`
|
||||
(the active agent sessions have no live control channel),
|
||||
`no_live_turn` (the stage has no model turn in flight),
|
||||
`no_such_stage` (the named stage is not running),
|
||||
`interrupt_refused` (no live agent stage, or several and none
|
||||
named).
|
||||
headers:
|
||||
x-request-id:
|
||||
$ref: "#/components/headers/XRequestId"
|
||||
|
|
@ -10178,6 +10210,32 @@ components:
|
|||
maxLength: 200
|
||||
example: code@2
|
||||
|
||||
RunControlAcknowledgement:
|
||||
description: |
|
||||
The worker's answer to a steer or an interrupt, as the endpoint's
|
||||
`202` body.
|
||||
type: object
|
||||
required:
|
||||
- outcome
|
||||
properties:
|
||||
outcome:
|
||||
$ref: "#/components/schemas/RunControlOutcome"
|
||||
stage:
|
||||
type: string
|
||||
description: >-
|
||||
The label of the stage the control was delivered to
|
||||
(`node@visit`, or `node/e<execution>@visit`); absent when the
|
||||
outcome is `pending`.
|
||||
RunControlOutcome:
|
||||
description: |
|
||||
What became of a forwarded control: `delivered` once the worker
|
||||
delivered it to its stage; `pending` when the worker gave no answer
|
||||
within the wait, in which case the run's event stream says what
|
||||
became of it (a `run.notice` record on refusal).
|
||||
type: string
|
||||
enum:
|
||||
- delivered
|
||||
- pending
|
||||
InterruptRunRequest:
|
||||
description: Request body for interrupting a live agent stage's model turn.
|
||||
type: object
|
||||
|
|
|
|||
|
|
@ -29,7 +29,10 @@
|
|||
//! resolves its stage the same way and stops the stage's current model
|
||||
//! turn, with the text of an `interrupt_then_steer` as the stage's next
|
||||
//! input, and is refused with a `run.notice` (`no_live_turn`,
|
||||
//! `no_such_stage`) when Petri refuses it. The
|
||||
//! `no_such_stage`) when Petri refuses it. A steer or an interrupt that
|
||||
//! carries a request id is acknowledged over the control channel with its
|
||||
//! outcome, delivered or refused with the notice's code, so the server can
|
||||
//! answer the caller in its own response. 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
|
||||
|
|
@ -66,10 +69,10 @@ use std::time::Instant;
|
|||
use anyhow::{Context, Result, anyhow};
|
||||
use fabro_auth::VaultCredentialSource;
|
||||
use fabro_client::{Client, ServerTarget};
|
||||
use fabro_interview::{ControlInterviewer, WorkerControlMessage};
|
||||
use fabro_interview::{ControlInterviewer, WorkerControlMessage, WorkerControlOutcome};
|
||||
use fabro_llm::credentials::{CredentialProvider, readiness};
|
||||
use fabro_petri::blobs::ClientBlobs;
|
||||
use fabro_petri::controls::{ControlError, RunControls, SteerError};
|
||||
use fabro_petri::controls::{RunControls, SteerError};
|
||||
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
|
||||
use fabro_petri::hooks::HooksSpec;
|
||||
use fabro_petri::interview::{Approval, FabroInterviewer};
|
||||
|
|
@ -137,11 +140,10 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
|
|||
// Fabro's own records of the run, over the client.
|
||||
let records: Arc<dyn PlatformRecords> =
|
||||
Arc::new(HttpPlatformRecords::new(worker.client.clone_for_reuse()));
|
||||
let petri_controls = Arc::new(PetriControls::new(
|
||||
run_id,
|
||||
controls.clone(),
|
||||
Arc::clone(&records),
|
||||
));
|
||||
let petri_controls = Arc::new(
|
||||
PetriControls::new(run_id, controls.clone(), Arc::clone(&records))
|
||||
.with_muted_acks(test_control_acks_muted()),
|
||||
);
|
||||
let mut control_manager = runner::spawn_worker_control_manager(
|
||||
worker.target.clone(),
|
||||
run_id,
|
||||
|
|
@ -301,10 +303,13 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
|
|||
/// Cancel and answers are applied by the channel itself, before a message
|
||||
/// reaches here.
|
||||
pub(super) struct PetriControls {
|
||||
run_id: RunId,
|
||||
controls: RunControls,
|
||||
run_id: RunId,
|
||||
controls: RunControls,
|
||||
/// Where a refused steer's notice goes.
|
||||
records: Arc<dyn PlatformRecords>,
|
||||
records: Arc<dyn PlatformRecords>,
|
||||
/// A test hook: the worker applies every control but acknowledges
|
||||
/// none, so the server's wait for an answer runs out.
|
||||
muted_acks: bool,
|
||||
}
|
||||
|
||||
impl PetriControls {
|
||||
|
|
@ -317,42 +322,52 @@ impl PetriControls {
|
|||
run_id,
|
||||
controls,
|
||||
records,
|
||||
muted_acks: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Apply every control but acknowledge none: a test's stand-in for a
|
||||
/// worker that never answers.
|
||||
#[must_use]
|
||||
pub(super) fn with_muted_acks(mut self, muted: bool) -> Self {
|
||||
self.muted_acks = muted;
|
||||
self
|
||||
}
|
||||
|
||||
/// The run's controls, for a test that reads the paused state back.
|
||||
#[cfg(test)]
|
||||
pub(super) fn controls(&self) -> &RunControls {
|
||||
&self.controls
|
||||
}
|
||||
|
||||
pub(super) async fn apply(&self, message: WorkerControlMessage) {
|
||||
match message {
|
||||
/// Apply the control. The outcome, for the channel to acknowledge when
|
||||
/// the control asked for one: `None` for a control that has no
|
||||
/// outcome to report (a pause, an ignored pair control) and when the
|
||||
/// acknowledgements are muted.
|
||||
pub(super) async fn apply(
|
||||
&self,
|
||||
message: WorkerControlMessage,
|
||||
) -> Option<WorkerControlOutcome> {
|
||||
let outcome = match message {
|
||||
WorkerControlMessage::RunPause => {
|
||||
info!(run_id = %self.run_id, "pause requested: admission is held");
|
||||
self.controls.pause();
|
||||
None
|
||||
}
|
||||
WorkerControlMessage::RunUnpause => {
|
||||
self.controls.unpause().await;
|
||||
info!(run_id = %self.run_id, "unpause recorded: admission is released");
|
||||
None
|
||||
}
|
||||
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");
|
||||
self.notice("steer_refused", error.to_string()).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
WorkerControlMessage::Interrupt { stage, actor } => {
|
||||
self.interrupt(stage.as_deref(), None, &actor).await;
|
||||
}
|
||||
WorkerControlMessage::InterruptThenSteer { text, stage, actor } => {
|
||||
self.interrupt(stage.as_deref(), Some(&text), &actor).await;
|
||||
WorkerControlMessage::Steer {
|
||||
text, stage, actor, ..
|
||||
} => Some(self.steer(stage.as_deref(), &text, &actor).await),
|
||||
WorkerControlMessage::Interrupt { stage, actor, .. } => {
|
||||
Some(self.interrupt(stage.as_deref(), None, &actor).await)
|
||||
}
|
||||
WorkerControlMessage::InterruptThenSteer {
|
||||
text, stage, actor, ..
|
||||
} => Some(self.interrupt(stage.as_deref(), Some(&text), &actor).await),
|
||||
WorkerControlMessage::PairStart { .. }
|
||||
| WorkerControlMessage::PairMessage { .. }
|
||||
| WorkerControlMessage::PairEnd { .. } => {
|
||||
|
|
@ -361,8 +376,36 @@ impl PetriControls {
|
|||
control = control_name(&message),
|
||||
"control has no Petri adapter yet and is ignored"
|
||||
);
|
||||
None
|
||||
}
|
||||
WorkerControlMessage::InterviewAnswer { .. } | WorkerControlMessage::RunCancel => None,
|
||||
};
|
||||
if self.muted_acks { None } else { outcome }
|
||||
}
|
||||
|
||||
/// Deliver `text` to the named stage, or to the run's one live agent
|
||||
/// stage. A refusal is a `run.notice` whose code says why
|
||||
/// (`no_such_stage` when the name is not running, `steer_refused`
|
||||
/// otherwise) and the same code and reason go back as the outcome.
|
||||
async fn steer(
|
||||
&self,
|
||||
stage: Option<&str>,
|
||||
text: &str,
|
||||
actor: &Principal,
|
||||
) -> WorkerControlOutcome {
|
||||
match self.controls.steer(stage, text).await {
|
||||
Ok(stage) => {
|
||||
info!(run_id = %self.run_id, stage, actor = ?actor, "steer delivered");
|
||||
WorkerControlOutcome::Delivered { stage: Some(stage) }
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(run_id = %self.run_id, stage, error = %error, "steer refused");
|
||||
self.refuse(
|
||||
error.code().unwrap_or("steer_refused"),
|
||||
refusal_message("Steer", stage, &error),
|
||||
)
|
||||
.await
|
||||
}
|
||||
WorkerControlMessage::InterviewAnswer { .. } | WorkerControlMessage::RunCancel => {}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -371,8 +414,13 @@ impl PetriControls {
|
|||
/// when the stage has no model turn in flight, `no_such_stage` when the
|
||||
/// name is not running, `interrupt_refused` otherwise) and whose message
|
||||
/// names the stage and the reason as Petri spells it, for the web and
|
||||
/// the CLI to show.
|
||||
async fn interrupt(&self, stage: Option<&str>, text: Option<&str>, actor: &Principal) {
|
||||
/// the CLI to show; the same code and reason go back as the outcome.
|
||||
async fn interrupt(
|
||||
&self,
|
||||
stage: Option<&str>,
|
||||
text: Option<&str>,
|
||||
actor: &Principal,
|
||||
) -> WorkerControlOutcome {
|
||||
match self.controls.interrupt(stage, text).await {
|
||||
Ok(stage) => {
|
||||
info!(
|
||||
|
|
@ -382,18 +430,29 @@ impl PetriControls {
|
|||
actor = ?actor,
|
||||
"interrupt delivered"
|
||||
);
|
||||
WorkerControlOutcome::Delivered { stage: Some(stage) }
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(run_id = %self.run_id, stage, error = %error, "interrupt refused");
|
||||
self.notice(
|
||||
interrupt_refusal_code(&error),
|
||||
interrupt_refusal_message(stage, &error),
|
||||
self.refuse(
|
||||
error.code().unwrap_or("interrupt_refused"),
|
||||
refusal_message("Interrupt", stage, &error),
|
||||
)
|
||||
.await;
|
||||
.await
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A refused control: its `run.notice` on the run, and the refusal as
|
||||
/// the outcome to acknowledge.
|
||||
async fn refuse(&self, code: &str, message: String) -> WorkerControlOutcome {
|
||||
self.notice(code, message.clone()).await;
|
||||
WorkerControlOutcome::Refused {
|
||||
code: code.to_string(),
|
||||
message,
|
||||
}
|
||||
}
|
||||
|
||||
/// A `run.notice` record on the run, so a refused control is visible in
|
||||
/// the run's stream and not only in the worker's log.
|
||||
async fn notice(&self, code: &str, message: String) {
|
||||
|
|
@ -408,26 +467,13 @@ impl PetriControls {
|
|||
}
|
||||
}
|
||||
|
||||
/// The notice code of a refused interrupt.
|
||||
fn interrupt_refusal_code(error: &SteerError) -> &'static str {
|
||||
match error {
|
||||
SteerError::Control(ControlError::NoLiveTurn) => "no_live_turn",
|
||||
SteerError::Control(ControlError::NoSuchStage(_)) => "no_such_stage",
|
||||
SteerError::NoLiveAgent
|
||||
| SteerError::SeveralLiveAgents(_)
|
||||
| SteerError::Control(ControlError::NotLive | ControlError::Finished) => {
|
||||
"interrupt_refused"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The notice message of a refused interrupt: the stage it named, and the
|
||||
/// reason as Petri's `ControlError` (or the resolution's own refusal)
|
||||
/// The message of a refused control: the control, the stage it named, and
|
||||
/// the reason as Petri's `ControlError` (or the resolution's own refusal)
|
||||
/// spells it.
|
||||
fn interrupt_refusal_message(stage: Option<&str>, error: &SteerError) -> String {
|
||||
fn refusal_message(control: &str, stage: Option<&str>, error: &SteerError) -> String {
|
||||
match stage {
|
||||
Some(stage) => format!("Interrupt of stage `{stage}` refused: {error}"),
|
||||
None => format!("Interrupt refused: {error}"),
|
||||
Some(stage) => format!("{control} of stage `{stage}` refused: {error}"),
|
||||
None => format!("{control} refused: {error}"),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -505,6 +551,16 @@ fn test_checkpoint_gates() -> Option<PathBuf> {
|
|||
std::env::var_os(EnvVars::FABRO_TEST_CHECKPOINT_GATES).map(PathBuf::from)
|
||||
}
|
||||
|
||||
/// Whether a test asked this worker to acknowledge no control, so the
|
||||
/// server's wait for an answer runs out.
|
||||
#[expect(
|
||||
clippy::disallowed_methods,
|
||||
reason = "the mute is a test-only process-env facade the server forwards by name"
|
||||
)]
|
||||
fn test_control_acks_muted() -> bool {
|
||||
std::env::var_os(EnvVars::FABRO_TEST_CONTROL_ACKS_MUTED).is_some_and(|value| value == "1")
|
||||
}
|
||||
|
||||
/// Fabro's run tools for the run's agent sessions, when the run's settings
|
||||
/// enable them and the worker token carries the scope; `None` otherwise.
|
||||
/// The server issues the scope from the same setting, so the two agree
|
||||
|
|
|
|||
|
|
@ -9,8 +9,8 @@ use fabro_config::Storage;
|
|||
use fabro_interview::{
|
||||
AnswerSubmission, ControlInterviewer, WORKER_CONTROL_INVALID_CURSOR_REASON,
|
||||
WORKER_CONTROL_PONG_TIMEOUT_REASON, WORKER_CONTROL_WS_LIVENESS_TIMEOUT,
|
||||
WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlDeliveryFrame, WorkerControlEnvelope,
|
||||
WorkerControlMessage,
|
||||
WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlAck, WorkerControlDeliveryFrame,
|
||||
WorkerControlEnvelope, WorkerControlMessage,
|
||||
};
|
||||
use fabro_manifest::SuppliedWorkflowVersionPackager;
|
||||
use fabro_petri::controls::RunControls;
|
||||
|
|
@ -584,7 +584,7 @@ async fn handle_worker_control_socket(
|
|||
last_liveness = Instant::now();
|
||||
let frame = serde_json::from_str::<WorkerControlDeliveryFrame>(text.as_str())
|
||||
.map_err(|err| WorkerControlConnectError::Other(anyhow::Error::new(err)))?;
|
||||
apply_worker_control_delivery_frame(
|
||||
let applied = apply_worker_control_delivery_frame(
|
||||
interviewer,
|
||||
cancel_token,
|
||||
controls,
|
||||
|
|
@ -592,6 +592,17 @@ async fn handle_worker_control_socket(
|
|||
frame,
|
||||
)
|
||||
.await;
|
||||
if let Some(ack) = applied.ack {
|
||||
// The answer goes back over the stream the
|
||||
// control came in on; a stream that is gone
|
||||
// reconnects, and the server's wait runs out.
|
||||
let text = serde_json::to_string(&ack)
|
||||
.map_err(|err| WorkerControlConnectError::Other(anyhow::Error::new(err)))?;
|
||||
socket
|
||||
.send(WebSocketMessage::Text(text.into()))
|
||||
.await
|
||||
.map_err(|err| WorkerControlConnectError::Other(anyhow::Error::new(err)))?;
|
||||
}
|
||||
}
|
||||
Ok(WebSocketMessage::Ping(payload)) => {
|
||||
last_liveness = Instant::now();
|
||||
|
|
@ -621,43 +632,59 @@ async fn handle_worker_control_socket(
|
|||
}
|
||||
}
|
||||
|
||||
/// What a delivery frame came to: whether it was applied (a duplicate is
|
||||
/// not), and the acknowledgement to send back when the control asked for
|
||||
/// one.
|
||||
#[derive(Debug, Default, PartialEq, Eq)]
|
||||
struct AppliedDelivery {
|
||||
applied: bool,
|
||||
ack: Option<WorkerControlAck>,
|
||||
}
|
||||
|
||||
async fn apply_worker_control_delivery_frame(
|
||||
interviewer: &ControlInterviewer,
|
||||
cancel_token: &CancellationToken,
|
||||
controls: &WorkerControls,
|
||||
applied_ids: &mut AppliedWorkerControlDeliveryIds,
|
||||
frame: WorkerControlDeliveryFrame,
|
||||
) -> bool {
|
||||
) -> AppliedDelivery {
|
||||
// Duplicate ids cannot reach us under normal operation: the server replays
|
||||
// strictly after the last applied id. Guard against a server-side bug or
|
||||
// reconnect race by ignoring recently-applied delivery ids.
|
||||
if applied_ids.contains(&frame.id) {
|
||||
return false;
|
||||
return AppliedDelivery::default();
|
||||
}
|
||||
let frame_id = frame.id;
|
||||
apply_worker_control_message(interviewer, cancel_token, controls, frame.envelope).await;
|
||||
let ack =
|
||||
apply_worker_control_message(interviewer, cancel_token, controls, frame.envelope).await;
|
||||
applied_ids.record(frame_id);
|
||||
true
|
||||
AppliedDelivery { applied: true, ack }
|
||||
}
|
||||
|
||||
/// Apply the control. The acknowledgement to send back when the control
|
||||
/// carries a request id and has an outcome to report.
|
||||
async fn apply_worker_control_message(
|
||||
interviewer: &ControlInterviewer,
|
||||
cancel_token: &CancellationToken,
|
||||
controls: &WorkerControls,
|
||||
message: WorkerControlEnvelope,
|
||||
) {
|
||||
match message.message {
|
||||
) -> Option<WorkerControlAck> {
|
||||
let request_id = message.request_id().map(str::to_owned);
|
||||
let outcome = match message.message {
|
||||
WorkerControlMessage::InterviewAnswer { qid, answer, actor } => {
|
||||
let _ = interviewer
|
||||
.submit(&qid, AnswerSubmission::new(answer.into(), actor))
|
||||
.await;
|
||||
None
|
||||
}
|
||||
WorkerControlMessage::RunCancel => {
|
||||
cancel_token.cancel();
|
||||
interviewer.interrupt_all().await;
|
||||
None
|
||||
}
|
||||
other => controls.apply(other).await,
|
||||
}
|
||||
};
|
||||
Some(WorkerControlAck::new(request_id?, outcome?))
|
||||
}
|
||||
|
||||
pub(super) fn set_worker_title(run_id: &RunId, phase: WorkerTitlePhase) {
|
||||
|
|
@ -746,9 +773,12 @@ mod tests {
|
|||
use fabro_client::ServerTarget;
|
||||
use fabro_config::Storage;
|
||||
use fabro_interview::{
|
||||
AnswerValue, ControlInterviewer, Interviewer, Question, WorkerControlEnvelope,
|
||||
AnswerValue, ControlInterviewer, Interviewer, Question, WorkerControlAck,
|
||||
WorkerControlEnvelope, WorkerControlOutcome,
|
||||
};
|
||||
use fabro_types::{QuestionType, fixtures};
|
||||
use fabro_petri::test_support::MemoryPlatformRecords;
|
||||
use fabro_store::PlatformRecord;
|
||||
use fabro_types::{Principal, QuestionType, SystemActorKind, fixtures};
|
||||
use fabro_vault::{SecretType, Vault};
|
||||
use tokio::time;
|
||||
use tokio_tungstenite::tungstenite::protocol::{Message as TestWebSocketMessage, Role};
|
||||
|
|
@ -767,13 +797,35 @@ mod tests {
|
|||
/// A run's controls over records kept in memory: what the channel
|
||||
/// tests drive.
|
||||
fn test_controls() -> WorkerControls {
|
||||
test_controls_over(Arc::new(MemoryPlatformRecords::new()))
|
||||
}
|
||||
|
||||
fn test_controls_over(records: Arc<MemoryPlatformRecords>) -> WorkerControls {
|
||||
Arc::new(PetriControls::new(
|
||||
fixtures::RUN_1,
|
||||
RunControls::new(),
|
||||
Arc::new(fabro_petri::test_support::MemoryPlatformRecords::new()),
|
||||
records,
|
||||
))
|
||||
}
|
||||
|
||||
/// The `run.notice` records of the test run, as `(code, message)`.
|
||||
fn notices(records: &MemoryPlatformRecords) -> Vec<(String, String)> {
|
||||
records
|
||||
.records(&fixtures::RUN_1)
|
||||
.into_iter()
|
||||
.filter_map(|stored| match stored.record {
|
||||
PlatformRecord::RunNotice(notice) => Some((notice.code, notice.message)),
|
||||
_ => None,
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn engine_actor() -> Principal {
|
||||
Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clone_sandbox_credentials_are_required_for_clone_based_providers() {
|
||||
use fabro_types::SandboxProviderKind;
|
||||
|
|
@ -920,6 +972,135 @@ mod tests {
|
|||
assert!(!controls.controls().is_paused());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_refused_steer_is_acknowledged_with_the_notice_code() {
|
||||
let interviewer = Arc::new(ControlInterviewer::new());
|
||||
let cancel_token = CancellationToken::new();
|
||||
let records = Arc::new(MemoryPlatformRecords::new());
|
||||
let controls = test_controls_over(Arc::clone(&records));
|
||||
|
||||
// No live agent: refused under the control's own code.
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::steer("hurry up", None, engine_actor()).with_request_id("req-1"),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
ack,
|
||||
Some(WorkerControlAck::new(
|
||||
"req-1",
|
||||
WorkerControlOutcome::Refused {
|
||||
code: "steer_refused".to_string(),
|
||||
message: "Steer refused: Run has no active steerable agent session."
|
||||
.to_string(),
|
||||
}
|
||||
))
|
||||
);
|
||||
|
||||
// A stage that is not running: `no_such_stage`, for a steer as for
|
||||
// an interrupt.
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::steer("hurry up", Some("work".to_string()), engine_actor())
|
||||
.with_request_id("req-2"),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
ack,
|
||||
Some(WorkerControlAck::new(
|
||||
"req-2",
|
||||
WorkerControlOutcome::Refused {
|
||||
code: "no_such_stage".to_string(),
|
||||
message: "Steer of stage `work` refused: no stage named `work` is running"
|
||||
.to_string(),
|
||||
}
|
||||
))
|
||||
);
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::interrupt(Some("work".to_string()), engine_actor())
|
||||
.with_request_id("req-3"),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
ack,
|
||||
Some(WorkerControlAck::new(
|
||||
"req-3",
|
||||
WorkerControlOutcome::Refused {
|
||||
code: "no_such_stage".to_string(),
|
||||
message: "Interrupt of stage `work` refused: no stage named `work` is running"
|
||||
.to_string(),
|
||||
}
|
||||
))
|
||||
);
|
||||
|
||||
// Each refusal is also a notice on the run, under the same code.
|
||||
assert_eq!(
|
||||
notices(&records)
|
||||
.iter()
|
||||
.map(|(code, _)| code.as_str())
|
||||
.collect::<Vec<_>>(),
|
||||
["steer_refused", "no_such_stage", "no_such_stage"]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_control_without_a_request_id_is_not_acknowledged() {
|
||||
let interviewer = Arc::new(ControlInterviewer::new());
|
||||
let cancel_token = CancellationToken::new();
|
||||
let records = Arc::new(MemoryPlatformRecords::new());
|
||||
let controls = test_controls_over(Arc::clone(&records));
|
||||
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::steer("hurry up", None, engine_actor()),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(ack, None);
|
||||
assert_eq!(notices(&records).len(), 1, "the refusal is still a notice");
|
||||
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::pause_run(),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(ack, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn muted_acknowledgements_apply_the_control_and_answer_nothing() {
|
||||
let interviewer = Arc::new(ControlInterviewer::new());
|
||||
let cancel_token = CancellationToken::new();
|
||||
let records = Arc::new(MemoryPlatformRecords::new());
|
||||
let controls: WorkerControls = Arc::new(
|
||||
PetriControls::new(fixtures::RUN_1, RunControls::new(), records.clone())
|
||||
.with_muted_acks(true),
|
||||
);
|
||||
|
||||
let ack = apply_worker_control_message(
|
||||
&interviewer,
|
||||
&cancel_token,
|
||||
&controls,
|
||||
WorkerControlEnvelope::interrupt(None, engine_actor()).with_request_id("req-1"),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(ack, None);
|
||||
assert_eq!(notices(&records), [(
|
||||
"interrupt_refused".to_string(),
|
||||
"Interrupt refused: Run has no active steerable agent session.".to_string()
|
||||
)]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn duplicate_delivery_ids_are_not_applied_twice() {
|
||||
let interviewer = Arc::new(ControlInterviewer::new());
|
||||
|
|
@ -940,6 +1121,7 @@ mod tests {
|
|||
frame.clone(),
|
||||
)
|
||||
.await
|
||||
.applied
|
||||
);
|
||||
assert!(
|
||||
!apply_worker_control_delivery_frame(
|
||||
|
|
@ -950,6 +1132,7 @@ mod tests {
|
|||
frame,
|
||||
)
|
||||
.await
|
||||
.applied
|
||||
);
|
||||
|
||||
assert_eq!(applied_ids.last_applied_id(), Some("local:1"));
|
||||
|
|
|
|||
|
|
@ -1,11 +1,17 @@
|
|||
use anyhow::{Result, bail};
|
||||
use fabro_api::types::{RunControlAcknowledgement, RunControlOutcome};
|
||||
use tokio::io::{AsyncReadExt as _, stdin};
|
||||
use tracing::info;
|
||||
|
||||
use crate::args::SteerArgs;
|
||||
use crate::command_context::CommandContext;
|
||||
|
||||
/// Send a steer, or with `--interrupt` an interrupt carrying the text, and
|
||||
/// say what the worker made of it: delivered to its stage, or pending when
|
||||
/// the worker gave no answer in time. A refusal is the command's error,
|
||||
/// with the reason the worker gave.
|
||||
pub(crate) async fn run(args: SteerArgs, base_ctx: &CommandContext) -> Result<()> {
|
||||
let printer = base_ctx.printer();
|
||||
let ctx = base_ctx.with_target(&args.server)?;
|
||||
let client = ctx.server().await?;
|
||||
let run_id = client.resolve_run(&args.run).await?.id;
|
||||
|
|
@ -33,8 +39,48 @@ pub(crate) async fn run(args: SteerArgs, base_ctx: &CommandContext) -> Result<()
|
|||
.filter(|stage| !stage.is_empty())
|
||||
.map(str::to_owned);
|
||||
info!(run_id = %run_id, interrupt = args.interrupt, stage = ?stage, "Sending steer");
|
||||
client
|
||||
let acknowledgement = client
|
||||
.steer_run(&run_id, text, args.interrupt, stage)
|
||||
.await?;
|
||||
let control = if args.interrupt { "Interrupt" } else { "Steer" };
|
||||
fabro_util::printerr!(printer, "{}", describe(control, &acknowledgement));
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// One line on what became of the control.
|
||||
fn describe(control: &str, acknowledgement: &RunControlAcknowledgement) -> String {
|
||||
match (&acknowledgement.outcome, acknowledgement.stage.as_deref()) {
|
||||
(RunControlOutcome::Delivered, Some(stage)) => {
|
||||
format!("{control} delivered to stage {stage}.")
|
||||
}
|
||||
(RunControlOutcome::Delivered, None) => format!("{control} delivered."),
|
||||
(RunControlOutcome::Pending, _) => format!(
|
||||
"{control} forwarded; the worker has not answered yet. The run's events say what \
|
||||
became of it."
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn the_outcome_is_described_in_one_line() {
|
||||
assert_eq!(
|
||||
describe("Steer", &RunControlAcknowledgement {
|
||||
outcome: RunControlOutcome::Delivered,
|
||||
stage: Some("work@1".to_string()),
|
||||
}),
|
||||
"Steer delivered to stage work@1."
|
||||
);
|
||||
assert_eq!(
|
||||
describe("Interrupt", &RunControlAcknowledgement {
|
||||
outcome: RunControlOutcome::Pending,
|
||||
stage: None,
|
||||
}),
|
||||
"Interrupt forwarded; the worker has not answered yet. The run's events say what \
|
||||
became of it."
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1585,12 +1585,16 @@ async fn mcp_interact_actions_resolve_selector_and_call_expected_endpoints() {
|
|||
when.method(POST)
|
||||
.path(format!("/api/v1/runs/{run_id}/steer"))
|
||||
.json_body(serde_json::json!({ "text": "continue", "interrupt": true }));
|
||||
then.status(202);
|
||||
then.status(202)
|
||||
.header("Content-Type", "application/json")
|
||||
.json_body(serde_json::json!({ "outcome": "delivered", "stage": "code@1" }));
|
||||
});
|
||||
let interrupt = server.mock(|when, then| {
|
||||
when.method(POST)
|
||||
.path(format!("/api/v1/runs/{run_id}/interrupt"));
|
||||
then.status(202);
|
||||
then.status(202)
|
||||
.header("Content-Type", "application/json")
|
||||
.json_body(serde_json::json!({ "outcome": "delivered", "stage": "code@1" }));
|
||||
});
|
||||
let cancel = server.mock(|when, then| {
|
||||
when.method(POST)
|
||||
|
|
|
|||
|
|
@ -89,6 +89,8 @@ pub(super) struct RunningServer {
|
|||
pub(super) api_base_url: String,
|
||||
/// The checkpoint gate directory the server forwards to its workers.
|
||||
gates_dir: PathBuf,
|
||||
/// Extra environment on the server process, kept for a relaunch.
|
||||
env: Vec<(String, String)>,
|
||||
}
|
||||
|
||||
impl RunningServer {
|
||||
|
|
@ -101,6 +103,16 @@ impl RunningServer {
|
|||
/// in the vault before the first launch, so the server and its workers
|
||||
/// see them from the start.
|
||||
pub(super) async fn start_with(settings: &str, secrets: &[(&str, &str)]) -> Self {
|
||||
Self::start_with_env(settings, secrets, &[]).await
|
||||
}
|
||||
|
||||
/// `start_with`, plus `env` on the server process: the test hooks the
|
||||
/// server forwards to its workers by name.
|
||||
pub(super) async fn start_with_env(
|
||||
settings: &str,
|
||||
secrets: &[(&str, &str)],
|
||||
env: &[(&str, &str)],
|
||||
) -> Self {
|
||||
let home_root = tempfile::tempdir_in("/tmp").expect("home tempdir");
|
||||
let storage_root = isolated_storage_dir();
|
||||
let storage_dir = storage_root.path().join("storage");
|
||||
|
|
@ -139,6 +151,10 @@ impl RunningServer {
|
|||
port,
|
||||
api_base_url: format!("http://127.0.0.1:{port}"),
|
||||
gates_dir,
|
||||
env: env
|
||||
.iter()
|
||||
.map(|(name, value)| ((*name).to_string(), (*value).to_string()))
|
||||
.collect(),
|
||||
};
|
||||
server.launch().await;
|
||||
server
|
||||
|
|
@ -158,6 +174,9 @@ impl RunningServer {
|
|||
self.home_root.path().join("fabro-home"),
|
||||
);
|
||||
cmd.env(EnvVars::FABRO_TEST_CHECKPOINT_GATES, &self.gates_dir);
|
||||
for (name, value) in &self.env {
|
||||
cmd.env(name, value);
|
||||
}
|
||||
cmd.args(["server", "start", "--foreground"])
|
||||
.arg("--storage-dir")
|
||||
.arg(&self.storage_dir)
|
||||
|
|
|
|||
|
|
@ -9,6 +9,11 @@
|
|||
//! interrupt of a gate stage is refused with `no_live_turn`; a run paused
|
||||
//! when its server and worker die resumes paused and goes on once unpaused.
|
||||
//!
|
||||
//! The worker answers each steer and interrupt over its control stream,
|
||||
//! and the endpoint's response is that answer: `202` with `delivered` and
|
||||
//! the stage, `409` with the refusal's code, or `202` with `pending` when
|
||||
//! the worker never answers (a test hook mutes the worker's answers).
|
||||
//!
|
||||
//! 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
|
||||
//! sandbox-driver host plugin, so the tests skip, and say why, when the
|
||||
|
|
@ -108,14 +113,40 @@ async fn steer(server: &RunningServer, run_id: &str, text: &str) {
|
|||
steer_stage(server, run_id, text, None).await;
|
||||
}
|
||||
|
||||
/// `POST /runs/{id}/steer` naming `stage`, or no stage.
|
||||
/// `POST /runs/{id}/steer` naming `stage`, or no stage: the worker
|
||||
/// answers `delivered`, to the stage it steered.
|
||||
async fn steer_stage(server: &RunningServer, run_id: &str, text: &str, stage: Option<&str>) {
|
||||
let (status, body) = steer_request(server, run_id, text, stage).await;
|
||||
assert_eq!(status, 202, "steer: {body}");
|
||||
assert_eq!(body["outcome"], "delivered", "steer: {body}");
|
||||
let delivered = body["stage"].as_str().unwrap_or_default();
|
||||
match stage {
|
||||
Some(stage) => assert_eq!(delivered, stage, "steer: {body}"),
|
||||
None => assert!(!delivered.is_empty(), "steer: {body}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// `POST /runs/{id}/steer` naming `stage`, or no stage; the status and
|
||||
/// body, for a steer the worker may refuse.
|
||||
async fn steer_request(
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
text: &str,
|
||||
stage: Option<&str>,
|
||||
) -> (u16, Value) {
|
||||
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}");
|
||||
control(server, run_id, "steer", Some(body)).await
|
||||
}
|
||||
|
||||
/// A refused control: 409, with the refusal's code and message in the
|
||||
/// error entry.
|
||||
fn assert_refused(status: u16, body: &Value, code: &str, message: &str) {
|
||||
assert_eq!(status, 409, "{body}");
|
||||
assert_eq!(body["errors"][0]["code"], code, "{body}");
|
||||
assert_eq!(body["errors"][0]["detail"], message, "{body}");
|
||||
}
|
||||
|
||||
/// `fabro steer <run> --stage <stage> <text>` against the server.
|
||||
|
|
@ -138,6 +169,11 @@ fn steer_by_cli(
|
|||
String::from_utf8_lossy(&output.stdout),
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
assert!(
|
||||
stderr.contains(&format!("Steer delivered to stage {stage}.")),
|
||||
"fabro steer reports the worker's answer\nstderr:\n{stderr}"
|
||||
);
|
||||
}
|
||||
|
||||
/// `fabro steer --interrupt <run> <text>` against the server: the stage's
|
||||
|
|
@ -160,6 +196,37 @@ fn interrupt_by_cli(
|
|||
String::from_utf8_lossy(&output.stdout),
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
assert!(
|
||||
stderr.contains("Interrupt delivered to stage work@1."),
|
||||
"fabro steer --interrupt reports the worker's answer\nstderr:\n{stderr}"
|
||||
);
|
||||
}
|
||||
|
||||
/// `fabro steer --interrupt <run> --stage <stage> <text>` for a stage the
|
||||
/// worker refuses: the command fails and prints the refusal.
|
||||
fn interrupt_refused_by_cli(
|
||||
context: &fabro_test::TestContext,
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
stage: &str,
|
||||
refusal: &str,
|
||||
) {
|
||||
let output = context
|
||||
.command()
|
||||
.args(["steer", "--server", &server.target(), run_id])
|
||||
.args(["--interrupt", "--stage", stage, "Stop."])
|
||||
.output()
|
||||
.expect("the steer command executes");
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
assert!(
|
||||
!output.status.success(),
|
||||
"fabro steer --interrupt of `{stage}` succeeded\nstderr:\n{stderr}"
|
||||
);
|
||||
assert!(
|
||||
stderr.contains(refusal),
|
||||
"fabro steer --interrupt prints the refusal `{refusal}`\nstderr:\n{stderr}"
|
||||
);
|
||||
}
|
||||
|
||||
/// The stream's `run.notice` records, as `(code, message)`, in order.
|
||||
|
|
@ -463,6 +530,15 @@ async fn a_steer_reaches_the_agent_stage_on_the_twin() {
|
|||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
wait_until_gate_is_polled(&gate);
|
||||
eprintln!("run {run_id}: the agent's tool is waiting on the gate");
|
||||
// A stage that is not running: refused in the response, and on the
|
||||
// stream as a notice under the same code.
|
||||
let (status, body) = steer_request(&server, &run_id, "Steer nobody.", Some("nope")).await;
|
||||
assert_refused(
|
||||
status,
|
||||
&body,
|
||||
"no_such_stage",
|
||||
"Steer of stage `nope` refused: no stage named `nope` is running",
|
||||
);
|
||||
steer(&server, &run_id, STEER).await;
|
||||
wait_for_stream_count(&server, &run_id, "control.requested", 1).await;
|
||||
eprintln!("run {run_id}: the steer is recorded");
|
||||
|
|
@ -478,6 +554,14 @@ async fn a_steer_reaches_the_agent_stage_on_the_twin() {
|
|||
server.stderr_text()
|
||||
);
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
assert_eq!(
|
||||
notices(&items),
|
||||
[(
|
||||
"no_such_stage".to_string(),
|
||||
"Steer of stage `nope` refused: no stage named `nope` is running".to_string()
|
||||
)],
|
||||
"the refused steer is the one notice"
|
||||
);
|
||||
|
||||
let delivery = items
|
||||
.iter()
|
||||
|
|
@ -609,8 +693,16 @@ async fn two_live_agent_stages_are_steered_apart_by_their_labels() {
|
|||
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;
|
||||
// Unnamed, the steer has two candidates and is refused with both
|
||||
// named: in the response, and on the stream.
|
||||
let (status, body) = steer_request(&server, &run_id, "Steer nobody.", None).await;
|
||||
assert_eq!(status, 409, "{body}");
|
||||
assert_eq!(body["errors"][0]["code"], "steer_refused", "{body}");
|
||||
let detail = body["errors"][0]["detail"].as_str().unwrap_or_default();
|
||||
assert!(
|
||||
detail.contains("a@1") && detail.contains("b@1"),
|
||||
"the refusal names both live stages: {body}"
|
||||
);
|
||||
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)
|
||||
|
|
@ -790,22 +882,10 @@ async fn an_interrupt_ends_the_turn_and_its_text_is_the_next_input() {
|
|||
server.shutdown();
|
||||
}
|
||||
|
||||
/// An interrupt of a stage with no model turn to stop: the gate the run is
|
||||
/// blocked on, named by its node, is refused by Petri with `no_live_turn`;
|
||||
/// unnamed, with no agent stage live, the worker refuses it with
|
||||
/// `interrupt_refused`. Both refusals are `run.notice` records on the
|
||||
/// stream naming the stage and Petri's reason, nothing is delivered, and
|
||||
/// the gate's question is untouched: its answer routes the run to its end.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn an_interrupt_of_a_gate_stage_is_refused_with_no_live_turn() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let server = RunningServer::start().await;
|
||||
let marker = context.temp_dir.join("yes.marker");
|
||||
let workspace = write_petri_workflow(
|
||||
&context,
|
||||
/// A workflow of one human gate: `yes` leaves `marker`.
|
||||
fn gate_workspace(context: &fabro_test::TestContext, marker: &Path) -> PathBuf {
|
||||
write_petri_workflow(
|
||||
context,
|
||||
&format!(
|
||||
"digraph Gate {{\n graph [goal=\"Ask before running\"]\n start [shape=Mdiamond]\n \
|
||||
exit [shape=Msquare]\n gate [shape=hexagon, label=\"Go?\", \
|
||||
|
|
@ -814,7 +894,31 @@ async fn an_interrupt_of_a_gate_stage_is_refused_with_no_live_turn() {
|
|||
No\"]\n yes -> exit\n}}\n",
|
||||
marker = marker.display()
|
||||
),
|
||||
);
|
||||
)
|
||||
}
|
||||
|
||||
const GATE_INTERRUPT_REFUSAL: &str =
|
||||
"Interrupt of stage `gate` refused: the stage has no model turn to interrupt";
|
||||
const UNNAMED_INTERRUPT_REFUSAL: &str =
|
||||
"Interrupt refused: Run has no active steerable agent session.";
|
||||
|
||||
/// An interrupt of a stage with no model turn to stop: the gate the run is
|
||||
/// blocked on, named by its node, is refused by Petri with `no_live_turn`;
|
||||
/// unnamed, with no agent stage live, the worker refuses it with
|
||||
/// `interrupt_refused`. Each refusal is the endpoint's own answer, a 409
|
||||
/// with the code and the reason, and `fabro steer` prints it; each is also
|
||||
/// a `run.notice` record on the stream naming the stage and Petri's
|
||||
/// reason. Nothing is delivered, and the gate's question is untouched: its
|
||||
/// answer routes the run to its end.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn an_interrupt_of_a_gate_stage_is_refused_with_no_live_turn() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let server = RunningServer::start().await;
|
||||
let marker = context.temp_dir.join("yes.marker");
|
||||
let workspace = gate_workspace(&context, &marker);
|
||||
let run_id = run_detached_with(&context, &server, &workspace, &[]);
|
||||
|
||||
let pending = wait_for_questions(&server, &run_id, 1).await;
|
||||
|
|
@ -829,25 +933,34 @@ async fn an_interrupt_of_a_gate_stage_is_refused_with_no_live_turn() {
|
|||
Some(json!({ "stage": "gate" })),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(status, 202, "interrupt: {body}");
|
||||
assert_refused(status, &body, "no_live_turn", GATE_INTERRUPT_REFUSAL);
|
||||
let (status, body) = control(&server, &run_id, "interrupt", None).await;
|
||||
assert_eq!(status, 202, "interrupt: {body}");
|
||||
let names = wait_for_stream_count(&server, &run_id, "run.notice", 2).await;
|
||||
assert_refused(
|
||||
status,
|
||||
&body,
|
||||
"interrupt_refused",
|
||||
UNNAMED_INTERRUPT_REFUSAL,
|
||||
);
|
||||
interrupt_refused_by_cli(&context, &server, &run_id, "gate", GATE_INTERRUPT_REFUSAL);
|
||||
let names = wait_for_stream_count(&server, &run_id, "run.notice", 3).await;
|
||||
assert_eq!(count_of(&names, "control.requested"), 0, "{names:?}");
|
||||
let refused = notices(&run_stream(&server, &run_id).await);
|
||||
assert_eq!(refused.len(), 2, "{refused:?}");
|
||||
// The notice names the stage and carries Petri's reason as it spells
|
||||
// it, so the web and the CLI can show both.
|
||||
assert_eq!(refused[0].0, "no_live_turn", "{refused:?}");
|
||||
assert_eq!(
|
||||
refused[0].1,
|
||||
"Interrupt of stage `gate` refused: the stage has no model turn to interrupt"
|
||||
);
|
||||
assert_eq!(refused[1].0, "interrupt_refused", "{refused:?}");
|
||||
assert_eq!(
|
||||
refused[1].1,
|
||||
"Interrupt refused: Run has no active steerable agent session."
|
||||
);
|
||||
assert_eq!(refused, [
|
||||
(
|
||||
"no_live_turn".to_string(),
|
||||
GATE_INTERRUPT_REFUSAL.to_string()
|
||||
),
|
||||
(
|
||||
"interrupt_refused".to_string(),
|
||||
UNNAMED_INTERRUPT_REFUSAL.to_string()
|
||||
),
|
||||
(
|
||||
"no_live_turn".to_string(),
|
||||
GATE_INTERRUPT_REFUSAL.to_string()
|
||||
),
|
||||
]);
|
||||
assert_eq!(run_status(&server, &run_id).await, "blocked");
|
||||
|
||||
answer(&server, &run_id, &question_id, json!({ "kind": "yes" })).await;
|
||||
|
|
@ -874,6 +987,60 @@ async fn an_interrupt_of_a_gate_stage_is_refused_with_no_live_turn() {
|
|||
server.shutdown();
|
||||
}
|
||||
|
||||
/// A worker that never answers a control: the endpoint waits its bound
|
||||
/// (5 s) and answers `202` with `pending`; the control was still applied,
|
||||
/// so its refusal is on the stream as a notice. The worker's answers are
|
||||
/// muted through the server's test hook, forwarded to the worker by name.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn a_control_the_worker_never_answers_is_pending() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let server =
|
||||
RunningServer::start_with_env("", &[], &[(EnvVars::FABRO_TEST_CONTROL_ACKS_MUTED, "1")])
|
||||
.await;
|
||||
let marker = context.temp_dir.join("yes.marker");
|
||||
let workspace = gate_workspace(&context, &marker);
|
||||
let run_id = run_detached_with(&context, &server, &workspace, &[]);
|
||||
|
||||
let pending = wait_for_questions(&server, &run_id, 1).await;
|
||||
let question_id = pending[0]["id"].as_str().expect("an id").to_string();
|
||||
eprintln!("run {run_id}: the gate is asking; interrupting it with the answers muted");
|
||||
|
||||
let asked = Instant::now();
|
||||
let (status, body) = control(
|
||||
&server,
|
||||
&run_id,
|
||||
"interrupt",
|
||||
Some(json!({ "stage": "gate" })),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(status, 202, "interrupt: {body}");
|
||||
assert_eq!(body, json!({ "outcome": "pending" }));
|
||||
assert!(
|
||||
asked.elapsed() >= Duration::from_secs(4),
|
||||
"the endpoint waited its bound for the answer: {:?}",
|
||||
asked.elapsed()
|
||||
);
|
||||
let refused = notices(&run_stream(&server, &run_id).await);
|
||||
assert_eq!(refused, [(
|
||||
"no_live_turn".to_string(),
|
||||
GATE_INTERRUPT_REFUSAL.to_string()
|
||||
)]);
|
||||
|
||||
answer(&server, &run_id, &question_id, json!({ "kind": "yes" })).await;
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"server stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert!(marker.exists(), "the answer routed the gate");
|
||||
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
|
||||
|
|
|
|||
|
|
@ -55,11 +55,13 @@ use fabro_db::DbPool;
|
|||
use fabro_environment::EnvironmentStore;
|
||||
use fabro_interview::{
|
||||
Answer, AnswerSubmission, ControlInterviewer, Question, WorkerControlEnvelope,
|
||||
WorkerControlOutcome,
|
||||
};
|
||||
use fabro_llm::credentials::CredentialProvider;
|
||||
use fabro_llm::lithos_catalog::Catalog;
|
||||
use fabro_llm::{ClientOptions, FabroClient};
|
||||
use fabro_mcp_store::McpServerStore;
|
||||
use fabro_petri::controls::{RunControls, SteerError};
|
||||
use fabro_petri::projector::Projector;
|
||||
use fabro_redact::redact_jsonl_line;
|
||||
use fabro_sandbox::details::sandbox_details;
|
||||
|
|
@ -151,7 +153,10 @@ use crate::request_id::{self, RequestId};
|
|||
use crate::run_files::{FilesInFlight, new_files_in_flight};
|
||||
use crate::server_secrets::ServerSecrets;
|
||||
use crate::spawn_env::apply_render_graph_env;
|
||||
use crate::worker_control::{LocalWorkerControlBus, WorkerControlBus, WorkerControlBusError};
|
||||
use crate::worker_control::{
|
||||
LocalWorkerControlBus, WORKER_CONTROL_ACK_WAIT, WorkerControlAcks, WorkerControlBus,
|
||||
WorkerControlBusError,
|
||||
};
|
||||
use crate::worker_runtime::{
|
||||
LocalWorkerRuntime, WorkerExit, WorkerLaunchSpec, WorkerRef, WorkerRuntime,
|
||||
};
|
||||
|
|
@ -326,9 +331,14 @@ enum RunAnswerTransport {
|
|||
Worker {
|
||||
run_id: RunId,
|
||||
bus: Arc<dyn WorkerControlBus>,
|
||||
/// Where the worker's answers to steers and interrupts arrive.
|
||||
acks: Arc<WorkerControlAcks>,
|
||||
},
|
||||
InProcess {
|
||||
interviewer: Arc<ControlInterviewer>,
|
||||
/// The run's controls, answered in place: the in-process run has
|
||||
/// no worker to forward a steer or an interrupt to.
|
||||
controls: RunControls,
|
||||
},
|
||||
}
|
||||
|
||||
|
|
@ -338,6 +348,46 @@ enum AnswerTransportError {
|
|||
Timeout,
|
||||
}
|
||||
|
||||
/// What a steer or an interrupt came to, as far as the caller is told.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
enum RunControlAnswer {
|
||||
/// The worker delivered it, to the stage named when it had one.
|
||||
Delivered { stage: Option<String> },
|
||||
/// The worker (or Petri) refused it: the code and the reason.
|
||||
Refused { code: String, message: String },
|
||||
/// The control was forwarded but no answer arrived within
|
||||
/// [`WORKER_CONTROL_ACK_WAIT`]; the run's stream says what became of it.
|
||||
Pending,
|
||||
}
|
||||
|
||||
impl From<WorkerControlOutcome> for RunControlAnswer {
|
||||
fn from(outcome: WorkerControlOutcome) -> Self {
|
||||
match outcome {
|
||||
WorkerControlOutcome::Delivered { stage } => Self::Delivered { stage },
|
||||
WorkerControlOutcome::Refused { code, message } => Self::Refused { code, message },
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl RunControlAnswer {
|
||||
/// The answer to a control the run's own controls settled in place:
|
||||
/// `control` names it in the refusal's message, `refused_code` is the
|
||||
/// code of a refusal whose reason has none of its own.
|
||||
fn from_controls(
|
||||
control: &str,
|
||||
refused_code: &str,
|
||||
result: Result<String, SteerError>,
|
||||
) -> Self {
|
||||
match result {
|
||||
Ok(stage) => Self::Delivered { stage: Some(stage) },
|
||||
Err(error) => Self::Refused {
|
||||
code: error.code().unwrap_or(refused_code).to_string(),
|
||||
message: format!("{control} refused: {error}"),
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl RunAnswerTransport {
|
||||
async fn publish_worker_control(
|
||||
run_id: RunId,
|
||||
|
|
@ -359,13 +409,34 @@ impl RunAnswerTransport {
|
|||
}
|
||||
}
|
||||
|
||||
/// Publish a control the worker answers: registered with the run's
|
||||
/// acknowledgements first, so the answer has a waiter, then published
|
||||
/// with the request id, then waited for. `Pending` when no answer
|
||||
/// arrives within [`WORKER_CONTROL_ACK_WAIT`].
|
||||
async fn publish_answered_control(
|
||||
run_id: RunId,
|
||||
bus: &Arc<dyn WorkerControlBus>,
|
||||
acks: &WorkerControlAcks,
|
||||
message: WorkerControlEnvelope,
|
||||
) -> Result<RunControlAnswer, AnswerTransportError> {
|
||||
let pending = acks.register(run_id);
|
||||
let message = message.with_request_id(pending.request_id.clone());
|
||||
Self::publish_worker_control(run_id, bus, message)
|
||||
.await
|
||||
.map_err(|err| Self::answer_error_from_bus(&err))?;
|
||||
Ok(acks
|
||||
.wait(pending)
|
||||
.await
|
||||
.map_or(RunControlAnswer::Pending, RunControlAnswer::from))
|
||||
}
|
||||
|
||||
async fn submit(
|
||||
&self,
|
||||
qid: &str,
|
||||
submission: AnswerSubmission,
|
||||
) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, .. } => {
|
||||
let message = WorkerControlEnvelope::interview_answer(qid.to_string(), submission);
|
||||
Self::publish_worker_control(*run_id, bus, message)
|
||||
.await
|
||||
|
|
@ -380,7 +451,7 @@ impl RunAnswerTransport {
|
|||
|
||||
async fn cancel_run(&self) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, .. } => {
|
||||
let message = WorkerControlEnvelope::cancel_run();
|
||||
Self::publish_worker_control(*run_id, bus, message)
|
||||
.await
|
||||
|
|
@ -394,51 +465,55 @@ impl RunAnswerTransport {
|
|||
}
|
||||
|
||||
/// 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.
|
||||
/// one live agent stage, and wait for its answer. The in-process path
|
||||
/// answers from the run's own controls at once.
|
||||
async fn steer(
|
||||
&self,
|
||||
text: String,
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
) -> Result<(), AnswerTransportError> {
|
||||
) -> Result<RunControlAnswer, AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, acks } => {
|
||||
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))
|
||||
Self::publish_answered_control(*run_id, bus, acks, message).await
|
||||
}
|
||||
Self::InProcess { .. } => Err(AnswerTransportError::Closed),
|
||||
Self::InProcess { controls, .. } => Ok(RunControlAnswer::from_controls(
|
||||
"Steer",
|
||||
"steer_refused",
|
||||
controls.steer(stage.as_deref(), &text).await,
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
/// Forward an interrupt to the worker, for the stage it names or the
|
||||
/// run's one live agent stage; `text`, when given, is the stage's next
|
||||
/// input.
|
||||
/// run's one live agent stage, and wait for its answer; `text`, when
|
||||
/// given, is the stage's next input.
|
||||
async fn interrupt(
|
||||
&self,
|
||||
stage: Option<String>,
|
||||
text: Option<String>,
|
||||
actor: Principal,
|
||||
) -> Result<(), AnswerTransportError> {
|
||||
) -> Result<RunControlAnswer, AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, acks } => {
|
||||
let message = match text {
|
||||
Some(text) => WorkerControlEnvelope::interrupt_then_steer(text, stage, actor),
|
||||
None => WorkerControlEnvelope::interrupt(stage, actor),
|
||||
};
|
||||
Self::publish_worker_control(*run_id, bus, message)
|
||||
.await
|
||||
.map_err(|err| Self::answer_error_from_bus(&err))
|
||||
Self::publish_answered_control(*run_id, bus, acks, message).await
|
||||
}
|
||||
Self::InProcess { .. } => Err(AnswerTransportError::Closed),
|
||||
Self::InProcess { controls, .. } => Ok(RunControlAnswer::from_controls(
|
||||
"Interrupt",
|
||||
"interrupt_refused",
|
||||
controls.interrupt(stage.as_deref(), text.as_deref()).await,
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn pause_run(&self) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, .. } => {
|
||||
let message = WorkerControlEnvelope::pause_run();
|
||||
Self::publish_worker_control(*run_id, bus, message)
|
||||
.await
|
||||
|
|
@ -450,7 +525,7 @@ impl RunAnswerTransport {
|
|||
|
||||
async fn unpause_run(&self) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
Self::Worker { run_id, bus, .. } => {
|
||||
let message = WorkerControlEnvelope::unpause_run();
|
||||
Self::publish_worker_control(*run_id, bus, message)
|
||||
.await
|
||||
|
|
@ -1036,6 +1111,8 @@ pub struct AppState {
|
|||
resource_sampler: resource_sampler::ResourceSampler,
|
||||
max_concurrent_runs: usize,
|
||||
pub(crate) worker_control_bus: Arc<dyn WorkerControlBus>,
|
||||
/// The steers and interrupts awaiting their worker's answer.
|
||||
pub(crate) worker_control_acks: Arc<WorkerControlAcks>,
|
||||
pub(crate) worker_runtime: Arc<dyn WorkerRuntime>,
|
||||
/// The Petri runs held open for workers over the API.
|
||||
pub(crate) petri_runs: PetriRuns,
|
||||
|
|
@ -2574,6 +2651,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
resource_sampler: resource_sampler::ResourceSampler::new(),
|
||||
max_concurrent_runs,
|
||||
worker_control_bus,
|
||||
worker_control_acks: Arc::new(WorkerControlAcks::new(WORKER_CONTROL_ACK_WAIT)),
|
||||
worker_runtime,
|
||||
petri_runs,
|
||||
petri_projector,
|
||||
|
|
@ -3025,6 +3103,7 @@ fn clear_live_run_state(run: &mut ManagedRun) {
|
|||
|
||||
fn cleanup_worker_control_bus_for_run(state: &AppState, run_id: RunId) {
|
||||
let bus = Arc::clone(&state.worker_control_bus);
|
||||
state.worker_control_acks.forget_run(run_id);
|
||||
tokio::spawn(async move {
|
||||
bus.cleanup_run(run_id).await;
|
||||
});
|
||||
|
|
@ -4004,6 +4083,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
|
|||
managed_run.answer_transport = Some(RunAnswerTransport::Worker {
|
||||
run_id,
|
||||
bus: Arc::clone(&state.worker_control_bus),
|
||||
acks: Arc::clone(&state.worker_control_acks),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,11 +5,15 @@ use axum::extract::State;
|
|||
use axum::http::StatusCode;
|
||||
use axum::response::{IntoResponse, Response};
|
||||
use axum::routing::post;
|
||||
use fabro_api::types::{InterruptRunRequest, SteerRunRequest};
|
||||
use fabro_api::types::{
|
||||
InterruptRunRequest, RunControlAcknowledgement, RunControlOutcome, SteerRunRequest,
|
||||
};
|
||||
use fabro_types::Principal;
|
||||
use fabro_workflow::run_status::RunStatus;
|
||||
|
||||
use super::super::{AnswerTransportError, AppState, durable_run_status, reject_if_archived};
|
||||
use super::super::{
|
||||
AnswerTransportError, AppState, RunControlAnswer, durable_run_status, reject_if_archived,
|
||||
};
|
||||
use crate::error::ApiError;
|
||||
use crate::principal_middleware::RequireRunManagementTarget;
|
||||
|
||||
|
|
@ -20,8 +24,11 @@ pub(super) fn routes() -> axum::Router<Arc<AppState>> {
|
|||
}
|
||||
|
||||
/// A control forwarded to the run's worker. The worker resolves the stage
|
||||
/// and delivers the control to Petri; what it cannot deliver it refuses on
|
||||
/// the run's stream as a `run.notice` whose code says why.
|
||||
/// and delivers the control to Petri, and answers over its control stream:
|
||||
/// the answer is this endpoint's response, 202 once delivered, 409 with
|
||||
/// the refusal's code when refused, and 202 `pending` when no answer came
|
||||
/// within the wait. What the worker cannot deliver it also refuses on the
|
||||
/// run's stream as a `run.notice` whose code says why.
|
||||
enum RunControlRequest {
|
||||
/// Guidance for a live agent stage's session, run as a follow-up turn.
|
||||
Steer {
|
||||
|
|
@ -202,7 +209,11 @@ async fn control_run(
|
|||
};
|
||||
|
||||
match result {
|
||||
Ok(()) => StatusCode::ACCEPTED.into_response(),
|
||||
Ok(RunControlAnswer::Delivered { stage }) => accepted(RunControlOutcome::Delivered, stage),
|
||||
Ok(RunControlAnswer::Pending) => accepted(RunControlOutcome::Pending, None),
|
||||
Ok(RunControlAnswer::Refused { code, message }) => {
|
||||
ApiError::with_code(StatusCode::CONFLICT, message, code).into_response()
|
||||
}
|
||||
Err(AnswerTransportError::Timeout) => ApiError::with_code(
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
"Worker control channel timed out.",
|
||||
|
|
@ -218,6 +229,16 @@ async fn control_run(
|
|||
}
|
||||
}
|
||||
|
||||
/// The 202 of a control the worker took: delivered to `stage`, or still
|
||||
/// pending its answer.
|
||||
fn accepted(outcome: RunControlOutcome, stage: Option<String>) -> Response {
|
||||
(
|
||||
StatusCode::ACCEPTED,
|
||||
Json(RunControlAcknowledgement { outcome, stage }),
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
|
||||
/// The 409 code of a control the run's status refuses.
|
||||
fn not_controllable_code(control: &RunControlRequest) -> &'static str {
|
||||
match control {
|
||||
|
|
|
|||
|
|
@ -5,9 +5,10 @@ use axum::extract::ws::{
|
|||
};
|
||||
use fabro_interview::{
|
||||
WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON,
|
||||
WORKER_CONTROL_WS_LIVENESS_TIMEOUT, WORKER_CONTROL_WS_PING_INTERVAL,
|
||||
WORKER_CONTROL_WS_LIVENESS_TIMEOUT, WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlAck,
|
||||
WorkerControlDeliveryFrame,
|
||||
};
|
||||
use fabro_types::RunId;
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use tokio::time::{self, Instant, MissedTickBehavior};
|
||||
|
||||
|
|
@ -15,7 +16,9 @@ use super::super::{
|
|||
ApiError, AppState, IntoResponse, Query, RequireWorkerRunScoped, Response, Router, State,
|
||||
StatusCode, get,
|
||||
};
|
||||
use crate::worker_control::{WorkerControlBusError, WorkerControlCursor, WorkerControlReceiver};
|
||||
use crate::worker_control::{
|
||||
WorkerControlAcks, WorkerControlBusError, WorkerControlCursor, WorkerControlReceiver,
|
||||
};
|
||||
|
||||
#[derive(Debug, serde::Deserialize)]
|
||||
struct WorkerControlStreamQuery {
|
||||
|
|
@ -65,7 +68,8 @@ async fn worker_control_stream(
|
|||
Err(err) => return worker_control_bus_error_response(&err),
|
||||
};
|
||||
|
||||
ws.on_upgrade(move |socket| worker_control_websocket(socket, receiver))
|
||||
let acks = Arc::clone(&state.worker_control_acks);
|
||||
ws.on_upgrade(move |socket| worker_control_websocket(socket, receiver, id, acks))
|
||||
}
|
||||
|
||||
fn worker_control_bus_error_response(err: &WorkerControlBusError) -> Response {
|
||||
|
|
@ -79,7 +83,15 @@ fn worker_control_bus_error_response(err: &WorkerControlBusError) -> Response {
|
|||
ApiError::new(status, err.to_string()).into_response()
|
||||
}
|
||||
|
||||
async fn worker_control_websocket(socket: WebSocket, mut receiver: WorkerControlReceiver) {
|
||||
/// The stream to one worker: deliveries go out as text frames; the text
|
||||
/// frames that come back are the worker's answers to the controls that
|
||||
/// asked for one, and settle the callers waiting on them.
|
||||
async fn worker_control_websocket(
|
||||
socket: WebSocket,
|
||||
mut receiver: WorkerControlReceiver,
|
||||
run_id: RunId,
|
||||
acks: Arc<WorkerControlAcks>,
|
||||
) {
|
||||
let (mut sender, mut receiver_ws) = socket.split();
|
||||
let mut ping_interval = time::interval(WORKER_CONTROL_WS_PING_INTERVAL);
|
||||
ping_interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
|
||||
|
|
@ -132,7 +144,11 @@ async fn worker_control_websocket(socket: WebSocket, mut receiver: WorkerControl
|
|||
return;
|
||||
}
|
||||
}
|
||||
Ok(WsMessage::Pong(_) | WsMessage::Text(_) | WsMessage::Binary(_)) => {
|
||||
Ok(WsMessage::Text(text)) => {
|
||||
last_liveness = Instant::now();
|
||||
receive_worker_control_ack(run_id, &acks, text.as_str());
|
||||
}
|
||||
Ok(WsMessage::Pong(_) | WsMessage::Binary(_)) => {
|
||||
last_liveness = Instant::now();
|
||||
}
|
||||
Ok(WsMessage::Close(_)) | Err(_) => return,
|
||||
|
|
@ -154,6 +170,31 @@ async fn worker_control_websocket(socket: WebSocket, mut receiver: WorkerControl
|
|||
}
|
||||
}
|
||||
|
||||
/// A text frame from the worker: an acknowledgement of a control, handed
|
||||
/// to the caller waiting on it. Anything else is logged and dropped; the
|
||||
/// stream stays up.
|
||||
fn receive_worker_control_ack(run_id: RunId, acks: &WorkerControlAcks, text: &str) {
|
||||
match serde_json::from_str::<WorkerControlAck>(text) {
|
||||
Ok(ack) => {
|
||||
let request_id = ack.request_id.clone();
|
||||
if !acks.resolve(run_id, ack) {
|
||||
tracing::debug!(
|
||||
run_id = %run_id,
|
||||
request_id,
|
||||
"worker control acknowledgement had no waiting caller"
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(error) => {
|
||||
tracing::debug!(
|
||||
run_id = %run_id,
|
||||
error = %error,
|
||||
"worker control stream carried a text frame that is not an acknowledgement"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn invalid_cursor_close_message() -> WsMessage {
|
||||
WsMessage::Close(Some(CloseFrame {
|
||||
code: close_code::POLICY,
|
||||
|
|
|
|||
|
|
@ -283,6 +283,8 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
|
|||
// records above already moved the live status to Running; a run that
|
||||
// ended meanwhile (cancelled while starting) takes no transport.
|
||||
let interviewer = Arc::new(ControlInterviewer::new());
|
||||
// The steer and interrupt endpoints reach these controls in place.
|
||||
let controls = RunControls::new();
|
||||
{
|
||||
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
||||
if let Some(managed_run) = runs
|
||||
|
|
@ -291,6 +293,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
|
|||
{
|
||||
managed_run.answer_transport = Some(RunAnswerTransport::InProcess {
|
||||
interviewer: Arc::clone(&interviewer),
|
||||
controls: controls.clone(),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -319,9 +322,10 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
|
|||
runtime: runtime_spec(&state, &eligible, dry_run),
|
||||
provider: run_state.spec.settings.run.environment.provider.clone(),
|
||||
cancel,
|
||||
// The in-process test path drives no pause or steer: the server's
|
||||
// transports for those name the worker.
|
||||
controls: RunControls::new(),
|
||||
// The in-process test path drives no pause: the server's transport
|
||||
// for it names the worker. A steer or an interrupt is answered in
|
||||
// place.
|
||||
controls,
|
||||
interviewer: Arc::new(petri_interviewer),
|
||||
observers,
|
||||
secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))),
|
||||
|
|
|
|||
|
|
@ -12,7 +12,8 @@ use fabro_automation::AutomationId;
|
|||
use fabro_config::bind::Bind;
|
||||
use fabro_config::{LlmLayer, RunLayer, ServerSettingsBuilder};
|
||||
use fabro_interview::{
|
||||
AnswerValue, WorkerControlDeliveryFrame, WorkerControlEnvelope, WorkerControlMessage,
|
||||
AnswerValue, WorkerControlAck, WorkerControlDeliveryFrame, WorkerControlEnvelope,
|
||||
WorkerControlMessage, WorkerControlOutcome,
|
||||
};
|
||||
use fabro_llm::lithos_catalog::Catalog;
|
||||
use fabro_store::platform_records::{
|
||||
|
|
@ -48,7 +49,8 @@ use crate::github_webhooks::compute_signature;
|
|||
use crate::jwt_auth::{AuthMode, ConfiguredAuth};
|
||||
use crate::test_support::*;
|
||||
use crate::worker_control::{
|
||||
LocalWorkerControlBus, WorkerControlBus, WorkerControlCursor, WorkerControlReceiver,
|
||||
LocalWorkerControlBus, WorkerControlAcks, WorkerControlBus, WorkerControlCursor,
|
||||
WorkerControlReceiver,
|
||||
};
|
||||
use crate::worker_runtime::{
|
||||
LocalWorkerRuntime, StartedWorker, WorkerLaunchSpec, WorkerRef, WorkerRuntime,
|
||||
|
|
@ -931,6 +933,51 @@ async fn worker_control_stream_after_subscription_delivers_only_later_frames() {
|
|||
assert_eq!(frame.envelope, expected);
|
||||
}
|
||||
|
||||
/// An acknowledgement the worker sends over its control stream settles the
|
||||
/// caller waiting on that request; one for another run's request does
|
||||
/// not.
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn worker_control_stream_acknowledgements_settle_the_waiting_caller() {
|
||||
let (state, app) = jwt_auth_app();
|
||||
let user_bearer = issue_test_user_jwt();
|
||||
let run_id = create_run_with_bearer(&app, &user_bearer).await;
|
||||
let worker_bearer = issue_test_worker_token(&run_id);
|
||||
let server = WorkerControlWsTestServer::spawn(app).await;
|
||||
let mut socket = connect_worker_control_ws(&server, run_id, &worker_bearer, None).await;
|
||||
|
||||
let pending = state.worker_control_acks.register(run_id);
|
||||
let foreign = state.worker_control_acks.register(fixtures::RUN_2);
|
||||
for request_id in [&foreign.request_id, &pending.request_id] {
|
||||
let ack = WorkerControlAck::new(request_id, WorkerControlOutcome::Refused {
|
||||
code: "no_live_turn".to_string(),
|
||||
message: "the stage has no model turn to interrupt".to_string(),
|
||||
});
|
||||
futures_util::SinkExt::send(
|
||||
&mut socket,
|
||||
WebSocketMessage::Text(serde_json::to_string(&ack).unwrap().into()),
|
||||
)
|
||||
.await
|
||||
.expect("the acknowledgement sends");
|
||||
}
|
||||
|
||||
let outcome = tokio::time::timeout(
|
||||
Duration::from_secs(2),
|
||||
state.worker_control_acks.wait(pending),
|
||||
)
|
||||
.await
|
||||
.expect("the caller is answered");
|
||||
assert_eq!(
|
||||
outcome,
|
||||
Some(WorkerControlOutcome::Refused {
|
||||
code: "no_live_turn".to_string(),
|
||||
message: "the stage has no model turn to interrupt".to_string(),
|
||||
})
|
||||
);
|
||||
// The other run's request was not this worker's to answer.
|
||||
assert_eq!(state.worker_control_acks.outstanding(), 1);
|
||||
drop(foreign);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn worker_control_stream_invalid_cursor_is_http_gone_before_upgrade() {
|
||||
let (_state, app) = jwt_auth_app();
|
||||
|
|
@ -2742,9 +2789,26 @@ impl WorkerRuntime for RecordingWorkerRuntime {
|
|||
}
|
||||
}
|
||||
|
||||
/// How long a test transport waits for a worker's answer: short, so a
|
||||
/// test whose worker never answers sees `pending` at once.
|
||||
const TEST_WORKER_CONTROL_ACK_WAIT: Duration = Duration::from_millis(100);
|
||||
|
||||
async fn worker_transport_with_receiver(
|
||||
run_id: RunId,
|
||||
) -> (RunAnswerTransport, WorkerControlReceiver) {
|
||||
let (transport, receiver, _) = worker_transport_with_acks(run_id).await;
|
||||
(transport, receiver)
|
||||
}
|
||||
|
||||
/// A worker transport over a private bus, with the acknowledgements a
|
||||
/// test answers through.
|
||||
async fn worker_transport_with_acks(
|
||||
run_id: RunId,
|
||||
) -> (
|
||||
RunAnswerTransport,
|
||||
WorkerControlReceiver,
|
||||
StdArc<WorkerControlAcks>,
|
||||
) {
|
||||
let bus = StdArc::new(LocalWorkerControlBus::new());
|
||||
let receiver = bus
|
||||
.subscribe(run_id, WorkerControlCursor::Start)
|
||||
|
|
@ -2753,8 +2817,50 @@ async fn worker_transport_with_receiver(
|
|||
// Ensure the subscription task is waiting before the test publishes.
|
||||
tokio::task::yield_now().await;
|
||||
let bus: StdArc<dyn WorkerControlBus> = bus;
|
||||
let transport = RunAnswerTransport::Worker { run_id, bus };
|
||||
(transport, receiver)
|
||||
let acks = StdArc::new(WorkerControlAcks::new(TEST_WORKER_CONTROL_ACK_WAIT));
|
||||
let transport = RunAnswerTransport::Worker {
|
||||
run_id,
|
||||
bus,
|
||||
acks: StdArc::clone(&acks),
|
||||
};
|
||||
(transport, receiver, acks)
|
||||
}
|
||||
|
||||
/// A worker that answers every control carrying a request id with
|
||||
/// `outcome`, as the real worker answers over its control stream. The
|
||||
/// deliveries it read, for the test to inspect.
|
||||
fn answering_worker(
|
||||
run_id: RunId,
|
||||
mut receiver: WorkerControlReceiver,
|
||||
acks: StdArc<WorkerControlAcks>,
|
||||
outcome: WorkerControlOutcome,
|
||||
) -> tokio::sync::mpsc::UnboundedReceiver<WorkerControlEnvelope> {
|
||||
let (seen_tx, seen_rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
tokio::spawn(async move {
|
||||
while let Some(Ok(delivery)) = receiver.recv().await {
|
||||
if let Some(request_id) = delivery.envelope.request_id() {
|
||||
acks.resolve(run_id, WorkerControlAck::new(request_id, outcome.clone()));
|
||||
}
|
||||
if seen_tx.send(delivery.envelope).is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
});
|
||||
seen_rx
|
||||
}
|
||||
|
||||
/// The published envelope without the request id the transport added, so
|
||||
/// a test can compare it with the constructor's.
|
||||
fn without_request_id(mut envelope: WorkerControlEnvelope) -> WorkerControlEnvelope {
|
||||
match &mut envelope.message {
|
||||
WorkerControlMessage::Steer { request_id, .. }
|
||||
| WorkerControlMessage::Interrupt { request_id, .. }
|
||||
| WorkerControlMessage::InterruptThenSteer { request_id, .. } => {
|
||||
*request_id = None;
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
envelope
|
||||
}
|
||||
|
||||
async fn recv_worker_control_envelope(
|
||||
|
|
@ -2787,17 +2893,75 @@ async fn worker_answer_transport_steer_publishes_plain_steer_message() {
|
|||
system_kind: SystemActorKind::Engine,
|
||||
};
|
||||
|
||||
transport
|
||||
// Nobody answers: the steer is pending once the wait runs out.
|
||||
let answer = transport
|
||||
.steer("try again".to_string(), None, actor.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(answer, RunControlAnswer::Pending);
|
||||
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(envelope.request_id().is_some(), "{envelope:?}");
|
||||
assert_eq!(
|
||||
recv_worker_control_envelope(&mut control_rx).await,
|
||||
without_request_id(envelope),
|
||||
WorkerControlEnvelope::steer("try again", None, actor)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_answer_transport_steer_returns_the_workers_answer() {
|
||||
let run_id = fixtures::RUN_1;
|
||||
let (transport, control_rx, acks) = worker_transport_with_acks(run_id).await;
|
||||
let actor = Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
};
|
||||
let mut seen = answering_worker(run_id, control_rx, acks, WorkerControlOutcome::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
});
|
||||
|
||||
let answer = transport
|
||||
.steer("try again".to_string(), None, actor)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(answer, RunControlAnswer::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
});
|
||||
let envelope = seen.recv().await.expect("the worker read the steer");
|
||||
assert!(matches!(
|
||||
envelope.message,
|
||||
WorkerControlMessage::Steer { ref text, .. } if text == "try again"
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn in_process_transport_answers_a_steer_and_an_interrupt_in_place() {
|
||||
let transport = RunAnswerTransport::InProcess {
|
||||
interviewer: StdArc::new(ControlInterviewer::new()),
|
||||
controls: RunControls::new(),
|
||||
};
|
||||
let actor = Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
};
|
||||
|
||||
// The run has no live agent stage: refused at once, no worker asked.
|
||||
let answer = transport
|
||||
.steer("try again".to_string(), None, actor.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(answer, RunControlAnswer::Refused {
|
||||
code: "steer_refused".to_string(),
|
||||
message: "Steer refused: Run has no active steerable agent session.".to_string(),
|
||||
});
|
||||
let answer = transport
|
||||
.interrupt(Some("work".to_string()), None, actor)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(answer, RunControlAnswer::Refused {
|
||||
code: "no_such_stage".to_string(),
|
||||
message: "Interrupt refused: no stage named `work` is running".to_string(),
|
||||
});
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_answer_transport_pause_and_unpause_publish_control_messages() {
|
||||
let (transport, mut control_rx) = worker_transport_with_receiver(fixtures::RUN_1).await;
|
||||
|
|
@ -8601,6 +8765,122 @@ async fn interrupt_of_an_unknown_run_returns_not_found() {
|
|||
assert_status!(response, StatusCode::NOT_FOUND).await;
|
||||
}
|
||||
|
||||
/// A steer the worker delivers: 202 with the outcome and the stage.
|
||||
#[tokio::test]
|
||||
async fn steer_delivered_by_the_worker_returns_accepted_with_its_stage() {
|
||||
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, acks) = worker_transport_with_acks(run_id).await;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, Some(transport));
|
||||
let _seen = answering_worker(run_id, control_rx, acks, WorkerControlOutcome::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
});
|
||||
|
||||
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"}"#))
|
||||
.unwrap();
|
||||
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::ACCEPTED);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(
|
||||
body,
|
||||
serde_json::json!({ "outcome": "delivered", "stage": "work@1" })
|
||||
);
|
||||
}
|
||||
|
||||
/// A control the worker refuses: 409 with the refusal's code and reason,
|
||||
/// for each code the worker can answer with.
|
||||
#[tokio::test]
|
||||
async fn control_refused_by_the_worker_returns_conflict_with_its_code() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let refusals = [
|
||||
(
|
||||
"steer",
|
||||
r#"{"text":"try again"}"#,
|
||||
"steer_refused",
|
||||
"Steer refused: Run has no active steerable agent session.",
|
||||
),
|
||||
(
|
||||
"steer",
|
||||
r#"{"text":"try again","stage":"nope"}"#,
|
||||
"no_such_stage",
|
||||
"Steer of stage `nope` refused: no stage named `nope` is running",
|
||||
),
|
||||
(
|
||||
"interrupt",
|
||||
r#"{"stage":"gate"}"#,
|
||||
"no_live_turn",
|
||||
"Interrupt of stage `gate` refused: the stage has no model turn to interrupt",
|
||||
),
|
||||
(
|
||||
"interrupt",
|
||||
r#"{"stage":"nope"}"#,
|
||||
"no_such_stage",
|
||||
"Interrupt of stage `nope` refused: no stage named `nope` is running",
|
||||
),
|
||||
(
|
||||
"interrupt",
|
||||
"{}",
|
||||
"interrupt_refused",
|
||||
"Interrupt refused: Run has no active steerable agent session.",
|
||||
),
|
||||
];
|
||||
for (index, (action, body, code, message)) in refusals.into_iter().enumerate() {
|
||||
let run_id = RunId::new();
|
||||
let (transport, control_rx, acks) = worker_transport_with_acks(run_id).await;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, Some(transport));
|
||||
let _seen = answering_worker(run_id, control_rx, acks, WorkerControlOutcome::Refused {
|
||||
code: code.to_string(),
|
||||
message: message.to_string(),
|
||||
});
|
||||
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/{action}")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(body))
|
||||
.unwrap();
|
||||
let response = app.clone().oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::CONFLICT, "refusal {index}");
|
||||
let response_body = body_json(response.into_body()).await;
|
||||
assert_eq!(response_body["errors"][0]["code"], code, "{response_body}");
|
||||
assert_eq!(
|
||||
response_body["errors"][0]["detail"], message,
|
||||
"{response_body}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// A worker that never answers: 202 `pending` once the wait runs out,
|
||||
/// with the control still forwarded.
|
||||
#[tokio::test]
|
||||
async fn interrupt_unanswered_by_the_worker_returns_accepted_pending() {
|
||||
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}/interrupt")))
|
||||
.body(Body::empty())
|
||||
.unwrap();
|
||||
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::ACCEPTED);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(body, serde_json::json!({ "outcome": "pending" }));
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(envelope.request_id().is_some(), "{envelope:?}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_without_a_worker_channel_returns_unavailable() {
|
||||
let state = test_app_state();
|
||||
|
|
|
|||
|
|
@ -80,6 +80,9 @@ const WORKER_ENV_ALLOWLIST: &[&str] = &[
|
|||
// A test's checkpoint gates: the worker's hooks hold at a named point
|
||||
// until the test releases them, so a crash can be placed there.
|
||||
EnvVars::FABRO_TEST_CHECKPOINT_GATES,
|
||||
// A test's mute on the worker's control acknowledgements, so the
|
||||
// server's wait for one runs out.
|
||||
EnvVars::FABRO_TEST_CONTROL_ACKS_MUTED,
|
||||
];
|
||||
|
||||
const RENDER_GRAPH_ENV_ALLOWLIST: &[&str] = &[EnvVars::PATH, EnvVars::HOME, EnvVars::TMPDIR];
|
||||
|
|
|
|||
179
lib/apps/fabro-server/src/worker_control/acks.rs
Normal file
179
lib/apps/fabro-server/src/worker_control/acks.rs
Normal file
|
|
@ -0,0 +1,179 @@
|
|||
//! The answers workers give to controls: a steer or an interrupt goes out
|
||||
//! over the control bus with a request id, and the worker acknowledges it
|
||||
//! over the control stream it arrived on ([`WorkerControlAck`]). The
|
||||
//! registry pairs each outstanding request with the caller waiting for its
|
||||
//! answer, for a bounded time; a request nobody answers in time is
|
||||
//! forgotten, and its caller told the answer is still pending.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::{Mutex, MutexGuard, PoisonError};
|
||||
use std::time::Duration;
|
||||
|
||||
use fabro_interview::{WorkerControlAck, WorkerControlOutcome};
|
||||
use fabro_types::RunId;
|
||||
use tokio::sync::oneshot;
|
||||
use tokio::time::timeout;
|
||||
|
||||
/// How long a control waits for the worker's answer before the caller is
|
||||
/// told it is pending.
|
||||
pub(crate) const WORKER_CONTROL_ACK_WAIT: Duration = Duration::from_secs(5);
|
||||
|
||||
/// One outstanding control: the run it went to, and who waits on it.
|
||||
struct PendingControl {
|
||||
run_id: RunId,
|
||||
answer: oneshot::Sender<WorkerControlOutcome>,
|
||||
}
|
||||
|
||||
/// The controls awaiting a worker's answer, by request id.
|
||||
pub(crate) struct WorkerControlAcks {
|
||||
pending: Mutex<HashMap<String, PendingControl>>,
|
||||
wait: Duration,
|
||||
}
|
||||
|
||||
/// A registered request: its id, to send with the control, and the answer
|
||||
/// to wait on.
|
||||
pub(crate) struct PendingAck {
|
||||
pub(crate) request_id: String,
|
||||
answer: oneshot::Receiver<WorkerControlOutcome>,
|
||||
}
|
||||
|
||||
impl WorkerControlAcks {
|
||||
/// A registry whose waits last `wait`.
|
||||
#[must_use]
|
||||
pub(crate) fn new(wait: Duration) -> Self {
|
||||
Self {
|
||||
pending: Mutex::new(HashMap::new()),
|
||||
wait,
|
||||
}
|
||||
}
|
||||
|
||||
/// Register a control about to go to `run_id`'s worker: the id to send
|
||||
/// with it, and the answer to wait on. Registered before the control
|
||||
/// is published, so an answer cannot arrive before anyone waits for it.
|
||||
pub(crate) fn register(&self, run_id: RunId) -> PendingAck {
|
||||
let request_id = ulid::Ulid::new().to_string();
|
||||
let (answer, receiver) = oneshot::channel();
|
||||
self.pending
|
||||
.lock()
|
||||
.unwrap_or_else(PoisonError::into_inner)
|
||||
.insert(request_id.clone(), PendingControl { run_id, answer });
|
||||
PendingAck {
|
||||
request_id,
|
||||
answer: receiver,
|
||||
}
|
||||
}
|
||||
|
||||
/// Wait for the answer to `pending`, at most the registry's wait:
|
||||
/// `None` when none arrives in time, after which the request is
|
||||
/// forgotten and a late answer is dropped.
|
||||
pub(crate) async fn wait(&self, pending: PendingAck) -> Option<WorkerControlOutcome> {
|
||||
let outcome = timeout(self.wait, pending.answer).await;
|
||||
match outcome {
|
||||
Ok(Ok(outcome)) => Some(outcome),
|
||||
// The sender was dropped: the request was forgotten, or the
|
||||
// registry was.
|
||||
Ok(Err(_)) => None,
|
||||
Err(_elapsed) => {
|
||||
self.lock().remove(&pending.request_id);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Answer the request `ack` names, when it is outstanding and went to
|
||||
/// `run_id`'s worker: an answer from another run's worker, or to a
|
||||
/// request already forgotten, is dropped. Whether an answer was
|
||||
/// delivered to a waiting caller.
|
||||
pub(crate) fn resolve(&self, run_id: RunId, ack: WorkerControlAck) -> bool {
|
||||
let mut pending = self.lock();
|
||||
let Some(entry) = pending.get(&ack.request_id) else {
|
||||
return false;
|
||||
};
|
||||
if entry.run_id != run_id {
|
||||
return false;
|
||||
}
|
||||
let Some(entry) = pending.remove(&ack.request_id) else {
|
||||
return false;
|
||||
};
|
||||
drop(pending);
|
||||
entry.answer.send(ack.outcome).is_ok()
|
||||
}
|
||||
|
||||
/// Forget every request outstanding on `run_id`: its callers are told
|
||||
/// the answer is pending at once.
|
||||
pub(crate) fn forget_run(&self, run_id: RunId) {
|
||||
self.lock().retain(|_, entry| entry.run_id != run_id);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn outstanding(&self) -> usize {
|
||||
self.lock().len()
|
||||
}
|
||||
|
||||
fn lock(&self) -> MutexGuard<'_, HashMap<String, PendingControl>> {
|
||||
self.pending.lock().unwrap_or_else(PoisonError::into_inner)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use fabro_types::fixtures;
|
||||
|
||||
use super::*;
|
||||
|
||||
fn delivered(request_id: &str) -> WorkerControlAck {
|
||||
WorkerControlAck::new(request_id, WorkerControlOutcome::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
})
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_answer_reaches_the_caller_waiting_on_its_request() {
|
||||
let acks = WorkerControlAcks::new(Duration::from_secs(1));
|
||||
let pending = acks.register(fixtures::RUN_1);
|
||||
assert!(acks.resolve(fixtures::RUN_1, delivered(&pending.request_id)));
|
||||
assert_eq!(
|
||||
acks.wait(pending).await,
|
||||
Some(WorkerControlOutcome::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
})
|
||||
);
|
||||
assert_eq!(acks.outstanding(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_request_nobody_answers_in_time_is_pending_and_forgotten() {
|
||||
let acks = WorkerControlAcks::new(Duration::from_millis(20));
|
||||
let pending = acks.register(fixtures::RUN_1);
|
||||
let request_id = pending.request_id.clone();
|
||||
assert_eq!(acks.wait(pending).await, None);
|
||||
assert_eq!(acks.outstanding(), 0);
|
||||
assert!(!acks.resolve(fixtures::RUN_1, delivered(&request_id)));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn another_runs_worker_cannot_answer_a_request() {
|
||||
let acks = WorkerControlAcks::new(Duration::from_millis(20));
|
||||
let pending = acks.register(fixtures::RUN_1);
|
||||
assert!(!acks.resolve(fixtures::RUN_2, delivered(&pending.request_id)));
|
||||
assert_eq!(acks.outstanding(), 1);
|
||||
assert_eq!(acks.wait(pending).await, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_unknown_request_id_is_dropped() {
|
||||
let acks = WorkerControlAcks::new(Duration::from_secs(1));
|
||||
assert!(!acks.resolve(fixtures::RUN_1, delivered("nobody")));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn forgetting_a_run_answers_its_callers_with_pending_at_once() {
|
||||
let acks = WorkerControlAcks::new(Duration::from_secs(30));
|
||||
let pending = acks.register(fixtures::RUN_1);
|
||||
let other = acks.register(fixtures::RUN_2);
|
||||
acks.forget_run(fixtures::RUN_1);
|
||||
assert_eq!(acks.outstanding(), 1);
|
||||
assert_eq!(acks.wait(pending).await, None);
|
||||
assert!(acks.resolve(fixtures::RUN_2, delivered(&other.request_id)));
|
||||
}
|
||||
}
|
||||
|
|
@ -1,6 +1,8 @@
|
|||
mod acks;
|
||||
mod bus;
|
||||
mod local;
|
||||
|
||||
pub(crate) use acks::{WORKER_CONTROL_ACK_WAIT, WorkerControlAcks};
|
||||
pub(crate) use bus::{
|
||||
WorkerControlBus, WorkerControlBusError, WorkerControlCursor, WorkerControlDelivery,
|
||||
WorkerControlMessageId, WorkerControlReceiver,
|
||||
|
|
|
|||
|
|
@ -78,6 +78,7 @@ impl WorkerControlEnvelope {
|
|||
text: text.into(),
|
||||
stage,
|
||||
actor,
|
||||
request_id: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
|
@ -86,7 +87,11 @@ impl WorkerControlEnvelope {
|
|||
pub fn interrupt(stage: Option<String>, actor: Principal) -> Self {
|
||||
Self {
|
||||
v: WORKER_CONTROL_PROTOCOL_VERSION,
|
||||
message: WorkerControlMessage::Interrupt { stage, actor },
|
||||
message: WorkerControlMessage::Interrupt {
|
||||
stage,
|
||||
actor,
|
||||
request_id: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -102,10 +107,41 @@ impl WorkerControlEnvelope {
|
|||
text: text.into(),
|
||||
stage,
|
||||
actor,
|
||||
request_id: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// The same control, asking the worker to acknowledge it: the worker
|
||||
/// answers a control that carries a request id with a
|
||||
/// [`WorkerControlAck`] naming the id, over the control stream it
|
||||
/// arrived on. Only a steer or an interrupt carries one; on any other
|
||||
/// control the id is dropped.
|
||||
#[must_use]
|
||||
pub fn with_request_id(mut self, id: impl Into<String>) -> Self {
|
||||
match &mut self.message {
|
||||
WorkerControlMessage::Steer { request_id, .. }
|
||||
| WorkerControlMessage::Interrupt { request_id, .. }
|
||||
| WorkerControlMessage::InterruptThenSteer { request_id, .. } => {
|
||||
*request_id = Some(id.into());
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
self
|
||||
}
|
||||
|
||||
/// The request id the control carries, when its sender asked for an
|
||||
/// acknowledgement.
|
||||
#[must_use]
|
||||
pub fn request_id(&self) -> Option<&str> {
|
||||
match &self.message {
|
||||
WorkerControlMessage::Steer { request_id, .. }
|
||||
| WorkerControlMessage::Interrupt { request_id, .. }
|
||||
| WorkerControlMessage::InterruptThenSteer { request_id, .. } => request_id.as_deref(),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn start_pair(
|
||||
run_id: RunId,
|
||||
|
|
@ -170,28 +206,37 @@ pub enum WorkerControlMessage {
|
|||
RunUnpause,
|
||||
#[serde(rename = "run.steer")]
|
||||
Steer {
|
||||
text: String,
|
||||
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<String>,
|
||||
actor: Principal,
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
/// Set when the sender waits for a [`WorkerControlAck`].
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
request_id: Option<String>,
|
||||
},
|
||||
#[serde(rename = "run.interrupt")]
|
||||
Interrupt {
|
||||
/// The stage whose model turn to stop (`node@visit`, or the node
|
||||
/// name); `None` interrupts the run's one live agent stage.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
/// Set when the sender waits for a [`WorkerControlAck`].
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
request_id: Option<String>,
|
||||
},
|
||||
#[serde(rename = "run.interrupt_then_steer")]
|
||||
InterruptThenSteer {
|
||||
text: String,
|
||||
text: String,
|
||||
/// The stage to interrupt and steer, as for `Interrupt`.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
/// Set when the sender waits for a [`WorkerControlAck`].
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
request_id: Option<String>,
|
||||
},
|
||||
#[serde(rename = "pair.start")]
|
||||
PairStart {
|
||||
|
|
@ -219,6 +264,43 @@ pub struct WorkerControlDeliveryFrame {
|
|||
pub envelope: WorkerControlEnvelope,
|
||||
}
|
||||
|
||||
/// The worker's answer to a control that carried a request id, sent as a
|
||||
/// text frame over the control stream the control arrived on: what
|
||||
/// became of it, so the server can answer the caller in its own response.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct WorkerControlAck {
|
||||
pub v: u8,
|
||||
pub request_id: String,
|
||||
pub outcome: WorkerControlOutcome,
|
||||
}
|
||||
|
||||
impl WorkerControlAck {
|
||||
#[must_use]
|
||||
pub fn new(request_id: impl Into<String>, outcome: WorkerControlOutcome) -> Self {
|
||||
Self {
|
||||
v: WORKER_CONTROL_PROTOCOL_VERSION,
|
||||
request_id: request_id.into(),
|
||||
outcome,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// What became of a control at the worker.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(tag = "kind", rename_all = "snake_case")]
|
||||
pub enum WorkerControlOutcome {
|
||||
/// The control reached its stage: the label of the stage it went to,
|
||||
/// when the control had one.
|
||||
Delivered {
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
stage: Option<String>,
|
||||
},
|
||||
/// The worker or Petri refused the control: the code the refusal is
|
||||
/// known by (`no_live_turn`, `no_such_stage`, `steer_refused`,
|
||||
/// `interrupt_refused`) and the reason as the worker spells it.
|
||||
Refused { code: String, message: String },
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(tag = "kind", rename_all = "snake_case")]
|
||||
pub enum WorkerControlAnswer {
|
||||
|
|
@ -422,6 +504,69 @@ mod tests {
|
|||
assert_eq!(parsed, end);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_request_id_rides_a_steer_or_an_interrupt_and_nothing_else() {
|
||||
let actor = Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
};
|
||||
let steer =
|
||||
WorkerControlEnvelope::steer("try again", None, actor.clone()).with_request_id("req-1");
|
||||
assert_eq!(steer.request_id(), Some("req-1"));
|
||||
let json = serde_json::to_string(&steer).unwrap();
|
||||
assert_eq!(
|
||||
json,
|
||||
r#"{"v":1,"type":"run.steer","text":"try again","actor":{"kind":"system","system_kind":"engine"},"request_id":"req-1"}"#
|
||||
);
|
||||
let parsed: WorkerControlEnvelope = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(parsed, steer);
|
||||
|
||||
let interrupt = WorkerControlEnvelope::interrupt(Some("code@2".to_string()), actor.clone())
|
||||
.with_request_id("req-2");
|
||||
assert_eq!(interrupt.request_id(), Some("req-2"));
|
||||
let interrupt_then_steer = WorkerControlEnvelope::interrupt_then_steer("stop", None, actor)
|
||||
.with_request_id("req-3");
|
||||
assert_eq!(interrupt_then_steer.request_id(), Some("req-3"));
|
||||
|
||||
let pause = WorkerControlEnvelope::pause_run().with_request_id("req-4");
|
||||
assert_eq!(pause.request_id(), None);
|
||||
assert_eq!(pause, WorkerControlEnvelope::pause_run());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_steer_without_a_request_id_still_parses() {
|
||||
let parsed: WorkerControlEnvelope = serde_json::from_str(
|
||||
r#"{"v":1,"type":"run.steer","text":"try again","actor":{"kind":"system","system_kind":"engine"}}"#,
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(parsed.request_id(), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn control_acks_round_trip_through_json() {
|
||||
let delivered = WorkerControlAck::new("req-1", WorkerControlOutcome::Delivered {
|
||||
stage: Some("work@1".to_string()),
|
||||
});
|
||||
let json = serde_json::to_string(&delivered).unwrap();
|
||||
assert_eq!(
|
||||
json,
|
||||
r#"{"v":1,"request_id":"req-1","outcome":{"kind":"delivered","stage":"work@1"}}"#
|
||||
);
|
||||
let parsed: WorkerControlAck = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(parsed, delivered);
|
||||
|
||||
let refused = WorkerControlAck::new("req-2", WorkerControlOutcome::Refused {
|
||||
code: "no_live_turn".to_string(),
|
||||
message: "the stage has no model turn to interrupt".to_string(),
|
||||
});
|
||||
let json = serde_json::to_string(&refused).unwrap();
|
||||
assert_eq!(
|
||||
json,
|
||||
r#"{"v":1,"request_id":"req-2","outcome":{"kind":"refused","code":"no_live_turn","message":"the stage has no model turn to interrupt"}}"#
|
||||
);
|
||||
let parsed: WorkerControlAck = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(parsed, refused);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delivery_frame_round_trips_through_json() {
|
||||
let frame = WorkerControlDeliveryFrame {
|
||||
|
|
|
|||
|
|
@ -227,8 +227,8 @@ pub use control::{ControlInterviewer, SubmitError};
|
|||
pub use control_protocol::{
|
||||
WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON,
|
||||
WORKER_CONTROL_PROTOCOL_VERSION, WORKER_CONTROL_WS_LIVENESS_TIMEOUT,
|
||||
WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlAnswer, WorkerControlDeliveryFrame,
|
||||
WorkerControlEnvelope, WorkerControlMessage,
|
||||
WORKER_CONTROL_WS_PING_INTERVAL, WorkerControlAck, WorkerControlAnswer,
|
||||
WorkerControlDeliveryFrame, WorkerControlEnvelope, WorkerControlMessage, WorkerControlOutcome,
|
||||
};
|
||||
pub use queue::QueueInterviewer;
|
||||
pub use recording::RecordingInterviewer;
|
||||
|
|
|
|||
|
|
@ -81,6 +81,25 @@ pub enum SteerError {
|
|||
Control(#[from] ControlError),
|
||||
}
|
||||
|
||||
impl SteerError {
|
||||
/// The code Fabro knows the refusal by, when the reason has one of its
|
||||
/// own: `no_live_turn` for a stage with no model turn in flight,
|
||||
/// `no_such_stage` for a name that is not running. `None` for a
|
||||
/// refusal named only by the control it refused (`steer_refused`,
|
||||
/// `interrupt_refused`): no live agent, several unnamed, a stage that
|
||||
/// ended, a run that finished.
|
||||
#[must_use]
|
||||
pub fn code(&self) -> Option<&'static str> {
|
||||
match self {
|
||||
Self::Control(ControlError::NoLiveTurn) => Some("no_live_turn"),
|
||||
Self::Control(ControlError::NoSuchStage(_)) => Some("no_such_stage"),
|
||||
Self::NoLiveAgent
|
||||
| Self::SeveralLiveAgents(_)
|
||||
| Self::Control(ControlError::NotLive | ControlError::Finished) => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 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)]
|
||||
|
|
@ -422,6 +441,24 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_refusal_has_a_code_when_its_reason_has_one() {
|
||||
assert_eq!(
|
||||
SteerError::Control(ControlError::NoLiveTurn).code(),
|
||||
Some("no_live_turn")
|
||||
);
|
||||
assert_eq!(
|
||||
SteerError::Control(ControlError::NoSuchStage("work".to_string())).code(),
|
||||
Some("no_such_stage")
|
||||
);
|
||||
assert_eq!(SteerError::NoLiveAgent.code(), None);
|
||||
assert_eq!(
|
||||
SteerError::SeveralLiveAgents(vec!["a@1".to_string()]).code(),
|
||||
None
|
||||
);
|
||||
assert_eq!(SteerError::Control(ControlError::Finished).code(), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_stage_label_names_its_node_visit_and_execution() {
|
||||
assert_eq!(parse_label("work@1"), Some(("work", None, 1)));
|
||||
|
|
|
|||
|
|
@ -129,12 +129,14 @@ impl FabroToolBackend for ClientBackend {
|
|||
|
||||
async fn interrupt_run(&self, run_id: &RunId) -> anyhow::Result<()> {
|
||||
self.ensure_run_scope(run_id)?;
|
||||
self.client.interrupt_run(run_id, None, None).await
|
||||
self.client.interrupt_run(run_id, None, None).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
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, None).await
|
||||
self.client.steer_run(run_id, text, interrupt, None).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn archive_run(&self, run_id: &RunId) -> anyhow::Result<Run> {
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ use std::sync::{Arc, RwLock};
|
|||
use anyhow::{Context as _, Result, anyhow, bail};
|
||||
use bytes::Bytes;
|
||||
use fabro_api::types;
|
||||
use fabro_api::types::RunControlAcknowledgement;
|
||||
use fabro_http::header::{ACCEPT, AUTHORIZATION, CONTENT_LENGTH, CONTENT_TYPE};
|
||||
use fabro_http::multipart::{Form, Part};
|
||||
use fabro_types::settings::run::MergeStrategy;
|
||||
|
|
@ -1145,7 +1146,7 @@ impl Client {
|
|||
run_id: &RunId,
|
||||
stage: Option<String>,
|
||||
text: Option<String>,
|
||||
) -> Result<()> {
|
||||
) -> Result<RunControlAcknowledgement> {
|
||||
let stage = stage
|
||||
.map(|stage| {
|
||||
types::InterruptRunRequestStage::try_from(stage)
|
||||
|
|
@ -1159,30 +1160,33 @@ impl Client {
|
|||
})
|
||||
.transpose()?;
|
||||
let body = types::InterruptRunRequest { stage, text };
|
||||
self.send_api(|client| {
|
||||
let body = body.clone();
|
||||
async move {
|
||||
client
|
||||
.interrupt_run()
|
||||
.id(run_id.to_string())
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
}
|
||||
})
|
||||
.await?;
|
||||
Ok(())
|
||||
let response = self
|
||||
.send_api(|client| {
|
||||
let body = body.clone();
|
||||
async move {
|
||||
client
|
||||
.interrupt_run()
|
||||
.id(run_id.to_string())
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
}
|
||||
})
|
||||
.await?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
/// Steer a run: the named stage (`node@visit`, or the node name), or
|
||||
/// the run's one live agent stage when `stage` is `None`.
|
||||
/// the run's one live agent stage when `stage` is `None`. The worker's
|
||||
/// answer: delivered, or pending when none came in time. A refusal is
|
||||
/// the error, with the refusal's code as its API failure code.
|
||||
pub async fn steer_run(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
text: String,
|
||||
interrupt: bool,
|
||||
stage: Option<String>,
|
||||
) -> Result<()> {
|
||||
) -> Result<RunControlAcknowledgement> {
|
||||
let stage = stage
|
||||
.map(|stage| {
|
||||
types::SteerRunRequestStage::try_from(stage)
|
||||
|
|
@ -1195,19 +1199,20 @@ impl Client {
|
|||
.stage(stage)
|
||||
.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(())
|
||||
let response = self
|
||||
.send_api(|client| {
|
||||
let body = body.clone();
|
||||
async move {
|
||||
client
|
||||
.steer_run()
|
||||
.id(run_id.to_string())
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
}
|
||||
})
|
||||
.await?;
|
||||
Ok(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn get_run_pair_status(&self, run_id: &RunId) -> Result<RunPairStatusResponse> {
|
||||
|
|
|
|||
|
|
@ -41,6 +41,10 @@ impl EnvVars {
|
|||
/// run's checkpoint at a named point (`fabro_petri::hooks`); unset
|
||||
/// outside tests.
|
||||
pub const FABRO_TEST_CHECKPOINT_GATES: &'static str = "FABRO_TEST_CHECKPOINT_GATES";
|
||||
/// `1` makes a Petri worker apply every control but acknowledge none,
|
||||
/// so a test sees the server's wait for an answer run out; unset
|
||||
/// outside tests.
|
||||
pub const FABRO_TEST_CONTROL_ACKS_MUTED: &'static str = "FABRO_TEST_CONTROL_ACKS_MUTED";
|
||||
pub const FABRO_VERBOSE: &'static str = "FABRO_VERBOSE";
|
||||
pub const FABRO_WEB_URL: &'static str = "FABRO_WEB_URL";
|
||||
pub const FABRO_WORKER_TOKEN: &'static str = "FABRO_WORKER_TOKEN";
|
||||
|
|
@ -248,6 +252,7 @@ mod tests {
|
|||
EnvVars::FABRO_TEST_DISABLE_SPA_ASSETS,
|
||||
EnvVars::FABRO_TEST_MODE,
|
||||
EnvVars::FABRO_TEST_CHECKPOINT_GATES,
|
||||
EnvVars::FABRO_TEST_CONTROL_ACKS_MUTED,
|
||||
EnvVars::FABRO_VERBOSE,
|
||||
EnvVars::FABRO_WEB_URL,
|
||||
EnvVars::FABRO_WORKER_TOKEN,
|
||||
|
|
|
|||
|
|
@ -386,7 +386,9 @@ models/run-commit-parent.ts
|
|||
models/run-commit-person.ts
|
||||
models/run-commit.ts
|
||||
models/run-commits-meta.ts
|
||||
models/run-control-acknowledgement.ts
|
||||
models/run-control-action.ts
|
||||
models/run-control-outcome.ts
|
||||
models/run-diff.ts
|
||||
models/run-environment-settings.ts
|
||||
models/run-error.ts
|
||||
|
|
|
|||
|
|
@ -42,6 +42,8 @@ import type { PreviewUrlRequest } from '../models';
|
|||
// @ts-ignore
|
||||
import type { PreviewUrlResponse } from '../models';
|
||||
// @ts-ignore
|
||||
import type { RunControlAcknowledgement } from '../models';
|
||||
// @ts-ignore
|
||||
import type { RunPairStatusResponse } from '../models';
|
||||
// @ts-ignore
|
||||
import type { SandboxDetails } from '../models';
|
||||
|
|
@ -424,7 +426,7 @@ export const HumanInTheLoopApiAxiosParamCreator = function (configuration?: Conf
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker; an interrupt the worker cannot deliver (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed) is refused on the run\'s event stream as a `run.notice` record whose code says why (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker stopped the turn, `409` with the refusal\'s code when the worker or Petri refused it (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed), and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused interrupt is also a `run.notice` record on the run\'s event stream under the same code (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
|
|
@ -795,7 +797,7 @@ export const HumanInTheLoopApiAxiosParamCreator = function (configuration?: Conf
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker; a control the worker cannot deliver (no live agent stage, several unnamed, a stage that is not running, or, for an interrupt, a stage with no model turn in flight) is refused on the run\'s event stream as a `run.notice` record whose code says why (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker delivered it, `409` with the refusal\'s code when the worker or Petri refused it, and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused control is also a `run.notice` record on the run\'s event stream under the same code (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
@ -1010,14 +1012,14 @@ export const HumanInTheLoopApiFp = function(configuration?: Configuration) {
|
|||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker; an interrupt the worker cannot deliver (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed) is refused on the run\'s event stream as a `run.notice` record whose code says why (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker stopped the turn, `409` with the refusal\'s code when the worker or Petri refused it (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed), and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused interrupt is also a `run.notice` record on the run\'s event stream under the same code (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<void>> {
|
||||
async interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<RunControlAcknowledgement>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.interruptRun(id, interruptRunRequest, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['HumanInTheLoopApi.interruptRun']?.[localVarOperationServerIndex]?.url;
|
||||
|
|
@ -1124,14 +1126,14 @@ export const HumanInTheLoopApiFp = function(configuration?: Configuration) {
|
|||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker; a control the worker cannot deliver (no live agent stage, several unnamed, a stage that is not running, or, for an interrupt, a stage with no model turn in flight) is refused on the run\'s event stream as a `run.notice` record whose code says why (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker delivered it, `409` with the refusal\'s code when the worker or Petri refused it, and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused control is also a `run.notice` record on the run\'s event stream under the same code (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async steerRun(id: string, steerRunRequest: SteerRunRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<void>> {
|
||||
async steerRun(id: string, steerRunRequest: SteerRunRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<RunControlAcknowledgement>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.steerRun(id, steerRunRequest, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['HumanInTheLoopApi.steerRun']?.[localVarOperationServerIndex]?.url;
|
||||
|
|
@ -1250,14 +1252,14 @@ export const HumanInTheLoopApiFactory = function (configuration?: Configuration,
|
|||
return localVarFp.getSandboxFile(id, path, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker; an interrupt the worker cannot deliver (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed) is refused on the run\'s event stream as a `run.notice` record whose code says why (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker stopped the turn, `409` with the refusal\'s code when the worker or Petri refused it (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed), and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused interrupt is also a `run.notice` record on the run\'s event stream under the same code (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): AxiosPromise<void> {
|
||||
interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): AxiosPromise<RunControlAcknowledgement> {
|
||||
return localVarFp.interruptRun(id, interruptRunRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
|
|
@ -1340,14 +1342,14 @@ export const HumanInTheLoopApiFactory = function (configuration?: Configuration,
|
|||
return localVarFp.startRunPair(id, pairStartRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker; a control the worker cannot deliver (no live agent stage, several unnamed, a stage that is not running, or, for an interrupt, a stage with no model turn in flight) is refused on the run\'s event stream as a `run.notice` record whose code says why (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker delivered it, `409` with the refusal\'s code when the worker or Petri refused it, and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused control is also a `run.notice` record on the run\'s event stream under the same code (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
steerRun(id: string, steerRunRequest: SteerRunRequest, options?: RawAxiosRequestConfig): AxiosPromise<void> {
|
||||
steerRun(id: string, steerRunRequest: SteerRunRequest, options?: RawAxiosRequestConfig): AxiosPromise<RunControlAcknowledgement> {
|
||||
return localVarFp.steerRun(id, steerRunRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
|
|
@ -1466,7 +1468,7 @@ export class HumanInTheLoopApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker; an interrupt the worker cannot deliver (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed) is refused on the run\'s event stream as a `run.notice` record whose code says why (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* Stop the current model turn of a live agent stage (the stage `stage` names, or the run\'s one live agent stage) and keep its session. With `text`, the text is the stage\'s next input; without, the stage waits for the next steer message before starting another model turn. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker stopped the turn, `409` with the refusal\'s code when the worker or Petri refused it (a stage with no model turn in flight, such as an agent between turns or a human gate; a stage that is not running; no live agent stage, or several unnamed), and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused interrupt is also a `run.notice` record on the run\'s event stream under the same code (`no_live_turn`, `no_such_stage`, `interrupt_refused`). A delivered interrupt is the stage\'s `control.requested` record with `$interrupt`, followed by an `attractor.turn.interrupted` progress record.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
|
|
@ -1564,7 +1566,7 @@ export class HumanInTheLoopApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker; a control the worker cannot deliver (no live agent stage, several unnamed, a stage that is not running, or, for an interrupt, a stage with no model turn in flight) is refused on the run\'s event stream as a `run.notice` record whose code says why (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* Send a mid-run steering message to a live agent stage of a running run: the stage `stage` names, or the run\'s one live agent stage. Without `interrupt`, the text is guidance for the stage\'s session, run as a follow-up turn once its current answer is reached. With `interrupt=true`, the stage\'s current model turn (the model request and the tool calls it is running) is stopped first, the session is kept, and the text is the stage\'s next input. The control is forwarded to the run\'s worker, and the worker\'s answer is this response: `202` with `outcome: delivered` (and the stage\'s label) once the worker delivered it, `409` with the refusal\'s code when the worker or Petri refused it, and `202` with `outcome: pending` when the worker gave no answer within the wait (5 s). A refused control is also a `run.notice` record on the run\'s event stream under the same code (`steer_refused`, `no_live_turn`, `no_such_stage`, `interrupt_refused`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
|
|||
|
|
@ -357,7 +357,9 @@ export * from './run-commit';
|
|||
export * from './run-commit-parent';
|
||||
export * from './run-commit-person';
|
||||
export * from './run-commits-meta';
|
||||
export * from './run-control-acknowledgement';
|
||||
export * from './run-control-action';
|
||||
export * from './run-control-outcome';
|
||||
export * from './run-diff';
|
||||
export * from './run-environment-settings';
|
||||
export * from './run-error';
|
||||
|
|
|
|||
29
lib/packages/fabro-api-client/src/models/run-control-acknowledgement.ts
generated
Normal file
29
lib/packages/fabro-api-client/src/models/run-control-acknowledgement.ts
generated
Normal 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.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunControlOutcome } from './run-control-outcome';
|
||||
|
||||
/**
|
||||
* The worker\'s answer to a steer or an interrupt, as the endpoint\'s `202` body.
|
||||
*/
|
||||
export interface RunControlAcknowledgement {
|
||||
'outcome': RunControlOutcome;
|
||||
/**
|
||||
* The label of the stage the control was delivered to (`node@visit`, or `node/e<execution>@visit`); absent when the outcome is `pending`.
|
||||
*/
|
||||
'stage'?: string;
|
||||
}
|
||||
26
lib/packages/fabro-api-client/src/models/run-control-outcome.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/run-control-outcome.ts
generated
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* What became of a forwarded control: `delivered` once the worker delivered it to its stage; `pending` when the worker gave no answer within the wait, in which case the run\'s event stream says what became of it (a `run.notice` record on refusal).
|
||||
*/
|
||||
|
||||
export const RunControlOutcome = {
|
||||
DELIVERED: 'delivered',
|
||||
PENDING: 'pending'
|
||||
} as const;
|
||||
|
||||
export type RunControlOutcome = typeof RunControlOutcome[keyof typeof RunControlOutcome];
|
||||
Loading…
Add table
Reference in a new issue