fix(supervisor): non-vacuous retry-budget test, clear stale uptime, sync gate reconfigure
This commit is contained in:
@@ -150,17 +150,20 @@ impl<S: Spawner> SupervisorTask<S> {
|
|||||||
}
|
}
|
||||||
Err(StartFailure::Spawn(err)) => {
|
Err(StartFailure::Spawn(err)) => {
|
||||||
warn!(name = %self.cfg.name, error = %err, "spawn failed");
|
warn!(name = %self.cfg.name, error = %err, "spawn failed");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
let _ = ack.send(StartAck::SpawnFailed(err.to_string()));
|
let _ = ack.send(StartAck::SpawnFailed(err.to_string()));
|
||||||
}
|
}
|
||||||
Err(StartFailure::WaitTimedOut) => {
|
Err(StartFailure::WaitTimedOut) => {
|
||||||
warn!(name = %self.cfg.name, "wait-for timed out");
|
warn!(name = %self.cfg.name, "wait-for timed out");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
let _ = ack.send(StartAck::SpawnFailed(
|
let _ = ack.send(StartAck::SpawnFailed(
|
||||||
"wait-for timed out".to_string(),
|
"wait-for timed out".to_string(),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
Err(StartFailure::Cancelled) => {
|
Err(StartFailure::Cancelled) => {
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Stopped);
|
self.set_state(ServerState::Stopped);
|
||||||
}
|
}
|
||||||
Err(StartFailure::Shutdown) => return,
|
Err(StartFailure::Shutdown) => return,
|
||||||
@@ -185,13 +188,16 @@ impl<S: Spawner> SupervisorTask<S> {
|
|||||||
Ok(c) => child = Some(c),
|
Ok(c) => child = Some(c),
|
||||||
Err(StartFailure::Spawn(err)) => {
|
Err(StartFailure::Spawn(err)) => {
|
||||||
warn!(name = %self.cfg.name, error = %err, "restart spawn failed");
|
warn!(name = %self.cfg.name, error = %err, "restart spawn failed");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
}
|
}
|
||||||
Err(StartFailure::WaitTimedOut) => {
|
Err(StartFailure::WaitTimedOut) => {
|
||||||
warn!(name = %self.cfg.name, "wait-for timed out");
|
warn!(name = %self.cfg.name, "wait-for timed out");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
}
|
}
|
||||||
Err(StartFailure::Cancelled) => {
|
Err(StartFailure::Cancelled) => {
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Stopped);
|
self.set_state(ServerState::Stopped);
|
||||||
}
|
}
|
||||||
Err(StartFailure::Shutdown) => return,
|
Err(StartFailure::Shutdown) => return,
|
||||||
@@ -299,13 +305,16 @@ impl<S: Spawner> SupervisorTask<S> {
|
|||||||
Ok(c) => child = Some(c),
|
Ok(c) => child = Some(c),
|
||||||
Err(StartFailure::Spawn(err)) => {
|
Err(StartFailure::Spawn(err)) => {
|
||||||
warn!(name = %self.cfg.name, error = %err, "restart spawn failed");
|
warn!(name = %self.cfg.name, error = %err, "restart spawn failed");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
}
|
}
|
||||||
Err(StartFailure::WaitTimedOut) => {
|
Err(StartFailure::WaitTimedOut) => {
|
||||||
warn!(name = %self.cfg.name, "wait-for timed out");
|
warn!(name = %self.cfg.name, "wait-for timed out");
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Failed);
|
self.set_state(ServerState::Failed);
|
||||||
}
|
}
|
||||||
Err(StartFailure::Cancelled) => {
|
Err(StartFailure::Cancelled) => {
|
||||||
|
self.started_at = None;
|
||||||
self.set_state(ServerState::Stopped);
|
self.set_state(ServerState::Stopped);
|
||||||
}
|
}
|
||||||
Err(StartFailure::Shutdown) => return,
|
Err(StartFailure::Shutdown) => return,
|
||||||
@@ -372,6 +381,12 @@ impl<S: Spawner> SupervisorTask<S> {
|
|||||||
}
|
}
|
||||||
Some(SupervisorCmd::Reconfigure { new, ack }) => {
|
Some(SupervisorCmd::Reconfigure { new, ack }) => {
|
||||||
self.cfg = *new;
|
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(());
|
let _ = ack.send(());
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
@@ -732,7 +747,7 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[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(
|
let cfg = cfg_with_wait(
|
||||||
"x",
|
"x",
|
||||||
WaitForConfig {
|
WaitForConfig {
|
||||||
@@ -761,12 +776,76 @@ mod tests {
|
|||||||
wait_for(&mut status_rx, ServerState::Waiting).await;
|
wait_for(&mut status_rx, ServerState::Waiting).await;
|
||||||
wait_for(&mut status_rx, ServerState::Failed).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!(
|
assert_eq!(
|
||||||
status_rx.borrow().restart_count,
|
status_rx.borrow().state,
|
||||||
0,
|
ServerState::Running,
|
||||||
"waiting must not consume the restart budget"
|
"the crash after two gate timeouts must still be within the retry budget"
|
||||||
);
|
);
|
||||||
|
|
||||||
|
drop(ctl2);
|
||||||
|
|
||||||
let (ack_tx, ack_rx) = oneshot::channel();
|
let (ack_tx, ack_rx) = oneshot::channel();
|
||||||
cmd_tx
|
cmd_tx
|
||||||
.send(SupervisorCmd::Shutdown { ack: ack_tx })
|
.send(SupervisorCmd::Shutdown { ack: ack_tx })
|
||||||
|
|||||||
Reference in New Issue
Block a user