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 23576c3ce..f98002bf6 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -21,8 +21,9 @@ //! 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 +//! through the API continues; pause and unpause (and `SIGUSR1`/`SIGUSR2`) +//! hold and release admission through the run's [`RunControls`]; a steer +//! goes to the run's one live //! agent stage, or is refused with a `run.notice` record saying why. The //! paused state is mirrored to Fabro's lifecycle: a `paused` lifecycle //! record when admission is held and `unpaused` when it is released, so @@ -125,12 +126,12 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { )); let cancel_token = CancellationToken::new(); - runner::install_signal_handlers(cancel_token.clone())?; + let controls = RunControls::new(); + runner::install_signal_handlers(cancel_token.clone(), controls.clone())?; let interviewer = Arc::new(ControlInterviewer::new()); // Fabro's own records of the run, over the client. let records: Arc = Arc::new(HttpPlatformRecords::new(worker.client.clone_for_reuse())); - let controls = RunControls::new(); let petri_controls = Arc::new(PetriControls::new( run_id, controls.clone(), diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index 9facc453f..a7c3839b8 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -13,6 +13,7 @@ use fabro_interview::{ WorkerControlMessage, }; use fabro_manifest::SuppliedWorkflowVersionPackager; +use fabro_petri::controls::RunControls; use fabro_tool::fabro_client::ClientBackend; use fabro_types::RunId; use fabro_vault::{SecretStore, Vault}; @@ -684,8 +685,13 @@ fn worker_title(run_id: &RunId, phase: WorkerTitlePhase) -> String { format!("fabro {short_id} {phase}") } -/// `SIGTERM` and `SIGINT` cancel the run, the way the server's cancel does. -pub(super) fn install_signal_handlers(cancel_token: CancellationToken) -> Result<()> { +/// `SIGTERM` and `SIGINT` cancel the run, the way the server's cancel does; +/// `SIGUSR1` pauses it and `SIGUSR2` unpauses it, the way the server's pause +/// and unpause do, through the run's controls. +pub(super) fn install_signal_handlers( + cancel_token: CancellationToken, + controls: RunControls, +) -> Result<()> { #[cfg(unix)] { let mut terminate = signal(SignalKind::terminate())?; @@ -702,6 +708,27 @@ pub(super) fn install_signal_handlers(cancel_token: CancellationToken) -> Result cancel_token.cancel(); } }); + + let mut pause = signal(SignalKind::user_defined1())?; + let pause_controls = controls.clone(); + tokio::spawn(async move { + while pause.recv().await.is_some() { + tracing::info!("SIGUSR1: pause requested; admission is held"); + pause_controls.pause(); + } + }); + + let mut unpause = signal(SignalKind::user_defined2())?; + tokio::spawn(async move { + while unpause.recv().await.is_some() { + controls.unpause().await; + tracing::info!("SIGUSR2: unpause recorded; admission is released"); + } + }); + } + #[cfg(not(unix))] + { + let _ = (cancel_token, controls); } Ok(()) @@ -729,8 +756,8 @@ mod tests { use super::super::petri_worker::PetriControls; use super::{ - AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket, - WorkerControls, WorkerTitlePhase, apply_worker_control_delivery_frame, + AppliedWorkerControlDeliveryIds, RunControls, WorkerControlConnectError, + WorkerControlSocket, 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, worker_title, @@ -742,7 +769,7 @@ mod tests { fn test_controls() -> WorkerControls { Arc::new(PetriControls::new( fixtures::RUN_1, - fabro_petri::controls::RunControls::new(), + RunControls::new(), Arc::new(fabro_petri::test_support::MemoryPlatformRecords::new()), )) } diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs b/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs index 248a74d87..2de57fde8 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri_controls.rs @@ -1,6 +1,7 @@ //! The run controls on a Petri run through a real server and its worker: //! a pause holds the next stage until the unpause and the API says -//! `paused` in between; a steer reaches the agent stage on the twin, which +//! `paused` in between; `SIGUSR1` and `SIGUSR2` on the worker do the same +//! without the API; a steer reaches the agent stage on the twin, which //! sees it in its next request, and the stream carries the control record; //! a run paused when its server and worker die resumes paused and goes on //! once unpaused. @@ -229,6 +230,84 @@ async fn a_pause_holds_the_next_stage_until_the_unpause() { server.shutdown(); } +/// `SIGUSR1` on the worker pauses the run the way the API's pause does, +/// and `SIGUSR2` unpauses it: `b` is held at admission in between, Petri's +/// records and Fabro's lifecycle both carry the pause and the unpause, and +/// no control request is recorded, since none went through the API. +#[tokio::test(flavor = "multi_thread")] +async fn the_user_signals_pause_and_unpause_the_worker() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let server = RunningServer::start().await; + let gate = context.temp_dir.join("a.gate"); + let marker = context.temp_dir.join("b.marker"); + let workspace = two_stage_workspace(&context, &gate, &marker); + let run_id = run_detached(&context, &server, &workspace); + + wait_for_status(&server, &run_id, &["running"]).await; + let worker = wait_for_worker(&run_id); + wait_until_gate_is_polled(&gate); + eprintln!("run {run_id}: a is waiting on the gate; sending SIGUSR1 to worker {worker}"); + + fabro_proc::sigusr1(worker); + wait_for_status(&server, &run_id, &["paused"]).await; + eprintln!("run {run_id} is paused"); + + std::fs::write(&gate, "go").expect("the gate opens"); + wait_for_stream_count(&server, &run_id, "step.finished", 2).await; + tokio::time::sleep(HOLD).await; + assert!(!marker.exists(), "b started while the run was paused"); + assert_eq!(run_status(&server, &run_id).await, "paused"); + assert!(pending_control(&server, &run_id).await.is_null()); + + fabro_proc::sigusr2(worker); + 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(), "b ran after the unpause"); + assert_petri_succeeded(&server, &run_id).await; + + for event in [ + "run.paused", + "run.unpaused", + "lifecycle:paused", + "lifecycle:unpaused", + ] { + assert_eq!(count_of(&names, event), 1, "{event}: {names:?}"); + } + for event in ["lifecycle:pause_requested", "lifecycle:unpause_requested"] { + assert_eq!( + count_of(&names, event), + 0, + "a signal is not an API request: {event}: {names:?}" + ); + } + let unpaused = names + .iter() + .position(|name| name == "run.unpaused") + .expect("the unpause is recorded"); + let b_started = names + .iter() + .enumerate() + .filter(|(_, name)| *name == "step.started") + .nth(2) + .map(|(index, _)| index) + .expect("b started"); + assert!( + unpaused < b_started, + "b started before the unpause: {names:?}" + ); + server.shutdown(); +} + /// A steer sent while the agent stage waits on a tool reaches its /// session: the twin sees the steer text in the follow-up request, the /// stream carries the `control.requested` record, and the run succeeds.