diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index 3ea33cd..c10017d 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -85,7 +85,7 @@ impl Hydrant { // 6. re-queue any repos that lost their backfill state, then start the retry worker #[cfg(feature = "indexer")] { - if let Err(e) = state + match state .db .run({ let state = state.clone(); @@ -93,8 +93,14 @@ impl Hydrant { }) .await { - error!(err = %e, "failed to queue gone backfills"); - state.db.poison.check_report(&e); + Ok(()) => {} + // an embedder shutting its runtime down mid-startup, and the next start + // queues them anyway + Err(e) if e.downcast_ref::().is_some() => {} + Err(e) => { + error!(err = %e, "failed to queue gone backfills"); + state.db.poison.check_report(&e); + } } std::thread::spawn({ diff --git a/src/db/mod.rs b/src/db/mod.rs index 1163088..d62d9e6 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -111,17 +111,26 @@ pub(crate) fn record_lock_index_for_did(did: &jacquard_common::types::did::Did) } } +/// [`Db::run`] work that never ran, because its runtime shut down first +#[derive(Debug, miette::Diagnostic, thiserror::Error)] +#[error("the runtime shut down before this database work ran")] +pub(crate) struct Cancelled; + impl Db { - /// runs synchronous database work without blocking the async runtime. + /// runs synchronous database work without blocking the async runtime. fails with + /// [`Cancelled`] if the runtime shut down before the work started pub async fn run(&self, f: F) -> Result where T: Send + 'static, F: FnOnce(&Db) -> Result + Send + 'static, { let db = self.clone(); - tokio::task::spawn_blocking(move || f(&db)) - .await - .into_diagnostic()? + match tokio::task::spawn_blocking(move || f(&db)).await { + Ok(result) => result, + // nothing holds the handle to abort it, so only a runtime shutting down cancels it + Err(e) if e.is_cancelled() => Err(Cancelled.into()), + Err(e) => Err(e).into_diagnostic(), + } } pub fn persist(&self) -> Result<()> { @@ -362,6 +371,22 @@ mod tests { .expect("a poisoned error should trip every clone"); } + #[test] + fn work_a_shutdown_cut_off_says_so() -> Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let db = Db::open(&Config { + database_path: tmp.path().to_path_buf(), + ..Config::default() + })?; + let runtime = tokio::runtime::Runtime::new().into_diagnostic()?; + let handle = runtime.handle().clone(); + runtime.shutdown_background(); + + let err = handle.block_on(db.run(|_| Ok(()))).unwrap_err(); + assert!(err.downcast_ref::().is_some(), "{err}"); + Ok(()) + } + #[tokio::test] async fn test_db_compact_concurrency_guard() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?;