mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Merge branch 'petri-followup-interrupt' into petri-integration
# Conflicts: # lib/apps/fabro-cli/src/commands/run/petri_stream.rs
This commit is contained in:
commit
696acc18d6
18 changed files with 982 additions and 148 deletions
|
|
@ -1719,11 +1719,19 @@ paths:
|
|||
tags: [Human-in-the-Loop]
|
||||
summary: Steer Run
|
||||
description: |
|
||||
Send a mid-run steering message to the live agent session(s) of a
|
||||
running run. Set `interrupt=true` to atomically interrupt the active
|
||||
steerable agent round first, then deliver this message as the next
|
||||
user turn. Without `interrupt=true`, the message is appended to the
|
||||
steering queue and may buffer until the next steerable agent session.
|
||||
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`).
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
requestBody:
|
||||
|
|
@ -2037,14 +2045,38 @@ paths:
|
|||
tags: [Human-in-the-Loop]
|
||||
summary: Interrupt Run
|
||||
description: |
|
||||
Interrupt the active steerable agent round without sending steering
|
||||
text. The agent keeps its steering lease and waits for a later steer
|
||||
message before starting another LLM round.
|
||||
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.
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
requestBody:
|
||||
required: false
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/InterruptRunRequest"
|
||||
responses:
|
||||
"202":
|
||||
description: Interrupt accepted and forwarded to the worker
|
||||
"400":
|
||||
description: Invalid request body
|
||||
headers:
|
||||
x-request-id:
|
||||
$ref: "#/components/headers/XRequestId"
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/ErrorResponse"
|
||||
"404":
|
||||
description: Run not found
|
||||
headers:
|
||||
|
|
@ -2057,9 +2089,7 @@ paths:
|
|||
"409":
|
||||
description: |
|
||||
Run is not currently interruptible. Returned when the run is in a
|
||||
terminal state, blocked (use the answer endpoint instead), has no
|
||||
active steerable agent session, or active agent sessions have no
|
||||
live control channel.
|
||||
terminal state or is not running yet.
|
||||
headers:
|
||||
x-request-id:
|
||||
$ref: "#/components/headers/XRequestId"
|
||||
|
|
@ -9941,10 +9971,10 @@ components:
|
|||
interrupt:
|
||||
type: boolean
|
||||
description: |
|
||||
When true, apply a worker-control interrupt first, then deliver
|
||||
this text as steering in the same control operation. When false
|
||||
(default), append to the steering queue and let the agent pick it
|
||||
up at the next turn boundary.
|
||||
When true, stop the stage's current model turn first and make
|
||||
this text its next input, in one control. When false (default),
|
||||
the text is guidance the agent runs as a follow-up turn once its
|
||||
current answer is reached.
|
||||
default: false
|
||||
stage:
|
||||
type: string
|
||||
|
|
@ -9957,6 +9987,28 @@ components:
|
|||
maxLength: 200
|
||||
example: code@2
|
||||
|
||||
InterruptRunRequest:
|
||||
description: Request body for interrupting a live agent stage's model turn.
|
||||
type: object
|
||||
properties:
|
||||
stage:
|
||||
type: string
|
||||
description: |
|
||||
The agent stage to interrupt: its stage identifier
|
||||
(`node_id@visit`) or its node name. Omit it to interrupt the
|
||||
run's one live agent stage.
|
||||
minLength: 1
|
||||
maxLength: 200
|
||||
example: code@2
|
||||
text:
|
||||
type: string
|
||||
description: |
|
||||
The stage's next input once its turn is stopped. Omit it to let
|
||||
the stage wait for the next steer message.
|
||||
minLength: 1
|
||||
maxLength: 8192
|
||||
example: Stop and summarize what you have so far.
|
||||
|
||||
StartRunRequest:
|
||||
description: Request body for starting or resuming a run.
|
||||
type: object
|
||||
|
|
|
|||
|
|
@ -476,6 +476,19 @@ pub(crate) fn format_pretty(
|
|||
}
|
||||
"step.progress.recorded" => format_progress(&ts, view, label, styles),
|
||||
"control.requested" => {
|
||||
if let Some(interrupt) = view.body()?.pointer("/ctl/deliver/$interrupt") {
|
||||
let text = interrupt
|
||||
.get("steer")
|
||||
.and_then(Value::as_str)
|
||||
.map(|text| format!(": {text}"))
|
||||
.unwrap_or_default();
|
||||
return Some(format!(
|
||||
"{ts} {} {}{}",
|
||||
styles.yellow.apply_to("\u{23f8} Interrupt"),
|
||||
styles.bold.apply_to(label),
|
||||
text,
|
||||
));
|
||||
}
|
||||
let answer = view.derived()?.get("answer")?;
|
||||
let value = answer
|
||||
.get("choice")
|
||||
|
|
@ -621,6 +634,14 @@ fn format_progress(ts: &str, view: PetriItem<'_>, label: &str, styles: &Styles)
|
|||
styles.dim.apply_to(format!("{count} {noun} available")),
|
||||
))
|
||||
}
|
||||
"attractor.turn.interrupted" => {
|
||||
let backend = custom.get("backend").and_then(Value::as_str).unwrap_or("?");
|
||||
Some(format!(
|
||||
"{ts} {} {}",
|
||||
styles.dim.apply_to("\u{21b3}"),
|
||||
styles.dim.apply_to(format!("turn interrupted ({backend})")),
|
||||
))
|
||||
}
|
||||
"attractor.checkout" => {
|
||||
let repository = custom
|
||||
.get("repository")
|
||||
|
|
@ -1188,6 +1209,57 @@ mod tests {
|
|||
assert!(line.contains("2 tools available"), "{line}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_interrupt_and_the_turn_it_stopped_render_by_kind() {
|
||||
let styles = Styles::new(false);
|
||||
let mut state = PrettyState::default();
|
||||
let interrupt = petri(
|
||||
8,
|
||||
json!({
|
||||
"origin": "external",
|
||||
"context": {"invocation": 0, "execution": 0},
|
||||
"subject": subject("work", "agent"),
|
||||
"record": {"seq": 20, "body": {"event": "control.requested", "firing": 2,
|
||||
"ctl": {"deliver": {"$interrupt": {"steer": "stop and summarize"}}}}},
|
||||
"derived": {"deliverable": true}
|
||||
}),
|
||||
);
|
||||
let line = format_pretty(&interrupt, &styles, &mut state).expect("an interrupt line");
|
||||
assert!(
|
||||
line.contains("\u{23f8} Interrupt work: stop and summarize"),
|
||||
"{line}"
|
||||
);
|
||||
|
||||
let plain = petri(
|
||||
9,
|
||||
json!({
|
||||
"origin": "external",
|
||||
"context": {"invocation": 0, "execution": 0},
|
||||
"subject": subject("work", "agent"),
|
||||
"record": {"seq": 21, "body": {"event": "control.requested", "firing": 2,
|
||||
"ctl": {"deliver": {"$interrupt": {}}}}},
|
||||
"derived": {"deliverable": true}
|
||||
}),
|
||||
);
|
||||
let line = format_pretty(&plain, &styles, &mut state).expect("an interrupt line");
|
||||
assert!(line.ends_with("\u{23f8} Interrupt work"), "{line}");
|
||||
|
||||
let stopped = petri(
|
||||
10,
|
||||
json!({
|
||||
"origin": "external",
|
||||
"context": {"invocation": 0, "execution": 0},
|
||||
"subject": subject("work", "agent"),
|
||||
"record": {"seq": 22, "body": {"event": "step.progress.recorded", "firing": 2,
|
||||
"ev": {"custom": {"kind": "attractor.turn.interrupted", "node": "work",
|
||||
"firing": 2, "attempt": 1, "backend": "api", "session": "s-1"}}}},
|
||||
"derived": {}
|
||||
}),
|
||||
);
|
||||
let line = format_pretty(&stopped, &styles, &mut state).expect("a stopped-turn line");
|
||||
assert!(line.contains("\u{21b3} turn interrupted (api)"), "{line}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_engine_finish_decides_the_exit_code() {
|
||||
let finished = petri(
|
||||
|
|
|
|||
|
|
@ -25,14 +25,18 @@
|
|||
//! hold and release admission through the run's [`RunControls`]; a steer
|
||||
//! goes to the agent stage it names (`node@visit`, or the node name) or,
|
||||
//! unnamed, to the run's one live agent stage, and is refused with a
|
||||
//! `run.notice` record saying why when neither resolves. The
|
||||
//! `run.notice` record saying why when neither resolves; an interrupt
|
||||
//! 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
|
||||
//! 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
|
||||
//! `run.paused` and `run.unpaused` records. A resumed run that was paused when
|
||||
//! its worker died comes back paused, and the mirror reports that too. The
|
||||
//! interrupt and pair controls have no Petri adapter yet and are ignored with a
|
||||
//! warning. A control channel that is lost for good cancels the run the same
|
||||
//! pair controls have no Petri adapter yet and are ignored with a warning.
|
||||
//! A control channel that is lost for good cancels the run the same
|
||||
//! way, and the worker exits with that loss as its error once the run has
|
||||
//! settled.
|
||||
//!
|
||||
|
|
@ -65,7 +69,7 @@ use fabro_client::{Client, ServerTarget};
|
|||
use fabro_interview::{ControlInterviewer, WorkerControlMessage};
|
||||
use fabro_llm::credentials::{CredentialProvider, readiness};
|
||||
use fabro_petri::blobs::ClientBlobs;
|
||||
use fabro_petri::controls::RunControls;
|
||||
use fabro_petri::controls::{ControlError, RunControls, SteerError};
|
||||
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
|
||||
use fabro_petri::hooks::HooksSpec;
|
||||
use fabro_petri::interview::{Approval, FabroInterviewer};
|
||||
|
|
@ -80,7 +84,7 @@ use fabro_store::platform_records::{
|
|||
PlatformRecord, RunLifecycleKind, RunLifecycleRecord, RunNoticeRecord,
|
||||
};
|
||||
use fabro_types::settings::run::{ApprovalMode, RunMode};
|
||||
use fabro_types::{FailureReason, RunId, RunNoticeLevel, RunStatus, SuccessReason};
|
||||
use fabro_types::{FailureReason, Principal, RunId, RunNoticeLevel, RunStatus, SuccessReason};
|
||||
use fabro_vault::Vault;
|
||||
use fabro_workflow::Error as WorkflowError;
|
||||
use fabro_workflow::services::FabroRunToolServices;
|
||||
|
|
@ -343,9 +347,13 @@ impl PetriControls {
|
|||
}
|
||||
}
|
||||
}
|
||||
WorkerControlMessage::Interrupt { .. }
|
||||
| WorkerControlMessage::InterruptThenSteer { .. }
|
||||
| WorkerControlMessage::PairStart { .. }
|
||||
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::PairStart { .. }
|
||||
| WorkerControlMessage::PairMessage { .. }
|
||||
| WorkerControlMessage::PairEnd { .. } => {
|
||||
warn!(
|
||||
|
|
@ -358,6 +366,29 @@ impl PetriControls {
|
|||
}
|
||||
}
|
||||
|
||||
/// Stop the named stage's model turn, `text` as its next input when
|
||||
/// given. A refusal is a `run.notice` whose code says why: `no_live_turn`
|
||||
/// when the stage has no model turn in flight, `no_such_stage` when the
|
||||
/// name is not running, `interrupt_refused` otherwise.
|
||||
async fn interrupt(&self, stage: Option<&str>, text: Option<&str>, actor: &Principal) {
|
||||
match self.controls.interrupt(stage, text).await {
|
||||
Ok(stage) => {
|
||||
info!(
|
||||
run_id = %self.run_id,
|
||||
stage,
|
||||
steered = text.is_some(),
|
||||
actor = ?actor,
|
||||
"interrupt delivered"
|
||||
);
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(run_id = %self.run_id, error = %error, "interrupt refused");
|
||||
self.notice(interrupt_refusal_code(&error), error.to_string())
|
||||
.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 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) {
|
||||
|
|
@ -372,6 +403,19 @@ 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 wire name of a control, for a log line.
|
||||
fn control_name(message: &WorkerControlMessage) -> &'static str {
|
||||
match message {
|
||||
|
|
|
|||
|
|
@ -808,7 +808,7 @@ async fn run_status_offline(server: &RunningServer) -> Option<String> {
|
|||
}
|
||||
|
||||
/// The run's pending questions, as the API lists them.
|
||||
async fn questions(server: &RunningServer, run_id: &str) -> Vec<serde_json::Value> {
|
||||
pub(super) async fn questions(server: &RunningServer, run_id: &str) -> Vec<serde_json::Value> {
|
||||
run_json(server, &format!("runs/{run_id}/questions")).await["data"]
|
||||
.as_array()
|
||||
.cloned()
|
||||
|
|
@ -816,7 +816,7 @@ async fn questions(server: &RunningServer, run_id: &str) -> Vec<serde_json::Valu
|
|||
}
|
||||
|
||||
/// Wait until `count` questions are pending at once.
|
||||
async fn wait_for_questions(
|
||||
pub(super) async fn wait_for_questions(
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
count: usize,
|
||||
|
|
@ -838,7 +838,12 @@ async fn wait_for_questions(
|
|||
/// Answer a question through the API, as the web app and the CLI do. The
|
||||
/// question id is Petri's (`gate#2`), so it travels as one percent-encoded
|
||||
/// path segment, as the generated clients send it.
|
||||
async fn answer(server: &RunningServer, run_id: &str, question_id: &str, body: serde_json::Value) {
|
||||
pub(super) async fn answer(
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
question_id: &str,
|
||||
body: serde_json::Value,
|
||||
) {
|
||||
let mut url = fabro_http::Url::parse(&format!(
|
||||
"{}/api/v1/runs/{run_id}/questions",
|
||||
server.api_base_url
|
||||
|
|
|
|||
|
|
@ -4,8 +4,10 @@
|
|||
//! without the API; a steer reaches the agent stage on the twin, which
|
||||
//! sees it in its next request, and the stream carries the control record;
|
||||
//! two live agent stages are steered apart by their stage labels, and an
|
||||
//! unnamed steer between them is refused; a run paused when its server and
|
||||
//! worker die resumes paused and goes on once unpaused.
|
||||
//! unnamed steer between them is refused; an interrupt during a long tool
|
||||
//! call ends the agent's turn and its text is the next input, while an
|
||||
//! 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 harness is `petri.rs`'s: a foreground server on disk storage, the
|
||||
//! run started with `fabro run --detach`, and the host scope through the
|
||||
|
|
@ -30,9 +32,9 @@ use fabro_test::{TwinScenario, TwinScenarios, TwinToolCall, test_context, twin_o
|
|||
use serde_json::{Value, json};
|
||||
|
||||
use super::petri::{
|
||||
RunningServer, count_of, host_plugin, run_detached, run_detached_with, run_json, run_status,
|
||||
run_stream, settled_stream, stream_names, wait_for_status, wait_for_worker,
|
||||
wait_until_gate_is_polled, write_petri_workflow,
|
||||
RunningServer, answer, count_of, host_plugin, run_detached, run_detached_with, run_json,
|
||||
run_status, run_stream, settled_stream, stream_names, wait_for_questions, wait_for_status,
|
||||
wait_for_worker, wait_until_gate_is_polled, write_petri_workflow,
|
||||
};
|
||||
use crate::support::TEST_DEV_TOKEN;
|
||||
|
||||
|
|
@ -50,6 +52,9 @@ const PROMPT_A: &str = "Alpha: wait for the gate, then report.";
|
|||
const PROMPT_B: &str = "Bravo: wait for the gate, then report.";
|
||||
const STEER_A: &str = "Steer alpha: mention the word lighthouse.";
|
||||
const STEER_B: &str = "Steer bravo: mention the word windmill.";
|
||||
/// The text an interrupt carries: the agent's next input once its turn
|
||||
/// is stopped.
|
||||
const INTERRUPT_STEER: &str = "Stop waiting and summarize what you have.";
|
||||
|
||||
/// Two command stages: `a` waits on `gate`, `b` leaves `marker`.
|
||||
fn two_stage_workspace(context: &fabro_test::TestContext, gate: &Path, marker: &Path) -> PathBuf {
|
||||
|
|
@ -135,6 +140,55 @@ fn steer_by_cli(
|
|||
);
|
||||
}
|
||||
|
||||
/// `fabro steer --interrupt <run> <text>` against the server: the stage's
|
||||
/// current turn is stopped and `text` is its next input.
|
||||
fn interrupt_by_cli(
|
||||
context: &fabro_test::TestContext,
|
||||
server: &RunningServer,
|
||||
run_id: &str,
|
||||
text: &str,
|
||||
) {
|
||||
let output = context
|
||||
.command()
|
||||
.args(["steer", "--server", &server.target(), run_id])
|
||||
.args(["--interrupt", text])
|
||||
.output()
|
||||
.expect("the steer command executes");
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"fabro steer --interrupt failed\nstdout:\n{}\nstderr:\n{}",
|
||||
String::from_utf8_lossy(&output.stdout),
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
}
|
||||
|
||||
/// The stream's `run.notice` records, as `(code, message)`, in order.
|
||||
fn notices(items: &[Value]) -> Vec<(String, String)> {
|
||||
items
|
||||
.iter()
|
||||
.filter(|item| item["kind"] == "platform")
|
||||
.map(|item| &item["item"]["record"])
|
||||
.filter(|record| record["kind"] == "run.notice")
|
||||
.map(|record| {
|
||||
(
|
||||
record["code"].as_str().unwrap_or_default().to_string(),
|
||||
record["message"].as_str().unwrap_or_default().to_string(),
|
||||
)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The stream's `step.progress.recorded` custom payloads of `kind`.
|
||||
fn progress_of_kind<'a>(items: &'a [Value], kind: &str) -> Vec<&'a Value> {
|
||||
items
|
||||
.iter()
|
||||
.map(|item| &item["item"]["record"]["body"])
|
||||
.filter(|body| body["event"] == "step.progress.recorded")
|
||||
.map(|body| &body["ev"]["custom"])
|
||||
.filter(|custom| custom["kind"] == kind)
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The twin's request inputs that carry `prompt`, in order.
|
||||
fn inputs_with(logs: &Value, prompt: &str) -> Vec<String> {
|
||||
logs["requests"]
|
||||
|
|
@ -622,6 +676,196 @@ async fn two_live_agent_stages_are_steered_apart_by_their_labels() {
|
|||
server.shutdown();
|
||||
}
|
||||
|
||||
/// An interrupt while the agent's tool call waits on a gate that never
|
||||
/// opens: `fabro steer --interrupt` stops the turn (the tool call is
|
||||
/// cancelled, no answer is reached), the session is kept, and the text is
|
||||
/// the agent's next input, which the twin answers. The stream carries the
|
||||
/// `control.requested` record with the `$interrupt` value and the stage's
|
||||
/// `attractor.turn.interrupted` report, and the run succeeds.
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn an_interrupt_ends_the_turn_and_its_text_is_the_next_input() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let context = test_context!();
|
||||
let twin = twin_openai().await;
|
||||
let namespace = format!("{}::{}", module_path!(), line!());
|
||||
let server = RunningServer::start_with(
|
||||
&format!(
|
||||
"\n[llm.providers.openai]\nbase_url = \"{}\"\n",
|
||||
twin.base_url
|
||||
),
|
||||
&[(EnvVars::OPENAI_API_KEY, &namespace)],
|
||||
)
|
||||
.await;
|
||||
// The gate is never opened: only the interrupt ends the tool call.
|
||||
let gate = context.temp_dir.join("interrupt.gate");
|
||||
TwinScenarios::new(namespace.clone())
|
||||
.scenario(
|
||||
TwinScenario::responses(MODEL)
|
||||
.input_contains(PROMPT)
|
||||
.tool_call(TwinToolCall::new(
|
||||
"shell",
|
||||
json!({ "command": format!("while [ ! -f {} ]; do sleep 0.05; done", gate.display()) }),
|
||||
)),
|
||||
)
|
||||
.scenario(
|
||||
TwinScenario::responses(MODEL)
|
||||
.input_contains(INTERRUPT_STEER)
|
||||
.text("Summary: I was waiting on the gate."),
|
||||
)
|
||||
.load(twin)
|
||||
.await;
|
||||
let workspace = write_petri_workflow(
|
||||
&context,
|
||||
&format!(
|
||||
"digraph Interrupt {{\n graph [goal=\"Wait then report\", default_max_retries=0]\n \
|
||||
start [shape=Mdiamond]\n exit [shape=Msquare]\n work [shape=box, \
|
||||
prompt=\"{PROMPT}\", max_retries=0]\n start -> work -> exit\n}}\n"
|
||||
),
|
||||
);
|
||||
let run_id = run_detached_with(&context, &server, &workspace, &[
|
||||
"--auto-approve",
|
||||
"--provider",
|
||||
"openai",
|
||||
"--model",
|
||||
MODEL,
|
||||
]);
|
||||
|
||||
wait_for_status(&server, &run_id, &["running"]).await;
|
||||
wait_until_gate_is_polled(&gate);
|
||||
eprintln!("run {run_id}: the agent's tool is waiting on the gate; interrupting");
|
||||
interrupt_by_cli(&context, &server, &run_id, INTERRUPT_STEER);
|
||||
wait_for_stream_count(&server, &run_id, "control.requested", 1).await;
|
||||
eprintln!("run {run_id}: the interrupt is recorded");
|
||||
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
let items = settled_stream(&server, &run_id).await;
|
||||
let names = stream_names(&items);
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"stream: {names:?}\nserver stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert!(!gate.exists(), "nothing opened the gate");
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
assert_eq!(notices(&items), Vec::new(), "nothing was refused");
|
||||
|
||||
let delivery = items
|
||||
.iter()
|
||||
.find(|item| item["item"]["record"]["body"]["event"] == "control.requested")
|
||||
.expect("the interrupt is in the stream");
|
||||
let record = &delivery["item"]["record"]["body"];
|
||||
assert_eq!(
|
||||
record["ctl"]["deliver"]["$interrupt"]["steer"], INTERRUPT_STEER,
|
||||
"the control record carries the interrupt and its text: {record}"
|
||||
);
|
||||
assert_eq!(
|
||||
delivery["item"]["derived"]["deliverable"], true,
|
||||
"the interrupt was delivered to a live firing: {delivery}"
|
||||
);
|
||||
let interrupted = progress_of_kind(&items, "attractor.turn.interrupted");
|
||||
assert_eq!(interrupted.len(), 1, "{names:?}");
|
||||
assert_eq!(interrupted[0]["node"], "work", "{}", interrupted[0]);
|
||||
assert_eq!(interrupted[0]["backend"], "api", "{}", interrupted[0]);
|
||||
|
||||
let logs = twin.request_logs(&namespace).await;
|
||||
let inputs = inputs_with(&logs, PROMPT);
|
||||
assert_eq!(
|
||||
inputs.len(),
|
||||
2,
|
||||
"the interrupted turn, then the steered one: {inputs:?}"
|
||||
);
|
||||
assert!(
|
||||
!inputs[0].contains(INTERRUPT_STEER),
|
||||
"the first request came before the interrupt: {}",
|
||||
inputs[0]
|
||||
);
|
||||
assert!(
|
||||
inputs[1].contains(INTERRUPT_STEER),
|
||||
"the next request carries the interrupt's text as its input: {}",
|
||||
inputs[1]
|
||||
);
|
||||
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, 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,
|
||||
&format!(
|
||||
"digraph Gate {{\n graph [goal=\"Ask before running\"]\n start [shape=Mdiamond]\n \
|
||||
exit [shape=Msquare]\n gate [shape=hexagon, label=\"Go?\", \
|
||||
question_type=\"yes_no\"]\n yes [shape=parallelogram, script=\"touch {marker}\"]\n \
|
||||
start -> gate\n gate -> yes [label=\"[Y] Yes\"]\n gate -> exit [label=\"[N] \
|
||||
No\"]\n yes -> exit\n}}\n",
|
||||
marker = marker.display()
|
||||
),
|
||||
);
|
||||
let run_id = run_detached_with(&context, &server, &workspace, &[]);
|
||||
|
||||
let pending = wait_for_questions(&server, &run_id, 1).await;
|
||||
assert_eq!(pending[0]["stage"], "gate@1", "{}", pending[0]);
|
||||
let question_id = pending[0]["id"].as_str().expect("an id").to_string();
|
||||
eprintln!("run {run_id}: the gate is asking; interrupting it");
|
||||
|
||||
let (status, body) = control(
|
||||
&server,
|
||||
&run_id,
|
||||
"interrupt",
|
||||
Some(json!({ "stage": "gate" })),
|
||||
)
|
||||
.await;
|
||||
assert_eq!(status, 202, "interrupt: {body}");
|
||||
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_eq!(count_of(&names, "control.requested"), 0, "{names:?}");
|
||||
let refused = notices(&run_stream(&server, &run_id).await);
|
||||
assert_eq!(refused.len(), 2, "{refused:?}");
|
||||
assert_eq!(refused[0].0, "no_live_turn", "{refused:?}");
|
||||
assert_eq!(refused[0].1, "the stage has no model turn to interrupt");
|
||||
assert_eq!(refused[1].0, "interrupt_refused", "{refused:?}");
|
||||
assert_eq!(refused[1].1, "Run has no active steerable agent session.");
|
||||
assert_eq!(run_status(&server, &run_id).await, "blocked");
|
||||
|
||||
answer(&server, &run_id, &question_id, json!({ "kind": "yes" })).await;
|
||||
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
|
||||
let items = settled_stream(&server, &run_id).await;
|
||||
let names = stream_names(&items);
|
||||
assert_eq!(
|
||||
status,
|
||||
"succeeded",
|
||||
"stream: {names:?}\nserver stderr:\n{}",
|
||||
server.stderr_text()
|
||||
);
|
||||
assert!(marker.exists(), "the answer routed the gate");
|
||||
assert_petri_succeeded(&server, &run_id).await;
|
||||
assert_eq!(
|
||||
count_of(&names, "control.requested"),
|
||||
1,
|
||||
"only the answer was delivered: {names:?}"
|
||||
);
|
||||
assert!(
|
||||
progress_of_kind(&items, "attractor.turn.interrupted").is_empty(),
|
||||
"no turn was stopped: {names:?}"
|
||||
);
|
||||
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
|
||||
|
|
|
|||
|
|
@ -413,6 +413,29 @@ impl RunAnswerTransport {
|
|||
}
|
||||
}
|
||||
|
||||
/// 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.
|
||||
async fn interrupt(
|
||||
&self,
|
||||
stage: Option<String>,
|
||||
text: Option<String>,
|
||||
actor: Principal,
|
||||
) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
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::InProcess { .. } => Err(AnswerTransportError::Closed),
|
||||
}
|
||||
}
|
||||
|
||||
async fn pause_run(&self) -> Result<(), AnswerTransportError> {
|
||||
match self {
|
||||
Self::Worker { run_id, bus } => {
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ use axum::extract::State;
|
|||
use axum::http::StatusCode;
|
||||
use axum::response::{IntoResponse, Response};
|
||||
use axum::routing::post;
|
||||
use fabro_api::types::SteerRunRequest;
|
||||
use fabro_api::types::{InterruptRunRequest, SteerRunRequest};
|
||||
use fabro_types::Principal;
|
||||
use fabro_workflow::run_status::RunStatus;
|
||||
|
||||
|
|
@ -19,11 +19,35 @@ pub(super) fn routes() -> axum::Router<Arc<AppState>> {
|
|||
.route("/runs/{id}/interrupt", post(interrupt_run))
|
||||
}
|
||||
|
||||
/// 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.
|
||||
enum RunControlRequest {
|
||||
/// Guidance for a live agent stage's session, run as a follow-up turn.
|
||||
Steer {
|
||||
text: String,
|
||||
stage: Option<String>,
|
||||
},
|
||||
/// Stop a live agent stage's current model turn and keep its session;
|
||||
/// `text`, when given, is the stage's next input.
|
||||
Interrupt {
|
||||
stage: Option<String>,
|
||||
text: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
impl RunControlRequest {
|
||||
fn name(&self) -> &'static str {
|
||||
match self {
|
||||
Self::Steer { .. } => "steer",
|
||||
Self::Interrupt { .. } => "interrupt",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A stage name, when given, must not be blank.
|
||||
fn blank_stage(stage: Option<&str>) -> bool {
|
||||
stage.is_some_and(|stage| stage.trim().is_empty())
|
||||
}
|
||||
|
||||
async fn steer_run(
|
||||
|
|
@ -43,37 +67,42 @@ async fn steer_run(
|
|||
return ApiError::bad_request("Steer text must not be empty.").into_response();
|
||||
}
|
||||
let stage = stage.map(String::from);
|
||||
if stage
|
||||
.as_deref()
|
||||
.is_some_and(|stage| stage.trim().is_empty())
|
||||
{
|
||||
if blank_stage(stage.as_deref()) {
|
||||
return ApiError::bad_request("Steer stage must not be empty.").into_response();
|
||||
}
|
||||
if interrupt {
|
||||
return interrupt_unsupported();
|
||||
}
|
||||
control_run(actor, state, id, RunControlRequest::Steer { text, stage }).await
|
||||
let control = if interrupt {
|
||||
RunControlRequest::Interrupt {
|
||||
stage,
|
||||
text: Some(text),
|
||||
}
|
||||
} else {
|
||||
RunControlRequest::Steer { text, stage }
|
||||
};
|
||||
control_run(actor, state, id, control).await
|
||||
}
|
||||
|
||||
/// Interrupting a live agent turn has no adapter over Petri's control
|
||||
/// service yet, which delivers a steer to a live stage and cancels a whole
|
||||
/// run but does not interrupt one stage's turn; the request is refused
|
||||
/// with that reason rather than accepted and dropped.
|
||||
/// Stop a live agent stage's current model turn. The body is optional: no
|
||||
/// body interrupts the run's one live agent stage and leaves it waiting
|
||||
/// for the next steer.
|
||||
async fn interrupt_run(
|
||||
RequireRunManagementTarget(_id, _actor): RequireRunManagementTarget,
|
||||
State(_state): State<Arc<AppState>>,
|
||||
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
|
||||
State(state): State<Arc<AppState>>,
|
||||
body: Option<Json<InterruptRunRequest>>,
|
||||
) -> Response {
|
||||
interrupt_unsupported()
|
||||
}
|
||||
|
||||
fn interrupt_unsupported() -> Response {
|
||||
ApiError::with_code(
|
||||
StatusCode::NOT_IMPLEMENTED,
|
||||
"Interrupting a run's agent turn is not supported: Petri's control service has no \
|
||||
per-stage interrupt yet. Steer the run without `interrupt`, or cancel it.",
|
||||
"interrupt_unsupported",
|
||||
)
|
||||
.into_response()
|
||||
let InterruptRunRequest { stage, text } = body.map(|Json(body)| body).unwrap_or_default();
|
||||
let stage = stage.map(String::from);
|
||||
if blank_stage(stage.as_deref()) {
|
||||
return ApiError::bad_request("Interrupt stage must not be empty.").into_response();
|
||||
}
|
||||
let text = text.map(String::from);
|
||||
if text.as_deref().is_some_and(|text| text.trim().is_empty()) {
|
||||
return ApiError::bad_request("Interrupt text must not be empty.").into_response();
|
||||
}
|
||||
control_run(actor, state, id, RunControlRequest::Interrupt {
|
||||
stage,
|
||||
text,
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
async fn control_run(
|
||||
|
|
@ -92,8 +121,13 @@ async fn control_run(
|
|||
let runs = state.runs.lock().expect("runs lock poisoned");
|
||||
match runs.get(&id) {
|
||||
Some(managed_run) => {
|
||||
match managed_run.status {
|
||||
RunStatus::Blocked { .. } => {
|
||||
match (&control, &managed_run.status) {
|
||||
// A blocked run may still have an agent stage running a
|
||||
// turn beside the question; the worker judges the
|
||||
// interrupt per stage and refuses the gate itself.
|
||||
(RunControlRequest::Interrupt { .. }, RunStatus::Blocked { .. })
|
||||
| (_, RunStatus::Running) => {}
|
||||
(RunControlRequest::Steer { .. }, RunStatus::Blocked { .. }) => {
|
||||
return ApiError::with_code(
|
||||
StatusCode::CONFLICT,
|
||||
"Run is blocked on a question; use the interview-answer endpoint \
|
||||
|
|
@ -102,25 +136,30 @@ async fn control_run(
|
|||
)
|
||||
.into_response();
|
||||
}
|
||||
RunStatus::Submitted
|
||||
| RunStatus::Pending { .. }
|
||||
| RunStatus::Runnable
|
||||
| RunStatus::Starting
|
||||
| RunStatus::Paused { .. } => {
|
||||
(
|
||||
_,
|
||||
RunStatus::Submitted
|
||||
| RunStatus::Pending { .. }
|
||||
| RunStatus::Runnable
|
||||
| RunStatus::Starting
|
||||
| RunStatus::Paused { .. },
|
||||
) => {
|
||||
return ApiError::with_code(
|
||||
StatusCode::CONFLICT,
|
||||
"Run is not currently running.",
|
||||
"run_not_steerable",
|
||||
not_controllable_code(&control),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
RunStatus::Failed { .. }
|
||||
| RunStatus::Succeeded { .. }
|
||||
| RunStatus::Removing
|
||||
| RunStatus::Dead => {
|
||||
(
|
||||
_,
|
||||
RunStatus::Failed { .. }
|
||||
| RunStatus::Succeeded { .. }
|
||||
| RunStatus::Removing
|
||||
| RunStatus::Dead,
|
||||
) => {
|
||||
return terminal_control_response(&control);
|
||||
}
|
||||
RunStatus::Running => {}
|
||||
}
|
||||
// Plain steers buffer in the worker hub when no agent session
|
||||
// is active; if active agents exist but none are steerable,
|
||||
|
|
@ -153,8 +192,14 @@ async fn control_run(
|
|||
.into_response();
|
||||
};
|
||||
|
||||
let RunControlRequest::Steer { text, stage } = control;
|
||||
let result = answer_transport.steer(text, stage, actor).await;
|
||||
let result = match control {
|
||||
RunControlRequest::Steer { text, stage } => {
|
||||
answer_transport.steer(text, stage, actor).await
|
||||
}
|
||||
RunControlRequest::Interrupt { stage, text } => {
|
||||
answer_transport.interrupt(stage, text, actor).await
|
||||
}
|
||||
};
|
||||
|
||||
match result {
|
||||
Ok(()) => StatusCode::ACCEPTED.into_response(),
|
||||
|
|
@ -173,11 +218,19 @@ async fn control_run(
|
|||
}
|
||||
}
|
||||
|
||||
fn terminal_control_response(_control: &RunControlRequest) -> Response {
|
||||
/// The 409 code of a control the run's status refuses.
|
||||
fn not_controllable_code(control: &RunControlRequest) -> &'static str {
|
||||
match control {
|
||||
RunControlRequest::Steer { .. } => "run_not_steerable",
|
||||
RunControlRequest::Interrupt { .. } => "run_not_interruptible",
|
||||
}
|
||||
}
|
||||
|
||||
fn terminal_control_response(control: &RunControlRequest) -> Response {
|
||||
ApiError::with_code(
|
||||
StatusCode::CONFLICT,
|
||||
"Run is no longer steerable.",
|
||||
"run_not_steerable",
|
||||
format!("Run no longer accepts a {}.", control.name()),
|
||||
not_controllable_code(control),
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8392,29 +8392,168 @@ async fn steer_with_active_non_steerable_session_returns_conflict() {
|
|||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn steer_with_interrupt_returns_unsupported() {
|
||||
async fn steer_with_interrupt_forwards_an_interrupt_then_steer_to_the_worker() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let run_id = fixtures::RUN_1;
|
||||
let (transport, _control_rx) = worker_transport_with_receiver(run_id).await;
|
||||
let (transport, mut control_rx) = worker_transport_with_receiver(run_id).await;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, Some(transport));
|
||||
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/steer")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(r#"{"text":"try again","interrupt":true}"#))
|
||||
.body(Body::from(
|
||||
r#"{"text":"stop and summarize","interrupt":true,"stage":"code@2"}"#,
|
||||
))
|
||||
.unwrap();
|
||||
|
||||
// A steer with an interrupt is not a control a Petri run takes.
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::NOT_IMPLEMENTED);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(body["errors"][0]["code"], "interrupt_unsupported");
|
||||
assert_status!(response, StatusCode::ACCEPTED).await;
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(
|
||||
matches!(
|
||||
envelope.message,
|
||||
WorkerControlMessage::InterruptThenSteer { ref text, ref stage, .. }
|
||||
if text == "stop and summarize" && stage.as_deref() == Some("code@2")
|
||||
),
|
||||
"{envelope:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_returns_unsupported() {
|
||||
async fn interrupt_forwards_the_stage_and_text_to_the_worker() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let run_id = fixtures::RUN_1;
|
||||
let (transport, mut control_rx) = worker_transport_with_receiver(run_id).await;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, Some(transport));
|
||||
|
||||
// No body: the run's one live agent stage, waiting for the next steer.
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/interrupt")))
|
||||
.body(Body::empty())
|
||||
.unwrap();
|
||||
let response = app.clone().oneshot(req).await.unwrap();
|
||||
assert_status!(response, StatusCode::ACCEPTED).await;
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(
|
||||
matches!(envelope.message, WorkerControlMessage::Interrupt {
|
||||
stage: None,
|
||||
..
|
||||
}),
|
||||
"{envelope:?}"
|
||||
);
|
||||
|
||||
// A stage alone: a plain interrupt of that stage.
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/interrupt")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(r#"{"stage":"code@2"}"#))
|
||||
.unwrap();
|
||||
let response = app.clone().oneshot(req).await.unwrap();
|
||||
assert_status!(response, StatusCode::ACCEPTED).await;
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(
|
||||
matches!(
|
||||
envelope.message,
|
||||
WorkerControlMessage::Interrupt { ref stage, .. } if stage.as_deref() == Some("code@2")
|
||||
),
|
||||
"{envelope:?}"
|
||||
);
|
||||
|
||||
// A stage and a text: the text is the stage's next input.
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/interrupt")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(
|
||||
r#"{"stage":"code@2","text":"stop and summarize"}"#,
|
||||
))
|
||||
.unwrap();
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_status!(response, StatusCode::ACCEPTED).await;
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(
|
||||
matches!(
|
||||
envelope.message,
|
||||
WorkerControlMessage::InterruptThenSteer { ref text, ref stage, .. }
|
||||
if text == "stop and summarize" && stage.as_deref() == Some("code@2")
|
||||
),
|
||||
"{envelope:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_with_a_blank_stage_or_text_returns_bad_request() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let run_id = fixtures::RUN_1;
|
||||
let (transport, _control_rx) = worker_transport_with_receiver(run_id).await;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, Some(transport));
|
||||
|
||||
for body in [r#"{"stage":" "}"#, r#"{"text":" "}"#] {
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/interrupt")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(body))
|
||||
.unwrap();
|
||||
let response = app.clone().oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::BAD_REQUEST, "{body}");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_of_a_blocked_run_is_forwarded_for_the_worker_to_judge() {
|
||||
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 mut runs = state.runs.lock().expect("runs lock poisoned");
|
||||
runs.get_mut(&run_id).unwrap().status = RunStatus::Blocked {
|
||||
blocked_reason: BlockedReason::HumanInputRequired,
|
||||
};
|
||||
}
|
||||
|
||||
// A steer of a blocked run goes to the answer endpoint; an interrupt
|
||||
// may still name an agent stage running beside the question, so the
|
||||
// worker decides, and refuses a gate with `no_live_turn` on the stream.
|
||||
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.clone().oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::CONFLICT);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(body["errors"][0]["code"], "use_answer_endpoint");
|
||||
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{run_id}/interrupt")))
|
||||
.header("content-type", "application/json")
|
||||
.body(Body::from(r#"{"stage":"gate@1"}"#))
|
||||
.unwrap();
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_status!(response, StatusCode::ACCEPTED).await;
|
||||
let envelope = recv_worker_control_envelope(&mut control_rx).await;
|
||||
assert!(
|
||||
matches!(
|
||||
envelope.message,
|
||||
WorkerControlMessage::Interrupt { ref stage, .. } if stage.as_deref() == Some("gate@1")
|
||||
),
|
||||
"{envelope:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_of_a_finished_run_returns_conflict() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let run_id = fixtures::RUN_1;
|
||||
|
|
@ -8441,12 +8580,44 @@ async fn interrupt_returns_unsupported() {
|
|||
.body(Body::empty())
|
||||
.unwrap();
|
||||
|
||||
// An interrupt is not a control a Petri run takes: the answer is
|
||||
// `unsupported`, whatever the run's state.
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::NOT_IMPLEMENTED);
|
||||
assert_eq!(response.status(), StatusCode::CONFLICT);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(body["errors"][0]["code"], "interrupt_unsupported");
|
||||
assert_eq!(body["errors"][0]["code"], "run_not_interruptible");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_of_an_unknown_run_returns_not_found() {
|
||||
let app = test_app_with();
|
||||
let missing_run_id = fixtures::RUN_64;
|
||||
|
||||
let req = Request::builder()
|
||||
.method("POST")
|
||||
.uri(api(&format!("/runs/{missing_run_id}/interrupt")))
|
||||
.body(Body::empty())
|
||||
.unwrap();
|
||||
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_status!(response, StatusCode::NOT_FOUND).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn interrupt_without_a_worker_channel_returns_unavailable() {
|
||||
let state = test_app_state();
|
||||
let app = crate::test_support::build_test_router(Arc::clone(&state));
|
||||
let run_id = fixtures::RUN_1;
|
||||
let _temp_dir = insert_running_control_run(&state, run_id, None);
|
||||
|
||||
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::SERVICE_UNAVAILABLE);
|
||||
let body = body_json(response.into_body()).await;
|
||||
assert_eq!(body["errors"][0]["code"], "worker_control_unavailable");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
|
|||
|
|
@ -83,19 +83,24 @@ impl WorkerControlEnvelope {
|
|||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn interrupt(actor: Principal) -> Self {
|
||||
pub fn interrupt(stage: Option<String>, actor: Principal) -> Self {
|
||||
Self {
|
||||
v: WORKER_CONTROL_PROTOCOL_VERSION,
|
||||
message: WorkerControlMessage::Interrupt { actor },
|
||||
message: WorkerControlMessage::Interrupt { stage, actor },
|
||||
}
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn interrupt_then_steer(text: impl Into<String>, actor: Principal) -> Self {
|
||||
pub fn interrupt_then_steer(
|
||||
text: impl Into<String>,
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
) -> Self {
|
||||
Self {
|
||||
v: WORKER_CONTROL_PROTOCOL_VERSION,
|
||||
message: WorkerControlMessage::InterruptThenSteer {
|
||||
text: text.into(),
|
||||
stage,
|
||||
actor,
|
||||
},
|
||||
}
|
||||
|
|
@ -173,9 +178,21 @@ pub enum WorkerControlMessage {
|
|||
actor: Principal,
|
||||
},
|
||||
#[serde(rename = "run.interrupt")]
|
||||
Interrupt { actor: Principal },
|
||||
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,
|
||||
},
|
||||
#[serde(rename = "run.interrupt_then_steer")]
|
||||
InterruptThenSteer { text: String, actor: Principal },
|
||||
InterruptThenSteer {
|
||||
text: String,
|
||||
/// The stage to interrupt and steer, as for `Interrupt`.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
stage: Option<String>,
|
||||
actor: Principal,
|
||||
},
|
||||
#[serde(rename = "pair.start")]
|
||||
PairStart {
|
||||
run_id: RunId,
|
||||
|
|
@ -320,7 +337,7 @@ mod tests {
|
|||
|
||||
#[test]
|
||||
fn interrupt_round_trips_through_json() {
|
||||
let envelope = WorkerControlEnvelope::interrupt(Principal::System {
|
||||
let envelope = WorkerControlEnvelope::interrupt(None, Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
});
|
||||
let json = serde_json::to_string(&envelope).unwrap();
|
||||
|
|
@ -334,14 +351,17 @@ mod tests {
|
|||
|
||||
#[test]
|
||||
fn interrupt_then_steer_round_trips_through_json() {
|
||||
let envelope =
|
||||
WorkerControlEnvelope::interrupt_then_steer("stop, do X instead", Principal::System {
|
||||
let envelope = WorkerControlEnvelope::interrupt_then_steer(
|
||||
"stop, do X instead",
|
||||
Some("code@2".to_string()),
|
||||
Principal::System {
|
||||
system_kind: SystemActorKind::Engine,
|
||||
});
|
||||
},
|
||||
);
|
||||
let json = serde_json::to_string(&envelope).unwrap();
|
||||
assert_eq!(
|
||||
json,
|
||||
r#"{"v":1,"type":"run.interrupt_then_steer","text":"stop, do X instead","actor":{"kind":"system","system_kind":"engine"}}"#
|
||||
r#"{"v":1,"type":"run.interrupt_then_steer","text":"stop, do X instead","stage":"code@2","actor":{"kind":"system","system_kind":"engine"}}"#
|
||||
);
|
||||
let parsed: WorkerControlEnvelope = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(parsed, envelope);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
//! The controls Fabro drives on a live Petri run: pause and unpause at
|
||||
//! admission, a steer into the run's agent stage, and cancel.
|
||||
//! admission, a steer into the run's agent stage, an interrupt of its
|
||||
//! current model turn, and cancel.
|
||||
//!
|
||||
//! [`RunControls`] is Petri's `ControlService` as the run's worker holds it:
|
||||
//! one per run, built before the run and handed to [`engine::run`] in its
|
||||
|
|
@ -26,6 +27,16 @@
|
|||
//! share one) or by the node's name; unnamed, it goes to the one live agent
|
||||
//! stage. With no live agent, several unnamed, or a name that is not running,
|
||||
//! it is refused with the reason, and nothing is recorded.
|
||||
//! - interrupt names its stage the way a steer does and stops the stage's
|
||||
//! current model turn (the model request and the tool calls it runs), keeping
|
||||
//! the session: the firing records `control.requested` with the
|
||||
//! `{"$interrupt": …}` value and the stage reports the stopped turn as
|
||||
//! `attractor.turn.interrupted`. The text given with the interrupt is the
|
||||
//! stage's next input; without one, the next steer is. A stage with no model
|
||||
//! turn in flight (an agent between turns, a gate, a command) refuses it with
|
||||
//! `NoLiveTurn`, and nothing is recorded. The check reads the service's
|
||||
//! live-turn set, which [`engine::run`] installs as a runtime capability
|
||||
//! beside the pause gate.
|
||||
//! - cancel is the caller's cancellation token ([`RunRequest::cancel`]); the
|
||||
//! service's own cancel is here for a host that holds only this.
|
||||
//!
|
||||
|
|
@ -41,7 +52,8 @@
|
|||
use std::collections::BTreeMap;
|
||||
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
|
||||
|
||||
use petri_execution::controls::{ControlError, ControlService};
|
||||
pub use petri_execution::controls::ControlError;
|
||||
use petri_execution::controls::{ControlService, LiveTurns};
|
||||
use petri_execution::{
|
||||
CoordinatorHandle, CoordinatorRecord, CoordinatorState, ExecutionId, ExecutionObserver,
|
||||
};
|
||||
|
|
@ -49,20 +61,22 @@ use petri_frontend_attractor::kinds::AGENT_KIND;
|
|||
use petri_runtime::driver::lifecycle::ExecutionHooks;
|
||||
use petri_runtime::engine::{EngineState, Event, EventRecord};
|
||||
use petri_runtime::ir::FiringId;
|
||||
use petri_runtime::steps::Interrupt;
|
||||
use tokio::sync::watch;
|
||||
|
||||
/// Why a steer was not delivered.
|
||||
/// Why a steer or an interrupt was not delivered.
|
||||
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
|
||||
pub enum SteerError {
|
||||
/// No agent stage is running: the same refusal the legacy server gave a
|
||||
/// control that needs a live agent session.
|
||||
#[error("Run has no active steerable agent session.")]
|
||||
NoLiveAgent,
|
||||
/// More than one agent stage is running and the steer names none, or
|
||||
/// More than one agent stage is running and the control names none, or
|
||||
/// names a label several live firings answer to.
|
||||
#[error("Run has several active agent stages ({}); the steer names none of them.", .0.join(", "))]
|
||||
#[error("Run has several active agent stages ({}); the control names none of them.", .0.join(", "))]
|
||||
SeveralLiveAgents(Vec<String>),
|
||||
/// The named stage is not running, or the run has ended.
|
||||
/// The named stage is not running, the stage has no model turn to
|
||||
/// interrupt, or the run has ended.
|
||||
#[error(transparent)]
|
||||
Control(#[from] ControlError),
|
||||
}
|
||||
|
|
@ -209,33 +223,65 @@ impl RunControls {
|
|||
/// a node name), or to the one live agent stage when `stage` is `None`.
|
||||
/// The label of the stage steered.
|
||||
pub async fn steer(&self, stage: Option<&str>, text: &str) -> Result<String, SteerError> {
|
||||
let live = self.agents().labelled();
|
||||
let ((execution, firing), label) = match stage {
|
||||
None => one_live_agent(live)?,
|
||||
Some(stage) => {
|
||||
let Some((node, execution, visit)) = parse_label(stage) else {
|
||||
// A node name: the service's own live-stage index.
|
||||
self.service.steer(stage, text).await?;
|
||||
return Ok(stage.to_owned());
|
||||
};
|
||||
let agents = self.agents();
|
||||
let matches = live
|
||||
.into_iter()
|
||||
.filter(|(key, _)| {
|
||||
let agent = &agents.firings[key];
|
||||
agent.node == node
|
||||
&& agent.visit == visit
|
||||
&& execution.is_none_or(|execution| key.0.raw() == execution)
|
||||
})
|
||||
.collect();
|
||||
drop(agents);
|
||||
labelled_agent(stage, matches)?
|
||||
}
|
||||
};
|
||||
let ((execution, firing), label) = self.resolve(stage)?;
|
||||
self.service.steer_firing(execution, firing, text).await?;
|
||||
Ok(label)
|
||||
}
|
||||
|
||||
/// Stop the current model turn of the stage `stage` names, or of the one
|
||||
/// live agent stage when `stage` is `None`, and keep its session. With
|
||||
/// `text`, the text is the stage's next input; without, the next steer
|
||||
/// is. Refused with [`ControlError::NoLiveTurn`] when the stage has no
|
||||
/// turn in flight (an agent between turns, or a stage that is not an
|
||||
/// agent). The label of the stage interrupted.
|
||||
pub async fn interrupt(
|
||||
&self,
|
||||
stage: Option<&str>,
|
||||
text: Option<&str>,
|
||||
) -> Result<String, SteerError> {
|
||||
let ((execution, firing), label) = self.resolve(stage)?;
|
||||
let interrupt = match text {
|
||||
Some(text) => Interrupt::and_steer(text),
|
||||
None => Interrupt::new(),
|
||||
};
|
||||
self.service
|
||||
.interrupt_firing(execution, firing, interrupt)
|
||||
.await?;
|
||||
Ok(label)
|
||||
}
|
||||
|
||||
/// The live firing a control goes to: the one `stage` names by label
|
||||
/// (`node@visit`, `node/e<execution>@visit`) or by node name, or the
|
||||
/// run's one live agent stage when `stage` is `None`.
|
||||
fn resolve(&self, stage: Option<&str>) -> Result<LabelledAgent, SteerError> {
|
||||
let live = self.agents().labelled();
|
||||
let Some(stage) = stage else {
|
||||
return one_live_agent(live);
|
||||
};
|
||||
let Some((node, execution, visit)) = parse_label(stage) else {
|
||||
// A node name: the service's own live-stage index, where a node
|
||||
// running in two executions keeps the latest.
|
||||
return match self.service.stage(stage) {
|
||||
Some(live) => Ok(((live.execution, live.firing), stage.to_owned())),
|
||||
None => Err(SteerError::Control(ControlError::NoSuchStage(
|
||||
stage.to_owned(),
|
||||
))),
|
||||
};
|
||||
};
|
||||
let agents = self.agents();
|
||||
let matches = live
|
||||
.into_iter()
|
||||
.filter(|(key, _)| {
|
||||
let agent = &agents.firings[key];
|
||||
agent.node == node
|
||||
&& agent.visit == visit
|
||||
&& execution.is_none_or(|execution| key.0.raw() == execution)
|
||||
})
|
||||
.collect();
|
||||
drop(agents);
|
||||
labelled_agent(stage, matches)
|
||||
}
|
||||
|
||||
/// Cancel the whole run politely; a second call reaches the kill tier.
|
||||
pub fn cancel(&self) -> Result<(), ControlError> {
|
||||
self.service.cancel()
|
||||
|
|
@ -246,6 +292,13 @@ impl RunControls {
|
|||
self.service.hooks(inner)
|
||||
}
|
||||
|
||||
/// The live-turn set the agent step marks while a model turn runs, for
|
||||
/// the runtime's capabilities. Without it installed no turn is ever
|
||||
/// live and every interrupt is refused.
|
||||
pub(crate) fn turns(&self) -> LiveTurns {
|
||||
self.service.turns()
|
||||
}
|
||||
|
||||
/// Hand the service the run's coordinator handle.
|
||||
pub(crate) fn wire(&self, handle: CoordinatorHandle) {
|
||||
self.service.wire(handle);
|
||||
|
|
@ -344,6 +397,31 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_interrupt_resolves_its_stage_the_way_a_steer_does() {
|
||||
let controls = RunControls::new();
|
||||
assert_eq!(
|
||||
controls.interrupt(None, None).await,
|
||||
Err(SteerError::NoLiveAgent)
|
||||
);
|
||||
assert_eq!(
|
||||
controls.interrupt(Some("work"), Some("stop")).await,
|
||||
Err(SteerError::Control(ControlError::NoSuchStage(
|
||||
"work".to_string()
|
||||
)))
|
||||
);
|
||||
assert_eq!(
|
||||
controls.interrupt(Some("work@1"), None).await,
|
||||
Err(SteerError::Control(ControlError::NoSuchStage(
|
||||
"work@1".to_string()
|
||||
)))
|
||||
);
|
||||
assert_eq!(
|
||||
ControlError::NoLiveTurn.to_string(),
|
||||
"the stage has no model turn to interrupt"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_stage_label_names_its_node_visit_and_execution() {
|
||||
assert_eq!(parse_label("work@1"), Some(("work", None, 1)));
|
||||
|
|
|
|||
|
|
@ -215,9 +215,12 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
|
|||
}
|
||||
// The pause gate goes outermost, over Fabro's hooks and Petri's own,
|
||||
// so a held attempt runs none of them until the unpause.
|
||||
// The live-turn set beside it: what an interrupt can reach.
|
||||
let controls = request.controls;
|
||||
let installed = runtime.installed_hooks();
|
||||
runtime = runtime.hooks(controls.hooks(installed));
|
||||
runtime = runtime
|
||||
.hooks(controls.hooks(installed))
|
||||
.capability(controls.turns());
|
||||
|
||||
let dispatcher = InterviewDispatcher::new(request.interviewer);
|
||||
let cancel = request.cancel.clone();
|
||||
|
|
|
|||
|
|
@ -129,7 +129,7 @@ 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).await
|
||||
self.client.interrupt_run(run_id, None, None).await
|
||||
}
|
||||
|
||||
async fn steer_run(&self, run_id: &RunId, text: String, interrupt: bool) -> anyhow::Result<()> {
|
||||
|
|
|
|||
|
|
@ -1128,9 +1128,40 @@ impl Client {
|
|||
convert_type(response.into_inner())
|
||||
}
|
||||
|
||||
pub async fn interrupt_run(&self, run_id: &RunId) -> Result<()> {
|
||||
self.send_api(|client| async move {
|
||||
client.interrupt_run().id(run_id.to_string()).send().await
|
||||
/// Interrupt a run's live agent stage: stop its current model turn and
|
||||
/// keep its session. The stage is the one `stage` names (`node@visit`,
|
||||
/// or the node name) or the run's one live agent stage; `text`, when
|
||||
/// given, is the stage's next input, else the stage waits for the next
|
||||
/// steer.
|
||||
pub async fn interrupt_run(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
stage: Option<String>,
|
||||
text: Option<String>,
|
||||
) -> Result<()> {
|
||||
let stage = stage
|
||||
.map(|stage| {
|
||||
types::InterruptRunRequestStage::try_from(stage)
|
||||
.map_err(|e| anyhow!("invalid interrupt stage: {e}"))
|
||||
})
|
||||
.transpose()?;
|
||||
let text = text
|
||||
.map(|text| {
|
||||
types::InterruptRunRequestText::try_from(text)
|
||||
.map_err(|e| anyhow!("invalid interrupt text: {e}"))
|
||||
})
|
||||
.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(())
|
||||
|
|
|
|||
|
|
@ -207,6 +207,7 @@ models/integration-connection-status.ts
|
|||
models/integration-provider.ts
|
||||
models/integration-status.ts
|
||||
models/integration-webhooks-settings.ts
|
||||
models/interrupt-run-request.ts
|
||||
models/interview-option.ts
|
||||
models/interview-provider-settings.ts
|
||||
models/interview-question-record.ts
|
||||
|
|
@ -487,7 +488,6 @@ models/server-sandbox-provider-settings.ts
|
|||
models/server-sandbox-settings.ts
|
||||
models/server-scheduler-settings.ts
|
||||
models/server-settings.ts
|
||||
models/server-slate-db-settings.ts
|
||||
models/server-storage-settings.ts
|
||||
models/server-web-settings.ts
|
||||
models/session-detail.ts
|
||||
|
|
|
|||
|
|
@ -24,6 +24,8 @@ import { BASE_PATH, COLLECTION_FORMATS, type RequestArgs, BaseAPI, RequiredError
|
|||
// @ts-ignore
|
||||
import type { ErrorResponse } from '../models';
|
||||
// @ts-ignore
|
||||
import type { InterruptRunRequest } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PaginatedApiQuestionList } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PairMessageRecord } from '../models';
|
||||
|
|
@ -422,13 +424,14 @@ export const HumanInTheLoopApiAxiosParamCreator = function (configuration?: Conf
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Interrupt the active steerable agent round without sending steering text. The agent keeps its steering lease and waits for a later steer message before starting another LLM round.
|
||||
* 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.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
interruptRun: async (id: string, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
interruptRun: async (id: string, interruptRunRequest?: InterruptRunRequest, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
// verify required parameter 'id' is not null or undefined
|
||||
assertParamExists('interruptRun', 'id', id)
|
||||
const localVarPath = `/api/v1/runs/{id}/interrupt`
|
||||
|
|
@ -450,11 +453,13 @@ export const HumanInTheLoopApiAxiosParamCreator = function (configuration?: Conf
|
|||
// http bearer authentication required
|
||||
await setBearerAuthToObject(localVarHeaderParameter, configuration)
|
||||
|
||||
localVarHeaderParameter['Content-Type'] = 'application/json';
|
||||
localVarHeaderParameter['Accept'] = 'application/json';
|
||||
|
||||
setSearchParams(localVarUrlObj, localVarQueryParameter);
|
||||
let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {};
|
||||
localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers};
|
||||
localVarRequestOptions.data = serializeDataIfNeeded(interruptRunRequest, localVarRequestOptions, configuration)
|
||||
|
||||
return {
|
||||
url: toPathString(localVarUrlObj),
|
||||
|
|
@ -790,7 +795,7 @@ export const HumanInTheLoopApiAxiosParamCreator = function (configuration?: Conf
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Send a mid-run steering message to the live agent session(s) of a running run. Set `interrupt=true` to atomically interrupt the active steerable agent round first, then deliver this message as the next user turn. Without `interrupt=true`, the message is appended to the steering queue and may buffer until the next steerable agent session.
|
||||
* 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`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
@ -1005,14 +1010,15 @@ export const HumanInTheLoopApiFp = function(configuration?: Configuration) {
|
|||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Interrupt the active steerable agent round without sending steering text. The agent keeps its steering lease and waits for a later steer message before starting another LLM round.
|
||||
* 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.
|
||||
* @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, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<void>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.interruptRun(id, options);
|
||||
async interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<void>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.interruptRun(id, interruptRunRequest, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['HumanInTheLoopApi.interruptRun']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
|
|
@ -1118,7 +1124,7 @@ 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 the live agent session(s) of a running run. Set `interrupt=true` to atomically interrupt the active steerable agent round first, then deliver this message as the next user turn. Without `interrupt=true`, the message is appended to the steering queue and may buffer until the next steerable agent session.
|
||||
* 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`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
@ -1244,14 +1250,15 @@ export const HumanInTheLoopApiFactory = function (configuration?: Configuration,
|
|||
return localVarFp.getSandboxFile(id, path, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Interrupt the active steerable agent round without sending steering text. The agent keeps its steering lease and waits for a later steer message before starting another LLM round.
|
||||
* 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.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
interruptRun(id: string, options?: RawAxiosRequestConfig): AxiosPromise<void> {
|
||||
return localVarFp.interruptRun(id, options).then((request) => request(axios, basePath));
|
||||
interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig): AxiosPromise<void> {
|
||||
return localVarFp.interruptRun(id, interruptRunRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Returns pending human-in-the-loop questions for a run. Questions are generated when the workflow needs user input to proceed.
|
||||
|
|
@ -1333,7 +1340,7 @@ export const HumanInTheLoopApiFactory = function (configuration?: Configuration,
|
|||
return localVarFp.startRunPair(id, pairStartRequest, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Send a mid-run steering message to the live agent session(s) of a running run. Set `interrupt=true` to atomically interrupt the active steerable agent round first, then deliver this message as the next user turn. Without `interrupt=true`, the message is appended to the steering queue and may buffer until the next steerable agent session.
|
||||
* 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`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
@ -1459,14 +1466,15 @@ export class HumanInTheLoopApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Interrupt the active steerable agent round without sending steering text. The agent keeps its steering lease and waits for a later steer message before starting another LLM round.
|
||||
* 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.
|
||||
* @summary Interrupt Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {InterruptRunRequest} [interruptRunRequest]
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
public interruptRun(id: string, options?: RawAxiosRequestConfig) {
|
||||
return HumanInTheLoopApiFp(this.configuration).interruptRun(id, options).then((request) => request(this.axios, this.basePath));
|
||||
public interruptRun(id: string, interruptRunRequest?: InterruptRunRequest, options?: RawAxiosRequestConfig) {
|
||||
return HumanInTheLoopApiFp(this.configuration).interruptRun(id, interruptRunRequest, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -1556,7 +1564,7 @@ export class HumanInTheLoopApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Send a mid-run steering message to the live agent session(s) of a running run. Set `interrupt=true` to atomically interrupt the active steerable agent round first, then deliver this message as the next user turn. Without `interrupt=true`, the message is appended to the steering queue and may buffer until the next steerable agent session.
|
||||
* 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`).
|
||||
* @summary Steer Run
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {SteerRunRequest} steerRunRequest
|
||||
|
|
|
|||
|
|
@ -177,6 +177,7 @@ export * from './integration-connection-status';
|
|||
export * from './integration-provider';
|
||||
export * from './integration-status';
|
||||
export * from './integration-webhooks-settings';
|
||||
export * from './interrupt-run-request';
|
||||
export * from './interview-option';
|
||||
export * from './interview-provider-settings';
|
||||
export * from './interview-question-record';
|
||||
|
|
|
|||
29
lib/packages/fabro-api-client/src/models/interrupt-run-request.ts
generated
Normal file
29
lib/packages/fabro-api-client/src/models/interrupt-run-request.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.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* Request body for interrupting a live agent stage\'s model turn.
|
||||
*/
|
||||
export interface InterruptRunRequest {
|
||||
/**
|
||||
* The agent stage to interrupt: its stage identifier (`node_id@visit`) or its node name. Omit it to interrupt the run\'s one live agent stage.
|
||||
*/
|
||||
'stage'?: string;
|
||||
/**
|
||||
* The stage\'s next input once its turn is stopped. Omit it to let the stage wait for the next steer message.
|
||||
*/
|
||||
'text'?: string;
|
||||
}
|
||||
|
|
@ -23,7 +23,7 @@ export interface SteerRunRequest {
|
|||
*/
|
||||
'text': string;
|
||||
/**
|
||||
* When true, apply a worker-control interrupt first, then deliver this text as steering in the same control operation. When false (default), append to the steering queue and let the agent pick it up at the next turn boundary.
|
||||
* When true, stop the stage\'s current model turn first and make this text its next input, in one control. When false (default), the text is guidance the agent runs as a follow-up turn once its current answer is reached.
|
||||
*/
|
||||
'interrupt'?: boolean;
|
||||
/**
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue