From aef7e759e4141aa755af98045d99bbae17236568 Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Sun, 9 Aug 2026 07:07:44 -0400 Subject: [PATCH 1/4] fix(ui): report session persistence failures on shutdown --- .../storage-shutdown-persistence.fixed.md | 1 + n00n-ui/src/app/mod.rs | 2 +- n00n-ui/src/app/tests.rs | 15 ++- n00n-ui/src/event_loop.rs | 34 +++-- n00n-ui/src/storage_writer.rs | 124 +++++++++++++++--- 5 files changed, 144 insertions(+), 32 deletions(-) create mode 100644 changelog.d/storage-shutdown-persistence.fixed.md diff --git a/changelog.d/storage-shutdown-persistence.fixed.md b/changelog.d/storage-shutdown-persistence.fixed.md new file mode 100644 index 000000000..351d88898 --- /dev/null +++ b/changelog.d/storage-shutdown-persistence.fixed.md @@ -0,0 +1 @@ +Report storage writer shutdown failures instead of silently accepting unpersisted session snapshots. diff --git a/n00n-ui/src/app/mod.rs b/n00n-ui/src/app/mod.rs index 729d5a679..0a6750d8c 100644 --- a/n00n-ui/src/app/mod.rs +++ b/n00n-ui/src/app/mod.rs @@ -218,7 +218,7 @@ impl Drop for TestStateDir { return; }; let writer_stopped = match Arc::try_unwrap(writer) { - Ok(writer) => writer.wait_for_shutdown(TEST_WRITER_DRAIN_TIMEOUT), + Ok(writer) => writer.wait_for_shutdown(TEST_WRITER_DRAIN_TIMEOUT).is_ok(), Err(writer) => { drop(writer); false diff --git a/n00n-ui/src/app/tests.rs b/n00n-ui/src/app/tests.rs index ba0e5aebe..46f52d0d2 100644 --- a/n00n-ui/src/app/tests.rs +++ b/n00n-ui/src/app/tests.rs @@ -2551,7 +2551,8 @@ fn drain_writer(app: App, writer: Arc) { Arc::try_unwrap(writer) .ok() .expect("app must hold the only other writer reference") - .shutdown(WRITER_DRAIN_TIMEOUT); + .shutdown(WRITER_DRAIN_TIMEOUT) + .unwrap(); } #[test] @@ -2730,7 +2731,8 @@ fn draw_failure_pending_submission_restores_fifo_images_and_control_after_restar Arc::try_unwrap(writer) .ok() .expect("test owns the storage writer") - .shutdown(WRITER_DRAIN_TIMEOUT); + .shutdown(WRITER_DRAIN_TIMEOUT) + .unwrap(); let writer = Arc::new(StorageWriter::new(dir.clone()).unwrap()); let mut restarted = build_app(dir.clone(), Arc::clone(&writer)); @@ -2769,7 +2771,8 @@ fn draw_failure_pending_submission_restores_fifo_images_and_control_after_restar Arc::try_unwrap(writer) .ok() .expect("test owns the restarted storage writer") - .shutdown(WRITER_DRAIN_TIMEOUT); + .shutdown(WRITER_DRAIN_TIMEOUT) + .unwrap(); } #[test] @@ -2831,7 +2834,8 @@ fn mcp_prompt_draw_failure_survives_restart_without_text_fallback() { Arc::try_unwrap(writer) .ok() .expect("test owns the storage writer") - .shutdown(WRITER_DRAIN_TIMEOUT); + .shutdown(WRITER_DRAIN_TIMEOUT) + .unwrap(); let writer = Arc::new(StorageWriter::new(dir.clone()).unwrap()); let mut restarted = build_app_with_mcp(dir.clone(), Arc::clone(&writer), mcp_reader); @@ -2865,7 +2869,8 @@ fn mcp_prompt_draw_failure_survives_restart_without_text_fallback() { Arc::try_unwrap(writer) .ok() .expect("test owns the restarted storage writer") - .shutdown(WRITER_DRAIN_TIMEOUT); + .shutdown(WRITER_DRAIN_TIMEOUT) + .unwrap(); } #[test] diff --git a/n00n-ui/src/event_loop.rs b/n00n-ui/src/event_loop.rs index 5f7caaa82..a1b2c1114 100644 --- a/n00n-ui/src/event_loop.rs +++ b/n00n-ui/src/event_loop.rs @@ -62,6 +62,8 @@ const PERIODIC_SAVE_INTERVAL: Duration = Duration::from_secs(1); const DRAIN_BUDGET: usize = 256; const AGENT_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(3); const STORAGE_WRITER_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5); +const STORAGE_WRITER_REFS_ERR: &str = + "storage writer has outstanding references, skipping graceful shutdown"; const DELETE_FOCUSED_ERR: &str = "cannot delete the focused session"; const NOT_LIVE_ERR: &str = "session not live"; const TEAM_TOOL_NAME: &str = "team"; @@ -606,8 +608,16 @@ impl<'t> EventLoop<'t> { // Fatal errors still save every session, shut down MCP transports // (terminating and reaping their child processes), and drain the // storage writer before the process exits. - let report = self.shutdown(); - result.map(|()| report) + let shutdown = self.shutdown(); + match result { + Ok(()) => shutdown, + Err(error) => { + if let Err(shutdown_error) = shutdown { + warn!(error = %shutdown_error, "shutdown after fatal error was incomplete"); + } + Err(error) + } + } } /// Wait for the next event from any source, or time out so animations @@ -1490,7 +1500,7 @@ impl<'t> EventLoop<'t> { } } - fn shutdown(mut self) -> ShutdownReport { + fn shutdown(mut self) -> Result { self.preserve_post_draw_submissions(); let exit = self.sessions[self.focused].app.exit_request; for rt in &self.sessions { @@ -1512,17 +1522,19 @@ impl<'t> EventLoop<'t> { smol::block_on(h.shutdown()); } crate::agent::join_all(agent_tasks, AGENT_SHUTDOWN_TIMEOUT); - match Arc::try_unwrap(self.ctx.storage_writer) { - Ok(writer) => writer.shutdown(STORAGE_WRITER_SHUTDOWN_TIMEOUT), - Err(_) => { - warn!("storage writer has outstanding references, skipping graceful shutdown"); - } - } - ShutdownReport { + let storage_result = match Arc::try_unwrap(self.ctx.storage_writer) { + Ok(writer) => writer + .shutdown(STORAGE_WRITER_SHUTDOWN_TIMEOUT) + .map_err(Into::into), + Err(_) => Err(eyre!(STORAGE_WRITER_REFS_ERR)), + }; + let report = ShutdownReport { exit, tabs, focused: self.focused, - } + }; + storage_result?; + Ok(report) } } diff --git a/n00n-ui/src/storage_writer.rs b/n00n-ui/src/storage_writer.rs index 5cdc49244..08300fe27 100644 --- a/n00n-ui/src/storage_writer.rs +++ b/n00n-ui/src/storage_writer.rs @@ -27,9 +27,20 @@ struct PendingSnapshot { type Pending = Arc>>; type FailedRevisions = HashMap; +#[derive(Debug, thiserror::Error)] +pub(crate) enum StorageWriterShutdownError { + #[error("storage writer stopped with {count} unpersisted snapshot(s)")] + UnpersistedSnapshots { count: usize }, + #[error("storage writer did not drain within {timeout:?}")] + Timeout { timeout: Duration }, + #[error("storage writer completion channel disconnected")] + Disconnected, +} + #[derive(Default)] struct RetryState { attempts: HashMap<(n00nId, u64), u32>, + exhausted: HashMap, } type DeleteCallback = Box) + Send>; @@ -53,7 +64,7 @@ enum Op { pub struct StorageWriter { pending: Pending, ops: flume::Sender, - done_rx: flume::Receiver<()>, + done_rx: flume::Receiver>, } impl StorageWriter { @@ -61,7 +72,7 @@ impl StorageWriter { let pending: Pending = Arc::default(); let writer_pending = Arc::clone(&pending); let (ops, ops_rx) = flume::unbounded::(); - let (done_tx, done_rx) = flume::bounded::<()>(1); + let (done_tx, done_rx) = flume::bounded::>(1); std::thread::Builder::new() .name("storage-writer".into()) @@ -103,9 +114,15 @@ impl StorageWriter { } } } - flush(&writer_pending, &mut logs, &mut durable_revisions, &dir); + let failed = flush(&writer_pending, &mut logs, &mut durable_revisions, &dir); drop(logs); - let _ = done_tx.send(()); + let unpersisted = retries.unpersisted_count(&failed, &durable_revisions); + let completion = if unpersisted == 0 { + Ok(()) + } else { + Err(unpersisted) + }; + let _ = done_tx.send(completion); })?; Ok(Self { @@ -159,16 +176,26 @@ impl StorageWriter { } } - pub fn shutdown(self, timeout: Duration) { - if !self.wait_for_shutdown(timeout) { - warn!("storage writer did not drain within {timeout:?}"); - } + pub(crate) fn shutdown(self, timeout: Duration) -> Result<(), StorageWriterShutdownError> { + self.wait_for_shutdown(timeout) } - pub(crate) fn wait_for_shutdown(self, timeout: Duration) -> bool { + pub(crate) fn wait_for_shutdown( + self, + timeout: Duration, + ) -> Result<(), StorageWriterShutdownError> { let Self { ops, done_rx, .. } = self; drop(ops); - done_rx.recv_timeout(timeout).is_ok() + match done_rx.recv_timeout(timeout) { + Ok(Ok(())) => Ok(()), + Ok(Err(count)) => Err(StorageWriterShutdownError::UnpersistedSnapshots { count }), + Err(flume::RecvTimeoutError::Timeout) => { + Err(StorageWriterShutdownError::Timeout { timeout }) + } + Err(flume::RecvTimeoutError::Disconnected) => { + Err(StorageWriterShutdownError::Disconnected) + } + } } } @@ -202,6 +229,10 @@ impl RetryState { .is_some_and(|snapshot| snapshot.revision == revision) { pending.remove(&id); + self.exhausted + .entry(id) + .and_modify(|current| *current = (*current).max(revision)) + .or_insert(revision); } } self.attempts.retain(|(id, revision), _| { @@ -210,6 +241,24 @@ impl RetryState { .is_some_and(|snapshot| snapshot.revision == *revision) }); } + + fn unpersisted_count( + &self, + failed: &FailedRevisions, + durable_revisions: &HashMap, + ) -> usize { + failed.len() + + self + .exhausted + .iter() + .filter(|(id, revision)| { + !failed.contains_key(id) + && durable_revisions + .get(id) + .is_none_or(|durable| durable < revision) + }) + .count() + } } fn flush_and_persist( @@ -402,7 +451,7 @@ mod tests { writer.send(Box::new(b.clone())); b.title = "renamed".into(); writer.send(Box::new(b)); - writer.shutdown(DRAIN_TIMEOUT); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); assert!(AppSession::load(a_id, &dir).is_ok()); assert_eq!(AppSession::load(b_id, &dir).unwrap().title, "renamed"); @@ -432,6 +481,21 @@ mod tests { assert!(lock(&pending).contains_key(&id)); } + #[test] + fn shutdown_does_not_report_success_with_unpersisted_snapshot() { + let (tmp, dir) = state_dir(); + let writer = StorageWriter::new(dir).unwrap(); + let session = AppSession::new("test-model", "/tmp/shutdown-failure"); + let id = session.id; + fs::create_dir_all(tmp.path().join(SESSIONS_DIR).join(format!("{id}.jsonl"))).unwrap(); + writer.send(Box::new(session)); + + assert!(matches!( + writer.wait_for_shutdown(DRAIN_TIMEOUT), + Err(StorageWriterShutdownError::UnpersistedSnapshots { count: 1 }) + )); + } + #[test] fn retry_exhaustion_removes_only_exact_failed_revisions() { let pending: Pending = Arc::default(); @@ -489,6 +553,36 @@ mod tests { assert!(!pending.contains_key(&exact_id)); } + #[test] + fn exhausted_retry_remains_a_shutdown_failure_until_durable() { + let pending: Pending = Arc::default(); + let mut retries = RetryState::default(); + let mut session = AppSession::new("test-model", "/tmp/exhausted"); + session.meta.revision = 4; + let id = session.id; + lock(&pending).insert( + id, + PendingSnapshot { + revision: 4, + session: Box::new(session), + }, + ); + let failures = HashMap::from([(id, 4)]); + for _ in 0..MAX_RETRY_ATTEMPTS { + retries.record_failures(&pending, failures.clone()); + } + + assert!(lock(&pending).is_empty()); + assert_eq!( + retries.unpersisted_count(&FailedRevisions::new(), &HashMap::new()), + 1 + ); + assert_eq!( + retries.unpersisted_count(&FailedRevisions::new(), &HashMap::from([(id, 4)]),), + 0 + ); + } + #[test] fn exhausted_retry_for_one_session_does_not_block_explicit_persist() { let (tmp, dir) = state_dir(); @@ -545,7 +639,7 @@ mod tests { assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); assert!(AppSession::load(id, &dir).is_ok()); - writer.shutdown(DRAIN_TIMEOUT); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); } #[test] @@ -563,7 +657,7 @@ mod tests { let mut second = AppSession::load(id, &dir).unwrap(); second.title = "same revision, new snapshot".into(); writer.send(Box::new(second)); - writer.shutdown(DRAIN_TIMEOUT); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); assert_eq!( AppSession::load(id, &dir).unwrap().title, @@ -591,7 +685,7 @@ mod tests { assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); assert_eq!(AppSession::load(id, &dir).unwrap().title, "periodic save"); - writer.shutdown(DRAIN_TIMEOUT); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); } #[test] @@ -605,7 +699,7 @@ mod tests { writer.delete(id, move |res| { let _ = done_tx.send(res); }); - writer.shutdown(DRAIN_TIMEOUT); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); assert!(done_rx.recv().unwrap().is_ok()); assert!(AppSession::load(id, &dir).is_err()); From 46de823a280d8a461d6b03f7a93afa20de1d44ff Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Sun, 9 Aug 2026 20:14:31 -0400 Subject: [PATCH 2/4] fix(ui): preserve storage writer ordering --- n00n-storage/src/sessions.rs | 37 +- n00n-ui/src/storage_writer.rs | 1245 +++++++++++++++++++++++---------- 2 files changed, 911 insertions(+), 371 deletions(-) diff --git a/n00n-storage/src/sessions.rs b/n00n-storage/src/sessions.rs index 74c9b9bee..0c35c66e1 100644 --- a/n00n-storage/src/sessions.rs +++ b/n00n-storage/src/sessions.rs @@ -3283,15 +3283,18 @@ where } /// # Errors - /// Returns `SessionError` if the session file cannot be found or removed. + /// Returns `SessionError` if the session file cannot be found or cleanup fails. pub fn delete_from(id: n00nId, dir: &Path) -> Result<(), SessionError> { let _lock = lock_openai_response_chain_in(dir, id)?; - let Some(path) = locate_session_file(dir, id) else { - return Err(StorageError::NotFound(id.to_string()).into()); - }; - try_remove(&path)?; + let path = locate_session_file(dir, id); + if let Some(path) = &path { + try_remove(path)?; + } try_remove(&openai_response_chain_path(dir, id))?; remove_from_cwd_index(dir, id)?; + if path.is_none() { + return Err(StorageError::NotFound(id.to_string()).into()); + } Ok(()) } @@ -4740,6 +4743,30 @@ mod tests { )); } + #[test] + fn delete_missing_primary_cleans_sidecar_and_cwd_index() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path(); + let mut session: TestSession = Session::new("m", "/orphaned"); + session.save_to(dir).unwrap(); + let sidecar = openai_response_chain_path(dir, session.id); + fs::write(&sidecar, b"orphaned").unwrap(); + fs::remove_file(jsonl_path(dir, session.id)).unwrap(); + + let error = TestSession::delete_from(session.id, dir).unwrap_err(); + + assert!(matches!( + error, + SessionError::Storage(StorageError::NotFound(_)) + )); + assert!(!sidecar.exists()); + assert!( + !load_cwd_index(dir) + .values() + .any(|value| *value == session.id.to_string()) + ); + } + #[test] fn title_unicode_safe() { let input = "あ".repeat(100); diff --git a/n00n-ui/src/storage_writer.rs b/n00n-ui/src/storage_writer.rs index 08300fe27..26593d8f8 100644 --- a/n00n-ui/src/storage_writer.rs +++ b/n00n-ui/src/storage_writer.rs @@ -1,31 +1,75 @@ //! Coalescing write-behind cache with incremental JSONL persistence. //! -//! Apps post session snapshots keyed by session id; the writer thread drains -//! the newest snapshot of every session per wake and performs O(delta) -//! appends. Deletes run on the same thread, so an append and a delete of the -//! same session can never race: a queued save cannot resurrect deleted files. +//! Callers only assign an ordering generation and enqueue commands. The single +//! writer thread owns all mutable persistence state, coalesces snapshots, and +//! performs filesystem I/O. Generations make command effects deterministic even +//! when concurrent callers enqueue in the opposite order. Per-session command +//! tracking retains ordering barriers only while an older reserved command can +//! still arrive, then collects them after that command or wake is consumed. -use std::collections::HashMap; +use std::collections::{BTreeSet, HashMap}; use std::io; use std::mem; use std::path::Path; use std::sync::{Arc, Mutex}; use std::time::Duration; -use n00n_storage::StateDir; use n00n_storage::id::n00nId; use n00n_storage::sessions::{SESSIONS_DIR, SessionError, SessionLog}; +use n00n_storage::{StateDir, StorageError}; use tracing::warn; use crate::AppSession; +const RETRY_DELAY: Duration = Duration::from_secs(1); +const MAX_RETRY_ATTEMPTS: u32 = 5; + +#[derive(Clone)] struct PendingSnapshot { + version: SnapshotVersion, + session: Arc, +} + +impl PendingSnapshot { + fn new(generation: u64, session: Box) -> Self { + Self { + version: SnapshotVersion { + generation, + revision: session.meta.revision, + }, + session: Arc::from(session), + } + } +} + +#[derive(Default)] +struct SnapshotInboxState { + snapshots: HashMap, + wake_queued: bool, +} + +type SnapshotInbox = Arc>; + +#[derive(Default)] +struct CommandTrackerState { + next_generation: u64, + outstanding: HashMap>, +} + +type CommandTracker = Arc>; + +#[derive(Clone, Debug, PartialEq, Eq)] +struct SnapshotVersion { + generation: u64, revision: u64, - session: Box, } -type Pending = Arc>>; -type FailedRevisions = HashMap; +struct FailedSnapshot { + version: SnapshotVersion, + error: Option, +} + +type FailedSnapshots = HashMap; #[derive(Debug, thiserror::Error)] pub(crate) enum StorageWriterShutdownError { @@ -40,48 +84,77 @@ pub(crate) enum StorageWriterShutdownError { #[derive(Default)] struct RetryState { attempts: HashMap<(n00nId, u64), u32>, - exhausted: HashMap, + exhausted: HashMap, +} + +#[derive(Default)] +struct WriterState { + pending: HashMap, + latest_generations: HashMap, + latest_snapshots: HashMap, + delete_generations: HashMap, + logs: HashMap, + durable_versions: HashMap, + retries: RetryState, } type DeleteCallback = Box) + Send>; type PersistCallback = Box) + Send>; -const RETRY_DELAY: Duration = Duration::from_secs(1); -const MAX_RETRY_ATTEMPTS: u32 = 5; - enum Op { Flush, + Wake, Persist { + generation: u64, session: Box, done: PersistCallback, }, Delete { id: n00nId, + generation: u64, done: DeleteCallback, }, + #[cfg(test)] + Pause { + entered: flume::Sender<()>, + release: flume::Receiver<()>, + }, + #[cfg(test)] + Inspect { + done: flume::Sender, + }, +} + +#[cfg(test)] +#[derive(Debug, PartialEq, Eq)] +struct WriterStateCounts { + latest_generations: usize, + delete_generations: usize, + outstanding_commands: usize, } pub struct StorageWriter { - pending: Pending, + tracker: CommandTracker, + inbox: SnapshotInbox, ops: flume::Sender, done_rx: flume::Receiver>, } impl StorageWriter { pub fn new(dir: StateDir) -> std::io::Result { - let pending: Pending = Arc::default(); - let writer_pending = Arc::clone(&pending); + let inbox = SnapshotInbox::default(); + let writer_inbox = Arc::clone(&inbox); + let tracker = CommandTracker::default(); + let writer_tracker = Arc::clone(&tracker); let (ops, ops_rx) = flume::unbounded::(); let (done_tx, done_rx) = flume::bounded::>(1); std::thread::Builder::new() .name("storage-writer".into()) .spawn(move || { - let mut logs: HashMap = HashMap::new(); - let mut durable_revisions: HashMap = HashMap::new(); - let mut retries = RetryState::default(); + let mut state = WriterState::default(); loop { - let op = if lock(&writer_pending).is_empty() { + let op = if state.pending.is_empty() { ops_rx.recv().ok() } else { match ops_rx.recv_timeout(RETRY_DELAY) { @@ -93,30 +166,60 @@ impl StorageWriter { let Some(op) = op else { break }; match op { Op::Flush => { - let failed = - flush(&writer_pending, &mut logs, &mut durable_revisions, &dir); - retries.record_failures(&writer_pending, failed); + state.flush(&dir); + } + Op::Wake => { + state.stage_snapshots(take_snapshots(&writer_inbox), &writer_tracker); + if ops_rx.is_empty() { + state.flush(&dir); + } + } + Op::Persist { + generation, + session, + done, + } => { + state.stage_snapshots(take_snapshots(&writer_inbox), &writer_tracker); + let id = session.id; + let result = state.persist(generation, session, &dir); + complete_command(&writer_tracker, id, generation); + state.collect_barriers(id, &writer_tracker); + done(result); + } + Op::Delete { + id, + generation, + done, + } => { + state.stage_snapshots(take_snapshots(&writer_inbox), &writer_tracker); + let result = state.delete(id, generation, &dir); + complete_command(&writer_tracker, id, generation); + state.collect_barriers(id, &writer_tracker); + done(result); } - Op::Persist { session, done } => { - done(flush_and_persist( - &writer_pending, - &mut logs, - &mut durable_revisions, - &mut retries, - &dir, - &session, - )); + #[cfg(test)] + Op::Pause { entered, release } => { + let _ = entered.send(()); + let _ = release.recv(); } - Op::Delete { id, done } => { - lock(&writer_pending).remove(&id); - logs.remove(&id); - done(AppSession::delete(id, &dir)); + #[cfg(test)] + Op::Inspect { done } => { + let outstanding_commands = lock_tracker(&writer_tracker) + .outstanding + .values() + .map(BTreeSet::len) + .sum(); + let _ = done.send(WriterStateCounts { + latest_generations: state.latest_generations.len(), + delete_generations: state.delete_generations.len(), + outstanding_commands, + }); } } } - let failed = flush(&writer_pending, &mut logs, &mut durable_revisions, &dir); - drop(logs); - let unpersisted = retries.unpersisted_count(&failed, &durable_revisions); + state.stage_snapshots(take_snapshots(&writer_inbox), &writer_tracker); + let failed = state.flush(&dir); + let unpersisted = state.retries.unpersisted_count(&failed); let completion = if unpersisted == 0 { Ok(()) } else { @@ -126,26 +229,17 @@ impl StorageWriter { })?; Ok(Self { - pending, + tracker, + inbox, ops, done_rx, }) } pub fn send(&self, session: Box) { - let mut pending = lock(&self.pending); - let was_empty = pending.is_empty(); - let revision = session.meta.revision; - let replace = pending - .get(&session.id) - .is_none_or(|current| current.revision <= revision); - if replace { - pending.insert(session.id, PendingSnapshot { revision, session }); - } - drop(pending); - if was_empty { - let _ = self.ops.send(Op::Flush); - } + let id = session.id; + let generation = reserve_command(&self.tracker, id); + self.enqueue_reserved_snapshot(generation, session); } /// Persists this snapshot before invoking `done` on the writer thread. @@ -154,28 +248,73 @@ impl StorageWriter { session: Box, done: impl FnOnce(Result<(), SessionError>) + Send + 'static, ) { + let id = session.id; + let generation = reserve_command(&self.tracker, id); let op = Op::Persist { + generation, session, done: Box::new(done), }; if let Err(flume::SendError(Op::Persist { done, .. })) = self.ops.send(op) { + complete_command(&self.tracker, id, generation); done(Err(writer_gone())); } } - /// Delete a session's files on the writer thread, discarding any pending - /// snapshot first. Runs after already-queued flushes; `done` fires on the - /// writer thread, so callers never block on disk. + /// Deletes a session on the writer thread after superseded commands have + /// been rejected by generation. The caller never waits for filesystem I/O. pub fn delete(&self, id: n00nId, done: impl FnOnce(Result<(), SessionError>) + Send + 'static) { + let generation = reserve_command(&self.tracker, id); let op = Op::Delete { id, + generation, done: Box::new(done), }; if let Err(flume::SendError(Op::Delete { done, .. })) = self.ops.send(op) { + complete_command(&self.tracker, id, generation); done(Err(writer_gone())); } } + #[cfg(test)] + fn enqueue_snapshot(&self, generation: u64, session: Box) { + reserve_explicit_command(&self.tracker, session.id, generation); + self.enqueue_reserved_snapshot(generation, session); + } + + fn enqueue_reserved_snapshot(&self, generation: u64, session: Box) { + let snapshot = PendingSnapshot::new(generation, session); + let id = snapshot.session.id; + let (should_wake, completed_generation) = { + let mut inbox = lock_inbox(&self.inbox); + let replace = inbox + .snapshots + .get(&id) + .is_none_or(|current| current.version.generation < generation); + if replace { + let replaced = inbox + .snapshots + .insert(id, snapshot) + .map(|snapshot| snapshot.version.generation); + let should_wake = if inbox.wake_queued { + false + } else { + inbox.wake_queued = true; + true + }; + (should_wake, replaced) + } else { + (false, Some(generation)) + } + }; + if let Some(completed_generation) = completed_generation { + complete_command(&self.tracker, id, completed_generation); + } + if should_wake && self.ops.send(Op::Wake).is_err() { + lock_inbox(&self.inbox).wake_queued = false; + } + } + pub(crate) fn shutdown(self, timeout: Duration) -> Result<(), StorageWriterShutdownError> { self.wait_for_shutdown(timeout) } @@ -199,21 +338,257 @@ impl StorageWriter { } } -fn lock(pending: &Pending) -> std::sync::MutexGuard<'_, HashMap> { - pending +fn lock_inbox(inbox: &SnapshotInbox) -> std::sync::MutexGuard<'_, SnapshotInboxState> { + inbox .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } +impl CommandTrackerState { + fn reserve(&mut self, id: n00nId) -> u64 { + self.next_generation += 1; + let generation = self.next_generation; + self.outstanding.entry(id).or_default().insert(generation); + generation + } -fn writer_gone() -> SessionError { - n00n_storage::StorageError::Io(io::Error::other("storage writer unavailable")).into() + #[cfg(test)] + fn reserve_explicit(&mut self, id: n00nId, generation: u64) { + self.next_generation = self.next_generation.max(generation); + self.outstanding.entry(id).or_default().insert(generation); + } + + fn complete(&mut self, id: n00nId, generation: u64) { + let remove_id = self.outstanding.get_mut(&id).is_some_and(|generations| { + generations.remove(&generation); + generations.is_empty() + }); + if remove_id { + self.outstanding.remove(&id); + } + } + + fn has_outstanding_through(&self, id: n00nId, generation: u64) -> bool { + self.outstanding + .get(&id) + .and_then(BTreeSet::first) + .is_some_and(|oldest| *oldest <= generation) + } +} + +fn lock_tracker(tracker: &CommandTracker) -> std::sync::MutexGuard<'_, CommandTrackerState> { + tracker + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn reserve_command(tracker: &CommandTracker, id: n00nId) -> u64 { + lock_tracker(tracker).reserve(id) +} + +#[cfg(test)] +fn reserve_explicit_command(tracker: &CommandTracker, id: n00nId, generation: u64) { + lock_tracker(tracker).reserve_explicit(id, generation); +} + +fn complete_command(tracker: &CommandTracker, id: n00nId, generation: u64) { + lock_tracker(tracker).complete(id, generation); +} + +fn take_snapshots(inbox: &SnapshotInbox) -> HashMap { + let mut inbox = lock_inbox(inbox); + inbox.wake_queued = false; + mem::take(&mut inbox.snapshots) +} + +impl WriterState { + fn stage_snapshots( + &mut self, + snapshots: HashMap, + tracker: &CommandTracker, + ) { + for snapshot in snapshots.into_values() { + let id = snapshot.session.id; + let generation = snapshot.version.generation; + self.stage_snapshot(snapshot); + complete_command(tracker, id, generation); + self.collect_barriers(id, tracker); + } + } + + fn collect_barriers(&mut self, id: n00nId, tracker: &CommandTracker) { + let tracker = lock_tracker(tracker); + if self + .latest_generations + .get(&id) + .is_some_and(|generation| !tracker.has_outstanding_through(id, *generation)) + { + self.latest_generations.remove(&id); + } + if self + .delete_generations + .get(&id) + .is_some_and(|generation| !tracker.has_outstanding_through(id, *generation)) + { + self.delete_generations.remove(&id); + } + } + + fn stage_snapshot(&mut self, snapshot: PendingSnapshot) -> Option { + let id = snapshot.session.id; + let generation = snapshot.version.generation; + if self + .latest_generations + .get(&id) + .is_some_and(|latest| *latest >= generation) + || self + .delete_generations + .get(&id) + .is_some_and(|deleted| *deleted >= generation) + { + return self + .pending + .get(&id) + .map(|snapshot| snapshot.version.clone()); + } + + self.latest_generations.insert(id, generation); + self.latest_snapshots.insert(id, snapshot.clone()); + let replace = self.pending.get(&id).is_none_or(|current| { + current.version.revision < snapshot.version.revision + || (current.version.revision == snapshot.version.revision + && current.version.generation < snapshot.version.generation) + }); + if replace { + self.pending.insert(id, snapshot); + } + self.pending.get(&id).map(|pending| pending.version.clone()) + } + + fn flush(&mut self, dir: &StateDir) -> FailedSnapshots { + self.flush_target(dir, None) + } + + fn flush_target( + &mut self, + dir: &StateDir, + target: Option<(n00nId, &SnapshotVersion)>, + ) -> FailedSnapshots { + let (failed, persisted) = flush( + &mut self.pending, + &mut self.logs, + &mut self.durable_versions, + dir, + target, + ); + for (id, version) in &persisted { + let snapshot_is_durable = self.durable_versions.get(id) == Some(version); + if snapshot_is_durable + && self + .latest_snapshots + .get(id) + .is_some_and(|snapshot| snapshot.version == *version) + { + self.latest_snapshots.remove(id); + } + } + self.retries.record_successes(&persisted); + self.retries.record_failures(&mut self.pending, &failed); + failed + } + + fn persist( + &mut self, + generation: u64, + session: Box, + dir: &StateDir, + ) -> Result<(), SessionError> { + let snapshot = PendingSnapshot::new(generation, session); + let id = snapshot.session.id; + let target = self.stage_snapshot(snapshot); + let mut failed = self.flush_target(dir, target.as_ref().map(|version| (id, version))); + if let Some(target) = target + && let Some(mut failure) = failed.remove(&id) + && failure.version == target + { + return Err(failure.error.take().unwrap_or_else(unpersisted_snapshot)); + } + if self.retries.exhausted.contains_key(&id) { + return Err(unpersisted_snapshot()); + } + Ok(()) + } + + fn delete(&mut self, id: n00nId, generation: u64, dir: &StateDir) -> Result<(), SessionError> { + let delete_generation = self + .delete_generations + .entry(id) + .and_modify(|current| *current = (*current).max(generation)) + .or_insert(generation); + let delete_generation = *delete_generation; + if self + .latest_generations + .get(&id) + .is_none_or(|latest| *latest < generation) + { + self.latest_generations.insert(id, generation); + } + self.pending.retain(|session_id, snapshot| { + *session_id != id || snapshot.version.generation > delete_generation + }); + self.latest_snapshots.retain(|session_id, snapshot| { + *session_id != id || snapshot.version.generation > delete_generation + }); + + let durable_is_newer = self + .durable_versions + .get(&id) + .is_some_and(|version| version.generation > delete_generation); + if durable_is_newer { + return Ok(()); + } + + let crossing_snapshot = self.latest_snapshots.get(&id).cloned(); + let crossing_target = crossing_snapshot + .as_ref() + .map(|snapshot| snapshot.version.clone()); + if let Some(snapshot) = crossing_snapshot { + self.pending.insert(id, snapshot); + } + + let primary_path = dir.path().join(SESSIONS_DIR).join(format!("{id}.jsonl")); + self.logs.remove(&id); + let delete_result = AppSession::delete(id, dir); + if !primary_path.exists() { + self.durable_versions.remove(&id); + self.retries.clear(id); + } + match delete_result { + Ok(()) | Err(SessionError::Storage(StorageError::NotFound(_))) => {} + Err(error) => return Err(error), + } + + let Some(target) = crossing_target else { + return Ok(()); + }; + let mut failed = self.flush_target(dir, Some((id, &target))); + if let Some(mut failure) = failed.remove(&id) { + return Err(failure.error.take().unwrap_or_else(unpersisted_snapshot)); + } + Ok(()) + } } impl RetryState { - fn record_failures(&mut self, pending: &Pending, failed: FailedRevisions) { - let mut pending = lock(pending); - for (id, revision) in failed { - let attempts = self.attempts.entry((id, revision)).or_default(); + fn record_failures( + &mut self, + pending: &mut HashMap, + failed: &FailedSnapshots, + ) { + for (id, failure) in failed { + let attempts = self + .attempts + .entry((*id, failure.version.generation)) + .or_default(); *attempts += 1; if *attempts < MAX_RETRY_ATTEMPTS { continue; @@ -221,129 +596,129 @@ impl RetryState { warn!( retry_count = *attempts, %id, - revision, + revision = failure.version.revision, "storage writer exhausted retry attempts, dropping snapshot" ); if pending - .get(&id) - .is_some_and(|snapshot| snapshot.revision == revision) + .get(id) + .is_some_and(|snapshot| snapshot.version == failure.version) { - pending.remove(&id); - self.exhausted - .entry(id) - .and_modify(|current| *current = (*current).max(revision)) - .or_insert(revision); + pending.remove(id); + self.exhausted.insert(*id, failure.version.clone()); } } - self.attempts.retain(|(id, revision), _| { + self.attempts.retain(|(session_id, generation), _| { pending - .get(id) - .is_some_and(|snapshot| snapshot.revision == *revision) + .get(session_id) + .is_some_and(|snapshot| snapshot.version.generation == *generation) }); } - fn unpersisted_count( - &self, - failed: &FailedRevisions, - durable_revisions: &HashMap, - ) -> usize { + fn record_successes(&mut self, persisted: &[(n00nId, SnapshotVersion)]) { + for (session_id, snapshot) in persisted { + let supersedes_exhausted = self.exhausted.get(session_id).is_some_and(|exhausted| { + exhausted.generation != snapshot.generation + && exhausted.revision <= snapshot.revision + }); + if supersedes_exhausted { + self.exhausted.remove(session_id); + } + } + } + + fn clear(&mut self, id: n00nId) { + self.attempts.retain(|(session_id, _), _| *session_id != id); + self.exhausted.remove(&id); + } + + fn unpersisted_count(&self, failed: &FailedSnapshots) -> usize { failed.len() + self .exhausted - .iter() - .filter(|(id, revision)| { - !failed.contains_key(id) - && durable_revisions - .get(id) - .is_none_or(|durable| durable < revision) - }) + .keys() + .filter(|id| !failed.contains_key(id)) .count() } } -fn flush_and_persist( - pending: &Pending, - logs: &mut HashMap, - durable_revisions: &mut HashMap, - retries: &mut RetryState, - dir: &StateDir, - session: &AppSession, -) -> Result<(), SessionError> { - let failed = flush(pending, logs, durable_revisions, dir); - retries.record_failures(pending, failed); - persist_session(pending, logs, durable_revisions, dir, session) +fn writer_gone() -> SessionError { + StorageError::Io(io::Error::other("storage writer unavailable")).into() +} + +fn unpersisted_snapshot() -> SessionError { + StorageError::Io(io::Error::other( + "newer session snapshot remains unpersisted", + )) + .into() } fn flush( - pending: &Pending, + pending: &mut HashMap, logs: &mut HashMap, - durable_revisions: &mut HashMap, + durable_versions: &mut HashMap, dir: &StateDir, -) -> FailedRevisions { - let mut pending_guard = lock(pending); - let batch = mem::take(&mut *pending_guard); + target: Option<(n00nId, &SnapshotVersion)>, +) -> (FailedSnapshots, Vec<(n00nId, SnapshotVersion)>) { + let batch = mem::take(pending); if batch.is_empty() { - return FailedRevisions::new(); + return (FailedSnapshots::new(), Vec::new()); } let sessions_dir = match dir.ensure_subdir(SESSIONS_DIR) { - Ok(d) => d, - Err(e) => { - warn!(error = %e, "failed to ensure sessions dir"); - let mut failed = FailedRevisions::with_capacity(batch.len()); + Ok(sessions_dir) => sessions_dir, + Err(error) => { + warn!(error = %error, "failed to ensure sessions dir"); + let mut target_error = Some(SessionError::from(error)); + let mut failed = FailedSnapshots::with_capacity(batch.len()); for (id, snapshot) in batch { - failed.insert(id, snapshot.revision); - let replace = pending_guard - .get(&id) - .is_none_or(|current| current.revision < snapshot.revision); - if replace { - pending_guard.insert(id, snapshot); - } + let error = if target.is_some_and(|(target_id, version)| { + target_id == id && *version == snapshot.version + }) { + target_error.take() + } else { + None + }; + failed.insert( + id, + FailedSnapshot { + version: snapshot.version.clone(), + error, + }, + ); + pending.insert(id, snapshot); } - return failed; + return (failed, Vec::new()); } }; - let mut failed = FailedRevisions::new(); - for snapshot in batch.into_values() { - if let Err(error) = write_session(&sessions_dir, logs, durable_revisions, &snapshot.session) - { - warn!(error = %error, id = %snapshot.session.id, "session write failed"); - let id = snapshot.session.id; - failed.insert(id, snapshot.revision); - let replace = pending_guard - .get(&id) - .is_none_or(|current| current.revision < snapshot.revision); - if replace { - pending_guard.insert(id, snapshot); - } - } - } - failed -} -fn persist_session( - pending: &Pending, - logs: &mut HashMap, - durable_revisions: &mut HashMap, - dir: &StateDir, - session: &AppSession, -) -> Result<(), SessionError> { - let mut pending_guard = lock(pending); - if pending_guard - .get(&session.id) - .is_some_and(|snapshot| snapshot.revision > session.meta.revision) - { - let snapshot = pending_guard.remove(&session.id).ok_or_else(|| { - n00n_storage::StorageError::Io(io::Error::other("pending snapshot disappeared")) - })?; - write_session( - &dir.ensure_subdir(SESSIONS_DIR)?, + let mut failed = FailedSnapshots::new(); + let mut persisted = Vec::new(); + for snapshot in batch.into_values() { + let id = snapshot.session.id; + match write_session( + &sessions_dir, logs, - durable_revisions, + durable_versions, + &snapshot.version, &snapshot.session, - )?; + ) { + Ok(()) => persisted.push((id, snapshot.version)), + Err(error) => { + warn!(error = %error, %id, "session write failed"); + let is_target = target.is_some_and(|(target_id, version)| { + target_id == id && *version == snapshot.version + }); + failed.insert( + id, + FailedSnapshot { + version: snapshot.version.clone(), + error: is_target.then_some(error), + }, + ); + pending.insert(id, snapshot); + } + } } - drop(pending_guard); - write_session_if_newer(logs, durable_revisions, dir, session) + (failed, persisted) } fn append_or_compact_result( @@ -357,56 +732,57 @@ fn append_or_compact_result( Err(error) => Err(error), } } + fn write_session( sessions_dir: &Path, logs: &mut HashMap, - durable_revisions: &mut HashMap, + durable_versions: &mut HashMap, + version: &SnapshotVersion, session: &AppSession, ) -> Result<(), SessionError> { - if durable_revisions + if durable_versions .get(&session.id) - .is_some_and(|revision| *revision > session.meta.revision) + .is_some_and(|durable| durable.revision > version.revision) { return Ok(()); } if let Some(log) = logs.get_mut(&session.id) { - if durable_revisions.get(&session.id) == Some(&session.meta.revision) { + if durable_versions + .get(&session.id) + .is_some_and(|durable| durable.revision == version.revision) + { log.compact(sessions_dir, session)?; } else { append_or_compact_result(log, sessions_dir, session)?; } - durable_revisions.insert(session.id, session.meta.revision); + durable_versions.insert(session.id, version.clone()); return Ok(()); } let (mut log, on_disk_revision) = open_or_create_log(sessions_dir, session)?; - if on_disk_revision > session.meta.revision { - durable_revisions.insert(session.id, on_disk_revision); + if on_disk_revision > version.revision { + durable_versions.insert( + session.id, + SnapshotVersion { + generation: 0, + revision: on_disk_revision, + }, + ); return Ok(()); } - if on_disk_revision == session.meta.revision { + if on_disk_revision == version.revision { log.compact(sessions_dir, session)?; } else { append_or_compact_result(&mut log, sessions_dir, session)?; } logs.insert(session.id, log); - durable_revisions.insert(session.id, session.meta.revision); + durable_versions.insert(session.id, version.clone()); Ok(()) } -fn write_session_if_newer( - logs: &mut HashMap, - durable_revisions: &mut HashMap, - dir: &StateDir, - session: &AppSession, -) -> Result<(), SessionError> { - let sessions_dir = dir.ensure_subdir(SESSIONS_DIR)?; - write_session(&sessions_dir, logs, durable_revisions, session) -} - fn open_or_create_log( sessions_dir: &Path, session: &AppSession, -) -> Result<(SessionLog, u64), n00n_storage::sessions::SessionError> { +) -> Result<(SessionLog, u64), SessionError> { let jsonl_path = sessions_dir.join(format!("{}.jsonl", session.id)); if jsonl_path.exists() { let id = session.id; @@ -428,9 +804,15 @@ fn open_or_create_log( mod tests { use super::*; use std::fs; + + use n00n_storage::sessions::lock_openai_response_chain; use tempfile::TempDir; const DRAIN_TIMEOUT: Duration = Duration::from_secs(30); + const BLOCKED_TIMEOUT: Duration = Duration::from_millis(100); + const NONBLOCKING_TIMEOUT: Duration = Duration::from_secs(2); + const OPENAI_RESPONSE_SUFFIX: &str = "openai-response.json"; + const STRESS_SNAPSHOT_COUNT: u64 = 10_000; fn state_dir() -> (TempDir, StateDir) { let tmp = TempDir::new().unwrap(); @@ -438,8 +820,42 @@ mod tests { (tmp, dir) } - /// Snapshots must coalesce per session id, not into one `latest` slot: - /// two racing sessions used to silently drop one. + fn pause_writer(writer: &StorageWriter) -> flume::Sender<()> { + let (entered_tx, entered_rx) = flume::bounded(1); + let (release_tx, release_rx) = flume::bounded(1); + writer + .ops + .send(Op::Pause { + entered: entered_tx, + release: release_rx, + }) + .unwrap(); + entered_rx.recv_timeout(DRAIN_TIMEOUT).unwrap(); + release_tx + } + + fn inspect_writer(writer: &StorageWriter) -> WriterStateCounts { + let (done_tx, done_rx) = flume::bounded(1); + writer.ops.send(Op::Inspect { done: done_tx }).unwrap(); + done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap() + } + + fn persist_and_wait(writer: &StorageWriter, session: AppSession) -> Result<(), SessionError> { + let (done_tx, done_rx) = flume::bounded(1); + writer.persist(Box::new(session), move |result| { + done_tx.send(result).unwrap(); + }); + done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap() + } + + fn delete_and_wait(writer: &StorageWriter, id: n00nId) -> Result<(), SessionError> { + let (done_tx, done_rx) = flume::bounded(1); + writer.delete(id, move |result| { + done_tx.send(result).unwrap(); + }); + done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap() + } + #[test] fn shutdown_drains_newest_snapshot_of_every_session() { let (_tmp, dir) = state_dir(); @@ -451,6 +867,7 @@ mod tests { writer.send(Box::new(b.clone())); b.title = "renamed".into(); writer.send(Box::new(b)); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); assert!(AppSession::load(a_id, &dir).is_ok()); @@ -458,250 +875,346 @@ mod tests { } #[test] - fn failed_flush_keeps_snapshot_for_retry() { - let (tmp, dir) = state_dir(); - let pending: Pending = Arc::default(); - let session = AppSession::new("test-model", "/tmp/retry"); + fn blocked_writer_coalesces_same_session_snapshots_and_persists_latest() { + let (_tmp, dir) = state_dir(); + let writer = StorageWriter::new(dir.clone()).unwrap(); + let release = pause_writer(&writer); + let mut session = AppSession::new("test-model", "/tmp/coalesced"); let id = session.id; - fs::create_dir_all(tmp.path().join(SESSIONS_DIR).join(format!("{id}.jsonl"))).unwrap(); - lock(&pending).insert( - id, - PendingSnapshot { - revision: session.meta.revision, - session: Box::new(session), - }, + + for revision in 1..=STRESS_SNAPSHOT_COUNT { + session.meta.revision = revision; + session.meta.input_draft = Some(format!("snapshot {revision}")); + writer.send(Box::new(session.clone())); + } + + let inbox = lock_inbox(&writer.inbox); + assert_eq!(inbox.snapshots.len(), 1); + assert!(inbox.wake_queued); + assert_eq!( + inbox.snapshots[&id].session.meta.input_draft.as_deref(), + Some("snapshot 10000") ); - let mut logs = HashMap::new(); - let mut durable_revisions = HashMap::new(); + drop(inbox); + assert_eq!(writer.ops.len(), 1); + + release.send(()).unwrap(); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); + let loaded = AppSession::load(id, &dir).unwrap(); + assert_eq!(loaded.meta.revision, STRESS_SNAPSHOT_COUNT); + assert_eq!(loaded.meta.input_draft.as_deref(), Some("snapshot 10000")); + } + + #[test] + fn inverse_delete_save_enqueue_order_obeys_generation() { + let (_tmp, dir) = state_dir(); + let writer = StorageWriter::new(dir.clone()).unwrap(); + let deleted = AppSession::new("test-model", "/tmp/delete-newer"); + let deleted_id = deleted.id; + let saved = AppSession::new("test-model", "/tmp/save-newer"); + let saved_id = saved.id; + persist_and_wait(&writer, deleted.clone()).unwrap(); + persist_and_wait(&writer, saved.clone()).unwrap(); + let release = pause_writer(&writer); + + reserve_explicit_command(&writer.tracker, deleted_id, 20); + writer + .ops + .send(Op::Delete { + id: deleted_id, + generation: 20, + done: Box::new(|result| assert!(result.is_ok())), + }) + .unwrap(); + writer.enqueue_snapshot(10, Box::new(deleted)); + let mut newer_saved = saved; + newer_saved.title = "newer save".into(); + writer.enqueue_snapshot(40, Box::new(newer_saved)); + reserve_explicit_command(&writer.tracker, saved_id, 30); + writer + .ops + .send(Op::Delete { + id: saved_id, + generation: 30, + done: Box::new(|result| assert!(result.is_ok())), + }) + .unwrap(); + release.send(()).unwrap(); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); + + assert!(AppSession::load(deleted_id, &dir).is_err()); assert_eq!( - flush(&pending, &mut logs, &mut durable_revisions, &dir), - HashMap::from([(id, 0)]) + AppSession::load(saved_id, &dir).unwrap().title, + "newer save" ); - assert!(lock(&pending).contains_key(&id)); } #[test] - fn shutdown_does_not_report_success_with_unpersisted_snapshot() { - let (tmp, dir) = state_dir(); - let writer = StorageWriter::new(dir).unwrap(); - let session = AppSession::new("test-model", "/tmp/shutdown-failure"); - let id = session.id; - fs::create_dir_all(tmp.path().join(SESSIONS_DIR).join(format!("{id}.jsonl"))).unwrap(); - writer.send(Box::new(session)); + fn delayed_delete_barrier_discards_pre_delete_higher_revision() { + let (_tmp, dir) = state_dir(); + let mut state = WriterState::default(); + let mut pre_delete = AppSession::new("test-model", "/tmp/delete-crossing"); + pre_delete.title = "pre-delete revision five".into(); + pre_delete.meta.revision = 5; + let id = pre_delete.id; + let mut post_delete = pre_delete.clone(); + post_delete.title = "post-delete revision four".into(); + post_delete.meta.revision = 4; - assert!(matches!( - writer.wait_for_shutdown(DRAIN_TIMEOUT), - Err(StorageWriterShutdownError::UnpersistedSnapshots { count: 1 }) - )); + state.stage_snapshot(PendingSnapshot::new(1, Box::new(pre_delete))); + assert!(state.flush(&dir).is_empty()); + state.stage_snapshot(PendingSnapshot::new(3, Box::new(post_delete))); + assert!(state.flush(&dir).is_empty()); + state.delete(id, 2, &dir).unwrap(); + + let loaded = AppSession::load(id, &dir).unwrap(); + assert_eq!(loaded.meta.revision, 4); + assert_eq!(loaded.title, "post-delete revision four"); } #[test] - fn retry_exhaustion_removes_only_exact_failed_revisions() { - let pending: Pending = Arc::default(); - let mut retries = RetryState::default(); - let mut old = AppSession::new("test-model", "/tmp/old"); - old.meta.revision = 1; - let old_id = old.id; - let mut exact = AppSession::new("test-model", "/tmp/exact"); - exact.meta.revision = 3; - let exact_id = exact.id; - let concurrent = AppSession::new("test-model", "/tmp/concurrent"); - let concurrent_id = concurrent.id; - lock(&pending).insert( - old_id, - PendingSnapshot { - revision: 1, - session: Box::new(old.clone()), - }, - ); - lock(&pending).insert( - exact_id, - PendingSnapshot { - revision: 3, - session: Box::new(exact), - }, - ); - let failures = HashMap::from([(old_id, 1), (exact_id, 3)]); - for _ in 1..MAX_RETRY_ATTEMPTS { - retries.record_failures(&pending, failures.clone()); - } + fn partial_delete_requeues_crossing_generation_lower_revision() { + let (_tmp, dir) = state_dir(); + let mut state = WriterState::default(); + let mut pre_delete = AppSession::new("test-model", "/tmp/partial-delete-crossing"); + pre_delete.title = "pre-delete revision five".into(); + pre_delete.meta.revision = 5; + let id = pre_delete.id; + let mut post_delete = pre_delete.clone(); + post_delete.title = "post-delete revision four".into(); + post_delete.meta.revision = 4; + let mut durable = pre_delete.clone(); + durable.title = "durable revision zero".into(); + durable.meta.revision = 0; - old.meta.revision = 2; - lock(&pending).insert( - old_id, - PendingSnapshot { - revision: 2, - session: Box::new(old), - }, - ); - lock(&pending).insert( - concurrent_id, - PendingSnapshot { - revision: concurrent.meta.revision, - session: Box::new(concurrent), - }, - ); - retries.record_failures(&pending, failures); + state.stage_snapshot(PendingSnapshot::new(0, Box::new(durable))); + assert!(state.flush(&dir).is_empty()); + state.stage_snapshot(PendingSnapshot::new(1, Box::new(pre_delete))); + state.stage_snapshot(PendingSnapshot::new(3, Box::new(post_delete))); + let sessions_dir = dir.ensure_subdir(SESSIONS_DIR).unwrap(); + let sidecar_path = sessions_dir.join(format!("{id}.{OPENAI_RESPONSE_SUFFIX}")); + fs::create_dir(&sidecar_path).unwrap(); - let pending = lock(&pending); - assert_eq!( - pending.get(&old_id).map(|snapshot| snapshot.revision), - Some(2) - ); - assert!(pending.contains_key(&concurrent_id)); - assert!(!pending.contains_key(&exact_id)); + assert!(state.delete(id, 2, &dir).is_err()); + assert!(!sessions_dir.join(format!("{id}.jsonl")).exists()); + assert_eq!(state.pending[&id].version.generation, 3); + assert_eq!(state.pending[&id].version.revision, 4); + + assert!(state.flush(&dir).is_empty()); + let loaded = AppSession::load(id, &dir).unwrap(); + assert_eq!(loaded.meta.revision, 4); + assert_eq!(loaded.title, "post-delete revision four"); } #[test] - fn exhausted_retry_remains_a_shutdown_failure_until_durable() { - let pending: Pending = Arc::default(); - let mut retries = RetryState::default(); - let mut session = AppSession::new("test-model", "/tmp/exhausted"); - session.meta.revision = 4; - let id = session.id; - lock(&pending).insert( - id, - PendingSnapshot { - revision: 4, - session: Box::new(session), - }, - ); - let failures = HashMap::from([(id, 4)]); - for _ in 0..MAX_RETRY_ATTEMPTS { - retries.record_failures(&pending, failures.clone()); + fn completed_unique_nonexistent_deletes_do_not_accumulate_tombstones() { + let (_tmp, dir) = state_dir(); + let writer = StorageWriter::new(dir).unwrap(); + let (done_tx, done_rx) = flume::unbounded(); + + for _ in 0..1_000 { + let done_tx = done_tx.clone(); + writer.delete(n00nId::generate(), move |result| { + done_tx.send(result).unwrap(); + }); + } + drop(done_tx); + for _ in 0..1_000 { + assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); } - assert!(lock(&pending).is_empty()); - assert_eq!( - retries.unpersisted_count(&FailedRevisions::new(), &HashMap::new()), - 1 - ); assert_eq!( - retries.unpersisted_count(&FailedRevisions::new(), &HashMap::from([(id, 4)]),), - 0 + inspect_writer(&writer), + WriterStateCounts { + latest_generations: 0, + delete_generations: 0, + outstanding_commands: 0, + } ); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); } #[test] - fn exhausted_retry_for_one_session_does_not_block_explicit_persist() { - let (tmp, dir) = state_dir(); - let pending: Pending = Arc::default(); - let failing = AppSession::new("test-model", "/tmp/failing"); - let failing_id = failing.id; - fs::create_dir_all( - tmp.path() - .join(SESSIONS_DIR) - .join(format!("{failing_id}.jsonl")), - ) - .unwrap(); - lock(&pending).insert( - failing_id, - PendingSnapshot { - revision: failing.meta.revision, - session: Box::new(failing), - }, - ); - let failures = HashMap::from([(failing_id, 0)]); - let mut retries = RetryState::default(); - for _ in 1..MAX_RETRY_ATTEMPTS { - retries.record_failures(&pending, failures.clone()); - } - let explicit = AppSession::new("test-model", "/tmp/explicit"); - let explicit_id = explicit.id; - let mut logs = HashMap::new(); - let mut durable_revisions = HashMap::new(); - - assert!( - flush_and_persist( - &pending, - &mut logs, - &mut durable_revisions, - &mut retries, - &dir, - &explicit, - ) - .is_ok() - ); - assert!(AppSession::load(explicit_id, &dir).is_ok()); + fn explicit_persist_preserves_ensure_subdir_io_error() { + let tmp = TempDir::new().unwrap(); + let state_path = tmp.path().join("state-file"); + fs::write(&state_path, b"not a directory").unwrap(); + let expected = fs::create_dir_all(state_path.join(SESSIONS_DIR)).unwrap_err(); + let writer = StorageWriter::new(StateDir::from_path(state_path)).unwrap(); + let session = AppSession::new("test-model", "/tmp/ensure-subdir-error"); + + let error = persist_and_wait(&writer, session).unwrap_err(); + let SessionError::Storage(StorageError::Io(error)) = error else { + panic!("expected storage I/O error"); + }; + assert_eq!(error.kind(), expected.kind()); + assert_eq!(error.raw_os_error(), expected.raw_os_error()); + assert!(matches!( + writer.shutdown(DRAIN_TIMEOUT), + Err(StorageWriterShutdownError::UnpersistedSnapshots { count: 1 }) + )); } #[test] - fn persist_reports_success_after_writing_snapshot() { + fn stale_equal_revision_persist_cannot_overwrite_newer_save() { let (_tmp, dir) = state_dir(); let writer = StorageWriter::new(dir.clone()).unwrap(); - let session = AppSession::new("test-model", "/tmp/persist"); - let id = session.id; + let mut stale = AppSession::new("test-model", "/tmp/equal-race"); + stale.title = "stale persist".into(); + let id = stale.id; + let mut newer = stale.clone(); + newer.title = "newer save".into(); + let release = pause_writer(&writer); let (done_tx, done_rx) = flume::bounded(1); - writer.persist(Box::new(session), move |result| { - let _ = done_tx.send(result); - }); + + reserve_explicit_command(&writer.tracker, id, 1); + writer.enqueue_snapshot(2, Box::new(newer)); + writer + .ops + .send(Op::Persist { + generation: 1, + session: Box::new(stale), + done: Box::new(move |result| done_tx.send(result).unwrap()), + }) + .unwrap(); + release.send(()).unwrap(); assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); - assert!(AppSession::load(id, &dir).is_ok()); writer.shutdown(DRAIN_TIMEOUT).unwrap(); + assert_eq!(AppSession::load(id, &dir).unwrap().title, "newer save"); } #[test] - fn equal_revision_different_snapshot_is_not_dropped() { - let (_tmp, dir) = state_dir(); - let writer = StorageWriter::new(dir.clone()).unwrap(); - let first = AppSession::new("test-model", "/tmp/equal"); - let id = first.id; + fn failed_stale_persist_retains_newer_snapshot_for_shutdown() { + let (tmp, dir) = state_dir(); + let writer = StorageWriter::new(dir).unwrap(); + let mut stale = AppSession::new("test-model", "/tmp/persist-failure"); + stale.title = "stale persist".into(); + let id = stale.id; + let mut newer = stale.clone(); + newer.title = "newer pending".into(); + fs::create_dir_all(tmp.path().join(SESSIONS_DIR).join(format!("{id}.jsonl"))).unwrap(); + let release = pause_writer(&writer); let (done_tx, done_rx) = flume::bounded(1); - writer.persist(Box::new(first), move |result| { - let _ = done_tx.send(result); + + reserve_explicit_command(&writer.tracker, id, 1); + writer.enqueue_snapshot(2, Box::new(newer)); + writer + .ops + .send(Op::Persist { + generation: 1, + session: Box::new(stale), + done: Box::new(move |result| done_tx.send(result).unwrap()), + }) + .unwrap(); + release.send(()).unwrap(); + + assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_err()); + assert!(matches!( + writer.shutdown(DRAIN_TIMEOUT), + Err(StorageWriterShutdownError::UnpersistedSnapshots { count: 1 }) + )); + } + + #[test] + fn retry_exhaustion_remains_accounted_at_shutdown() { + let (tmp, dir) = state_dir(); + let mut state = WriterState::default(); + let session = AppSession::new("test-model", "/tmp/exhausted"); + let id = session.id; + fs::create_dir_all(tmp.path().join(SESSIONS_DIR).join(format!("{id}.jsonl"))).unwrap(); + state.stage_snapshot(PendingSnapshot::new(1, Box::new(session))); + + for _ in 0..MAX_RETRY_ATTEMPTS { + state.flush(&dir); + } + + assert!(state.pending.is_empty()); + assert_eq!(state.retries.unpersisted_count(&FailedSnapshots::new()), 1); + } + + #[test] + fn public_operations_do_not_block_behind_writer_filesystem_io() { + let (_tmp, dir) = state_dir(); + let writer = Arc::new(StorageWriter::new(dir.clone()).unwrap()); + let session = AppSession::new("test-model", "/tmp/nonblocking"); + let id = session.id; + persist_and_wait(&writer, session).unwrap(); + let response_lock = lock_openai_response_chain(&dir, id).unwrap(); + let (delete_tx, delete_rx) = flume::bounded(1); + writer.delete(id, move |result| delete_tx.send(result).unwrap()); + let barrier = AppSession::new("test-model", "/tmp/io-barrier"); + let (barrier_tx, barrier_rx) = flume::bounded(1); + writer.persist(Box::new(barrier), move |result| { + barrier_tx.send(result).unwrap(); }); - assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); + assert!(matches!( + barrier_rx.recv_timeout(BLOCKED_TIMEOUT), + Err(flume::RecvTimeoutError::Timeout) + )); - let mut second = AppSession::load(id, &dir).unwrap(); - second.title = "same revision, new snapshot".into(); - writer.send(Box::new(second)); - writer.shutdown(DRAIN_TIMEOUT).unwrap(); + let caller_writer = Arc::clone(&writer); + let (returned_tx, returned_rx) = flume::bounded(1); + let caller = std::thread::spawn(move || { + caller_writer.send(Box::new(AppSession::new( + "test-model", + "/tmp/nonblocking-save", + ))); + caller_writer.delete(n00nId::generate(), |_| {}); + returned_tx.send(()).unwrap(); + }); + returned_rx.recv_timeout(NONBLOCKING_TIMEOUT).unwrap(); + caller.join().unwrap(); - assert_eq!( - AppSession::load(id, &dir).unwrap().title, - "same revision, new snapshot" - ); + drop(response_lock); + assert!(delete_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); + assert!(barrier_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); + Arc::into_inner(writer) + .unwrap() + .shutdown(DRAIN_TIMEOUT) + .unwrap(); } #[test] - fn persist_cannot_overwrite_newer_periodic_snapshot() { + fn partial_delete_closes_cached_log_before_later_persist() { let (_tmp, dir) = state_dir(); let writer = StorageWriter::new(dir.clone()).unwrap(); - let mut older = AppSession::new("test-model", "/tmp/race"); - older.meta.revision = 1; - older.title = "submission".into(); - let id = older.id; - let mut newer = older.clone(); - newer.meta.revision = 2; - newer.title = "periodic save".into(); - writer.send(Box::new(newer)); + let mut session = AppSession::new("test-model", "/tmp/partial-delete"); + let id = session.id; + persist_and_wait(&writer, session.clone()).unwrap(); + let sessions_dir = dir.ensure_subdir(SESSIONS_DIR).unwrap(); + let log_path = sessions_dir.join(format!("{id}.jsonl")); + let sidecar_path = sessions_dir.join(format!("{id}.{OPENAI_RESPONSE_SUFFIX}")); + fs::create_dir(&sidecar_path).unwrap(); - let (done_tx, done_rx) = flume::bounded(1); - writer.persist(Box::new(older), move |result| { - let _ = done_tx.send(result); - }); + assert!(delete_and_wait(&writer, id).is_err()); + assert!(!log_path.exists()); + fs::remove_dir(sidecar_path).unwrap(); - assert!(done_rx.recv_timeout(DRAIN_TIMEOUT).unwrap().is_ok()); - assert_eq!(AppSession::load(id, &dir).unwrap().title, "periodic save"); + session.meta.input_draft = Some("persisted after partial delete".into()); + session.meta.revision += 1; + persist_and_wait(&writer, session).unwrap(); + assert_eq!( + AppSession::load(id, &dir) + .unwrap() + .meta + .input_draft + .as_deref(), + Some("persisted after partial delete") + ); writer.shutdown(DRAIN_TIMEOUT).unwrap(); } #[test] - fn delete_discards_pending_snapshot() { + fn persist_reports_success_after_writing_snapshot() { let (_tmp, dir) = state_dir(); let writer = StorageWriter::new(dir.clone()).unwrap(); - let session = AppSession::new("test-model", "/tmp/c"); + let session = AppSession::new("test-model", "/tmp/persist"); let id = session.id; - writer.send(Box::new(session)); - let (done_tx, done_rx) = flume::bounded(1); - writer.delete(id, move |res| { - let _ = done_tx.send(res); - }); - writer.shutdown(DRAIN_TIMEOUT).unwrap(); - assert!(done_rx.recv().unwrap().is_ok()); - assert!(AppSession::load(id, &dir).is_err()); + assert!(persist_and_wait(&writer, session).is_ok()); + assert!(AppSession::load(id, &dir).is_ok()); + writer.shutdown(DRAIN_TIMEOUT).unwrap(); } } From f99565d6bbfb37b0674d79795144cef85f2b1bba Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Sun, 9 Aug 2026 21:21:07 -0400 Subject: [PATCH 3/4] fix(ui): resolve shutdown review findings --- n00n-ui/src/app/mod.rs | 8 +++++++- n00n-ui/src/storage_writer.rs | 33 +++++++++++++++++++++++++++++++++ 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/n00n-ui/src/app/mod.rs b/n00n-ui/src/app/mod.rs index 0a6750d8c..933df0728 100644 --- a/n00n-ui/src/app/mod.rs +++ b/n00n-ui/src/app/mod.rs @@ -218,7 +218,13 @@ impl Drop for TestStateDir { return; }; let writer_stopped = match Arc::try_unwrap(writer) { - Ok(writer) => writer.wait_for_shutdown(TEST_WRITER_DRAIN_TIMEOUT).is_ok(), + Ok(writer) => match writer.wait_for_shutdown(TEST_WRITER_DRAIN_TIMEOUT) { + Ok(()) => true, + Err(error) => { + tracing::warn!(%error, "test storage writer shutdown failed"); + false + } + }, Err(writer) => { drop(writer); false diff --git a/n00n-ui/src/storage_writer.rs b/n00n-ui/src/storage_writer.rs index 26593d8f8..b5ea3e154 100644 --- a/n00n-ui/src/storage_writer.rs +++ b/n00n-ui/src/storage_writer.rs @@ -1134,6 +1134,39 @@ mod tests { assert_eq!(state.retries.unpersisted_count(&FailedSnapshots::new()), 1); } + #[test] + fn successful_delete_clears_exhausted_retry_state() { + let (_tmp, dir) = state_dir(); + let mut state = WriterState::default(); + let mut session = AppSession::new("test-model", "/tmp/delete-exhausted"); + let id = session.id; + state.stage_snapshot(PendingSnapshot::new(1, Box::new(session.clone()))); + assert!(state.flush(&dir).is_empty()); + + session.meta.revision += 1; + let failed_snapshot = PendingSnapshot::new(2, Box::new(session)); + let failed_version = failed_snapshot.version.clone(); + state.stage_snapshot(failed_snapshot); + let failed = FailedSnapshots::from([( + id, + FailedSnapshot { + version: failed_version.clone(), + error: None, + }, + )]); + for _ in 0..MAX_RETRY_ATTEMPTS { + state.retries.record_failures(&mut state.pending, &failed); + } + assert_eq!(state.retries.exhausted.get(&id), Some(&failed_version)); + assert_eq!(state.retries.unpersisted_count(&FailedSnapshots::new()), 1); + + state.delete(id, 3, &dir).unwrap(); + + assert!(!state.retries.exhausted.contains_key(&id)); + assert_eq!(state.retries.unpersisted_count(&FailedSnapshots::new()), 0); + assert!(AppSession::load(id, &dir).is_err()); + } + #[test] fn public_operations_do_not_block_behind_writer_filesystem_io() { let (_tmp, dir) = state_dir(); From 92aa1b5b2241c53b2230a7ed9ae9fea54b10ba28 Mon Sep 17 00:00:00 2001 From: w0wl0lxd Date: Sun, 9 Aug 2026 21:55:24 -0400 Subject: [PATCH 4/4] fix(ui): preserve higher revision on delete --- n00n-ui/src/storage_writer.rs | 30 ++++++++++++++++++++++++++++-- 1 file changed, 28 insertions(+), 2 deletions(-) diff --git a/n00n-ui/src/storage_writer.rs b/n00n-ui/src/storage_writer.rs index b5ea3e154..ea74dedd7 100644 --- a/n00n-ui/src/storage_writer.rs +++ b/n00n-ui/src/storage_writer.rs @@ -453,6 +453,12 @@ impl WriterState { self.latest_generations.insert(id, generation); self.latest_snapshots.insert(id, snapshot.clone()); + self.queue_pending_snapshot(snapshot); + self.pending.get(&id).map(|pending| pending.version.clone()) + } + + fn queue_pending_snapshot(&mut self, snapshot: PendingSnapshot) { + let id = snapshot.session.id; let replace = self.pending.get(&id).is_none_or(|current| { current.version.revision < snapshot.version.revision || (current.version.revision == snapshot.version.revision @@ -461,7 +467,6 @@ impl WriterState { if replace { self.pending.insert(id, snapshot); } - self.pending.get(&id).map(|pending| pending.version.clone()) } fn flush(&mut self, dir: &StateDir) -> FailedSnapshots { @@ -552,7 +557,7 @@ impl WriterState { .as_ref() .map(|snapshot| snapshot.version.clone()); if let Some(snapshot) = crossing_snapshot { - self.pending.insert(id, snapshot); + self.queue_pending_snapshot(snapshot); } let primary_path = dir.path().join(SESSIONS_DIR).join(format!("{id}.jsonl")); @@ -973,6 +978,27 @@ mod tests { assert_eq!(loaded.title, "post-delete revision four"); } + #[test] + fn delayed_delete_preserves_highest_revision_post_delete_snapshot() { + let (_tmp, dir) = state_dir(); + let mut state = WriterState::default(); + let mut revision_nine = AppSession::new("test-model", "/tmp/delete-post-delete-order"); + revision_nine.title = "generation five revision nine".into(); + revision_nine.meta.revision = 9; + let id = revision_nine.id; + let mut revision_four = revision_nine.clone(); + revision_four.title = "generation seven revision four".into(); + revision_four.meta.revision = 4; + + state.stage_snapshot(PendingSnapshot::new(5, Box::new(revision_nine))); + state.stage_snapshot(PendingSnapshot::new(7, Box::new(revision_four))); + state.delete(id, 3, &dir).unwrap(); + + let loaded = AppSession::load(id, &dir).unwrap(); + assert_eq!(loaded.meta.revision, 9); + assert_eq!(loaded.title, "generation five revision nine"); + } + #[test] fn partial_delete_requeues_crossing_generation_lower_revision() { let (_tmp, dir) = state_dir();