diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 9c9bb0f70..d72260531 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -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, diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index a28b82ded..14acc08c4 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -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, + run_control: Arc, + }, + Petri(Arc), +} + pub(super) struct WorkerControlManagerHandle { first_connection: Option>>, fatal: Option>, @@ -403,8 +416,7 @@ pub(super) fn spawn_worker_control_manager( worker_token: String, interviewer: Arc, cancel_token: CancellationToken, - steering_hub: Arc, - run_control: Arc, + 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, cancel_token: CancellationToken, - steering_hub: Arc, - run_control: Arc, + controls: WorkerControls, done: CancellationToken, first_tx: oneshot::Sender>, fatal_tx: oneshot::Sender, @@ -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) -> 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, ) diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 4563305b5..3a4a0a85f 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -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, 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))), diff --git a/lib/components/fabro-petri/src/controls.rs b/lib/components/fabro-petri/src/controls.rs new file mode 100644 index 000000000..d056eb6d1 --- /dev/null +++ b/lib/components/fabro-petri/src/controls.rs @@ -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), + /// 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, + 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>, +} + +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 { + self.service.paused_changes() + } + + /// The names of the agent stages running now. + #[must_use] + pub fn live_agents(&self) -> Vec { + 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 { + 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 { + 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 { + 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()); + } +} diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index cfb6e06f8..a1ed6daff 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -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, /// The caller's observers of every record, registered ahead of the @@ -193,6 +201,11 @@ pub async fn run(request: RunRequest) -> Result { if let Some(hooks) = &fabro_hooks { runtime = runtime.hooks(Arc::clone(hooks) as Arc); } + // 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 { 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 { })); }; let mut observers = request.observers; + observers.push(controls.observer()); observers.push(Arc::new(dispatcher.clone())); let result = match request.execution { Execution::Start(graphs) => { diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index fe3a64df2..1b9308bda 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -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; diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 3f227d033..cd7b2bd6a 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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, diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index 987b4dff8..95e2adeb5 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -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,