Wire pause, unpause and steer into the Petri worker

A Petri run answered only cancel and answers; pause, unpause and steer
were ignored with a warning. `fabro_petri::controls::RunControls` now
wraps Petri's `ControlService` per run: `engine::run` installs its pause
gate over the run's hooks, observes the run through it and wires it to
the coordinator, on a start and a resume alike, so a run paused when its
worker died resumes paused.

The worker's control channel takes a `WorkerControls` enum: the legacy
hub and pause flag, or the Petri run's controls. On Petri, `run.pause`
holds admission, `run.unpause` releases it once the record is durable,
and `run.steer` goes to the one live agent stage (Fabro's steer names no
stage); with none or several it is refused with a `run.notice` record.
The paused state is mirrored to Fabro's lifecycle as `run.paused` and
`run.unpaused` events, so the server's live status and the projection
follow Petri's own records.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 07:49:54 -04:00
parent 806ac90c98
commit 5093265efb
No known key found for this signature in database
8 changed files with 498 additions and 79 deletions

View file

@ -17,15 +17,24 @@
//! (`run.starting`, `run.running`, then `run.completed` or `run.failed`)
//! through the client, as the legacy worker does.
//!
//! Of the server's controls, cancel and answers are wired: the control
//! channel's cancel and `SIGTERM`/`SIGINT` fire one token, which cancels
//! Petri's root invocation politely, and an `interview.answer` message
//! reaches the control interviewer the run's questions wait on
//! (`fabro_petri::interview`), so a human gate answered through the API
//! continues. Pause, unpause and steer are received and ignored with a
//! warning until their Petri adapters land. 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.
//! The server's controls arrive over the control channel and go to Petri
//! through [`PetriControls`]: cancel (and `SIGTERM`/`SIGINT`) fires one
//! token, which cancels Petri's root invocation politely; an
//! `interview.answer` message reaches the control interviewer the run's
//! questions wait on (`fabro_petri::interview`), so a human gate answered
//! through the API continues; pause and unpause hold and release admission
//! through the run's [`RunControls`]; a steer goes to the run's one live
//! agent stage, or is refused with a `run.notice` record saying why. The
//! paused state is mirrored to Fabro's lifecycle as the legacy worker
//! reported it: a `run.paused` lifecycle event when admission is held and
//! `run.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; the
//! `SIGUSR1`/`SIGUSR2` pause signals reach only the legacy executor. 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.
//!
//! Fabro's hooks ride the run with their platform records over the same
//! client: the checkpoint commit in the run's host workspace before every
@ -52,9 +61,10 @@ use std::time::Instant;
use anyhow::{Context, Result, anyhow, bail};
use fabro_auth::VaultCredentialSource;
use fabro_client::{Client, ServerTarget};
use fabro_interview::ControlInterviewer;
use fabro_interview::{ControlInterviewer, WorkerControlMessage};
use fabro_llm::credentials::{CredentialProvider, readiness};
use fabro_petri::blobs::ClientBlobs;
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::interview::{Approval, EventSinkQuestions, FabroInterviewer};
@ -66,18 +76,19 @@ use fabro_petri::{HttpRunStore, admission};
use fabro_static::EnvVars;
use fabro_store::RunProjection;
use fabro_types::settings::run::{ApprovalMode, RunMode};
use fabro_types::{FailureReason, RunId, RunTiming, StageOutcome, SuccessReason};
use fabro_types::{FailureReason, RunId, RunNoticeLevel, RunTiming, StageOutcome, SuccessReason};
use fabro_vault::Vault;
use fabro_workflow::Error as WorkflowError;
use fabro_workflow::event::{self as workflow_event, Emitter, Event, RunEventSink};
use fabro_workflow::event::{self as workflow_event, Event, RunEventSink};
use fabro_workflow::run_control::RunControlState;
use fabro_workflow::runtime_store::RunStoreHandle;
use fabro_workflow::services::FabroRunToolServices;
use tokio::sync::RwLock as AsyncRwLock;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use super::runner::{self, WorkerTitlePhase};
use super::runner::{self, WorkerControls, WorkerTitlePhase};
use crate::args::RunWorkerMode;
use crate::command_context;
@ -123,27 +134,25 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
let run_control = RunControlState::new();
runner::install_signal_handlers(Arc::clone(&run_control), cancel_token.clone())?;
let interviewer = Arc::new(ControlInterviewer::new());
let emitter = Arc::new(Emitter::new(run_id));
let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::clone(&emitter)));
let sink = RunEventSink::map(
runner::stamp_system_worker,
RunEventSink::backend(worker.run_store.clone()),
);
let controls = RunControls::new();
let petri_controls = Arc::new(PetriControls {
run_id,
controls: controls.clone(),
sink: sink.clone(),
});
let mut control_manager = runner::spawn_worker_control_manager(
worker.target.clone(),
run_id,
worker.worker_token.to_owned(),
Arc::clone(&interviewer),
cancel_token.clone(),
steering_hub,
run_control,
WorkerControls::Petri(petri_controls),
);
control_manager.wait_for_first_connection().await?;
warn!(
run_id = %run_id,
"a Petri run answers cancel and questions only: pause, unpause and steer are not wired \
yet and are ignored"
);
let sink = RunEventSink::map(
runner::stamp_system_worker,
RunEventSink::backend(worker.run_store.clone()),
);
let approval = if worker.run_state.spec.settings.run.execution.approval == ApprovalMode::Auto {
Approval::Auto
} else {
@ -206,6 +215,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
.provider
.clone(),
cancel: cancel_token.clone(),
controls: controls.clone(),
interviewer: Arc::new(petri_interviewer),
observers,
secrets: Some(Arc::new(secrets)),
@ -215,6 +225,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
))),
hooks: Some(hooks),
};
let paused_mirror = mirror_paused_state(run_id, &controls, sink.clone());
let run = Box::pin(engine::run(request));
tokio::pin!(run);
let mut control_lost = None;
@ -231,6 +242,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
}
};
control_manager.finish();
paused_mirror.abort();
let timing = RunTiming {
wall_time_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
@ -285,6 +297,117 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
}
}
/// The Petri run's controls as the worker's control channel drives them.
/// Cancel and answers are applied by the channel itself, before a message
/// reaches here.
pub(super) struct PetriControls {
run_id: RunId,
controls: RunControls,
/// Where a refused steer's notice goes.
sink: RunEventSink,
}
impl PetriControls {
pub(super) async fn apply(&self, message: WorkerControlMessage) {
match message {
WorkerControlMessage::RunPause => {
info!(run_id = %self.run_id, "pause requested: admission is held");
self.controls.pause();
}
WorkerControlMessage::RunUnpause => {
self.controls.unpause().await;
info!(run_id = %self.run_id, "unpause recorded: admission is released");
}
WorkerControlMessage::Steer { text, actor } => {
match self.controls.steer(None, &text).await {
Ok(node) => {
info!(run_id = %self.run_id, node, actor = ?actor, "steer delivered");
}
Err(error) => {
warn!(run_id = %self.run_id, error = %error, "steer refused");
self.notice("steer_refused", error.to_string()).await;
}
}
}
WorkerControlMessage::Interrupt { .. }
| WorkerControlMessage::InterruptThenSteer { .. }
| WorkerControlMessage::PairStart { .. }
| WorkerControlMessage::PairMessage { .. }
| WorkerControlMessage::PairEnd { .. } => {
warn!(
run_id = %self.run_id,
control = control_name(&message),
"control has no Petri adapter yet and is ignored"
);
}
WorkerControlMessage::InterviewAnswer { .. } | WorkerControlMessage::RunCancel => {}
}
}
/// 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) {
let event = Event::RunNotice {
level: RunNoticeLevel::Warn,
code: code.to_string(),
message,
exec_output_tail: None,
};
if let Err(error) =
workflow_event::append_event_to_sink(&self.sink, &self.run_id, &event).await
{
warn!(run_id = %self.run_id, error = %error, "the control notice was not recorded");
}
}
}
/// The wire name of a control, for a log line.
fn control_name(message: &WorkerControlMessage) -> &'static str {
match message {
WorkerControlMessage::InterviewAnswer { .. } => "interview.answer",
WorkerControlMessage::RunCancel => "run.cancel",
WorkerControlMessage::RunPause => "run.pause",
WorkerControlMessage::RunUnpause => "run.unpause",
WorkerControlMessage::Steer { .. } => "run.steer",
WorkerControlMessage::Interrupt { .. } => "run.interrupt",
WorkerControlMessage::InterruptThenSteer { .. } => "run.interrupt_then_steer",
WorkerControlMessage::PairStart { .. } => "pair.start",
WorkerControlMessage::PairMessage { .. } => "pair.message",
WorkerControlMessage::PairEnd { .. } => "pair.end",
}
}
/// Mirror the run's paused state to Fabro's lifecycle: `run.paused` when
/// admission is held (a pause, or a resume that came back paused) and
/// `run.unpaused` when it is released, each once per change, with the
/// worker's title alongside. Aborted with the run.
fn mirror_paused_state(
run_id: RunId,
controls: &RunControls,
sink: RunEventSink,
) -> JoinHandle<()> {
let mut changes = controls.paused_changes();
tokio::spawn(async move {
let mut last = *changes.borrow_and_update();
while changes.changed().await.is_ok() {
let paused = *changes.borrow_and_update();
if paused == last {
continue;
}
last = paused;
let (event, phase) = if paused {
(Event::RunPaused, WorkerTitlePhase::Paused)
} else {
(Event::RunUnpaused, WorkerTitlePhase::Running)
};
if let Err(error) = workflow_event::append_event_to_sink(&sink, &run_id, &event).await {
warn!(run_id = %run_id, error = %error, "the paused state was not reported");
}
runner::set_worker_title(&run_id, phase);
}
})
}
/// A test's checkpoint gate directory, when the server forwarded one.
#[expect(
clippy::disallowed_methods,

