diff --git a/CHANGELOG.md b/CHANGELOG.md index e953f93..783586c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Preserve purgatory and rejected-event recovery across abrupt termination. + Their snapshots now remain after restore and are atomically replaced every + 60 seconds, bounding crash loss to one interval instead of consuming the + only durable copy at startup. - Treat cumulative retained-subscription byte refusals as capacity signals, not temporary query-rate episodes. The sync client learns the disclosed cap, rebuilds persistent coverage within it while reserving one maximum transient diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index a33703f..b4f1444 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -68,9 +68,11 @@ runtime: - Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51) - Set up shared storage (LMDB or Memory), purgatory, sync manager, and background maintenance tasks +- Atomically checkpoint purgatory and rejected-event recovery state every 60 + seconds without consuming the checkpoint during restore - Serve HTTP + WebSocket until a caller-supplied shutdown future - resolves, then persist state (purgatory, rejected-events cache, - placeholder ref cleanup) and stop background tasks + resolves, then stop background mutation, persist a final state snapshot + (purgatory and rejected-events cache), and clean up placeholder refs **Key Dependencies:** diff --git a/docs/explanation/purgatory-design.md b/docs/explanation/purgatory-design.md index aa66744..bc98052 100644 --- a/docs/explanation/purgatory-design.md +++ b/docs/explanation/purgatory-design.md @@ -39,11 +39,18 @@ This ensures we only serve announcements for repos that actually have content. ## Key Design Principles -### 1. Graceful-Shutdown Persistence +### 1. Crash-Safe Checkpoint Persistence -Purgatory state is **saved to disk on graceful shutdown** and **restored on startup**. This preserves in-flight work across planned restarts (deployments, reboots). +Purgatory state is atomically checkpointed every 60 seconds, saved once more +on graceful shutdown, and restored on startup. This preserves in-flight work +across planned restarts and bounds state lost to an abrupt process or machine +stop to the checkpoint interval. -On `SIGINT` / Ctrl-C, `main.rs` calls `purgatory.save_to_disk()` before exiting. On startup, if the state file exists, `purgatory.restore_from_disk()` is called before the server begins accepting connections. +The restored file remains in place until a complete newer snapshot is durable; +restore never consumes the only crash-safe copy. Snapshot replacement writes a +temporary file, syncs it, renames it over the checkpoint, and syncs the parent +directory. On `SIGINT` / Ctrl-C, `RelayServer` first stops background mutation +and then writes the final checkpoint before exiting. **What is persisted:** @@ -55,13 +62,9 @@ On `SIGINT` / Ctrl-C, `main.rs` calls `purgatory.save_to_disk()` before exiting. | `expired_events` | ✅ Yes | Prevents re-sync loops after restart | | `sync_queue` | ❌ No | Rebuilt automatically after restore | -**What is NOT persisted (unclean shutdown):** - -On a crash or `SIGKILL`, the state file is not written. In that case: - -- Events are still on other relays (can be re-submitted) -- Git data can be re-pushed -- 30-minute expiry means data is transient anyway +**Unclean shutdown:** mutations after the most recent 60-second checkpoint can +be lost. The previous complete checkpoint remains valid even if the process is +terminated while its replacement is being written. **State file location:** `/purgatory-state.json` diff --git a/src/atomic_file.rs b/src/atomic_file.rs new file mode 100644 index 0000000..ee2ceb3 --- /dev/null +++ b/src/atomic_file.rs @@ -0,0 +1,34 @@ +//! Crash-safe replacement of small state snapshots. + +use std::fs::File; +use std::io::{self, Write}; +use std::path::Path; + +/// Replace `path` only after the complete new contents are durable. +pub(crate) fn write(path: &Path, contents: &[u8]) -> io::Result<()> { + let parent = path.parent().ok_or_else(|| { + io::Error::new(io::ErrorKind::InvalidInput, "snapshot path has no parent") + })?; + let mut temporary = tempfile::NamedTempFile::new_in(parent)?; + temporary.write_all(contents)?; + temporary.as_file_mut().sync_all()?; + temporary.persist(path).map_err(|error| error.error)?; + File::open(parent)?.sync_all()?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn atomically_replaces_an_existing_snapshot() { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("state.json"); + std::fs::write(&path, b"old").unwrap(); + + write(&path, b"complete new snapshot").unwrap(); + + assert_eq!(std::fs::read(path).unwrap(), b"complete new snapshot"); + } +} diff --git a/src/lib.rs b/src/lib.rs index 6dd989d..ce5620d 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,3 +1,4 @@ +mod atomic_file; pub mod audit_cleanup; pub mod cleanup_empty_repos; pub mod config; diff --git a/src/purgatory/mod.rs b/src/purgatory/mod.rs index 27a125d..d2d9932 100644 --- a/src/purgatory/mod.rs +++ b/src/purgatory/mod.rs @@ -6,7 +6,8 @@ //! //! ## Architecture //! -//! - **In-memory only**: Data is lost on restart (acceptable per spec) +//! - **Crash-safe checkpoints**: In-memory state is periodically snapshotted +//! and restored across graceful or abrupt restarts //! - **Thread-safe**: Uses DashMap for concurrent access from multiple handlers //! - **Automatic expiry**: Entries expire after 30 minutes by default //! - **Separate stores**: State events and PR events use different indexing strategies @@ -1645,11 +1646,12 @@ impl Purgatory { expired_events, }; - // Serialize to JSON and write to file + // Replace the previous checkpoint only after the new snapshot is + // complete and durable. An abrupt stop must leave one valid version. let json = serde_json::to_string_pretty(&state)?; - std::fs::write(path, json)?; + crate::atomic_file::write(path, json.as_bytes())?; - tracing::info!( + tracing::debug!( path = %path.display(), announcements = state.announcement_purgatory.len(), state_events = state.state_events.len(), @@ -1667,8 +1669,9 @@ impl Purgatory { /// the current purgatory instance. Adjusts time-based fields to account for downtime /// between save and restore. /// - /// After successful restore, the state file is deleted to prevent accidental - /// double-restore. + /// The checkpoint remains after restore until a newer periodic or shutdown + /// snapshot atomically replaces it. Restoring it more than once is safe: + /// addressable entries replace their in-memory keys. /// /// # Arguments /// * `path` - Path to the saved state file @@ -1816,10 +1819,6 @@ impl Purgatory { "Restored purgatory state from disk" ); - // Delete state file after successful restore - std::fs::remove_file(path)?; - tracing::debug!(path = %path.display(), "Deleted state file after restore"); - Ok(()) } } @@ -2491,8 +2490,8 @@ async fn test_save_and_restore_state_events() { let purgatory2 = Purgatory::new(PathBuf::new()); purgatory2.restore_from_disk(&state_file).unwrap(); - // Verify file was deleted after restore - assert!(!state_file.exists()); + // The last durable checkpoint remains available after restore. + assert!(state_file.exists()); // Verify state events were restored let (_, state_count, _) = purgatory2.count(); @@ -2504,6 +2503,11 @@ async fn test_save_and_restore_state_events() { // Verify event IDs match let restored_ids: Vec = restored_entries.iter().map(|e| e.event.id).collect(); assert!(restored_ids.contains(&event1_id)); + + // A second startup before the next checkpoint restores the same state. + let purgatory3 = Purgatory::new(PathBuf::new()); + purgatory3.restore_from_disk(&state_file).unwrap(); + assert_eq!(purgatory3.count().1, 2); assert!(restored_ids.contains(&event2_id)); // Verify identifiers and authors match @@ -3026,8 +3030,8 @@ async fn test_file_cleanup_after_successful_restore() { let purgatory2 = Purgatory::new(PathBuf::new()); purgatory2.restore_from_disk(&state_file).unwrap(); - // File should be deleted after successful restore - assert!(!state_file.exists()); + // The checkpoint remains crash-safe after successful restore. + assert!(state_file.exists()); } #[tokio::test] @@ -3068,8 +3072,8 @@ async fn test_save_and_restore_announcement_events() { let purgatory2 = Purgatory::new(PathBuf::new()); purgatory2.restore_from_disk(&state_file).unwrap(); - // File should be deleted after restore - assert!(!state_file.exists()); + // The checkpoint remains crash-safe after restore. + assert!(state_file.exists()); // Verify announcement was restored let (ann_count, _, _) = purgatory2.count(); diff --git a/src/server.rs b/src/server.rs index 66f635c..4119ace 100644 --- a/src/server.rs +++ b/src/server.rs @@ -34,6 +34,10 @@ use crate::{ sync::{naughty_list::NaughtyListTracker, rejected_index::RejectedEventsIndex, SyncManager}, }; +/// Limits recovery work lost to an abrupt process or machine stop without +/// turning every in-memory mutation into synchronous disk I/O. +const SYNC_STATE_CHECKPOINT_INTERVAL: Duration = Duration::from_secs(60); + /// A fully-wired relay that has bound its listener but not yet started /// accepting connections. /// @@ -253,6 +257,32 @@ impl RelayServer { sync_manager.run().await; })); + // Retain a recent durable copy after startup. Restore deliberately + // leaves the prior checkpoint in place, and these atomic replacements + // bound any later crash loss to one interval. + let checkpoint_purgatory = purgatory.clone(); + let checkpoint_rejected = rejected_events_index.clone(); + let checkpoint_root = PathBuf::from(config.effective_git_data_path()); + background_tasks.push(tokio::spawn(async move { + let first = tokio::time::Instant::now() + SYNC_STATE_CHECKPOINT_INTERVAL; + let mut interval = tokio::time::interval_at(first, SYNC_STATE_CHECKPOINT_INTERVAL); + loop { + interval.tick().await; + let purgatory_path = checkpoint_root.join("purgatory-state.json"); + if let Err(error) = checkpoint_purgatory.save_to_disk(&purgatory_path) { + warn!(%error, "Failed to checkpoint purgatory state"); + } + let rejected_path = checkpoint_root.join("rejected-events-cache.json"); + if let Err(error) = checkpoint_rejected.save_to_disk(&rejected_path) { + warn!(%error, "Failed to checkpoint rejected-events cache"); + } + } + })); + info!( + interval_secs = SYNC_STATE_CHECKPOINT_INTERVAL.as_secs(), + "Crash-safe sync-state checkpoint task started" + ); + // Spawn background cleanup task for purgatory entries (60s interval) let cleanup_purgatory = purgatory.clone(); background_tasks.push(tokio::spawn(async move { @@ -367,9 +397,9 @@ impl RelayServer { /// itself fails), then persist state and tear down background tasks. /// /// The shutdown sequence mirrors the binary's signal handling: stop the - /// holding-cleanup task gracefully, save purgatory state and the - /// rejected-events cache to disk, remove placeholder `refs/nostr/` - /// refs, and abort the remaining background loops. + /// holding-cleanup task gracefully, stop background mutation, save + /// purgatory state and the rejected-events cache to disk, and remove + /// placeholder `refs/nostr/` refs. pub async fn run_until(self, shutdown: impl Future) -> Result<()> { info!("Starting HTTP server on {}", self.config.bind_address); @@ -394,6 +424,14 @@ impl RelayServer { } }; + // Stop and join all state-mutating background loops before taking the + // final snapshot. This also prevents an older periodic checkpoint from + // racing and replacing the shutdown snapshot after it is written. + for task in self.background_tasks { + task.abort(); + let _ = task.await; + } + self.deletion_cleanup.shutdown().await; // Save purgatory state to disk @@ -426,12 +464,6 @@ impl RelayServer { git::cleanup_placeholder_refs(&self.git_data_path, &placeholder_ids); } - // Abort detached background loops so the host process does not - // accumulate live tasks (and the state they capture) per instance. - for task in self.background_tasks { - task.abort(); - } - result } } diff --git a/src/sync/rejected_index.rs b/src/sync/rejected_index.rs index c9a3262..1e21590 100644 --- a/src/sync/rejected_index.rs +++ b/src/sync/rejected_index.rs @@ -1092,9 +1092,10 @@ impl RejectedEventsIndex { }, }; - // Serialize to JSON and write to file + // Replace the previous checkpoint only after the new snapshot is + // complete and durable. An abrupt stop must leave one valid version. let json = serde_json::to_string_pretty(&state)?; - std::fs::write(path, json)?; + crate::atomic_file::write(path, json.as_bytes())?; Ok(()) } @@ -1103,7 +1104,9 @@ impl RejectedEventsIndex { /// /// Loads the serialized state from disk and populates both hot cache and cold index. /// Adjusts all timestamps by adding the downtime duration (time since save) to maintain - /// correct expiry behavior. Deletes the state file after successful restore. + /// correct expiry behavior. The checkpoint remains in place until a newer + /// periodic or shutdown snapshot atomically replaces it, so another crash + /// before graceful shutdown cannot erase all recovery history. /// /// # Arguments /// @@ -1192,14 +1195,13 @@ impl RejectedEventsIndex { unrecoverable_entries.insert(event_id, entry); } - // Release locks before deleting file + // Release locks after the complete snapshot has been restored. Keep the + // checkpoint: it is the last crash-safe state until the periodic writer + // replaces it. drop(hot_entries); drop(cold_entries); drop(unrecoverable_entries); - // Delete the state file after successful restore - std::fs::remove_file(path)?; - Ok(()) } } @@ -1774,13 +1776,19 @@ mod tests { RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); index2.restore_from_disk(&state_path).unwrap(); - // Verify state file was deleted after restore - assert!(!state_path.exists()); + // The last durable checkpoint remains available after restore. + assert!(state_path.exists()); // Verify hot cache restored assert_eq!(index2.hot_cache_len(), 1); assert!(index2.hot_cache.contains(&event.id)); + // A second startup before the next checkpoint restores the same state. + let index3 = + RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); + index3.restore_from_disk(&state_path).unwrap(); + assert!(index3.hot_cache.contains(&event.id)); + // Verify cold index restored assert_eq!(index2.cold_index_len(), 1); assert!(index2.cold_index.contains(&event.id)); @@ -1963,8 +1971,8 @@ mod tests { RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); index2.restore_from_disk(&state_path).unwrap(); - // File should be deleted after successful restore - assert!(!state_path.exists()); + // The checkpoint remains crash-safe after successful restore. + assert!(state_path.exists()); } #[tokio::test]