mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Route SIGUSR1 and SIGUSR2 in the Petri worker to the run's controls
The worker's signal handlers pause and unpause the run through its `RunControls`, the same path the server's pause and unpause take, so a signal holds admission, records Petri's `run.paused` and the lifecycle mirror, and releases it on the unpause. The legacy pause state the plan named no longer exists in the tree; nothing was left to delete. A controls scenario sends both signals to a real worker and reads the records back. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
a4af6fac92
commit
c730aca50c
3 changed files with 117 additions and 10 deletions
|
|
@ -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<dyn PlatformRecords> =
|
||||
Arc::new(HttpPlatformRecords::new(worker.client.clone_for_reuse()));
|
||||
let controls = RunControls::new();
|
||||
let petri_controls = Arc::new(PetriControls::new(
|
||||
run_id,
|
||||
controls.clone(),
|
||||
|
|
|
|||
|
|
@ -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()),
|
||||
))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue