use std::future::Future;
use std::sync::OnceLock;
use std::time::{Duration, Instant};
pub const EXIT_MARGIN: Duration = Duration::from_secs(1);
static SHUTDOWN_REQUESTED: OnceLock<Instant> = OnceLock::new();
pub fn note_shutdown_requested() {
let _ = SHUTDOWN_REQUESTED.set(Instant::now());
}
pub fn teardown_budget(requested: Option<Instant>, now: Instant, grace: Duration) -> Duration {
let start = requested.unwrap_or(now);
let spent = now.saturating_duration_since(start);
grace.saturating_sub(EXIT_MARGIN).saturating_sub(spent)
}
pub fn run_main<F, T>(fut: F) -> anyhow::Result<T>
where
F: Future<Output = T>,
{
block_on_bounded(fut, || {
teardown_budget(
SHUTDOWN_REQUESTED.get().copied(),
Instant::now(),
trusty_common::shutdown::termination_grace(),
)
})
}
pub fn block_on_bounded<F, T>(fut: F, budget: impl FnOnce() -> Duration) -> anyhow::Result<T>
where
F: Future<Output = T>,
{
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
let out = rt.block_on(fut);
let budget = budget();
let started = Instant::now();
rt.shutdown_timeout(budget);
if started.elapsed() >= budget {
tracing::warn!(
budget_ms = budget.as_millis(),
"#8314: runtime teardown spent its budget; blocking work still \
running (a redb write or commit) was abandoned. redb commits are \
atomic, so the store reopens at its last committed state"
);
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_parked_blocking_task_does_not_hold_the_process_open() {
let (done_tx, done_rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let (_never_tx, never_rx) = std::sync::mpsc::channel::<()>();
let out = block_on_bounded(
async move {
tokio::task::spawn_blocking(move || {
let _ = never_rx.recv();
});
7_u8
},
|| Duration::from_millis(200),
);
let _ = done_tx.send(out.ok());
});
let out = done_rx
.recv_timeout(Duration::from_secs(5))
.expect("runtime teardown waited on a parked blocking task (#8314)");
assert_eq!(out, Some(7));
}
#[test]
fn a_finite_commit_longer_than_a_second_completes_before_exit() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
let committed = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&committed);
run_main(async move {
tokio::task::spawn_blocking(move || {
std::thread::sleep(Duration::from_millis(1_500));
flag.store(true, Ordering::SeqCst);
});
})
.expect("runtime builds");
assert!(
committed.load(Ordering::SeqCst),
"teardown abandoned a commit that would have finished (#8314)"
);
}
#[test]
fn teardown_budget_spends_what_the_grace_window_has_left() {
let grace = Duration::from_secs(60);
let signal = Instant::now();
let cases = [
(None, Duration::ZERO, Duration::from_secs(59)),
(
Some(signal),
Duration::from_secs(50),
Duration::from_secs(9),
),
(Some(signal), Duration::from_secs(59), Duration::ZERO),
(Some(signal), Duration::from_secs(70), Duration::ZERO),
];
for (requested, elapsed, want) in cases {
let now = signal + elapsed;
assert_eq!(
teardown_budget(requested, now, grace),
want,
"requested={requested:?} elapsed={elapsed:?}"
);
}
}
}