Our Personal Data Server from scratch!
0

Configure Feed

Select the types of activity you want to include in your feed.

fix(store): unblock eventlog sync&freeze when writer dies

Lewis: May this revision serve well! <lu5a@proton.me>

Lewis (Jun 2, 2026, 5:22 PM +0300) 0b5fae79 54244075

+149 -14
+50
crates/tranquil-store/src/eventlog/commit_loop.rs
··· 21 21 const REORDER_TIMEOUT: Duration = Duration::from_millis(100); 22 22 const GAP_ABANDON_TIMEOUT: Duration = Duration::from_secs(5); 23 23 24 + pub(super) fn writer_terminated() -> io::Error { 25 + io::Error::other("eventlog writer thread terminated") 26 + } 27 + 24 28 pub struct FreezeResponse { 25 29 pub synced_through: EventSequence, 26 30 pub segment_id: SegmentId, ··· 42 46 pub struct WriterNotify { 43 47 synced_seq: AtomicU64, 44 48 poisoned: AtomicBool, 49 + terminated: AtomicBool, 45 50 mutex: Mutex<()>, 46 51 cond: Condvar, 47 52 } ··· 128 133 Self { 129 134 synced_seq: AtomicU64::new(initial_synced), 130 135 poisoned: AtomicBool::new(false), 136 + terminated: AtomicBool::new(false), 131 137 mutex: Mutex::new(()), 132 138 cond: Condvar::new(), 133 139 } 140 + } 141 + 142 + pub(super) fn is_terminated(&self) -> bool { 143 + self.terminated.load(Ordering::Acquire) 134 144 } 135 145 136 146 pub fn wait_for_sync(&self, target: EventSequence) -> io::Result<()> { ··· 143 153 if self.poisoned.load(Ordering::Acquire) { 144 154 return Err(io::Error::other("eventlog writer poisoned")); 145 155 } 156 + if self.terminated.load(Ordering::Acquire) { 157 + return Err(io::Error::other("eventlog writer thread terminated")); 158 + } 146 159 147 160 let deadline = Instant::now() + SYNC_TIMEOUT; 148 161 let mut guard = self.mutex.lock(); ··· 153 166 } 154 167 if self.poisoned.load(Ordering::Acquire) { 155 168 return Err(io::Error::other("eventlog writer poisoned")); 169 + } 170 + if self.terminated.load(Ordering::Acquire) { 171 + return Err(io::Error::other("eventlog writer thread terminated")); 156 172 } 157 173 158 174 let now = Instant::now(); ··· 175 191 176 192 fn poison(&self) { 177 193 self.poisoned.store(true, Ordering::Release); 194 + let _guard = self.mutex.lock(); 195 + self.cond.notify_all(); 196 + } 197 + 198 + fn terminate(&self) { 199 + self.terminated.store(true, Ordering::Release); 178 200 let _guard = self.mutex.lock(); 179 201 self.cond.notify_all(); 180 202 } ··· 372 394 } 373 395 } 374 396 397 + fn fail_abandoned_request(request: WriterRequest) { 398 + match request { 399 + WriterRequest::SyncBarrier { response } => { 400 + let _ = response.send(Err(writer_terminated())); 401 + } 402 + WriterRequest::Freeze { response, .. } => { 403 + let _ = response.send(Err(writer_terminated())); 404 + } 405 + WriterRequest::Append(_) | WriterRequest::Shutdown => {} 406 + } 407 + } 408 + 409 + struct TerminateOnDrop<'a> { 410 + notify: &'a WriterNotify, 411 + receiver: &'a flume::Receiver<WriterRequest>, 412 + } 413 + 414 + impl Drop for TerminateOnDrop<'_> { 415 + fn drop(&mut self) { 416 + self.notify.terminate(); 417 + self.receiver.drain().for_each(fail_abandoned_request); 418 + } 419 + } 420 + 375 421 fn writer_loop<S: StorageIO>( 376 422 receiver: &flume::Receiver<WriterRequest>, 377 423 writer: &mut EventLogWriter<S>, 378 424 ctx: &WriterCtx<'_, S>, 379 425 ) { 380 426 let _close = CloseOnDrop(ctx.pending_bytes); 427 + let _terminate = TerminateOnDrop { 428 + notify: ctx.notify, 429 + receiver, 430 + }; 381 431 let mut reorder = ReorderBuffer::new(writer.current_seq().next().raw()); 382 432 383 433 loop {
+26 -10
crates/tranquil-store/src/eventlog/mod.rs
··· 26 26 use crate::fsync_order::PostBlockstoreHook; 27 27 use crate::io::StorageIO; 28 28 29 - use commit_loop::{CommitThread, FreezeResponse, PendingBytesBudget, WriterNotify, WriterRequest}; 29 + use commit_loop::{ 30 + CommitThread, FreezeResponse, PendingBytesBudget, WriterNotify, WriterRequest, 31 + writer_terminated, 32 + }; 30 33 31 34 pub use bridge::{DeferredBroadcast, EventLogBridge}; 32 35 pub use manager::{SEGMENT_FILE_EXTENSION, SegmentManager, parse_segment_id, segment_path}; ··· 51 54 52 55 const DEFAULT_BROADCAST_BUFFER: usize = 16384; 53 56 pub const DEFAULT_PENDING_BYTES_BUDGET: u64 = 1024 * 1024 * 1024; 57 + const WRITER_RESPONSE_POLL: Duration = Duration::from_millis(100); 54 58 55 59 pub struct EventWithMutations { 56 60 pub event: SequencedEvent, ··· 187 191 self.commit_thread 188 192 .sender() 189 193 .send(WriterRequest::Append(event)) 190 - .map_err(|_| io::Error::other("eventlog writer thread terminated")) 194 + .map_err(|_| writer_terminated()) 191 195 } 192 196 193 197 pub fn append_event( ··· 261 265 self.commit_thread 262 266 .sender() 263 267 .send(WriterRequest::SyncBarrier { response: resp_tx }) 264 - .map_err(|_| io::Error::other("eventlog writer thread terminated"))?; 265 - resp_rx 266 - .recv() 267 - .map_err(|_| io::Error::other("eventlog writer thread terminated"))? 268 + .map_err(|_| writer_terminated())?; 269 + self.await_writer_response(&resp_rx) 270 + } 271 + 272 + fn await_writer_response<T>(&self, resp_rx: &flume::Receiver<io::Result<T>>) -> io::Result<T> { 273 + loop { 274 + match resp_rx.recv_timeout(WRITER_RESPONSE_POLL) { 275 + Ok(result) => return result, 276 + Err(flume::RecvTimeoutError::Disconnected) => { 277 + return Err(writer_terminated()); 278 + } 279 + Err(flume::RecvTimeoutError::Timeout) => { 280 + if self.notify.is_terminated() { 281 + return Err(writer_terminated()); 282 + } 283 + } 284 + } 285 + } 268 286 } 269 287 270 288 pub fn get_events_since( ··· 432 450 response: resp_tx, 433 451 resume: resume_rx, 434 452 }) 435 - .map_err(|_| io::Error::other("eventlog writer thread terminated"))?; 453 + .map_err(|_| writer_terminated())?; 436 454 437 - let freeze_resp: FreezeResponse = resp_rx 438 - .recv() 439 - .map_err(|_| io::Error::other("eventlog writer thread terminated"))??; 455 + let freeze_resp: FreezeResponse = self.await_writer_response(&resp_rx)?; 440 456 441 457 let all_segments = self.manager.list_segments()?; 442 458 let sealed_segments: Vec<SegmentId> = all_segments
+73 -4
crates/tranquil-store/tests/sim_cross_store.rs
··· 51 51 uid 52 52 } 53 53 54 - fn append_event(stores: &TestStores, idx: u64) { 55 - let did = test_did(idx); 56 - let event = SequencedEvent { 54 + fn make_commit_event(did: &Did, idx: u64) -> SequencedEvent { 55 + SequencedEvent { 57 56 seq: SequenceNumber::from_raw(0), 58 57 did: did.clone(), 59 58 created_at: chrono::Utc::now(), ··· 68 67 active: None, 69 68 status: None, 70 69 rev: Some(format!("rev{idx}")), 71 - }; 70 + } 71 + } 72 + 73 + fn append_event(stores: &TestStores, idx: u64) { 74 + let did = test_did(idx); 75 + let event = make_commit_event(&did, idx); 72 76 stores 73 77 .eventlog 74 78 .append_event(&did, RepoEventType::Commit, &event) ··· 988 992 fn sim_cross_store_coordinated_commit_under_faults() { 989 993 sim_seed_range().into_par_iter().for_each(|seed| { 990 994 run_cross_store_fault_scenario(seed); 995 + }); 996 + } 997 + 998 + fn open_fault_eventlog( 999 + dir: &std::path::Path, 1000 + sim: &std::sync::Arc<SimulatedIO>, 1001 + ) -> std::sync::Arc<EventLog<std::sync::Arc<SimulatedIO>>> { 1002 + std::sync::Arc::new( 1003 + EventLog::open( 1004 + EventLogConfig { 1005 + segments_dir: dir.join("eventlog/segments"), 1006 + ..EventLogConfig::default() 1007 + }, 1008 + std::sync::Arc::clone(sim), 1009 + ) 1010 + .unwrap(), 1011 + ) 1012 + } 1013 + 1014 + fn run_eventlog_commit_loop_under_faults(seed: u64) { 1015 + let dir = tempfile::TempDir::new().unwrap(); 1016 + std::fs::create_dir_all(dir.path().join("eventlog/segments")).unwrap(); 1017 + let fault = sim_fault_for(seed); 1018 + let sim = std::sync::Arc::new(SimulatedIO::new(seed, fault)); 1019 + let did = test_did(seed); 1020 + let event_count = (seed % 40) + 20; 1021 + 1022 + let mut acked_sync_seq: u64 = 0; 1023 + { 1024 + sim.set_pristine_mode(true); 1025 + let eventlog = open_fault_eventlog(dir.path(), &sim); 1026 + sim.set_pristine_mode(false); 1027 + 1028 + (0..event_count).for_each(|i| { 1029 + if eventlog 1030 + .append_event(&did, RepoEventType::Commit, &make_commit_event(&did, i)) 1031 + .is_ok() 1032 + && let Ok(result) = eventlog.sync() 1033 + { 1034 + acked_sync_seq = acked_sync_seq.max(result.synced_through.raw()); 1035 + } 1036 + }); 1037 + 1038 + sim.crash(); 1039 + drop(eventlog); 1040 + } 1041 + 1042 + sim.set_pristine_mode(true); 1043 + let reopened = open_fault_eventlog(dir.path(), &sim); 1044 + let recovered = reopened.max_seq().raw(); 1045 + assert!( 1046 + recovered >= acked_sync_seq, 1047 + "seed={seed} fault={fault:?}: events acknowledged synced through {acked_sync_seq} must survive crash, recovered max_seq={recovered}" 1048 + ); 1049 + assert!( 1050 + recovered <= event_count, 1051 + "seed={seed} fault={fault:?}: recovery invented events beyond the {event_count} appended, recovered max_seq={recovered}" 1052 + ); 1053 + let _ = reopened.shutdown(); 1054 + } 1055 + 1056 + #[test] 1057 + fn sim_eventlog_commit_loop_durability_under_faults() { 1058 + sim_seed_range().into_par_iter().for_each(|seed| { 1059 + run_eventlog_commit_loop_under_faults(seed); 991 1060 }); 992 1061 } 993 1062