From bf0d28fed8a8106cf1aacee41eb663c468ea4ba2 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 5 Oct 2026 00:34:39 +0300 Subject: [PATCH] [db] report a poisoned database through run instead of exiting Signed-off-by: dawn <90008@klbr.net> --- AGENTS.md | 2 ++ docs/embedding.md | 1 + src/api/debug.rs | 4 +-- src/backfill/manager.rs | 6 ++-- src/backfill/worker.rs | 6 ++-- src/control/hydrant.rs | 43 ++++++++++++++++++++---- src/control/hydrant/run.rs | 32 ++++++++++-------- src/control/mod.rs | 2 +- src/control/stream/indexer.rs | 2 +- src/db/keyspaces.rs | 1 + src/db/mod.rs | 61 +++++++++++++++++++++++++++++------ src/db/open.rs | 2 ++ src/db/schema.rs | 14 ++++---- src/ingest/indexer/shard.rs | 5 +-- src/main.rs | 12 +++++-- 15 files changed, 145 insertions(+), 48 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index cab2f16..8f2565f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -72,6 +72,8 @@ Non-rust programs embed hydrant through the C api in `ffi/` (`hydrant-ffi`, a wo `Hydrant::shutdown` waits for the database to close, so every std thread hydrant starts has to end once `AppState::stop` is set: timer loops sleep on `state.stop.wait(..)` instead of `thread::sleep`, and channel workers check `state.stop.is_set()` per message. A thread that doesn't keeps the database open and makes shutdown time out. +Hydrant never exits the process, since it may be embedded. Pass fjall errors to `db.poison.check(..)` (or `check_report(..)` for a `miette::Report`): once fjall says the database is poisoned, `Hydrant::run` ends with `DatabasePoisoned`, and only the binary turns that into exit code 10. + ## General conventions ### Correctness over convenience diff --git a/docs/embedding.md b/docs/embedding.md index f8bcbde..2d0453b 100644 --- a/docs/embedding.md +++ b/docs/embedding.md @@ -36,6 +36,7 @@ the socket is created owner-only (`0600`), so only the embedding program's user - `hydrant_wait` blocks until hydrant stops: when it fails (for example when its socket is already in use) it gives the reason, and after `hydrant_shutdown` it returns 0. - `hydrant_shutdown` stops hydrant and waits for its database to close, which gives its memory back. anything it was in the middle of is cut off like a crash, and the next start picks up from its saved cursors. a hydrant that failed still holds its database until it's shut down. - `hydrant_free` frees the handle once nothing is blocked in `hydrant_wait` on it. +- if a disk write fails, fjall refuses every write after it and `hydrant_wait` reports the database as poisoned. shut down and start again, which recovers it, and look at the disk if it keeps happening. - only one hydrant can use a database at a time, a second `hydrant_start` on it fails until the first one is shut down. ## go diff --git a/src/api/debug.rs b/src/api/debug.rs index fe339ac..5c76fca 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -269,9 +269,9 @@ pub async fn handle_debug_get( let partition = req.partition.clone(); let value = state .db - .run(move |_| { + .run(move |db| { ks.get(key) - .inspect_err(crate::db::check_poisoned) + .inspect_err(|e| db.poison.check(e)) .into_diagnostic() }) .await diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index 5c715da..2b9c6eb 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -86,7 +86,7 @@ pub fn queue_gone_backfills(state: &Arc) -> Result<()> { Ok(false) => {} Err(e) => { error!(did = %did, err = %e, "failed to queue gone repo"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } } } @@ -108,7 +108,7 @@ pub fn retry_worker(state: Arc) { Ok(t) => t, Err(e) => { error!(err = %e, "failed to get resync state"); - db::check_poisoned(&e); + state.db.poison.check(&e); continue; } }; @@ -165,7 +165,7 @@ pub fn retry_worker(state: Arc) { Ok(false) => {} Err(e) => { error!(did = %did, err = %e, "failed to queue retry"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } } } diff --git a/src/backfill/worker.rs b/src/backfill/worker.rs index 3baa850..b316e47 100644 --- a/src/backfill/worker.rs +++ b/src/backfill/worker.rs @@ -2,7 +2,7 @@ use crate::backfill::admission::BackfillAdmission; use crate::backfill::client::ThrottledHttpClient; use crate::backfill::error::BackfillError; use crate::config::BackfillStrategy; -use crate::db::{self, types::TrimmedDid}; +use crate::db::types::TrimmedDid; use crate::ingest::indexer::IndexerTx; use crate::net::public_http_with_timeout; use crate::state::AppState; @@ -99,7 +99,7 @@ impl BackfillWorker { Ok(kv) => kv, Err(e) => { error!(err = %e, "failed to read pending entry"); - db::check_poisoned(&e); + self.state.db.poison.check(&e); continue; } }; @@ -177,7 +177,7 @@ impl BackfillWorker { } } if let BackfillError::Generic(report) = &e { - db::check_poisoned_report(report); + state.db.poison.check_report(report); } } diff --git a/src/control/hydrant.rs b/src/control/hydrant.rs index 219555c..a571a59 100644 --- a/src/control/hydrant.rs +++ b/src/control/hydrant.rs @@ -32,6 +32,13 @@ use super::seed; pub(crate) mod run; +/// [`Hydrant::run`] fails with this once fjall refuses writes because a disk write +/// failed. every later write fails the same way, so the database has to be closed +/// and opened again, and looked at if it keeps happening. +#[derive(Debug, miette::Diagnostic, thiserror::Error)] +#[error("the database is poisoned: a disk write failed, so fjall refuses all writes")] +pub struct DatabasePoisoned; + /// the top-level handle to a hydrant instance. /// /// `Hydrant` is cheaply cloneable. all sub-handles share the same underlying state. @@ -444,7 +451,13 @@ mod tests { } /// starts a hydrant that never touches the network on its own runtime. - fn start_offline(path: &std::path::Path) -> (tokio::runtime::Runtime, Hydrant) { + fn start_offline( + path: &std::path::Path, + ) -> ( + tokio::runtime::Runtime, + Hydrant, + tokio::task::JoinHandle>, + ) { let runtime = tokio::runtime::Runtime::new().unwrap(); let config = Config { relays: Vec::new(), @@ -454,24 +467,24 @@ mod tests { ..test_config(path) }; let hydrant = runtime.block_on(Hydrant::new(config)).unwrap(); - runtime.spawn({ + let run = runtime.spawn({ let hydrant = hydrant.clone(); async move { hydrant.run()?.await } }); // long enough for run to have spawned its threads std::thread::sleep(Duration::from_millis(300)); - (runtime, hydrant) + (runtime, hydrant, run) } #[test] fn shutdown_closes_the_database_for_the_next_start() { let tmp = tempdir().unwrap(); - let (runtime, hydrant) = start_offline(tmp.path()); + let (runtime, hydrant, _) = start_offline(tmp.path()); runtime.shutdown_background(); hydrant.shutdown(Duration::from_secs(10)).unwrap(); // fjall locks the database, so this only opens if the first one let go - let (runtime, hydrant) = start_offline(tmp.path()); + let (runtime, hydrant, _) = start_offline(tmp.path()); runtime.shutdown_background(); hydrant.shutdown(Duration::from_secs(10)).unwrap(); } @@ -479,7 +492,7 @@ mod tests { #[test] fn shutdown_fails_while_a_clone_holds_the_database() { let tmp = tempdir().unwrap(); - let (runtime, hydrant) = start_offline(tmp.path()); + let (runtime, hydrant, _) = start_offline(tmp.path()); runtime.shutdown_background(); let forgotten = hydrant.clone(); @@ -489,4 +502,22 @@ mod tests { // the clone still works and is the last one, so it can finish the job forgotten.shutdown(Duration::from_secs(10)).unwrap(); } + + #[test] + fn run_fails_once_the_database_is_poisoned() { + let tmp = tempdir().unwrap(); + let (runtime, hydrant, run) = start_offline(tmp.path()); + hydrant.state.db.poison.check(&fjall::Error::Poisoned); + + let err = runtime + .block_on(async { tokio::time::timeout(Duration::from_secs(10), run).await }) + .expect("run kept going on a poisoned database") + .unwrap() + .unwrap_err(); + assert!(err.downcast_ref::().is_some(), "{err}"); + + // and it can still be shut down, so the database can be reopened + runtime.shutdown_background(); + hydrant.shutdown(Duration::from_secs(10)).unwrap(); + } } diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index d76b20a..b4635d9 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -9,7 +9,7 @@ use url::Url; use super::super::firehose::FirehoseShared; use super::super::seed; -use super::Hydrant; +use super::{DatabasePoisoned, Hydrant}; use crate::config::SignatureVerification; use crate::db::{self, load_persisted_firehose_sources}; use crate::state::AppState; @@ -94,7 +94,7 @@ impl Hydrant { .await { error!(err = %e, "failed to queue gone backfills"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } std::thread::spawn({ @@ -532,16 +532,22 @@ impl Hydrant { // fatal_rx.changed() returns Err and we return Ok(()). drop(fatal_tx); - loop { - match fatal_rx.changed().await { - Ok(()) => { - if let Some(result) = fatal_rx.borrow().clone() { - return result.map_err(|s| miette::miette!("{s}")); + let fatal = async move { + loop { + match fatal_rx.changed().await { + Ok(()) => { + if let Some(result) = fatal_rx.borrow().clone() { + return result.map_err(|s| miette::miette!("{s}")); + } } + // all fatal_tx clones dropped: all tasks finished cleanly + Err(_) => return Ok(()), } - // all fatal_tx clones dropped: all tasks finished cleanly - Err(_) => return Ok(()), } + }; + tokio::select! { + r = fatal => r, + _ = state.db.poison.wait() => Err(DatabasePoisoned.into()), } }; Ok(fut) @@ -556,7 +562,7 @@ fn persist(state: &AppState) { && let Err(e) = db::set_firehose_cursor(&state.db, relay, seq) { error!(relay = %relay, err = %e, "failed to save cursor"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } true }); @@ -565,14 +571,14 @@ fn persist(state: &AppState) { Ok(watermark) => watermark, Err(e) => { error!(err = %e, "failed to checkpoint count deltas"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); None } }; if let Err(e) = state.db.persist() { error!(err = %e, "db persist failed"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } else { let watermark = checkpoint_watermark .map(Ok) @@ -581,7 +587,7 @@ fn persist(state: &AppState) { Ok(watermark) => state.db.mark_count_checkpoint_persisted(watermark), Err(e) => { error!(err = %e, "failed to load durable count checkpoint watermark"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); } } } diff --git a/src/control/mod.rs b/src/control/mod.rs index 4fbeb2f..8b1cccb 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -38,9 +38,9 @@ pub use repos::{ListedRecord, Record, RecordList, RepoHandle, RepoInfo, ReposCon pub use db::DbControl; pub use hosts::{ApiBind, ApiBinds, Host}; -pub use hydrant::Hydrant; #[cfg(feature = "indexer")] pub use hydrant::ScannedBlock; +pub use hydrant::{DatabasePoisoned, Hydrant}; pub use stats::StatsResponse; #[cfg(feature = "indexer_stream")] diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index 7e59cd4..eb4a64a 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -213,7 +213,7 @@ pub(crate) fn resolve_event_body( Ok(body) => body, Err(e) => { error!(err = %e, id, "cant resolve event record body"); - db::check_poisoned_report(&e); + state.db.poison.check_report(&e); None } } diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index cff56e0..d87d97a 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -38,6 +38,7 @@ pub struct OpenCx<'a> { pub(super) compression: &'a dyn Fn(&str, i32) -> CompressionType, /// names opened so far, checked against the registry at the end of open. pub(super) opened: RefCell>, + pub(super) poison: super::Poison, } impl OpenCx<'_> { diff --git a/src/db/mod.rs b/src/db/mod.rs index d65bcd3..dde5bc4 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -11,6 +11,7 @@ use std::sync::atomic::AtomicU64; #[cfg(test)] use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; +use tokio::sync::watch; use url::Url; pub mod compaction; @@ -84,6 +85,7 @@ pub struct Db { count_delta_gc_watermark: Arc, count_delta_in_flight: Arc>>, pub(crate) compaction_running: Arc, + pub(crate) poison: Poison, #[cfg(test)] pub(crate) persist_failures: Arc, /// 256 lock-sharded mutexes keyed by one byte of the trimmed DID. this is @@ -242,18 +244,39 @@ pub fn deser_repo_state<'b>(bytes: &'b [u8]) -> Result> { rmp_serde::from_slice(bytes).into_diagnostic() } -pub fn check_poisoned(e: &fjall::Error) { - if matches!(e, fjall::Error::Poisoned) { - error!("!!! DATABASE POISONED !!! exiting"); - std::process::exit(10); +/// fjall refuses every write once a disk write failed, and the first error saying so +/// ends [`Hydrant::run`](crate::control::Hydrant::run) instead of the whole process, +/// since hydrant may share it with the program embedding it. +#[derive(Clone)] +pub struct Poison(watch::Sender); + +impl Default for Poison { + fn default() -> Self { + Self(watch::Sender::new(false)) } } -pub fn check_poisoned_report(e: &miette::Report) { - let Some(err) = e.downcast_ref::() else { - return; - }; - self::check_poisoned(err); +impl Poison { + pub fn check(&self, e: &fjall::Error) { + if matches!(e, fjall::Error::Poisoned) && !self.0.send_replace(true) { + error!("!!! DATABASE POISONED !!! stopping"); + } + } + + pub fn check_report(&self, e: &miette::Report) { + // into_diagnostic hides the fjall error behind a private wrapper that only + // passes its message on, so a downcast never finds it + let poisoned = fjall::Error::Poisoned.to_string(); + if e.chain().any(|cause| cause.to_string() == poisoned) { + self.check(&fjall::Error::Poisoned); + } + } + + /// resolves once the database is poisoned. + pub async fn wait(&self) { + // only errs once the sender is dropped, and `self` holds it + let _ = self.0.subscribe().wait_for(|poisoned| *poisoned).await; + } } /// Load the persisted `(day, count)` pair for the daily PDS add counter. @@ -317,8 +340,28 @@ pub fn load_persisted_firehose_sources( mod tests { use super::*; use crate::config::Config; + use std::sync::atomic::Ordering; + #[tokio::test] + async fn only_a_poisoned_error_trips_the_poison() { + let poison = Poison::default(); + let wait = std::time::Duration::from_millis(50); + poison.check(&fjall::Error::Locked); + poison.check_report(&miette::miette!("not a fjall error")); + assert!(tokio::time::timeout(wait, poison.wait()).await.is_err()); + + // wrapped the way callsites hand it over, through into_diagnostic and context + let report = Err::<(), _>(fjall::Error::Poisoned) + .into_diagnostic() + .wrap_err("committing a batch") + .unwrap_err(); + poison.clone().check_report(&report); + tokio::time::timeout(std::time::Duration::from_secs(5), poison.wait()) + .await + .expect("a poisoned error should trip every clone"); + } + #[tokio::test] async fn test_db_compact_concurrency_guard() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; diff --git a/src/db/open.rs b/src/db/open.rs index cccb34f..a49a4ef 100644 --- a/src/db/open.rs +++ b/src/db/open.rs @@ -41,6 +41,7 @@ impl Db { cfg, compression: &get_compression, opened: std::cell::RefCell::new(Vec::new()), + poison: super::Poison::default(), }; let this = Self::assemble_keyspaces_and_verify(&cx, count_delta_gc_watermark)?; @@ -197,6 +198,7 @@ impl Db { count_delta_gc_watermark, count_delta_in_flight: Arc::new(Mutex::new(BTreeSet::new())), compaction_running: Arc::new(std::sync::atomic::AtomicBool::new(false)), + poison: cx.poison.clone(), #[cfg(test)] persist_failures: Arc::new(AtomicUsize::new(0)), #[cfg(feature = "indexer")] diff --git a/src/db/schema.rs b/src/db/schema.rs index b4a41a4..dad69ee 100644 --- a/src/db/schema.rs +++ b/src/db/schema.rs @@ -65,10 +65,10 @@ pub trait Schema { /// typed handle over a [`Keyspace`], parameterized by its [`Schema`]. /// -/// zero-cost: one `Keyspace` handle plus a marker. cloning is as cheap as -/// cloning the underlying handle. +/// one `Keyspace` handle plus the database's poison flag, both cheap to clone. pub struct Ks { inner: Keyspace, + poison: super::Poison, _s: PhantomData S>, } @@ -76,6 +76,7 @@ impl Clone for Ks { fn clone(&self) -> Self { Self { inner: self.inner.clone(), + poison: self.poison.clone(), _s: PhantomData, } } @@ -89,6 +90,7 @@ impl Ks { cx.record_opened(S::NAME); Ok(Self { inner, + poison: cx.poison.clone(), _s: PhantomData, }) } @@ -109,14 +111,14 @@ impl Ks { /// callsites route through these. #[inline] pub fn get>(&self, key: K) -> fjall::Result> { - self.inner.get(key).inspect_err(super::check_poisoned) + self.inner.get(key).inspect_err(|e| self.poison.check(e)) } #[inline] pub fn contains_key>(&self, key: K) -> fjall::Result { self.inner .contains_key(key) - .inspect_err(super::check_poisoned) + .inspect_err(|e| self.poison.check(e)) } #[inline] @@ -127,12 +129,12 @@ impl Ks { ) -> fjall::Result<()> { self.inner .insert(key, value) - .inspect_err(super::check_poisoned) + .inspect_err(|e| self.poison.check(e)) } #[inline] pub fn remove>(&self, key: K) -> fjall::Result<()> { - self.inner.remove(key).inspect_err(super::check_poisoned) + self.inner.remove(key).inspect_err(|e| self.poison.check(e)) } } diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 6e2c1f0..3d082d0 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -6,7 +6,7 @@ use std::sync::atomic::Ordering::SeqCst; use tokio::runtime::Handle as TokioHandle; use tracing::{debug, error, warn}; -use crate::db::{self, Txn, keys, ser_repo_meta}; +use crate::db::{Txn, keys, ser_repo_meta}; use crate::ingest::stream::types::AccountStatus; use crate::ingest::stream::{Account, Commit, Identity}; use crate::ingest::validation; @@ -223,7 +223,7 @@ impl FirehoseWorker { Ok(RepoProcessResult::NeedsBackfill(None)) => {} Err(e) => { if let IngestError::Generic(ref r) = e { - db::check_poisoned_report(r); + state.db.poison.check_report(r); } error!(err = %e, "error processing commit"); try_persist(&mut ctx, &commit); @@ -571,6 +571,7 @@ mod tests { use super::*; use crate::{ config::Config, + db, ingest::stream::{Datetime, RepoOp, RepoOpAction}, }; diff --git a/src/main.rs b/src/main.rs index f45ec09..a128d5d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,6 +1,6 @@ use futures::FutureExt; use hydrant::config::Config; -use hydrant::control::{ApiBind, ApiBinds, Hydrant}; +use hydrant::control::{ApiBind, ApiBinds, DatabasePoisoned, Hydrant}; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr}; mod allocator; @@ -68,9 +68,17 @@ async fn main() -> miette::Result<()> { .then(|| hydrant.serve_debug(app.debug_port).boxed()) .unwrap_or_else(|| std::future::pending().boxed()); - tokio::select! { + let result = tokio::select! { r = hydrant.run()? => r, r = api_fut => r, r = debug_fut => r, + }; + // a supervisor can tell this apart from other failures by the exit code + if let Err(e) = &result + && e.downcast_ref::().is_some() + { + eprintln!("{e:?}"); + std::process::exit(10); } + result } -- 2.51.2