diff --git a/src/error.rs b/src/error.rs index 84b5cb9..d368d9e 100644 --- a/src/error.rs +++ b/src/error.rs @@ -17,4 +17,23 @@ pub enum Error { Other(String), } +impl Error { + /// True if this error indicates the database is in an unrecoverable state. + /// + /// `fjall::Error::Poisoned` is fjall's signal that a prior flush/commit + /// failed and subsequent writes can't be trusted — fjall's own docs say + /// to crash the application. `Unrecoverable` is a similar terminal state. + /// Callers should force-exit the process rather than attempt graceful + /// shutdown, since the blocking thread pool may be stuck on the same + /// underlying failure. + pub fn is_db_fatal(&self) -> bool { + matches!( + self, + Error::Storage(StorageError::Fjall( + fjall::Error::Poisoned | fjall::Error::Unrecoverable + )) + ) + } +} + pub type Result = std::result::Result; diff --git a/src/main.rs b/src/main.rs index 130d257..7f402d2 100644 --- a/src/main.rs +++ b/src/main.rs @@ -121,8 +121,7 @@ struct Args { max_deep_crawl_workers: usize, } -#[tokio::main] -async fn main() -> Result<()> { +fn main() { rustls::crypto::aws_lc_rs::default_provider() .install_default() .expect("failed to install rustls crypto provider"); @@ -131,6 +130,31 @@ async fn main() -> Result<()> { .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) .init(); + let rt = tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build() + .expect("failed to build tokio runtime"); + + let result = rt.block_on(run()); + + // Force-shutdown the runtime after a bounded wait. Without this, a + // `spawn_blocking` task genuinely stuck in fjall (e.g. after Poisoned) + // holds a blocking-pool thread, and since those threads are non-daemon + // they'd prevent process exit. `shutdown_timeout` detaches any remaining + // tasks after the deadline; the explicit `process::exit` below then + // guarantees we don't wait for detached blocking threads either. + rt.shutdown_timeout(Duration::from_secs(10)); + + match result { + Ok(()) => std::process::exit(0), + Err(e) => { + eprintln!("fatal: {e}"); + std::process::exit(1); + } + } +} + +async fn run() -> Result<()> { let args = Args::parse(); let subscribe_host = args @@ -377,12 +401,22 @@ async fn main() -> Result<()> { /// Flatten a task join result into an optional error. /// Panics (JoinError) are treated as errors. +/// +/// If the error indicates an unrecoverable database state +/// ([`Error::is_db_fatal`]), this immediately force-exits the process rather +/// than returning. Graceful shutdown isn't safe in that state because other +/// tasks may be stuck in blocking fjall calls that will never return. fn into_error(r: std::result::Result, tokio::task::JoinError>) -> Option { - match r { - Ok(Ok(())) => None, - Ok(Err(e)) => Some(e), - Err(e) => Some(Error::TaskPanic(e)), + let err = match r { + Ok(Ok(())) => return None, + Ok(Err(e)) => e, + Err(e) => Error::TaskPanic(e), + }; + if err.is_db_fatal() { + eprintln!("FATAL: database poisoned, force-exiting: {err}"); + std::process::exit(2); } + Some(err) } fn install_metrics(addr: SocketAddr) -> Result<()> { diff --git a/src/sync/resync/dispatcher.rs b/src/sync/resync/dispatcher.rs index d3a1cc9..722e723 100644 --- a/src/sync/resync/dispatcher.rs +++ b/src/sync/resync/dispatcher.rs @@ -294,7 +294,13 @@ pub async fn run( } Ok(Ok(None)) => break, // queue empty or all ready items busy Ok(Err(e)) => { - error!(error = %e, "error claiming resync job; pausing"); + let wrapped = crate::error::Error::from(e); + if wrapped.is_db_fatal() { + error!(error = %wrapped, + "claim_resync hit unrecoverable db state; exiting dispatcher"); + return Err(wrapped); + } + error!(error = %wrapped, "error claiming resync job; pausing"); break; } }