diff --git a/crates/session-store/src/worker_metadata.rs b/crates/session-store/src/worker_metadata.rs index 7fd8cdf5..4491f393 100644 --- a/crates/session-store/src/worker_metadata.rs +++ b/crates/session-store/src/worker_metadata.rs @@ -829,6 +829,7 @@ where #[cfg(test)] mod tests { use super::*; + use crate::{LogEntry, Store}; #[test] fn worker_metadata_manifest_snapshot_roundtrips() { @@ -1050,6 +1051,111 @@ mod tests { assert_eq!(restored.reclaimed_children[0].scope_delegated, vec![scope]); } + #[test] + fn staged_segment_is_invisible_until_cas_and_reopen_selects_committed_history() { + let temp = tempfile::tempdir().unwrap(); + let sessions = temp.path().join("sessions"); + let workers = temp.path().join("workers"); + let open = || { + CombinedStore::new( + crate::FsStore::new(&sessions).unwrap(), + FsWorkerStore::new(&workers).unwrap(), + ) + }; + let store = open(); + let session_id = crate::new_session_id(); + let old_segment_id = crate::new_segment_id(); + let new_segment_id = crate::new_segment_id(); + let entry = |label: &str| LogEntry::Extension { + ts: 1, + domain: label.into(), + payload: serde_json::json!({}), + }; + store + .create_segment(session_id, old_segment_id, &[entry("old-history")]) + .unwrap(); + store + .write(&WorkerMetadata::new( + "agent", + Some(WorkerActiveSegmentRef::active_segment( + session_id, + old_segment_id, + )), + )) + .unwrap(); + store + .create_segment(session_id, new_segment_id, &[entry("new-history")]) + .unwrap(); + drop(store); + + let reopened = open(); + assert_eq!( + reopened + .read_by_name("agent") + .unwrap() + .unwrap() + .active + .unwrap() + .segment_id, + Some(old_segment_id) + ); + assert!( + reopened + .compare_and_swap_active( + "agent", + &WorkerActiveSegmentRef::active_segment(session_id, old_segment_id), + WorkerActiveSegmentRef::active_segment(session_id, new_segment_id), + ) + .unwrap() + ); + drop(reopened); + + let reopened = open(); + assert_eq!( + reopened + .read_by_name("agent") + .unwrap() + .unwrap() + .active + .unwrap() + .segment_id, + Some(new_segment_id) + ); + assert!(matches!( + reopened.read_all(session_id, new_segment_id).unwrap().as_slice(), + [LogEntry::Extension { domain, .. }] if domain == "new-history" + )); + } + + #[test] + fn aggregate_store_uses_expected_old_segment_cas() { + let temp = tempfile::tempdir().unwrap(); + let store = WorkerAggregateStore::new(temp.path(), "agent").unwrap(); + let session_id = crate::new_session_id(); + let old = WorkerActiveSegmentRef::active_segment(session_id, crate::new_segment_id()); + store + .write(&WorkerMetadata::new("agent", Some(old.clone()))) + .unwrap(); + assert!( + store + .compare_and_swap_active( + "agent", + &old, + WorkerActiveSegmentRef::active_segment(session_id, crate::new_segment_id()), + ) + .unwrap() + ); + assert!( + !store + .compare_and_swap_active( + "agent", + &old, + WorkerActiveSegmentRef::active_segment(session_id, crate::new_segment_id()), + ) + .unwrap() + ); + } + #[test] fn combined_store_delegates_atomic_active_segment_cas() { let temp = tempfile::tempdir().unwrap(); diff --git a/crates/tui/src/app.rs b/crates/tui/src/app.rs index be192281..b7900da5 100644 --- a/crates/tui/src/app.rs +++ b/crates/tui/src/app.rs @@ -1521,9 +1521,9 @@ impl App { } => { self.rewind_refresh_fence = false; self.pending_submissions = session.pending_submissions.clone(); + self.apply_worker_state_snapshot(&state); self.restore_snapshot(&session, greeting, in_flight); self.replace_internal_worker_snapshots(internal_workers); - self.apply_worker_state_snapshot(&state); } Event::InternalWorker { worker, @@ -4339,16 +4339,24 @@ mod completion_flow_tests { #[test] fn snapshot_restores_and_runtime_clear_removes_compaction_progress() { let mut app = App::new("test".into()); - app.worker_state.state = protocol::WorkerState::Busy( - protocol::WorkerBusyState::Maintenance(protocol::WorkerMaintenanceState::Compacting), - ); - app.apply_in_flight_snapshot(InFlightSnapshot { - compaction: Some(protocol::InFlightCompaction { - phase: protocol::CompactionPhase::Summarizing, - started_at_ms: 100, - trigger: protocol::CompactionTrigger::Manual, - }), - ..InFlightSnapshot::default() + assert_eq!(app.worker_state.state, protocol::WorkerState::Idle); + let mut state = protocol::WorkerStateSnapshot::initial(2); + state.state = protocol::WorkerState::Busy(protocol::WorkerBusyState::Maintenance( + protocol::WorkerMaintenanceState::Compacting, + )); + app.handle_worker_event(Event::Snapshot { + session: public_session(Vec::new()), + greeting: test_greeting(), + state, + in_flight: InFlightSnapshot { + compaction: Some(protocol::InFlightCompaction { + phase: protocol::CompactionPhase::Summarizing, + started_at_ms: 100, + trigger: protocol::CompactionTrigger::Manual, + }), + ..InFlightSnapshot::default() + }, + internal_workers: Vec::new(), }); assert_eq!(compact_block_count(&app), 0); assert_eq!(