返回 CodeWhale
persistence_actor.rs
根目录 / crates / tui / src / tui / persistence_actor.rs
1 //! Dedicated persistence actor for session save / checkpoint I/O.
2 //!
3 //! ## Motivation
4 //!
5 //! Before this module, `persist_checkpoint` and `persist_session_snapshot` ran
6 //! synchronously on the tokio worker thread that drives the TUI event loop.
7 //! Each call serialised all API messages to JSON, wrote a temp file, and
8 //! renamed it atomically — blocking keyboard input for the duration.
9 //! `save_session` additionally called `cleanup_old_sessions`, which listed all
10 //! session files, parsed metadata from every one, sorted, and deleted the
11 //! oldest — scaling O(session-bytes + file-count) with every turn.
12 //!
13 //! ## Design
14 //!
15 //! - **One dedicated tokio task** owns disk I/O. The UI only sends requests;
16 //! keystrokes never wait for writes.
17 //! - **Latest-wins coalescing per session**: when multiple `SaveCheckpoint`,
18 //! `SessionSnapshot`, or offline-queue requests pile up before the actor's
19 //! next write cycle, only the most recent one per session is written.
20 //! Checkpoints and clears are keyed by session id, so concurrent sessions
21 //! never coalesce into (or clear) each other's slot.
22 //! - **Durability reporting**: `FlushAndReport` returns accumulated results;
23 //! cycles without a listener log failures instead of discarding them.
24 //! - **Bounded command channel with sender-side coalescing** (#6212):
25 //! `try_send` absorbs each request into a shared latest-wins state at send
26 //! time and wakes the actor through a small bounded channel, so a paused
27 //! consumer retains one snapshot per session instead of one per send.
28 //! Queued snapshots retain only the canonical journal; legacy `messages`
29 //! are derived at the disk boundary instead of doubling every paused
30 //! request.
31
32 use std::collections::{BTreeMap, BTreeSet};
33 use std::sync::{Arc, OnceLock};
34
35 use tokio::sync::{mpsc, oneshot};
36
37 use crate::session_manager::{OfflineQueueLease, OfflineQueueState, SavedSession, SessionManager};
38 use crate::utils::spawn_supervised;
39
40 // ---------------------------------------------------------------------------
41 // Request type
42 // ---------------------------------------------------------------------------
43
44 /// Persistence work item sent to the actor.
45 #[derive(Debug)]
46 pub enum PersistRequest {
47 /// Write a crash-recovery checkpoint (in-flight turn state) to the
48 /// session's own file (`checkpoints/<session_id>.json`).
49 SaveCheckpoint { session: SavedSession },
50 /// Write a full session snapshot (completed turn, durable save).
51 SessionSnapshot(SavedSession),
52 /// Compound completion commit: write the completed session snapshot,
53 /// and only if that write succeeds, clear that same session's
54 /// crash-recovery checkpoint. A failed snapshot write RETAINS the
55 /// checkpoint as the only surviving recovery record, and the clear is
56 /// scoped to the committed session's id — it can never remove another
57 /// session's checkpoint. Turn completion must send this instead of a
58 /// `SessionSnapshot` + `ClearCheckpoint` pair, which the actor could
59 /// otherwise apply with the clear first, erasing the recovery record
60 /// before the snapshot safely landed.
61 CompletedCommit { session: SavedSession },
62 /// Write queued/draft offline input for crash recovery.
63 OfflineQueue {
64 state: OfflineQueueState,
65 lease: Arc<OfflineQueueLease>,
66 },
67 /// Remove the queued/draft offline input file.
68 ClearOfflineQueue {
69 /// Captures the exact owner and retains its exclusive editor lease
70 /// until the removal finishes. An unowned clear is unrepresentable.
71 lease: Arc<OfflineQueueLease>,
72 },
73 /// Remove one session's crash-recovery checkpoint file. Scoped: cannot
74 /// remove another session's checkpoint.
75 ClearCheckpoint { session_id: String },
76 /// Flush all pending work now and report durability results through
77 /// `reply`. The report aggregates every write/removal result since the
78 /// previous report (including background write cycles) — errors are
79 /// collected and surfaced, never discarded.
80 FlushAndReport { reply: oneshot::Sender<FlushReport> },
81 /// Graceful shutdown — flush pending writes, then exit the actor loop.
82 Shutdown,
83 }
84
85 /// Aggregated durability results: how many writes/removals completed and
86 /// which failed (labelled by what was being persisted, with the I/O error
87 /// kind).
88 #[derive(Debug, Default)]
89 pub struct FlushReport {
90 pub completed: usize,
91 pub failures: Vec<(String, std::io::ErrorKind)>,
92 }
93
94 impl FlushReport {
95 /// Upper bound on retained failure entries when accumulating across
96 /// write cycles. Every failure is logged at the cycle it happened, so
97 /// dropping older-than-bound entries from the reply loses no evidence.
98 const MAX_ACCUMULATED_FAILURES: usize = 256;
99
100 fn merge(&mut self, other: FlushReport) {
101 self.completed += other.completed;
102 self.failures.extend(other.failures);
103 if self.failures.len() > Self::MAX_ACCUMULATED_FAILURES {
104 let excess = self.failures.len() - Self::MAX_ACCUMULATED_FAILURES;
105 self.failures.drain(..excess);
106 }
107 }
108 }
109
110 #[derive(Debug)]
111 enum PendingOfflineQueue {
112 Save {
113 state: Box<OfflineQueueState>,
114 lease: Arc<OfflineQueueLease>,
115 },
116 Clear {
117 lease: Arc<OfflineQueueLease>,
118 },
119 }
120
121 // ---------------------------------------------------------------------------
122 // Handle (held by the TUI)
123 // ---------------------------------------------------------------------------
124
125 /// Control commands the actor reacts to. Work itself never crosses the
126 /// channel: requests are coalesced into the shared [`PendingState`] at send
127 /// time, so a paused consumer retains at most the latest request per session
128 /// instead of every snapshot ever sent (#6212).
129 enum ActorCommand {
130 /// The shared pending state has work the actor has not taken yet.
131 WorkReady,
132 FlushAndReport {
133 reply: oneshot::Sender<FlushReport>,
134 },
135 Shutdown,
136 }
137
138 /// Command-channel capacity. `WorkReady` is deduplicated by the `notified`
139 /// flag, so only `FlushAndReport`/`Shutdown` can occupy slots unplanned; the
140 /// capacity exists so those never observe a full channel in practice.
141 const ACTOR_COMMAND_CAPACITY: usize = 8;
142
143 /// The coalescing state shared between senders and the actor. Senders absorb
144 /// under the lock; the actor `take`s (swap to empty) under the same lock and
145 /// flushes outside it, so disk I/O never blocks a sender.
146 #[derive(Debug, Default)]
147 struct SharedPending {
148 pending: PendingState,
149 /// Whether a `WorkReady` is already queued (or being queued) and no
150 /// `take_pending` has observed it since. Reset only by `take_pending` —
151 /// the receiver side — so the flag always reflects the channel the actor
152 /// drains.
153 notified: bool,
154 }
155
156 #[derive(Clone)]
157 struct PersistRequestSender {
158 shared: Arc<std::sync::Mutex<SharedPending>>,
159 cmd_tx: mpsc::Sender<ActorCommand>,
160 health: SessionSaveHealth,
161 }
162
163 impl std::fmt::Debug for PersistRequestSender {
164 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
165 f.debug_struct("PersistRequestSender")
166 .finish_non_exhaustive()
167 }
168 }
169
170 struct PersistRequestReceiver {
171 shared: Arc<std::sync::Mutex<SharedPending>>,
172 cmd_rx: mpsc::Receiver<ActorCommand>,
173 health: SessionSaveHealth,
174 }
175
176 /// Single construction seam for the production persistence request channel.
177 ///
178 /// The ignored backlog measurement uses this same factory with the receiver
179 /// deliberately paused, so the measurement always characterizes whatever
180 /// representation the seam actually retains.
181 fn persistence_request_channel() -> (PersistRequestSender, PersistRequestReceiver) {
182 let (cmd_tx, cmd_rx) = mpsc::channel(ACTOR_COMMAND_CAPACITY);
183 let shared = Arc::new(std::sync::Mutex::new(SharedPending::default()));
184 let health = SessionSaveHealth::default();
185 (
186 PersistRequestSender {
187 shared: Arc::clone(&shared),
188 cmd_tx,
189 health: health.clone(),
190 },
191 PersistRequestReceiver {
192 shared,
193 cmd_rx,
194 health,
195 },
196 )
197 }
198
199 impl PersistRequestReceiver {
200 /// Await the next actor command. `None` means every sender is gone.
201 async fn recv(&mut self) -> Option<ActorCommand> {
202 self.cmd_rx.recv().await
203 }
204
205 /// Atomically take everything coalesced so far. Resets the `notified`
206 /// flag so the next absorbed request queues a fresh `WorkReady`.
207 fn take_pending(&mut self) -> PendingState {
208 let mut guard = self
209 .shared
210 .lock()
211 .unwrap_or_else(|poisoned| poisoned.into_inner());
212 guard.notified = false;
213 std::mem::take(&mut guard.pending)
214 }
215 }
216
217 /// Lightweight handle that the UI holds to queue persistence work.
218 #[derive(Debug, Clone)]
219 pub struct PersistActorHandle {
220 tx: PersistRequestSender,
221 }
222
223 impl PersistActorHandle {
224 /// Whether the latest save of each session document landed.
225 pub(crate) fn session_save_health(&self) -> SaveHealthReading {
226 self.tx.health.reading()
227 }
228
229 /// Queue a persistence request without blocking. The request is
230 /// coalesced into the shared pending state immediately (latest-wins per
231 /// session), so repeated snapshots of one session retain only the
232 /// newest. Returns `false` when the actor is already shut down.
233 pub fn try_send(&self, mut request: PersistRequest) -> bool {
234 match &mut request {
235 PersistRequest::SaveCheckpoint { session }
236 | PersistRequest::SessionSnapshot(session)
237 | PersistRequest::CompletedCommit { session } => {
238 session.compact_for_persistence_queue();
239 }
240 _ => {}
241 }
242 let control = {
243 let mut guard = self
244 .tx
245 .shared
246 .lock()
247 .unwrap_or_else(|poisoned| poisoned.into_inner());
248 guard.pending.absorb(request)
249 };
250 match control {
251 Control::Continue => {
252 // A WorkReady the actor has not taken yet is already queued;
253 // that take will observe this request too. `notified` resets
254 // only when the actor takes, so it always mirrors the
255 // channel the actor drains.
256 {
257 let mut guard = self
258 .tx
259 .shared
260 .lock()
261 .unwrap_or_else(|poisoned| poisoned.into_inner());
262 if guard.notified {
263 return true;
264 }
265 guard.notified = true;
266 }
267 if self.tx.cmd_tx.try_send(ActorCommand::WorkReady).is_ok() {
268 true
269 } else {
270 // Roll the flag back so later sends re-attempt (and
271 // re-fail) honestly instead of riding a dead wake.
272 let mut guard = self
273 .tx
274 .shared
275 .lock()
276 .unwrap_or_else(|poisoned| poisoned.into_inner());
277 guard.notified = false;
278 false
279 }
280 }
281 Control::Flush(reply) => self
282 .tx
283 .cmd_tx
284 .try_send(ActorCommand::FlushAndReport { reply })
285 .is_ok(),
286 Control::Shutdown => self.tx.cmd_tx.try_send(ActorCommand::Shutdown).is_ok(),
287 }
288 }
289 }
290
291 // ---------------------------------------------------------------------------
292 // Global singleton (avoid threading through App)
293 // ---------------------------------------------------------------------------
294
295 static ACTOR_TX: OnceLock<PersistActorHandle> = OnceLock::new();
296
297 /// Initialise the global persistence actor handle. Must be called once at
298 /// startup, before the event loop starts.
299 pub fn init_actor(handle: PersistActorHandle) {
300 let _ = ACTOR_TX.set(handle);
301 }
302
303 /// Queue a persistence request through the global handle. When the request
304 /// cannot be queued — actor not initialised yet (tests, early startup) or
305 /// already shut down — the drop is logged instead of discarded silently, so
306 /// lost session/work-graph state is diagnosable after the fact.
307 pub fn persist(request: PersistRequest) {
308 let label = request_label(&request);
309 if try_persist(request) {
310 return;
311 }
312 if ACTOR_TX.get().is_some() {
313 tracing::warn!(
314 request = label,
315 "persistence request dropped: actor channel is closed (shutdown already happened)"
316 );
317 } else {
318 tracing::debug!(
319 request = label,
320 "persistence request dropped: actor not initialised yet"
321 );
322 }
323 }
324
325 /// Order synchronous lifecycle saves after all previously queued snapshots.
326 /// Refuse a single-thread runtime rather than deadlocking its persistence task.
327 pub(crate) fn flush_before_transition() -> Result<(), String> {
328 let Some(handle) = ACTOR_TX.get() else {
329 return Ok(());
330 };
331 let runtime = tokio::runtime::Handle::try_current().ok();
332 if runtime
333 .as_ref()
334 .is_some_and(|r| r.runtime_flavor() != tokio::runtime::RuntimeFlavor::MultiThread)
335 {
336 return Err("session transition requires an asynchronous persistence barrier".into());
337 }
338 let (reply, receiver) = oneshot::channel();
339 if !handle.try_send(PersistRequest::FlushAndReport { reply }) {
340 return Err("session transition could not queue persistence barrier".into());
341 }
342 let receive = || receiver.blocking_recv();
343 let report = if runtime.is_some() {
344 tokio::task::block_in_place(receive)
345 } else {
346 receive()
347 }
348 .map_err(|_| "session persistence stopped before the transition".to_string())?;
349 if !report.failures.is_empty() {
350 return Err(format!(
351 "session transition refused after persistence failures: {:?}",
352 report.failures
353 ));
354 }
355 Ok(())
356 }
357
358 fn request_label(request: &PersistRequest) -> &'static str {
359 match request {
360 PersistRequest::SaveCheckpoint { .. } => "SaveCheckpoint",
361 PersistRequest::SessionSnapshot(_) => "SessionSnapshot",
362 PersistRequest::CompletedCommit { .. } => "CompletedCommit",
363 PersistRequest::OfflineQueue { .. } => "OfflineQueue",
364 PersistRequest::ClearOfflineQueue { .. } => "ClearOfflineQueue",
365 PersistRequest::ClearCheckpoint { .. } => "ClearCheckpoint",
366 PersistRequest::FlushAndReport { .. } => "FlushAndReport",
367 PersistRequest::Shutdown => "Shutdown",
368 }
369 }
370
371 /// Queue persistence and report whether the actor accepted ownership. Work
372 /// Graph projections use this acknowledgement as their publish boundary.
373 /// [`PersistActorHandle::session_save_health`] of the global actor; `None`
374 /// before it starts.
375 pub(crate) fn session_save_health() -> Option<SaveHealthReading> {
376 ACTOR_TX.get().map(PersistActorHandle::session_save_health)
377 }
378
379 pub fn try_persist(request: PersistRequest) -> bool {
380 ACTOR_TX
381 .get()
382 .is_some_and(|handle| handle.try_send(request))
383 }
384
385 // ---------------------------------------------------------------------------
386 // Actor spawn
387 // ---------------------------------------------------------------------------
388
389 /// Spawn the persistence actor task and return a handle for the caller to
390 /// store and initialise.
391 ///
392 /// The returned handle should be passed to [`init_actor`] so that the
393 /// `persist()` free function can reach it from anywhere in the TUI.
394 pub fn spawn_persistence_actor(
395 manager: SessionManager,
396 ) -> (PersistActorHandle, tokio::task::JoinHandle<()>) {
397 let (tx, mut rx) = persistence_request_channel();
398 let handle = PersistActorHandle { tx };
399
400 let task = spawn_supervised(
401 "persistence-actor",
402 std::panic::Location::caller(),
403 async move {
404 let mut unreported = FlushReport::default();
405
406 // Flush pending work, log new failures, and fold the cycle's
407 // results into the unreported accumulator.
408 fn flush_cycle(
409 manager: &SessionManager,
410 pending: &mut PendingState,
411 unreported: &mut FlushReport,
412 health: &SessionSaveHealth,
413 ) {
414 let cycle = flush_inner(manager, pending, health);
415 log_flush_failures(&cycle);
416 unreported.merge(cycle);
417 }
418
419 // Work is coalesced at send time into the shared pending state;
420 // every command handler takes whatever has accumulated and
421 // flushes it outside the sender lock.
422 while let Some(command) = rx.recv().await {
423 let mut pending = rx.take_pending();
424 match command {
425 ActorCommand::WorkReady => {
426 if !pending.is_empty() {
427 flush_cycle(&manager, &mut pending, &mut unreported, &rx.health);
428 }
429 }
430 ActorCommand::FlushAndReport { reply } => {
431 flush_cycle(&manager, &mut pending, &mut unreported, &rx.health);
432 let _ = reply.send(std::mem::take(&mut unreported));
433 }
434 ActorCommand::Shutdown => {
435 flush_cycle(&manager, &mut pending, &mut unreported, &rx.health);
436 return;
437 }
438 }
439 }
440 // Every sender is gone — final flush and exit. Each absorb
441 // guarantees a queued WorkReady, so this is normally empty.
442 let mut pending = rx.take_pending();
443 flush_cycle(&manager, &mut pending, &mut unreported, &rx.health);
444 },
445 );
446
447 (handle, task)
448 }
449
450 /// Coalesced work waiting for the next write cycle.
451 #[derive(Debug, Default)]
452 struct PendingState {
453 /// Latest-wins per session id. Crash checkpoints are keyed per session
454 /// (mirroring `sessions` below) so concurrent sessions can interleave
455 /// saves and clears without clobbering each other.
456 checkpoints: BTreeMap<String, SavedSession>,
457 /// Session ids whose checkpoint file should be removed.
458 checkpoint_clears: BTreeSet<String>,
459 /// Latest-wins per session id. Coalescing into one global slot can
460 /// drop session A when an immediate `/new` queues session B before
461 /// the actor drains.
462 sessions: BTreeMap<String, SavedSession>,
463 /// Compound completion commits, latest-wins per session id: the
464 /// completed snapshot body to write, followed by that session's own
465 /// checkpoint clear — the clear only if the write succeeded. Kept
466 /// separate from `sessions` so a plain snapshot can never be drained
467 /// as a completion (or vice versa) and the clear intent stays bound to
468 /// exactly the session that completed.
469 completed_commits: BTreeMap<String, SavedSession>,
470 /// Latest-wins per session id, for the same reason `sessions` above is:
471 /// a single global slot dropped session A's queued text when session B
472 /// queued before the actor drained, which defeats the per-session file
473 /// naming entirely. Each pending request retains its editor lease, so a
474 /// window changing session cannot release ownership ahead of its writes.
475 offline_queue: BTreeMap<String, PendingOfflineQueue>,
476 }
477
478 /// What the actor loop should do after absorbing a request.
479 enum Control {
480 Continue,
481 Flush(oneshot::Sender<FlushReport>),
482 Shutdown,
483 }
484
485 impl PendingState {
486 /// True when nothing coalesced is waiting. An empty `WorkReady` take is
487 /// skipped instead of running a no-op flush cycle.
488 fn is_empty(&self) -> bool {
489 self.checkpoints.is_empty()
490 && self.checkpoint_clears.is_empty()
491 && self.sessions.is_empty()
492 && self.completed_commits.is_empty()
493 && self.offline_queue.is_empty()
494 }
495
496 fn absorb(&mut self, req: PersistRequest) -> Control {
497 match req {
498 PersistRequest::SaveCheckpoint { session } => {
499 // Last-writer-wins per session: a fresh checkpoint supersedes
500 // a pending clear for the same session so the two never both
501 // apply in one drain (which previously cleared then re-wrote
502 // the stale checkpoint, undoing the clear).
503 let id = session.metadata.id.clone();
504 self.checkpoint_clears.remove(&id);
505 // A new in-flight checkpoint means newer turn work started
506 // after the completion it would have cleared: the compound's
507 // clear-on-success intent is stale now (it would erase the
508 // fresher recovery record), so drop the pending compound and
509 // keep only the newer checkpoint body.
510 if self.completed_commits.remove(&id).is_some() {
511 tracing::debug!(
512 session_id = %id,
513 "pending completed commit superseded by a newer in-flight checkpoint"
514 );
515 }
516 self.checkpoints.insert(id, session);
517 }
518 PersistRequest::SessionSnapshot(session) => {
519 // A newer full snapshot of a session with a pending
520 // completed commit refreshes the commit's body (and keeps
521 // its clear-on-success intent): the in-flight checkpoint it
522 // guards is already captured by the newer completed state.
523 let id = session.metadata.id.clone();
524 if let Some(pending) = self.completed_commits.get_mut(&id) {
525 *pending = session;
526 } else {
527 self.sessions.insert(id, session);
528 }
529 }
530 PersistRequest::CompletedCommit { session } => {
531 let id = session.metadata.id.clone();
532 // The compound owns this session's snapshot and clear: a
533 // pending plain snapshot is superseded, a pending standalone
534 // clear would erase the recovery record even when the save
535 // fails, and a pending checkpoint write would re-create the
536 // record after the compound cleared it.
537 self.sessions.remove(&id);
538 self.checkpoint_clears.remove(&id);
539 self.checkpoints.remove(&id);
540 self.completed_commits.insert(id, session);
541 }
542 PersistRequest::OfflineQueue { state, lease } => {
543 self.offline_queue.insert(
544 lease.session_id().to_string(),
545 PendingOfflineQueue::Save {
546 state: Box::new(state),
547 lease,
548 },
549 );
550 }
551 PersistRequest::ClearOfflineQueue { lease } => {
552 // A clear supersedes a pending save for its OWN session only.
553 self.offline_queue.insert(
554 lease.session_id().to_string(),
555 PendingOfflineQueue::Clear { lease },
556 );
557 }
558 PersistRequest::ClearCheckpoint { session_id } => {
559 // A clear supersedes a pending checkpoint write for the same
560 // session only — other sessions' pending work is untouched.
561 // An explicit clear is also a user-owned boundary (e.g.
562 // `/new`): it supersedes a pending compound for the same
563 // session so the discarded session is not re-written.
564 self.checkpoints.remove(&session_id);
565 self.completed_commits.remove(&session_id);
566 self.checkpoint_clears.insert(session_id);
567 }
568 PersistRequest::FlushAndReport { reply } => return Control::Flush(reply),
569 PersistRequest::Shutdown => return Control::Shutdown,
570 }
571 Control::Continue
572 }
573 }
574
575 /// Write all pending work to disk, draining `pending`. Every write and
576 /// removal result is collected into the returned [`FlushReport`] — failures
577 /// are reported, never silently discarded.
578 ///
579 /// Ordering is durability-critical: every session snapshot write (plain or
580 /// completion commit) happens BEFORE any checkpoint clear. A completion
581 /// commit clears its session's checkpoint only after that session's own
582 /// write succeeded, so a failed save always leaves the crash-recovery
583 /// checkpoint in place.
584 fn flush_inner(
585 manager: &SessionManager,
586 pending: &mut PendingState,
587 health: &SessionSaveHealth,
588 ) -> FlushReport {
589 let mut report = FlushReport::default();
590 let mut record = |what: String, result: std::io::Result<()>| match result {
591 Ok(()) => report.completed += 1,
592 Err(err) => report.failures.push((what, err.kind())),
593 };
594
595 // Document saves are bounded (#6842): an oversized journal is archived
596 // and pruned here, on the actor, never on the UI loop.
597 for (session_id, session) in std::mem::take(&mut pending.sessions) {
598 let result = manager.save_session_bounded(session).map(|_| ());
599 health.record(&session_id, &result);
600 record(format!("session:{session_id}"), result);
601 }
602 for (session_id, session) in std::mem::take(&mut pending.completed_commits) {
603 let commit_result = manager.save_session_bounded(session);
604 health.record(&session_id, &commit_result);
605 let save_succeeded = commit_result.is_ok();
606 record(
607 format!("completed-commit:{session_id}"),
608 commit_result.map(|_| ()),
609 );
610 if save_succeeded {
611 // Only the committed session's own checkpoint is cleared, and
612 // only because its snapshot safely landed. A failure above
613 // retains the checkpoint as the sole recovery record.
614 record(
615 format!("clear-checkpoint:{session_id}"),
616 manager.clear_session_checkpoint(&session_id),
617 );
618 }
619 }
620 for session_id in std::mem::take(&mut pending.checkpoint_clears) {
621 record(
622 format!("clear-checkpoint:{session_id}"),
623 manager.clear_session_checkpoint(&session_id),
624 );
625 }
626 for (session_id, session) in std::mem::take(&mut pending.checkpoints) {
627 record(
628 format!("checkpoint:{session_id}"),
629 manager.save_checkpoint_owned(session).map(|_| ()),
630 );
631 }
632 for (_, request) in std::mem::take(&mut pending.offline_queue) {
633 match request {
634 PendingOfflineQueue::Save { state, lease } => record(
635 "offline-queue".to_string(),
636 manager
637 .save_offline_queue_state(&state, Some(lease.session_id()))
638 .map(|_| ()),
639 ),
640 PendingOfflineQueue::Clear { lease } => record(
641 "clear-offline-queue".to_string(),
642 manager.clear_offline_queue_state_for(lease.session_id()),
643 ),
644 }
645 }
646 report
647 }
648
649 /// Which session documents' latest save failed. Only document saves count: a
650 /// failed checkpoint, queue or cleanup write does not lose the conversation.
651 /// A later successful save of the same session clears its entry, so a
652 /// transient failure does not leave a standing alarm.
653 #[derive(Debug, Clone, Default)]
654 pub(crate) struct SessionSaveHealth(Arc<std::sync::Mutex<SaveHealthState>>);
655
656 #[derive(Debug, Default)]
657 struct SaveHealthState {
658 /// Bumped whenever `failing` changes, so a poller can tell what is new.
659 generation: u64,
660 failing: BTreeMap<String, std::io::ErrorKind>,
661 }
662
663 /// One reading of [`SessionSaveHealth`].
664 #[derive(Debug, Clone, PartialEq, Eq)]
665 pub(crate) struct SaveHealthReading {
666 pub(crate) generation: u64,
667 /// A session whose latest save failed, and how, while any has.
668 pub(crate) failing: Option<(String, std::io::ErrorKind)>,
669 }
670
671 impl SessionSaveHealth {
672 fn record<T>(&self, session_id: &str, result: &std::io::Result<T>) {
673 let mut state = self
674 .0
675 .lock()
676 .unwrap_or_else(|poisoned| poisoned.into_inner());
677 let changed = match result {
678 Ok(_) => state.failing.remove(session_id).is_some(),
679 Err(error) => {
680 state.failing.insert(session_id.to_string(), error.kind()) != Some(error.kind())
681 }
682 };
683 if changed {
684 state.generation += 1;
685 }
686 }
687
688 fn reading(&self) -> SaveHealthReading {
689 let state = self
690 .0
691 .lock()
692 .unwrap_or_else(|poisoned| poisoned.into_inner());
693 SaveHealthReading {
694 generation: state.generation,
695 failing: state
696 .failing
697 .iter()
698 .next()
699 .map(|(id, kind)| (id.clone(), *kind)),
700 }
701 }
702 }
703
704 /// Surface flush failures in the log for write cycles that have no caller
705 /// waiting on a [`FlushReport`].
706 fn log_flush_failures(report: &FlushReport) {
707 for (what, kind) in &report.failures {
708 tracing::warn!(
709 target: "persistence",
710 what = %what,
711 error_kind = ?kind,
712 "persistence write failed",
713 );
714 }
715 }
716
717 #[cfg(test)]
718 #[path = "persistence_actor/tests.rs"]
719 mod backlog_measurement_tests;
720
721 #[cfg(test)]
722 mod tests {
723 use super::*;
724 use std::time::Duration;
725
726 use crate::session_manager::{OfflineQueueState, QueuedSessionMessage};
727
728 async fn wait_until(mut predicate: impl FnMut() -> bool) {
729 let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
730 loop {
731 if predicate() {
732 return;
733 }
734 assert!(
735 tokio::time::Instant::now() < deadline,
736 "timed out waiting for persistence actor"
737 );
738 tokio::time::sleep(Duration::from_millis(10)).await;
739 }
740 }
741
742 #[tokio::test]
743 async fn two_sessions_queueing_before_a_drain_both_survive() {
744 // Per-session FILENAMES are not enough on their own: the actor
745 // coalesces pending work before those names are ever used, and the
746 // queue used one global slot while its checkpoint/session neighbours
747 // were already keyed per session. Session A's unsent text was
748 // therefore dropped whenever session B queued first.
749 let tmp = tempfile::tempdir().expect("tempdir");
750 let sessions_dir = tmp.path().join("sessions");
751 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
752 let (handle, task) = spawn_persistence_actor(manager);
753
754 let queue_manager = SessionManager::new(sessions_dir.clone()).expect("queue manager");
755 let lease_a = queue_manager
756 .acquire_offline_queue_lease("session-A")
757 .expect("lease A");
758 let lease_b = queue_manager
759 .acquire_offline_queue_lease("session-B")
760 .expect("lease B");
761 for (session, body) in [("session-A", "text from A"), ("session-B", "text from B")] {
762 let state = OfflineQueueState {
763 messages: vec![QueuedSessionMessage {
764 display: body.to_string(),
765 skill_instruction: None,
766 skill_provenance: None,
767 }],
768 ..OfflineQueueState::default()
769 };
770 handle.try_send(PersistRequest::OfflineQueue {
771 state,
772 lease: Arc::clone(if session == "session-A" {
773 &lease_a
774 } else {
775 &lease_b
776 }),
777 });
778 }
779
780 let checkpoints = sessions_dir.join("checkpoints");
781 for (session, body) in [("session-A", "text from A"), ("session-B", "text from B")] {
782 let path = checkpoints.join(format!("{session}.offline_queue.json"));
783 // wait_until panics on timeout, which is the failure signal: a
784 // coalesced-away queue never appears.
785 wait_until(|| std::fs::read_to_string(&path).is_ok_and(|f| f.contains(body))).await;
786 }
787
788 // A clear names its own session and must not touch the other's.
789 handle.try_send(PersistRequest::ClearOfflineQueue {
790 lease: Arc::clone(&lease_a),
791 });
792 let a = checkpoints.join("session-A.offline_queue.json");
793 wait_until(|| !a.exists()).await;
794 assert!(
795 checkpoints.join("session-B.offline_queue.json").exists(),
796 "clearing one session must not delete another session's queued text"
797 );
798
799 drop(handle);
800 let _ = task.await;
801 }
802
803 #[tokio::test]
804 async fn actor_persists_and_clears_offline_queue_requests() {
805 let tmp = tempfile::tempdir().expect("tempdir");
806 let sessions_dir = tmp.path().join("sessions");
807 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
808 // The queue is keyed per session now (#5715-adjacent data-loss fix):
809 // two concurrent instances used to share one global file and the loser
810 // lost its unsent text. The request below carries session-A.
811 let queue_path = sessions_dir
812 .join("checkpoints")
813 .join("session-A.offline_queue.json");
814 let lease = manager
815 .acquire_offline_queue_lease("session-A")
816 .expect("queue lease");
817 let (handle, task) = spawn_persistence_actor(manager);
818
819 let state = OfflineQueueState {
820 messages: vec![QueuedSessionMessage {
821 display: "queued from enter".to_string(),
822 skill_instruction: None,
823 skill_provenance: None,
824 }],
825 ..OfflineQueueState::default()
826 };
827
828 handle.try_send(PersistRequest::OfflineQueue {
829 state,
830 lease: Arc::clone(&lease),
831 });
832 wait_until(|| {
833 std::fs::read_to_string(&queue_path)
834 .is_ok_and(|body| body.contains("queued from enter"))
835 })
836 .await;
837
838 handle.try_send(PersistRequest::ClearOfflineQueue {
839 lease: Arc::clone(&lease),
840 });
841 wait_until(|| !queue_path.exists()).await;
842 handle.try_send(PersistRequest::Shutdown);
843 task.await.expect("persistence actor join");
844 }
845
846 #[tokio::test]
847 async fn shutdown_wait_flushes_queued_session_before_returning() {
848 let tmp = tempfile::tempdir().expect("tempdir");
849 let sessions_dir = tmp.path().join("sessions");
850 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
851 let verification_manager = SessionManager::new(sessions_dir).expect("verification manager");
852 let session = crate::session_manager::create_saved_session_with_mode(
853 &[],
854 "deepseek-v4-pro",
855 tmp.path(),
856 0,
857 None,
858 Some("agent"),
859 );
860 let session_id = session.metadata.id.clone();
861 let (handle, task) = spawn_persistence_actor(manager);
862
863 handle.try_send(PersistRequest::SessionSnapshot(session));
864 handle.try_send(PersistRequest::Shutdown);
865 task.await.expect("persistence actor join");
866
867 let loaded = verification_manager
868 .load_session(&session_id)
869 .expect("shutdown must flush queued session");
870 assert_eq!(loaded.metadata.id, session_id);
871 }
872
873 #[tokio::test]
874 async fn shutdown_flushes_latest_snapshot_for_each_session_id() {
875 let tmp = tempfile::tempdir().expect("tempdir");
876 let sessions_dir = tmp.path().join("sessions");
877 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
878 let verification_manager = SessionManager::new(sessions_dir).expect("verification manager");
879 let mut first = crate::session_manager::create_saved_session_with_mode(
880 &[],
881 "deepseek-v4-pro",
882 tmp.path(),
883 0,
884 None,
885 Some("agent"),
886 );
887 first.metadata.title = "Session A".to_string();
888 let mut second = crate::session_manager::create_saved_session_with_mode(
889 &[],
890 "deepseek-v4-pro",
891 tmp.path(),
892 0,
893 None,
894 Some("agent"),
895 );
896 second.metadata.title = "Session B".to_string();
897 let first_id = first.metadata.id.clone();
898 let second_id = second.metadata.id.clone();
899 let (handle, task) = spawn_persistence_actor(manager);
900
901 handle.try_send(PersistRequest::SessionSnapshot(first));
902 handle.try_send(PersistRequest::SessionSnapshot(second));
903 handle.try_send(PersistRequest::Shutdown);
904 task.await.expect("persistence actor join");
905
906 assert_eq!(
907 verification_manager
908 .load_session(&first_id)
909 .expect("session A flushed")
910 .metadata
911 .title,
912 "Session A"
913 );
914 assert_eq!(
915 verification_manager
916 .load_session(&second_id)
917 .expect("session B flushed")
918 .metadata
919 .title,
920 "Session B"
921 );
922 }
923
924 #[tokio::test]
925 async fn interleaved_checkpoint_saves_and_clears_stay_per_session() {
926 let tmp = tempfile::tempdir().expect("tempdir");
927 let sessions_dir = tmp.path().join("sessions");
928 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
929 let verification_manager = SessionManager::new(sessions_dir).expect("verification manager");
930 let first = crate::session_manager::create_saved_session_with_mode(
931 &[],
932 "deepseek-v4-pro",
933 tmp.path(),
934 0,
935 None,
936 Some("agent"),
937 );
938 let second = crate::session_manager::create_saved_session_with_mode(
939 &[],
940 "deepseek-v4-pro",
941 tmp.path(),
942 0,
943 None,
944 Some("agent"),
945 );
946 let first_id = first.metadata.id.clone();
947 let second_id = second.metadata.id.clone();
948 let (handle, task) = spawn_persistence_actor(manager);
949
950 // Interleave: save A, save B, clear A — all coalesced into one drain.
951 handle.try_send(PersistRequest::SaveCheckpoint { session: first });
952 handle.try_send(PersistRequest::SaveCheckpoint { session: second });
953 handle.try_send(PersistRequest::ClearCheckpoint {
954 session_id: first_id.clone(),
955 });
956 handle.try_send(PersistRequest::Shutdown);
957 task.await.expect("persistence actor join");
958
959 assert!(
960 verification_manager
961 .load_session_checkpoint(&first_id)
962 .expect("load first checkpoint")
963 .is_none(),
964 "cleared session must have no checkpoint file"
965 );
966 let survivor = verification_manager
967 .load_session_checkpoint(&second_id)
968 .expect("load second checkpoint")
969 .expect("second session's checkpoint must survive an unrelated clear");
970 assert_eq!(survivor.metadata.id, second_id);
971 }
972
973 #[tokio::test]
974 async fn flush_and_report_returns_completed_counts() {
975 let tmp = tempfile::tempdir().expect("tempdir");
976 let sessions_dir = tmp.path().join("sessions");
977 let manager = SessionManager::new(sessions_dir).expect("manager");
978 let session = crate::session_manager::create_saved_session_with_mode(
979 &[],
980 "deepseek-v4-pro",
981 tmp.path(),
982 0,
983 None,
984 Some("agent"),
985 );
986 let (handle, task) = spawn_persistence_actor(manager);
987
988 handle.try_send(PersistRequest::SaveCheckpoint { session });
989 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
990 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
991 let report = reply_rx.await.expect("flush report reply");
992 // Whether the checkpoint was written by an earlier background cycle
993 // or by this flush, the accumulated report must count it and show no
994 // failures — and the actor keeps running afterwards.
995 assert!(report.completed >= 1, "checkpoint write must be counted");
996 assert!(report.failures.is_empty(), "no failures expected");
997 handle.try_send(PersistRequest::Shutdown);
998 task.await.expect("persistence actor join");
999 }
1000
1001 /// The health the UI polls follows real saves through the actor: a failed
1002 /// checkpoint does not count as a lost conversation, a failed document
1003 /// save does, and a later successful save of that session clears it.
1004 #[tokio::test]
1005 async fn session_save_health_follows_document_saves_through_the_actor() {
1006 let tmp = tempfile::tempdir().expect("tempdir");
1007 let sessions_dir = tmp.path().join("sessions");
1008 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1009 std::fs::write(sessions_dir.join("checkpoints"), b"not a directory")
1010 .expect("block checkpoints dir");
1011 let session = crate::session_manager::create_saved_session_with_mode(
1012 &[],
1013 "deepseek-v4-pro",
1014 tmp.path(),
1015 0,
1016 None,
1017 Some("agent"),
1018 );
1019 let session_id = session.metadata.id.clone();
1020 // A directory where the document goes makes its save fail.
1021 let blocker = sessions_dir.join(format!("{session_id}.json"));
1022 std::fs::create_dir(&blocker).expect("block session document");
1023 let (handle, task) = spawn_persistence_actor(manager);
1024 let flush = |handle: &PersistActorHandle| {
1025 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1026 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1027 reply_rx
1028 };
1029
1030 handle.try_send(PersistRequest::SaveCheckpoint {
1031 session: session.clone(),
1032 });
1033 flush(&handle).await.expect("flush");
1034 assert_eq!(handle.session_save_health().failing, None);
1035
1036 handle.try_send(PersistRequest::SessionSnapshot(session.clone()));
1037 flush(&handle).await.expect("flush");
1038 let failed = handle.session_save_health();
1039 assert_eq!(
1040 failed.failing.as_ref().map(|(id, _)| id.as_str()),
1041 Some(session_id.as_str())
1042 );
1043
1044 std::fs::remove_dir(&blocker).expect("unblock");
1045 handle.try_send(PersistRequest::SessionSnapshot(session));
1046 flush(&handle).await.expect("flush");
1047 let healed = handle.session_save_health();
1048 assert_eq!(healed.failing, None, "a later save clears the failure");
1049 assert!(healed.generation > failed.generation);
1050
1051 handle.try_send(PersistRequest::Shutdown);
1052 task.await.expect("persistence actor join");
1053 }
1054
1055 #[tokio::test]
1056 async fn flush_and_report_propagates_write_failures() {
1057 let tmp = tempfile::tempdir().expect("tempdir");
1058 let sessions_dir = tmp.path().join("sessions");
1059 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1060 // Occupy the checkpoints directory path with a regular file so every
1061 // checkpoint write deterministically fails on all platforms.
1062 std::fs::write(sessions_dir.join("checkpoints"), b"not a directory")
1063 .expect("block checkpoints dir");
1064 let session = crate::session_manager::create_saved_session_with_mode(
1065 &[],
1066 "deepseek-v4-pro",
1067 tmp.path(),
1068 0,
1069 None,
1070 Some("agent"),
1071 );
1072 let session_id = session.metadata.id.clone();
1073 let (handle, task) = spawn_persistence_actor(manager);
1074
1075 handle.try_send(PersistRequest::SaveCheckpoint { session });
1076 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1077 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1078 let report = reply_rx.await.expect("flush report reply");
1079
1080 assert!(
1081 report
1082 .failures
1083 .iter()
1084 .any(|(what, _)| what == &format!("checkpoint:{session_id}")),
1085 "failed checkpoint write must be reported, got: {:?}",
1086 report.failures
1087 );
1088 handle.try_send(PersistRequest::Shutdown);
1089 task.await.expect("persistence actor join");
1090 }
1091
1092 /// Pre-write a crash-recovery checkpoint file for `session_id` directly,
1093 /// so completion-commit tests can assert on its survival without needing
1094 /// a prior in-flight turn.
1095 fn seed_checkpoint_file(
1096 sessions_dir: &std::path::Path,
1097 session_id: &str,
1098 ) -> std::path::PathBuf {
1099 let path = sessions_dir
1100 .join("checkpoints")
1101 .join(format!("{session_id}.json"));
1102 std::fs::create_dir_all(path.parent().expect("checkpoints parent")).expect("mkdir");
1103 std::fs::write(&path, "{}").expect("seed checkpoint file");
1104 path
1105 }
1106
1107 /// Deterministically fail session-file saves only: a directory at the
1108 /// session's own `<id>.json` path makes `save_session` fail while the
1109 /// checkpoints directory stays fully usable.
1110 fn block_session_file(sessions_dir: &std::path::Path, session_id: &str) {
1111 std::fs::create_dir_all(sessions_dir.join(format!("{session_id}.json")))
1112 .expect("block session file path");
1113 }
1114
1115 #[tokio::test]
1116 async fn completed_commit_preserves_checkpoint_when_session_save_fails() {
1117 let tmp = tempfile::tempdir().expect("tempdir");
1118 let sessions_dir = tmp.path().join("sessions");
1119 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1120 let session = crate::session_manager::create_saved_session_with_mode(
1121 &[],
1122 "deepseek-v4-pro",
1123 tmp.path(),
1124 0,
1125 None,
1126 Some("agent"),
1127 );
1128 let session_id = session.metadata.id.clone();
1129 let checkpoint_path = seed_checkpoint_file(&sessions_dir, &session_id);
1130 block_session_file(&sessions_dir, &session_id);
1131 let (handle, task) = spawn_persistence_actor(manager);
1132
1133 handle.try_send(PersistRequest::CompletedCommit { session });
1134 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1135 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1136 let report = reply_rx.await.expect("flush report reply");
1137
1138 assert!(
1139 report
1140 .failures
1141 .iter()
1142 .any(|(what, _)| what == &format!("completed-commit:{session_id}")),
1143 "failed session save must be reported, got: {:?}",
1144 report.failures
1145 );
1146 assert!(
1147 !report
1148 .failures
1149 .iter()
1150 .any(|(what, _)| what == &format!("clear-checkpoint:{session_id}")),
1151 "the checkpoint clear must not be attempted after a failed save"
1152 );
1153 assert!(
1154 checkpoint_path.exists(),
1155 "a failed session save must retain the crash-recovery checkpoint"
1156 );
1157 handle.try_send(PersistRequest::Shutdown);
1158 task.await.expect("persistence actor join");
1159 }
1160
1161 #[tokio::test]
1162 async fn completed_commit_clears_only_after_session_save() {
1163 let tmp = tempfile::tempdir().expect("tempdir");
1164 let sessions_dir = tmp.path().join("sessions");
1165 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1166 let verification_manager = SessionManager::new(sessions_dir.clone()).expect("verify");
1167 let session = crate::session_manager::create_saved_session_with_mode(
1168 &[],
1169 "deepseek-v4-pro",
1170 tmp.path(),
1171 0,
1172 None,
1173 Some("agent"),
1174 );
1175 let session_id = session.metadata.id.clone();
1176 let checkpoint_path = seed_checkpoint_file(&sessions_dir, &session_id);
1177 let (handle, task) = spawn_persistence_actor(manager);
1178
1179 handle.try_send(PersistRequest::CompletedCommit { session });
1180 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1181 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1182 let report = reply_rx.await.expect("flush report reply");
1183
1184 assert!(
1185 report.failures.is_empty(),
1186 "commit save and its clear must both succeed, got: {:?}",
1187 report.failures
1188 );
1189 let saved = verification_manager
1190 .load_session(&session_id)
1191 .expect("completed session must be saved before its checkpoint is cleared");
1192 assert_eq!(saved.metadata.id, session_id);
1193 assert!(
1194 !checkpoint_path.exists(),
1195 "the checkpoint is cleared only after the session save succeeded"
1196 );
1197 handle.try_send(PersistRequest::Shutdown);
1198 task.await.expect("persistence actor join");
1199 }
1200
1201 #[tokio::test]
1202 async fn completed_commit_never_clears_another_sessions_checkpoint() {
1203 let tmp = tempfile::tempdir().expect("tempdir");
1204 let sessions_dir = tmp.path().join("sessions");
1205 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1206 let verification_manager = SessionManager::new(sessions_dir.clone()).expect("verify");
1207 let committed = crate::session_manager::create_saved_session_with_mode(
1208 &[],
1209 "deepseek-v4-pro",
1210 tmp.path(),
1211 0,
1212 None,
1213 Some("agent"),
1214 );
1215 let inflight = crate::session_manager::create_saved_session_with_mode(
1216 &[],
1217 "deepseek-v4-pro",
1218 tmp.path(),
1219 0,
1220 None,
1221 Some("agent"),
1222 );
1223 let committed_id = committed.metadata.id.clone();
1224 let inflight_id = inflight.metadata.id.clone();
1225 let committed_checkpoint = seed_checkpoint_file(&sessions_dir, &committed_id);
1226 // The concurrent session's checkpoint must carry a loadable body, so
1227 // seed it through the real manager instead of a placeholder file.
1228 let inflight_checkpoint = verification_manager
1229 .save_checkpoint(&inflight)
1230 .expect("seed the concurrent session's checkpoint");
1231 let (handle, task) = spawn_persistence_actor(manager);
1232
1233 handle.try_send(PersistRequest::CompletedCommit { session: committed });
1234 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1235 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1236 let report = reply_rx.await.expect("flush report reply");
1237 assert!(report.failures.is_empty(), "{:?}", report.failures);
1238
1239 assert!(
1240 !committed_checkpoint.exists(),
1241 "the committed session's own checkpoint must be cleared"
1242 );
1243 let survivor = verification_manager
1244 .load_session_checkpoint(&inflight_id)
1245 .expect("load the concurrent session's checkpoint")
1246 .expect("another session's checkpoint must never be cleared");
1247 assert_eq!(survivor.metadata.id, inflight_id);
1248 assert!(inflight_checkpoint.exists());
1249 handle.try_send(PersistRequest::Shutdown);
1250 task.await.expect("persistence actor join");
1251 }
1252
1253 #[tokio::test]
1254 async fn shutdown_preserves_inflight_checkpoint() {
1255 let tmp = tempfile::tempdir().expect("tempdir");
1256 let sessions_dir = tmp.path().join("sessions");
1257 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1258 let verification_manager = SessionManager::new(sessions_dir.clone()).expect("verify");
1259 let session = crate::session_manager::create_saved_session_with_mode(
1260 &[],
1261 "deepseek-v4-pro",
1262 tmp.path(),
1263 0,
1264 None,
1265 Some("agent"),
1266 );
1267 let session_id = session.metadata.id.clone();
1268 let checkpoint_path = seed_checkpoint_file(&sessions_dir, &session_id);
1269 let (handle, task) = spawn_persistence_actor(manager);
1270
1271 handle.try_send(PersistRequest::SaveCheckpoint { session });
1272 handle.try_send(PersistRequest::Shutdown);
1273 task.await.expect("persistence actor join");
1274
1275 assert!(
1276 checkpoint_path.exists(),
1277 "shutdown must never unconditionally clear an in-flight checkpoint"
1278 );
1279 let recovered = verification_manager
1280 .load_session_checkpoint(&session_id)
1281 .expect("load checkpoint after shutdown")
1282 .expect("in-flight work must survive shutdown for recovery review");
1283 assert_eq!(recovered.metadata.id, session_id);
1284 }
1285
1286 #[tokio::test]
1287 async fn newer_inflight_checkpoint_supersedes_pending_completed_commit() {
1288 let tmp = tempfile::tempdir().expect("tempdir");
1289 let sessions_dir = tmp.path().join("sessions");
1290 let manager = SessionManager::new(sessions_dir.clone()).expect("manager");
1291 let completed = crate::session_manager::create_saved_session_with_mode(
1292 &[],
1293 "deepseek-v4-pro",
1294 tmp.path(),
1295 0,
1296 None,
1297 Some("agent"),
1298 );
1299 let mut inflight = crate::session_manager::create_saved_session_with_mode(
1300 &[],
1301 "deepseek-v4-pro",
1302 tmp.path(),
1303 0,
1304 None,
1305 Some("agent"),
1306 );
1307 // The new turn reuses the same session id: a fresh in-flight
1308 // checkpoint arrives while the previous completion is still queued.
1309 inflight.metadata.id = completed.metadata.id.clone();
1310 let session_id = completed.metadata.id.clone();
1311 let (handle, task) = spawn_persistence_actor(manager);
1312
1313 handle.try_send(PersistRequest::CompletedCommit { session: completed });
1314 handle.try_send(PersistRequest::SaveCheckpoint { session: inflight });
1315 let (reply_tx, reply_rx) = tokio::sync::oneshot::channel();
1316 handle.try_send(PersistRequest::FlushAndReport { reply: reply_tx });
1317 let report = reply_rx.await.expect("flush report reply");
1318
1319 assert!(
1320 !report
1321 .failures
1322 .iter()
1323 .any(|(what, _)| what.starts_with("clear-checkpoint:")),
1324 "the newer in-flight checkpoint must not be cleared, got: {:?}",
1325 report.failures
1326 );
1327 let checkpoint = std::fs::read_to_string(
1328 sessions_dir
1329 .join("checkpoints")
1330 .join(format!("{session_id}.json")),
1331 )
1332 .expect("the newer checkpoint must survive the drain");
1333 assert!(
1334 !checkpoint.is_empty(),
1335 "the newer checkpoint body must be on disk"
1336 );
1337 assert!(
1338 report.completed >= 1,
1339 "the newer checkpoint write must be counted, got: {report:?}"
1340 );
1341 handle.try_send(PersistRequest::Shutdown);
1342 task.await.expect("persistence actor join");
1343 }
1344 #[test]
1345 fn offline_queue_editor_lease_survives_until_pending_write_finishes() {
1346 let directory = tempfile::tempdir().expect("queue fixture");
1347 let manager = SessionManager::new(directory.path().join("sessions")).expect("manager");
1348 let lease = manager
1349 .acquire_offline_queue_lease("session-A")
1350 .expect("first editor");
1351 let mut pending = PendingState::default();
1352 pending.absorb(PersistRequest::OfflineQueue {
1353 state: OfflineQueueState {
1354 draft: Some(QueuedSessionMessage {
1355 display: "last edited draft".into(),
1356 skill_instruction: None,
1357 skill_provenance: None,
1358 }),
1359 ..OfflineQueueState::default()
1360 },
1361 lease: Arc::clone(&lease),
1362 });
1363 drop(lease); // The old window changed session before the actor ran.
1364 assert!(manager.acquire_offline_queue_lease("session-A").is_err());
1365 let report = flush_inner(&manager, &mut pending, &SessionSaveHealth::default());
1366 assert!(report.failures.is_empty(), "draft write failed: {report:?}");
1367 assert_eq!(report.completed, 1);
1368 let _next_editor = manager
1369 .acquire_offline_queue_lease("session-A")
1370 .expect("released after write");
1371 assert_eq!(
1372 manager
1373 .load_offline_queue_state("session-A")
1374 .unwrap()
1375 .unwrap()
1376 .draft
1377 .unwrap()
1378 .display,
1379 "last edited draft"
1380 );
1381 }
1382 }
1383
1383 lines RUST