View file

@ -44,7 +44,7 @@ use tokio_tungstenite::tungstenite::protocol::{self, Message as WebSocketMessage
use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async, tungstenite};
use tokio_util::sync::CancellationToken;
use super::petri_worker::{self, PetriWorker};
use super::petri_worker::{self, PetriControls, PetriWorker};
use crate::args::RunWorkerMode;
use crate::shared::github::build_github_credentials;
use crate::{command_context, server_client};
@ -131,8 +131,10 @@ pub(crate) async fn execute(
worker_token.to_owned(),
Arc::clone(&interviewer),
cancel_token.clone(),
Arc::clone(&steering_hub),
Arc::clone(&run_control),
WorkerControls::Legacy {
steering_hub: Arc::clone(&steering_hub),
run_control: Arc::clone(&run_control),
},
))
};
if let Some(control_manager) = &mut control_manager {
@ -303,6 +305,17 @@ impl AppliedWorkerControlDeliveryIds {
}
}
/// Where the run's pause, unpause, steer and pair controls go: the legacy
/// executor's hub and pause flag, or the Petri run's controls. Cancel and
/// answers are applied by the channel itself, the same way for both.
pub(super) enum WorkerControls {
Legacy {
steering_hub: Arc<fabro_workflow::SteeringHub>,
run_control: Arc<RunControlState>,
},
Petri(Arc<PetriControls>),
}
pub(super) struct WorkerControlManagerHandle {
first_connection: Option<oneshot::Receiver<Result<()>>>,
fatal: Option<oneshot::Receiver<anyhow::Error>>,
@ -403,8 +416,7 @@ pub(super) fn spawn_worker_control_manager(
worker_token: String,
interviewer: Arc<ControlInterviewer>,
cancel_token: CancellationToken,
steering_hub: Arc<fabro_workflow::SteeringHub>,
run_control: Arc<RunControlState>,
controls: WorkerControls,
) -> WorkerControlManagerHandle {
let (first_tx, first_rx) = oneshot::channel();
let (fatal_tx, fatal_rx) = oneshot::channel();
@ -417,8 +429,7 @@ pub(super) fn spawn_worker_control_manager(
worker_token,
interviewer,
cancel_token,
steering_hub,
run_control,
controls,
task_done,
first_tx,
fatal_tx,
@ -443,8 +454,7 @@ async fn run_worker_control_manager(
worker_token: String,
interviewer: Arc<ControlInterviewer>,
cancel_token: CancellationToken,
steering_hub: Arc<fabro_workflow::SteeringHub>,
run_control: Arc<RunControlState>,
controls: WorkerControls,
done: CancellationToken,
first_tx: oneshot::Sender<Result<()>>,
fatal_tx: oneshot::Sender<anyhow::Error>,
@ -485,8 +495,7 @@ async fn run_worker_control_manager(
&mut socket,
&interviewer,
&cancel_token,
&steering_hub,
&run_control,
&controls,
&mut applied_ids,
&done,
)
@ -656,8 +665,7 @@ async fn handle_worker_control_socket(
socket: &mut WorkerControlSocket,
interviewer: &ControlInterviewer,
cancel_token: &CancellationToken,
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
controls: &WorkerControls,
applied_ids: &mut AppliedWorkerControlDeliveryIds,
done: &CancellationToken,
) -> Result<(), WorkerControlConnectError> {
@ -703,8 +711,7 @@ async fn handle_worker_control_socket(
apply_worker_control_delivery_frame(
interviewer,
cancel_token,
steering_hub,
run_control,
controls,
applied_ids,
frame,
)
@ -741,8 +748,7 @@ async fn handle_worker_control_socket(
async fn apply_worker_control_delivery_frame(
interviewer: &ControlInterviewer,
cancel_token: &CancellationToken,
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
controls: &WorkerControls,
applied_ids: &mut AppliedWorkerControlDeliveryIds,
frame: WorkerControlDeliveryFrame,
) -> bool {
@ -753,14 +759,7 @@ async fn apply_worker_control_delivery_frame(
return false;
}
let frame_id = frame.id;
apply_worker_control_message(
interviewer,
cancel_token,
steering_hub,
run_control,
frame.envelope,
)
.await;
apply_worker_control_message(interviewer, cancel_token, controls, frame.envelope).await;
applied_ids.record(frame_id);
true
}
@ -768,8 +767,7 @@ async fn apply_worker_control_delivery_frame(
async fn apply_worker_control_message(
interviewer: &ControlInterviewer,
cancel_token: &CancellationToken,
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
controls: &WorkerControls,
message: WorkerControlEnvelope,
) {
match message.message {
@ -782,6 +780,24 @@ async fn apply_worker_control_message(
cancel_token.cancel();
interviewer.interrupt_all().await;
}
other => match controls {
WorkerControls::Legacy {
steering_hub,
run_control,
} => apply_legacy_control(steering_hub, run_control, other),
WorkerControls::Petri(petri) => petri.apply(other).await,
},
}
}
/// The legacy executor's pause flag and steering hub.
fn apply_legacy_control(
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
message: WorkerControlMessage,
) {
match message {
WorkerControlMessage::InterviewAnswer { .. } | WorkerControlMessage::RunCancel => {}
WorkerControlMessage::RunPause => {
run_control.request_pause();
}
@ -1212,11 +1228,11 @@ mod tests {
use super::{
AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket,
WorkerTitlePhase, apply_worker_control_delivery_frame, apply_worker_control_message,
build_worker_control_stream_request, connect_worker_control_stream,
handle_worker_control_socket, initial_worker_title_phase, load_worker_vault,
next_worker_control_reconnect_backoff, stamp_system_worker, worker_title,
worker_title_phase_for_event,
WorkerControls, WorkerTitlePhase, apply_worker_control_delivery_frame,
apply_worker_control_message, build_worker_control_stream_request,
connect_worker_control_stream, handle_worker_control_socket, initial_worker_title_phase,
load_worker_vault, next_worker_control_reconnect_backoff, stamp_system_worker,
worker_title, worker_title_phase_for_event,
};
use crate::args::RunWorkerMode;
@ -1225,6 +1241,13 @@ mod tests {
Arc::new(fabro_workflow::SteeringHub::new(emitter))
}
fn test_controls(run_control: &Arc<RunControlState>) -> WorkerControls {
WorkerControls::Legacy {
steering_hub: test_steering_hub(),
run_control: Arc::clone(run_control),
}
}
#[test]
fn clone_sandbox_credentials_are_required_for_clone_based_providers() {
use fabro_types::SandboxProviderKind;
@ -1466,12 +1489,11 @@ mod tests {
let ask_interviewer = Arc::clone(&interviewer);
let answer_task = tokio::spawn(async move { ask_interviewer.ask(question).await });
let hub = test_steering_hub();
let controls = test_controls(&run_control);
apply_worker_control_message(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
WorkerControlEnvelope::interview_answer(
"q-1",
fabro_interview::AnswerSubmission::system(
@ -1498,12 +1520,11 @@ mod tests {
let answer_task = tokio::spawn(async move { ask_interviewer.ask(question).await });
tokio::task::yield_now().await;
let hub = test_steering_hub();
let controls = test_controls(&run_control);
apply_worker_control_message(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
WorkerControlEnvelope::cancel_run(),
)
.await;
@ -1518,13 +1539,12 @@ mod tests {
let interviewer = Arc::new(ControlInterviewer::new());
let cancel_token = CancellationToken::new();
let run_control = RunControlState::new();
let hub = test_steering_hub();
let controls = test_controls(&run_control);
apply_worker_control_message(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
WorkerControlEnvelope::pause_run(),
)
.await;
@ -1533,8 +1553,7 @@ mod tests {
apply_worker_control_message(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
WorkerControlEnvelope::unpause_run(),
)
.await;
@ -1546,7 +1565,7 @@ mod tests {
let interviewer = Arc::new(ControlInterviewer::new());
let cancel_token = CancellationToken::new();
let run_control = RunControlState::new();
let hub = test_steering_hub();
let controls = test_controls(&run_control);
let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
let frame = fabro_interview::WorkerControlDeliveryFrame {
id: "local:1".to_string(),
@ -1557,8 +1576,7 @@ mod tests {
apply_worker_control_delivery_frame(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
&mut applied_ids,
frame.clone(),
)
@ -1568,8 +1586,7 @@ mod tests {
!apply_worker_control_delivery_frame(
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
&mut applied_ids,
frame,
)
@ -1655,8 +1672,8 @@ mod tests {
let mut socket = WorkerControlSocket::Test(Box::new(worker_ws));
let interviewer = Arc::new(ControlInterviewer::new());
let cancel_token = CancellationToken::new();
let hub = test_steering_hub();
let run_control = RunControlState::new();
let controls = test_controls(&run_control);
let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
let done = CancellationToken::new();
@ -1665,8 +1682,7 @@ mod tests {
&mut socket,
&interviewer,
&cancel_token,
&hub,
&run_control,
&controls,
&mut applied_ids,
&done,
)

View file

@ -36,6 +36,7 @@ use fabro_config::{Home, SettingsLayer, Storage};
use fabro_interview::ControlInterviewer;
use fabro_llm::selection;
use fabro_petri::check::{self, Bundle, CheckError, CheckRequest, Diagnostic, Launch};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::interview::{Approval, DatabaseQuestions, FabroInterviewer};
@ -403,6 +404,9 @@ 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(),
interviewer: Arc::new(petri_interviewer),
observers,
secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))),

View file

@ -0,0 +1,253 @@
//! The controls Fabro drives on a live Petri run: pause and unpause at
//! admission, a steer into the run's agent stage, 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
//! [`RunRequest`], which installs the service's pause gate over the run's
//! hooks (Fabro's own [`FabroHooks`] over Petri's local hook service),
//! observes the run through it, and wires it to the coordinator once the
//! coordinator exists. A start and a resume install it the same way, so a
//! run that was paused when its worker died resumes paused: the service is
//! handed the replayed coordinator state before the first attempt is
//! admitted, and admission stays held until an unpause arrives through the
//! new worker's control channel.
//!
//! What each control does, and what the run's record says of it:
//!
//! - pause holds every attempt not yet admitted, at once; the coordinator
//! records `run.paused`. Running work continues to its end.
//! - unpause records `run.unpaused` first and releases admission once the
//! record is durable, so a crash between the two resumes paused.
//! - steer delivers a text to a live agent stage as guidance for its session:
//! the stage's firing records `control.requested` with the `{"$steer": …}`
//! value, and the agent runs the text as a follow-up turn once its current
//! answer is reached. Fabro's steer names no stage, so the steer goes to the
//! one live agent stage; with none, or several, it is refused with the
//! reason, and nothing is recorded.
//! - cancel is the caller's cancellation token ([`RunRequest::cancel`]); the
//! service's own cancel is here for a host that holds only this.
//!
//! The paused state is published as a watch ([`RunControls::paused_changes`])
//! so the worker can mirror it to Fabro's lifecycle (`run.paused` and
//! `run.unpaused` lifecycle events), including the flip a resume makes.
//!
//! [`engine::run`]: crate::engine::run
//! [`RunRequest`]: crate::engine::RunRequest
//! [`RunRequest::cancel`]: crate::engine::RunRequest::cancel
//! [`FabroHooks`]: crate::hooks::FabroHooks
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use petri_execution::controls::{ControlError, ControlService};
use petri_execution::{
CoordinatorHandle, CoordinatorRecord, CoordinatorState, ExecutionId, ExecutionObserver,
};
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 tokio::sync::watch;
/// Why a steer 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.
#[error("Run has several active agent stages ({}); the steer names none.", .0.join(", "))]
SeveralLiveAgents(Vec<String>),
/// The named stage is not running, or the run has ended.
#[error(transparent)]
Control(#[from] ControlError),
}
/// The live agent firings, by node name: what a steer that names no stage
/// is routed by.
#[derive(Default)]
struct LiveAgents {
stages: BTreeMap<String, (ExecutionId, FiringId)>,
firings: BTreeMap<(ExecutionId, FiringId), String>,
}
/// One run's controls. Clone freely: every clone drives the same service.
#[derive(Clone)]
pub struct RunControls {
service: ControlService,
agents: Arc<Mutex<LiveAgents>>,
}
impl Default for RunControls {
fn default() -> Self {
Self::new()
}
}
impl RunControls {
#[must_use]
pub fn new() -> Self {
Self {
service: ControlService::new(),
agents: Arc::new(Mutex::new(LiveAgents::default())),
}
}
/// Hold every attempt not yet admitted. Running work is not interrupted.
pub fn pause(&self) {
self.service.pause();
}
/// Release held and future attempts, once the unpause is durable.
pub async fn unpause(&self) {
self.service.unpause().await;
}
#[must_use]
pub fn is_paused(&self) -> bool {
self.service.is_paused()
}
/// Every change of the paused state, the flip a resume makes included.
#[must_use]
pub fn paused_changes(&self) -> watch::Receiver<bool> {
self.service.paused_changes()
}
/// The names of the agent stages running now.
#[must_use]
pub fn live_agents(&self) -> Vec<String> {
self.agents().stages.keys().cloned().collect()
}
/// Deliver `text` to the named agent stage, or to the one live agent
/// stage when `node` is `None`. The name of the stage steered.
pub async fn steer(&self, node: Option<&str>, text: &str) -> Result<String, SteerError> {
let node = if let Some(node) = node {
node.to_owned()
} else {
let mut live = self.live_agents();
match live.len() {
0 => return Err(SteerError::NoLiveAgent),
1 => live.remove(0),
_ => return Err(SteerError::SeveralLiveAgents(live)),
}
};
self.service.steer(&node, text).await?;
Ok(node)
}
/// Cancel the whole run politely; a second call reaches the kill tier.
pub fn cancel(&self) -> Result<(), ControlError> {
self.service.cancel()
}
/// The pause gate over `inner`, for the runtime.
pub(crate) fn hooks(&self, inner: Option<Arc<dyn ExecutionHooks>>) -> Arc<dyn ExecutionHooks> {
self.service.hooks(inner)
}
/// Hand the service the run's coordinator handle.
pub(crate) fn wire(&self, handle: CoordinatorHandle) {
self.service.wire(handle);
}
/// The observer to register on the run: the service's own, which keeps
/// the live firing of every stage and the paused state across a
/// resume, and the live agent stages beside it.
pub(crate) fn observer(&self) -> Arc<dyn ExecutionObserver> {
Arc::new(self.clone())
}
fn agents(&self) -> MutexGuard<'_, LiveAgents> {
self.agents.lock().unwrap_or_else(PoisonError::into_inner)
}
}
impl ExecutionObserver for RunControls {
fn on_engine_record(
&self,
execution: ExecutionId,
record: &EventRecord,
recorded_at: u64,
state: &EngineState,
) {
self.service
.on_engine_record(execution, record, recorded_at, state);
match &record.event {
Event::StepStarted { firing, .. } => {
let Some(node) = state
.firing_node(*firing)
.and_then(|id| state.graph().node(id))
else {
return;
};
if node.step.kind != AGENT_KIND {
return;
}
let name = node.name.to_string();
let mut agents = self.agents();
agents.stages.insert(name.clone(), (execution, *firing));
agents.firings.insert((execution, *firing), name);
}
Event::StepFinished { firing, .. } => {
let mut agents = self.agents();
if let Some(name) = agents.firings.remove(&(execution, *firing)) {
if agents.stages.get(&name) == Some(&(execution, *firing)) {
agents.stages.remove(&name);
}
}
}
_ => {}
}
}
fn on_lifecycle(&self, record: &CoordinatorRecord) {
self.service.on_lifecycle(record);
}
fn on_resumed(&self, state: &CoordinatorState) {
self.service.on_resumed(state);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn a_steer_with_no_live_agent_is_refused_with_the_legacy_reason() {
let controls = RunControls::new();
assert_eq!(
controls.steer(None, "hurry up").await,
Err(SteerError::NoLiveAgent)
);
assert_eq!(
SteerError::NoLiveAgent.to_string(),
"Run has no active steerable agent session."
);
}
#[tokio::test]
async fn a_steer_to_a_named_stage_that_is_not_running_is_refused() {
let controls = RunControls::new();
assert_eq!(
controls.steer(Some("work"), "hurry up").await,
Err(SteerError::Control(ControlError::NoSuchStage(
"work".to_string()
)))
);
}
#[test]
fn a_pause_holds_before_the_run_is_wired() {
let controls = RunControls::new();
let mut changes = controls.paused_changes();
assert!(!controls.is_paused());
controls.pause();
assert!(controls.is_paused());
assert!(changes.has_changed().expect("the sender is alive"));
assert!(*changes.borrow_and_update());
}
}

View file

@ -26,7 +26,11 @@
//! give the run: Petri's local hook service for `[[run.hooks]]`, no host
//! tools, and `Retention::Always` for every workspace, Fabro's default.
//! Cancellation rides the caller's token: when it fires, the root
//! invocation is cancelled politely and Petri records why.
//! invocation is cancelled politely and Petri records why. The run's other
//! controls (pause, unpause, steer) are the caller's [`RunControls`]: its
//! pause gate is installed over the run's hooks, it observes the run, and
//! it is wired to the coordinator with the interviewer, on a start and on
//! a resume alike, so a run that was paused resumes paused.
//!
//! A resume here is Petri's own: the run continues from its records, and
//! sandbox leases are reconciled by label. What the workspaces look like
@ -57,6 +61,7 @@ use tracing::{debug, info, warn};
use crate::admission::AdmittedGraphs;
use crate::blobs::{Blobs, RunBlobs};
use crate::controls::RunControls;
use crate::hooks::{FabroHooks, HooksSpec};
use crate::runtime::RuntimeSpec;
use crate::secrets::SharedSecrets;
@ -88,6 +93,9 @@ pub struct RunRequest {
pub provider: SandboxProviderKind,
/// Fires to cancel the run.
pub cancel: CancellationToken,
/// The run's pause, unpause and steer controls, which the caller keeps
/// a clone of to drive them while the run is live.
pub controls: RunControls,
/// Where the run's questions go.
pub interviewer: Arc<dyn Interviewer>,
/// The caller's observers of every record, registered ahead of the
@ -193,6 +201,11 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
if let Some(hooks) = &fabro_hooks {
runtime = runtime.hooks(Arc::clone(hooks) as Arc<dyn ExecutionHooks>);
}
// The pause gate goes outermost, over Fabro's hooks and Petri's own,
// so a held attempt runs none of them until the unpause.
let controls = request.controls;
let installed = runtime.installed_hooks();
runtime = runtime.hooks(controls.hooks(installed));
let dispatcher = InterviewDispatcher::new(request.interviewer);
let cancel = request.cancel.clone();
@ -202,6 +215,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
if let Some(hooks) = &fabro_hooks {
hooks.attach(handle.clone());
}
controls.wire(handle.clone());
cancel_task = Some(tokio::spawn(async move {
cancel.cancelled().await;
info!("cancelling the Petri run");
@ -209,6 +223,7 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
}));
};
let mut observers = request.observers;
observers.push(controls.observer());
observers.push(Arc::new(dispatcher.clone()));
let result = match request.execution {
Execution::Start(graphs) => {

View file

@ -44,7 +44,9 @@
//! - [`platform_records`]: Fabro's platform records as the adapters reach them,
//! in the server's database or over its API from a worker;
//! - [`host_tools`]: Fabro's run tools on every native agent session of a run,
//! through Petri's `HostTools` capability.
//! through Petri's `HostTools` capability;
//! - [`controls`]: the controls Fabro drives on a live run (pause, unpause,
//! steer, cancel), over Petri's control service.
//!
//! The Petri packages are pinned by revision in the workspace `Cargo.toml`
//! under `petri_*` keys.
@ -53,6 +55,7 @@ pub mod admission;
pub mod blobs;
pub mod check;
pub mod checkpoint;
pub mod controls;
pub mod engine;
pub mod hooks;
pub mod host_tools;

View file

@ -24,6 +24,7 @@ use fabro_checkpoint::author::GitAuthor;
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::checkpoint::{CHECKPOINT_FAILED_CLASS, CheckpointKey, RunWorkspaces};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::platform_records::PlatformRecords;
@ -142,6 +143,7 @@ impl Harness {
runtime: RuntimeSpec::default(),
provider: SandboxProviderKind::LOCAL,
cancel: CancellationToken::new(),
controls: RunControls::new(),
interviewer,
observers,
secrets: None,
@ -558,6 +560,7 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() {
},
provider: SandboxProviderKind::LOCAL,
cancel: CancellationToken::new(),
controls: RunControls::new(),
interviewer,
observers,
secrets: None,

View file

@ -15,6 +15,7 @@ use std::time::{Duration, Instant};
use fabro_petri::admission::AdmittedGraphs;
use fabro_petri::check::{self, Bundle, CheckRequest, Launch};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{Execution, RunRequest};
use fabro_petri::interview::{Approval, FabroInterviewer, QuestionNotice, QuestionSink};
use fabro_petri::runtime::RuntimeSpec;
@ -113,6 +114,7 @@ pub(crate) fn run_request(
runtime,
provider: SandboxProviderKind::LOCAL,
cancel: CancellationToken::new(),
controls: RunControls::new(),
observers: vec![interviewer.observer()],
interviewer: Arc::new(interviewer),
secrets: None,