From a533da8e80f988693a872c79478d4175e0770010 Mon Sep 17 00:00:00 2001 From: Anders Olsson Date: Fri, 7 Aug 2026 16:00:42 +0200 Subject: [PATCH] fix(supervisor): non-vacuous retry-budget test, clear stale uptime, sync gate reconfigure --- crates/xy-supervisor/src/supervisor.rs | 87 ++++++++++++++++++++++++-- 1 file changed, 83 insertions(+), 4 deletions(-) diff --git a/crates/xy-supervisor/src/supervisor.rs b/crates/xy-supervisor/src/supervisor.rs index 0518dd0..227934d 100644 --- a/crates/xy-supervisor/src/supervisor.rs +++ b/crates/xy-supervisor/src/supervisor.rs @@ -150,17 +150,20 @@ impl SupervisorTask { } Err(StartFailure::Spawn(err)) => { warn!(name = %self.cfg.name, error = %err, "spawn failed"); + self.started_at = None; self.set_state(ServerState::Failed); let _ = ack.send(StartAck::SpawnFailed(err.to_string())); } Err(StartFailure::WaitTimedOut) => { warn!(name = %self.cfg.name, "wait-for timed out"); + self.started_at = None; self.set_state(ServerState::Failed); let _ = ack.send(StartAck::SpawnFailed( "wait-for timed out".to_string(), )); } Err(StartFailure::Cancelled) => { + self.started_at = None; self.set_state(ServerState::Stopped); } Err(StartFailure::Shutdown) => return, @@ -185,13 +188,16 @@ impl SupervisorTask { Ok(c) => child = Some(c), Err(StartFailure::Spawn(err)) => { warn!(name = %self.cfg.name, error = %err, "restart spawn failed"); + self.started_at = None; self.set_state(ServerState::Failed); } Err(StartFailure::WaitTimedOut) => { warn!(name = %self.cfg.name, "wait-for timed out"); + self.started_at = None; self.set_state(ServerState::Failed); } Err(StartFailure::Cancelled) => { + self.started_at = None; self.set_state(ServerState::Stopped); } Err(StartFailure::Shutdown) => return, @@ -299,13 +305,16 @@ impl SupervisorTask { Ok(c) => child = Some(c), Err(StartFailure::Spawn(err)) => { warn!(name = %self.cfg.name, error = %err, "restart spawn failed"); + self.started_at = None; self.set_state(ServerState::Failed); } Err(StartFailure::WaitTimedOut) => { warn!(name = %self.cfg.name, "wait-for timed out"); + self.started_at = None; self.set_state(ServerState::Failed); } Err(StartFailure::Cancelled) => { + self.started_at = None; self.set_state(ServerState::Stopped); } Err(StartFailure::Shutdown) => return, @@ -372,6 +381,12 @@ impl SupervisorTask { } Some(SupervisorCmd::Reconfigure { new, ack }) => { self.cfg = *new; + self.backoff = + Backoff::new(self.cfg.restart.backoff_initial, self.cfg.restart.backoff_max); + self.retry_window = RetryWindow::new( + Duration::from_secs(60), + self.cfg.restart.max_retries_per_minute, + ); let _ = ack.send(()); None } @@ -732,7 +747,7 @@ mod tests { } #[tokio::test] - async fn a_timed_out_gate_fails_without_spawning_or_touching_the_retry_budget() { + async fn a_timed_out_gate_fails_without_spawning() { let cfg = cfg_with_wait( "x", WaitForConfig { @@ -761,12 +776,76 @@ mod tests { wait_for(&mut status_rx, ServerState::Waiting).await; wait_for(&mut status_rx, ServerState::Failed).await; + let (ack_tx, ack_rx) = oneshot::channel(); + cmd_tx + .send(SupervisorCmd::Shutdown { ack: ack_tx }) + .await + .unwrap(); + ack_rx.await.unwrap(); + h.await.unwrap(); + } + + #[tokio::test] + async fn repeated_gate_timeouts_do_not_pollute_the_retry_window() { + let tmp = tempfile::tempdir().unwrap(); + let file = tmp.path().join("sock"); + + let mut cfg = cfg("x", RestartPolicy::Always, 2); + + cfg.wait_for = Some(WaitForConfig { + condition: WaitCondition::Path(file.clone()), + timeout: Duration::from_millis(20), + interval: Duration::from_millis(5), + }); + + let (first, mut ctl) = MockChild::new(1); + let (second, ctl2) = MockChild::new(2); + let queue = Arc::new(Mutex::new(vec![first, second])); + let spawner = QueueSpawner { queue }; + + let (status_tx, mut status_rx) = watch::channel(initial_status(&cfg)); + let (cmd_tx, cmd_rx) = mpsc::channel(8); + let task = SupervisorTask::new(cfg, sink("x"), spawner, status_tx, cmd_rx); + let h = tokio::spawn(task.run()); + + let (ack_tx, _ack_rx) = oneshot::channel(); + cmd_tx + .send(SupervisorCmd::Start { ack: ack_tx }) + .await + .unwrap(); + wait_for(&mut status_rx, ServerState::Waiting).await; + wait_for(&mut status_rx, ServerState::Failed).await; + + let (ack_tx, _ack_rx) = oneshot::channel(); + cmd_tx + .send(SupervisorCmd::Start { ack: ack_tx }) + .await + .unwrap(); + wait_for(&mut status_rx, ServerState::Waiting).await; + wait_for(&mut status_rx, ServerState::Failed).await; + + std::fs::write(&file, b"").unwrap(); + + let (ack_tx, ack_rx) = oneshot::channel(); + cmd_tx + .send(SupervisorCmd::Start { ack: ack_tx }) + .await + .unwrap(); + assert_eq!(ack_rx.await.unwrap(), StartAck::Started); + wait_for(&mut status_rx, ServerState::Running).await; + + ctl.exit_tx.take().unwrap().send(Some(1)).unwrap(); + + wait_for_pid(&mut status_rx, 2).await; + assert_eq!( - status_rx.borrow().restart_count, - 0, - "waiting must not consume the restart budget" + status_rx.borrow().state, + ServerState::Running, + "the crash after two gate timeouts must still be within the retry budget" ); + drop(ctl2); + let (ack_tx, ack_rx) = oneshot::channel(); cmd_tx .send(SupervisorCmd::Shutdown { ack: ack_tx })