1use std::time::Duration;
19
20use crate::{
21 error::EventStoreError,
22 markers::{DataCursorSnapshot, HiFiMarker, MarkerBackend, MarkerGap, StreamDictEntry},
23};
24
25pub const DEFAULT_MARKER_CHANNEL_CAPACITY: usize = 10_000;
27pub const DEFAULT_MARKER_MAX_BATCH: usize = 100;
29pub const DEFAULT_MARKER_MAX_LATENCY: Duration = Duration::from_millis(5);
31
32#[derive(Clone, Debug)]
34pub struct MarkerWriterConfig {
35 pub channel_capacity: usize,
37 pub max_batch: usize,
39 pub max_latency: Duration,
41}
42
43impl Default for MarkerWriterConfig {
44 fn default() -> Self {
45 Self {
46 channel_capacity: DEFAULT_MARKER_CHANNEL_CAPACITY,
47 max_batch: DEFAULT_MARKER_MAX_BATCH,
48 max_latency: DEFAULT_MARKER_MAX_LATENCY,
49 }
50 }
51}
52
53#[derive(Debug, Clone, PartialEq, Eq)]
55pub enum MarkerMsg {
56 Snapshot(DataCursorSnapshot),
58 HiFi(HiFiMarker),
60 Dict(StreamDictEntry),
62 Close,
64 #[doc(hidden)]
65 GapThen { gap: MarkerGap, msg: Box<Self> },
66}
67
68#[cfg(not(madsim))]
69mod imp {
70 use std::{
71 fmt::Debug,
72 sync::{
73 atomic::{AtomicBool, AtomicU64, Ordering},
74 mpsc::{self, RecvTimeoutError, SyncSender, TrySendError},
75 },
76 thread::{self, JoinHandle},
77 time::Instant,
78 };
79
80 use nautilus_core::time::AtomicTime;
81 use parking_lot::Mutex;
82
83 use super::{
84 EventStoreError, MarkerBackend, MarkerGap, MarkerMsg, MarkerWriterConfig, StreamDictEntry,
85 };
86 use crate::{
87 manifest::RunStatus,
88 markers::{
89 MarkerGapReason, compute_dict_hash, compute_gap_hash, compute_hifi_hash,
90 compute_marker_hash,
91 },
92 };
93
94 const MARKER_WRITER_THREAD_NAME: &str = "event-store-marker-writer";
95
96 pub struct MarkerWriter {
98 tx: Option<SyncSender<MarkerMsg>>,
99 handle: Option<JoinHandle<()>>,
100 last_submitted_seq: AtomicU64,
101 dropped: Mutex<Option<DroppedRange>>,
102 closed: AtomicBool,
103 }
104
105 #[derive(Debug, Clone, Copy)]
106 struct DroppedRange {
107 from_marker_seq: u64,
108 to_marker_seq: u64,
109 }
110
111 impl DroppedRange {
112 const fn new(marker_seq: u64) -> Self {
113 Self {
114 from_marker_seq: marker_seq,
115 to_marker_seq: marker_seq,
116 }
117 }
118
119 fn extend(&mut self, marker_seq: u64) {
120 if marker_seq < self.from_marker_seq {
121 self.from_marker_seq = marker_seq;
122 }
123
124 if marker_seq > self.to_marker_seq {
125 self.to_marker_seq = marker_seq;
126 }
127 }
128
129 const fn gap(self, reason: MarkerGapReason) -> MarkerGap {
130 MarkerGap {
131 from_marker_seq: self.from_marker_seq,
132 to_marker_seq: self.to_marker_seq,
133 reason,
134 }
135 }
136 }
137
138 impl Debug for MarkerWriter {
139 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
140 f.debug_struct(stringify!(MarkerWriter))
141 .field(
142 "last_submitted_seq",
143 &self.last_submitted_seq.load(Ordering::Acquire),
144 )
145 .field("closed", &self.closed.load(Ordering::Acquire))
146 .field("tx_attached", &self.tx.is_some())
147 .finish_non_exhaustive()
148 }
149 }
150
151 impl MarkerWriter {
152 pub fn spawn(
159 backend: Box<dyn MarkerBackend + Send>,
160 _clock: &'static AtomicTime,
161 config: MarkerWriterConfig,
162 ) -> Result<Self, EventStoreError> {
163 backend.manifest()?;
164 let (tx, rx) = mpsc::sync_channel::<MarkerMsg>(config.channel_capacity.max(1));
168 let config_for_thread = config;
169
170 let handle = thread::Builder::new()
171 .name(MARKER_WRITER_THREAD_NAME.to_string())
172 .spawn(move || run(backend, &rx, &config_for_thread))
173 .map_err(|e| {
174 EventStoreError::Backend(format!("spawn marker writer thread: {e}"))
175 })?;
176
177 Ok(Self {
178 tx: Some(tx),
179 handle: Some(handle),
180 last_submitted_seq: AtomicU64::new(0),
181 dropped: Mutex::new(None),
182 closed: AtomicBool::new(false),
183 })
184 }
185
186 pub fn submit(&self, msg: MarkerMsg, marker_seq: u64) -> Result<bool, EventStoreError> {
196 if self.closed.load(Ordering::Acquire) {
197 return Err(EventStoreError::Closed);
198 }
199
200 if matches!(
201 msg,
202 MarkerMsg::Dict(_) | MarkerMsg::Close | MarkerMsg::GapThen { .. }
203 ) {
204 return Err(EventStoreError::Backend(
205 "submit accepts Snapshot or HiFi marker messages".to_string(),
206 ));
207 }
208
209 let tx = self.tx.as_ref().ok_or(EventStoreError::Closed)?;
210 let mut dropped = self.dropped.lock();
211 let outbound = if let Some(range) = *dropped {
212 MarkerMsg::GapThen {
213 gap: range.gap(MarkerGapReason::Overflow),
214 msg: Box::new(msg),
215 }
216 } else {
217 msg
218 };
219
220 match tx.try_send(outbound) {
221 Ok(()) => {
222 *dropped = None;
223 self.last_submitted_seq.store(marker_seq, Ordering::Release);
224 Ok(true)
225 }
226 Err(TrySendError::Full(_)) => {
227 extend_dropped_range(&mut dropped, marker_seq);
228 Ok(false)
229 }
230 Err(TrySendError::Disconnected(_)) => {
231 self.closed.store(true, Ordering::Release);
232 Err(EventStoreError::Closed)
233 }
234 }
235 }
236
237 pub fn put_dict(&self, entry: StreamDictEntry) -> Result<bool, EventStoreError> {
246 if self.closed.load(Ordering::Acquire) {
247 return Err(EventStoreError::Closed);
248 }
249
250 let tx = self.tx.as_ref().ok_or(EventStoreError::Closed)?;
251 if tx.send(MarkerMsg::Dict(entry)).is_ok() {
252 Ok(true)
253 } else {
254 self.closed.store(true, Ordering::Release);
255 Err(EventStoreError::Closed)
256 }
257 }
258
259 pub fn close(mut self) {
261 self.close_inner();
262 }
263
264 fn close_inner(&mut self) {
268 self.closed.store(true, Ordering::Release);
269
270 if let Some(tx) = self.tx.take() {
271 let close_msg = {
272 let mut dropped = self.dropped.lock();
273 if let Some(range) = dropped.take() {
274 MarkerMsg::GapThen {
275 gap: range.gap(MarkerGapReason::WriterClosed),
276 msg: Box::new(MarkerMsg::Close),
277 }
278 } else {
279 MarkerMsg::Close
280 }
281 };
282 let _ = tx.send(close_msg);
283 drop(tx);
284 }
285
286 if let Some(handle) = self.handle.take()
287 && handle.join().is_err()
288 {
289 log::error!("Marker writer thread panicked");
290 }
291 }
292 }
293
294 impl Drop for MarkerWriter {
295 fn drop(&mut self) {
296 self.close_inner();
297 }
298 }
299
300 fn extend_dropped_range(dropped: &mut Option<DroppedRange>, marker_seq: u64) {
301 if let Some(range) = dropped {
302 range.extend(marker_seq);
303 } else {
304 *dropped = Some(DroppedRange::new(marker_seq));
305 }
306 }
307
308 fn run(
309 mut backend: Box<dyn MarkerBackend + Send>,
310 rx: &mpsc::Receiver<MarkerMsg>,
311 config: &MarkerWriterConfig,
312 ) {
313 let max_batch = config.max_batch.max(1);
314 let mut batch = Vec::with_capacity(max_batch);
315
316 while let Ok(first) = rx.recv() {
317 let mut should_close = push_batch(&mut batch, first);
318 let mut disconnected = false;
319 let started = Instant::now();
320
321 while !should_close && batch.len() < max_batch {
322 let Some(remaining) = config.max_latency.checked_sub(started.elapsed()) else {
323 break;
324 };
325
326 if remaining.is_zero() {
327 break;
328 }
329
330 match rx.recv_timeout(remaining) {
331 Ok(msg) => should_close = push_batch(&mut batch, msg),
332 Err(RecvTimeoutError::Timeout) => break,
333 Err(RecvTimeoutError::Disconnected) => {
334 disconnected = true;
335 break;
336 }
337 }
338 }
339
340 if let Err(e) = write_batch(backend.as_mut(), batch.drain(..)) {
341 log::error!(
342 "Marker writer fail-stopped, marker capture is disabled for the rest of the run: {e}"
343 );
344 return;
345 }
346
347 if should_close {
348 if let Err(e) = backend.seal(RunStatus::Ended) {
349 log::error!("Failed to seal marker run on close: {e}");
350 }
351 return;
352 }
353
354 if disconnected {
355 return;
356 }
357 }
358 }
359
360 fn push_batch(batch: &mut Vec<MarkerMsg>, msg: MarkerMsg) -> bool {
361 match msg {
362 MarkerMsg::Close => true,
363 other => {
364 let should_close = closes_after_msg(&other);
365 batch.push(other);
366 should_close
367 }
368 }
369 }
370
371 fn closes_after_msg(msg: &MarkerMsg) -> bool {
372 match msg {
373 MarkerMsg::Close => true,
374 MarkerMsg::GapThen { msg, .. } => closes_after_msg(msg),
375 MarkerMsg::Snapshot(_) | MarkerMsg::HiFi(_) | MarkerMsg::Dict(_) => false,
376 }
377 }
378
379 fn write_batch(
380 backend: &mut dyn MarkerBackend,
381 batch: impl IntoIterator<Item = MarkerMsg>,
382 ) -> Result<(), EventStoreError> {
383 for msg in batch {
384 write_msg(backend, msg)?;
385 }
386 Ok(())
387 }
388
389 fn write_msg(backend: &mut dyn MarkerBackend, msg: MarkerMsg) -> Result<(), EventStoreError> {
390 match msg {
391 MarkerMsg::Snapshot(snapshot) => {
392 backend.append_snapshot(&snapshot, compute_marker_hash(&snapshot))
393 }
394 MarkerMsg::HiFi(marker) => backend.append_hifi(&marker, compute_hifi_hash(&marker)),
395 MarkerMsg::Dict(entry) => backend.put_dict(&entry, compute_dict_hash(&entry)),
396 MarkerMsg::GapThen { gap, msg } => {
397 backend.append_gap(&gap, compute_gap_hash(&gap))?;
398 write_msg(backend, *msg)
399 }
400 MarkerMsg::Close => Ok(()),
401 }
402 }
403}
404
405#[cfg(madsim)]
406mod imp {
407 use std::{
408 fmt::Debug,
409 sync::atomic::{AtomicU64, Ordering},
410 };
411
412 use nautilus_core::time::AtomicTime;
413 use parking_lot::Mutex;
414
415 use super::{EventStoreError, MarkerBackend, MarkerMsg, MarkerWriterConfig, StreamDictEntry};
416 use crate::{
417 manifest::RunStatus,
418 markers::{compute_dict_hash, compute_hifi_hash, compute_marker_hash},
419 };
420
421 pub struct MarkerWriter {
423 inner: Mutex<Inner>,
424 last_submitted_seq: AtomicU64,
425 }
426
427 struct Inner {
428 backend: Box<dyn MarkerBackend + Send>,
429 closed: bool,
430 }
431
432 impl Debug for MarkerWriter {
433 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
434 f.debug_struct(stringify!(MarkerWriter))
435 .field(
436 "last_submitted_seq",
437 &self.last_submitted_seq.load(Ordering::Acquire),
438 )
439 .finish_non_exhaustive()
440 }
441 }
442
443 impl MarkerWriter {
444 pub fn spawn(
450 backend: Box<dyn MarkerBackend + Send>,
451 _clock: &'static AtomicTime,
452 _config: MarkerWriterConfig,
453 ) -> Result<Self, EventStoreError> {
454 backend.manifest()?;
455 Ok(Self {
456 inner: Mutex::new(Inner {
457 backend,
458 closed: false,
459 }),
460 last_submitted_seq: AtomicU64::new(0),
461 })
462 }
463
464 pub fn submit(&self, msg: MarkerMsg, marker_seq: u64) -> Result<bool, EventStoreError> {
471 let mut inner = self.inner.lock();
472
473 if inner.closed {
474 return Err(EventStoreError::Closed);
475 }
476
477 match msg {
478 MarkerMsg::Snapshot(snapshot) => inner
479 .backend
480 .append_snapshot(&snapshot, compute_marker_hash(&snapshot))?,
481 MarkerMsg::HiFi(marker) => inner
482 .backend
483 .append_hifi(&marker, compute_hifi_hash(&marker))?,
484 MarkerMsg::Dict(_) | MarkerMsg::Close | MarkerMsg::GapThen { .. } => {
485 return Err(EventStoreError::Backend(
486 "submit accepts Snapshot or HiFi marker messages".to_string(),
487 ));
488 }
489 }
490
491 self.last_submitted_seq.store(marker_seq, Ordering::Release);
492 Ok(true)
493 }
494
495 #[expect(
501 clippy::needless_pass_by_value,
502 reason = "matches the threaded writer API, which transfers ownership"
503 )]
504 pub fn put_dict(&self, entry: StreamDictEntry) -> Result<bool, EventStoreError> {
505 let mut inner = self.inner.lock();
506
507 if inner.closed {
508 return Err(EventStoreError::Closed);
509 }
510
511 inner.backend.put_dict(&entry, compute_dict_hash(&entry))?;
512 Ok(true)
513 }
514
515 pub fn close(self) {
517 let mut inner = self.inner.lock();
518
519 if !inner.closed {
520 let _ = inner.backend.seal(RunStatus::Ended);
521 inner.closed = true;
522 }
523 }
524 }
525}
526
527pub use imp::MarkerWriter;
528
529#[cfg(test)]
530#[cfg(not(madsim))]
531mod tests {
532 use std::{
533 sync::{
534 Arc,
535 atomic::{AtomicUsize, Ordering},
536 },
537 time::{Duration, Instant},
538 };
539
540 use nautilus_core::{UnixNanos, time::get_atomic_clock_static};
541 use parking_lot::Mutex;
542 use rstest::rstest;
543
544 use super::{super::test_support::SharedMemoryMarker, *};
545 use crate::{
546 error::EventStoreError,
547 manifest::RunStatus,
548 markers::{
549 DataClass, MarkerGap, MarkerGapReason, MarkerManifest, MemoryMarkerBackend,
550 StreamCursor, StreamDictEntry,
551 },
552 };
553
554 fn manifest(run_id: &str) -> MarkerManifest {
555 MarkerManifest {
556 run_id: run_id.to_string(),
557 enabled_classes: vec![DataClass::Quote, DataClass::Trade],
558 high_fidelity: true,
559 snapshot_count: 0,
560 hifi_count: 0,
561 gap_count: 0,
562 dict_count: 0,
563 status: RunStatus::Running,
564 }
565 }
566
567 fn snapshot(marker_seq: u64) -> DataCursorSnapshot {
568 DataCursorSnapshot {
569 marker_seq,
570 event_seq_before: marker_seq.saturating_sub(1),
571 ts_init: UnixNanos::from(1_700_000_000_000_000_000 + marker_seq),
572 advanced: vec![StreamCursor {
573 slot: 0,
574 ts_init_hi: UnixNanos::from(1_700_000_000_000_000_000 + marker_seq),
575 count: marker_seq,
576 }],
577 }
578 }
579
580 fn hifi(marker_seq: u64) -> HiFiMarker {
581 HiFiMarker {
582 marker_seq,
583 event_seq_before: marker_seq.saturating_sub(1),
584 slot: 0,
585 ts_event: UnixNanos::from(1_700_000_000_000_000_100 + marker_seq),
586 ts_init: UnixNanos::from(1_700_000_000_000_000_200 + marker_seq),
587 same_ts_ordinal: 0,
588 record_fingerprint: [7u8; 32],
589 }
590 }
591
592 #[derive(Debug)]
593 struct BlockingMarkerBackend {
594 inner: Arc<Mutex<MemoryMarkerBackend>>,
595 gate: Arc<(Mutex<bool>, parking_lot::Condvar)>,
596 appends_seen: Arc<AtomicUsize>,
597 }
598
599 impl BlockingMarkerBackend {
600 fn new(
601 inner: Arc<Mutex<MemoryMarkerBackend>>,
602 gate: Arc<(Mutex<bool>, parking_lot::Condvar)>,
603 appends_seen: Arc<AtomicUsize>,
604 ) -> Self {
605 Self {
606 inner,
607 gate,
608 appends_seen,
609 }
610 }
611
612 fn wait_for_release(&self) {
613 let (lock, cvar) = &*self.gate;
614 let mut released = lock.lock();
615
616 while !*released {
617 cvar.wait(&mut released);
618 }
619 }
620 }
621
622 impl MarkerBackend for BlockingMarkerBackend {
623 fn open_run(&mut self, _: MarkerManifest) -> Result<(), EventStoreError> {
624 unreachable!("test wrapper does not forward open_run")
625 }
626
627 fn append_snapshot(
628 &mut self,
629 snapshot: &DataCursorSnapshot,
630 hash: [u8; 32],
631 ) -> Result<(), EventStoreError> {
632 self.appends_seen.fetch_add(1, Ordering::SeqCst);
633 self.wait_for_release();
634 self.inner.lock().append_snapshot(snapshot, hash)
635 }
636
637 fn append_hifi(
638 &mut self,
639 marker: &HiFiMarker,
640 hash: [u8; 32],
641 ) -> Result<(), EventStoreError> {
642 self.appends_seen.fetch_add(1, Ordering::SeqCst);
643 self.wait_for_release();
644 self.inner.lock().append_hifi(marker, hash)
645 }
646
647 fn append_gap(&mut self, gap: &MarkerGap, hash: [u8; 32]) -> Result<(), EventStoreError> {
648 self.appends_seen.fetch_add(1, Ordering::SeqCst);
649 self.wait_for_release();
650 self.inner.lock().append_gap(gap, hash)
651 }
652
653 fn put_dict(
654 &mut self,
655 entry: &StreamDictEntry,
656 hash: [u8; 32],
657 ) -> Result<(), EventStoreError> {
658 self.inner.lock().put_dict(entry, hash)
659 }
660
661 fn scan_snapshots(&self) -> Result<Vec<DataCursorSnapshot>, EventStoreError> {
662 self.inner.lock().scan_snapshots()
663 }
664
665 fn scan_hifi(&self) -> Result<Vec<HiFiMarker>, EventStoreError> {
666 self.inner.lock().scan_hifi()
667 }
668
669 fn scan_gaps(&self) -> Result<Vec<MarkerGap>, EventStoreError> {
670 self.inner.lock().scan_gaps()
671 }
672
673 fn scan_dict(&self) -> Result<Vec<StreamDictEntry>, EventStoreError> {
674 self.inner.lock().scan_dict()
675 }
676
677 fn seal(&mut self, status: RunStatus) -> Result<(), EventStoreError> {
678 self.inner.lock().seal(status)
679 }
680
681 fn manifest(&self) -> Result<MarkerManifest, EventStoreError> {
682 self.inner.lock().manifest()
683 }
684 }
685
686 fn wait_until(mut predicate: impl FnMut() -> bool, label: &str) {
687 let start = Instant::now();
688
689 while !predicate() {
690 assert!(
691 start.elapsed() < Duration::from_millis(500),
692 "timed out waiting for {label}"
693 );
694 std::thread::sleep(Duration::from_millis(5));
695 }
696 }
697
698 #[rstest]
699 fn submitted_snapshots_reach_backend() {
700 let (wrapper, shared) = SharedMemoryMarker::new();
701 shared
702 .lock()
703 .open_run(manifest("run-snapshots"))
704 .expect("open marker run");
705 let config = MarkerWriterConfig {
706 channel_capacity: 16,
707 max_batch: 2,
708 max_latency: Duration::from_secs(30),
709 };
710
711 let writer = MarkerWriter::spawn(Box::new(wrapper), get_atomic_clock_static(), config)
712 .expect("spawn marker writer");
713 let s1 = snapshot(1);
714 let s2 = snapshot(2);
715
716 assert!(
717 writer
718 .submit(MarkerMsg::Snapshot(s1.clone()), s1.marker_seq)
719 .expect("submit first")
720 );
721 assert!(
722 writer
723 .submit(MarkerMsg::Snapshot(s2.clone()), s2.marker_seq)
724 .expect("submit second")
725 );
726 writer.close();
727
728 let backend = shared.lock();
729 assert_eq!(
730 backend.scan_snapshots().expect("scan snapshots"),
731 vec![s1, s2]
732 );
733 }
734
735 #[rstest]
736 fn put_dict_waits_for_capacity_and_persists_metadata() {
737 let inner = Arc::new(Mutex::new(MemoryMarkerBackend::new()));
738 inner
739 .lock()
740 .open_run(manifest("run-dict-capacity"))
741 .expect("open marker run");
742 let gate = Arc::new((Mutex::new(false), parking_lot::Condvar::new()));
743 let appends_seen = Arc::new(AtomicUsize::new(0));
744 let backend = BlockingMarkerBackend::new(
745 Arc::clone(&inner),
746 Arc::clone(&gate),
747 Arc::clone(&appends_seen),
748 );
749
750 let writer = MarkerWriter::spawn(
751 Box::new(backend),
752 get_atomic_clock_static(),
753 MarkerWriterConfig {
754 channel_capacity: 1,
755 max_batch: 1,
756 max_latency: Duration::from_secs(30),
757 },
758 )
759 .expect("spawn marker writer");
760
761 let first = snapshot(1);
762 assert!(
763 writer
764 .submit(MarkerMsg::Snapshot(first.clone()), first.marker_seq)
765 .expect("submit first")
766 );
767 wait_until(
768 || appends_seen.load(Ordering::SeqCst) == 1,
769 "writer to block in backend append",
770 );
771
772 let second = snapshot(2);
773 assert!(
774 writer
775 .submit(MarkerMsg::Snapshot(second.clone()), second.marker_seq)
776 .expect("submit second")
777 );
778
779 let entry = StreamDictEntry {
780 slot: 1,
781 data_cls: DataClass::Trade,
782 identifier: "BTCUSDT.BINANCE".to_string(),
783 };
784 let gate_for_release = Arc::clone(&gate);
785
786 let release = std::thread::spawn(move || {
787 std::thread::sleep(Duration::from_millis(20));
788 let (lock, cvar) = &*gate_for_release;
789 *lock.lock() = true;
790 cvar.notify_all();
791 });
792
793 assert!(writer.put_dict(entry.clone()).expect("put dict"));
794 release.join().expect("release gate");
795 writer.close();
796
797 let backend = inner.lock();
798 assert_eq!(backend.scan_dict().expect("scan dict"), vec![entry]);
799 assert_eq!(
800 backend
801 .scan_snapshots()
802 .expect("scan snapshots")
803 .into_iter()
804 .map(|snapshot| snapshot.marker_seq)
805 .collect::<Vec<_>>(),
806 vec![1, 2]
807 );
808 }
809
810 #[rstest]
811 fn spawn_requires_open_marker_run() {
812 let backend = MemoryMarkerBackend::new();
813
814 let err = MarkerWriter::spawn(
815 Box::new(backend),
816 get_atomic_clock_static(),
817 MarkerWriterConfig::default(),
818 )
819 .expect_err("spawn must reject a backend with no open run");
820
821 match err {
822 EventStoreError::Backend(msg) => {
823 assert!(msg.contains("no run open"), "msg was: {msg}");
824 }
825 other => panic!("expected Backend, was {other:?}"),
826 }
827 }
828
829 #[rstest]
830 fn latency_window_flushes_marker_before_close() {
831 let (wrapper, shared) = SharedMemoryMarker::new();
832 shared
833 .lock()
834 .open_run(manifest("run-latency"))
835 .expect("open marker run");
836
837 let writer = MarkerWriter::spawn(
838 Box::new(wrapper),
839 get_atomic_clock_static(),
840 MarkerWriterConfig {
841 channel_capacity: 16,
842 max_batch: 100,
843 max_latency: Duration::from_millis(20),
844 },
845 )
846 .expect("spawn marker writer");
847 let snap = snapshot(1);
848
849 assert!(
850 writer
851 .submit(MarkerMsg::Snapshot(snap.clone()), snap.marker_seq)
852 .expect("submit")
853 );
854 wait_until(
855 || shared.lock().scan_snapshots().expect("scan").len() == 1,
856 "latency flush",
857 );
858 writer.close();
859
860 assert_eq!(
861 shared.lock().scan_snapshots().expect("scan snapshots"),
862 vec![snap]
863 );
864 }
865
866 #[rstest]
867 fn overflow_drops_and_records_single_gap() {
868 let inner = Arc::new(Mutex::new(MemoryMarkerBackend::new()));
869 inner
870 .lock()
871 .open_run(manifest("run-overflow"))
872 .expect("open marker run");
873 let gate = Arc::new((Mutex::new(false), parking_lot::Condvar::new()));
874 let appends_seen = Arc::new(AtomicUsize::new(0));
875 let backend = BlockingMarkerBackend::new(
876 Arc::clone(&inner),
877 Arc::clone(&gate),
878 Arc::clone(&appends_seen),
879 );
880
881 let writer = MarkerWriter::spawn(
882 Box::new(backend),
883 get_atomic_clock_static(),
884 MarkerWriterConfig {
885 channel_capacity: 1,
886 max_batch: 1,
887 max_latency: Duration::from_secs(30),
888 },
889 )
890 .expect("spawn marker writer");
891
892 let first = snapshot(1);
893 assert!(
894 writer
895 .submit(MarkerMsg::Snapshot(first.clone()), first.marker_seq)
896 .expect("submit first")
897 );
898 wait_until(
899 || appends_seen.load(Ordering::SeqCst) == 1,
900 "writer to block in backend append",
901 );
902
903 let second = snapshot(2);
904 assert!(
905 writer
906 .submit(MarkerMsg::Snapshot(second.clone()), second.marker_seq)
907 .expect("submit second")
908 );
909
910 let start = Instant::now();
911
912 for marker_seq in 3..=5 {
913 let dropped = snapshot(marker_seq);
914 assert!(
915 !writer
916 .submit(MarkerMsg::Snapshot(dropped), marker_seq)
917 .expect("submit drop")
918 );
919 }
920 assert!(
921 start.elapsed() < Duration::from_millis(100),
922 "overflow submits must not block"
923 );
924
925 let (lock, cvar) = &*gate;
926 *lock.lock() = true;
927 cvar.notify_all();
928
929 wait_until(
930 || inner.lock().scan_snapshots().expect("scan").len() == 2,
931 "first two snapshots to drain",
932 );
933
934 let sixth = snapshot(6);
935 assert!(
936 writer
937 .submit(MarkerMsg::Snapshot(sixth.clone()), sixth.marker_seq)
938 .expect("submit after overflow")
939 );
940 writer.close();
941
942 let backend = inner.lock();
943 let gaps = backend.scan_gaps().expect("scan gaps");
944 assert_eq!(
945 gaps,
946 vec![MarkerGap {
947 from_marker_seq: 3,
948 to_marker_seq: 5,
949 reason: MarkerGapReason::Overflow,
950 }]
951 );
952 assert_eq!(
953 backend
954 .scan_snapshots()
955 .expect("scan snapshots")
956 .into_iter()
957 .map(|snapshot| snapshot.marker_seq)
958 .collect::<Vec<_>>(),
959 vec![1, 2, 6]
960 );
961 }
962
963 #[rstest]
964 fn overflow_gap_precedes_next_hifi_marker() {
965 let inner = Arc::new(Mutex::new(MemoryMarkerBackend::new()));
966 inner
967 .lock()
968 .open_run(manifest("run-overflow-hifi"))
969 .expect("open marker run");
970 let gate = Arc::new((Mutex::new(false), parking_lot::Condvar::new()));
971 let appends_seen = Arc::new(AtomicUsize::new(0));
972 let backend = BlockingMarkerBackend::new(
973 Arc::clone(&inner),
974 Arc::clone(&gate),
975 Arc::clone(&appends_seen),
976 );
977
978 let writer = MarkerWriter::spawn(
979 Box::new(backend),
980 get_atomic_clock_static(),
981 MarkerWriterConfig {
982 channel_capacity: 1,
983 max_batch: 1,
984 max_latency: Duration::from_secs(30),
985 },
986 )
987 .expect("spawn marker writer");
988
989 let first = snapshot(1);
990 assert!(
991 writer
992 .submit(MarkerMsg::Snapshot(first.clone()), first.marker_seq)
993 .expect("submit first")
994 );
995 wait_until(
996 || appends_seen.load(Ordering::SeqCst) == 1,
997 "writer to block in backend append",
998 );
999
1000 let second = snapshot(2);
1001 assert!(
1002 writer
1003 .submit(MarkerMsg::Snapshot(second.clone()), second.marker_seq)
1004 .expect("submit second")
1005 );
1006
1007 for marker_seq in 3..=5 {
1008 let dropped = snapshot(marker_seq);
1009 assert!(
1010 !writer
1011 .submit(MarkerMsg::Snapshot(dropped), marker_seq)
1012 .expect("submit drop")
1013 );
1014 }
1015
1016 let (lock, cvar) = &*gate;
1017 *lock.lock() = true;
1018 cvar.notify_all();
1019
1020 wait_until(
1021 || inner.lock().scan_snapshots().expect("scan").len() == 2,
1022 "first two snapshots to drain",
1023 );
1024
1025 let marker = hifi(6);
1026 assert!(
1027 writer
1028 .submit(MarkerMsg::HiFi(marker.clone()), marker.marker_seq)
1029 .expect("submit hifi after overflow")
1030 );
1031 writer.close();
1032
1033 let backend = inner.lock();
1034 assert_eq!(
1035 backend.scan_gaps().expect("scan gaps"),
1036 vec![MarkerGap {
1037 from_marker_seq: 3,
1038 to_marker_seq: 5,
1039 reason: MarkerGapReason::Overflow,
1040 }]
1041 );
1042 assert_eq!(backend.scan_hifi().expect("scan hifi"), vec![marker]);
1043 }
1044
1045 #[rstest]
1046 fn close_records_pending_drop_as_writer_closed_gap() {
1047 let inner = Arc::new(Mutex::new(MemoryMarkerBackend::new()));
1048 inner
1049 .lock()
1050 .open_run(manifest("run-close-gap"))
1051 .expect("open marker run");
1052 let gate = Arc::new((Mutex::new(false), parking_lot::Condvar::new()));
1053 let appends_seen = Arc::new(AtomicUsize::new(0));
1054 let backend = BlockingMarkerBackend::new(
1055 Arc::clone(&inner),
1056 Arc::clone(&gate),
1057 Arc::clone(&appends_seen),
1058 );
1059
1060 let writer = MarkerWriter::spawn(
1061 Box::new(backend),
1062 get_atomic_clock_static(),
1063 MarkerWriterConfig {
1064 channel_capacity: 1,
1065 max_batch: 1,
1066 max_latency: Duration::from_secs(30),
1067 },
1068 )
1069 .expect("spawn marker writer");
1070
1071 let first = snapshot(1);
1072 assert!(
1073 writer
1074 .submit(MarkerMsg::Snapshot(first.clone()), first.marker_seq)
1075 .expect("submit first")
1076 );
1077 wait_until(
1078 || appends_seen.load(Ordering::SeqCst) == 1,
1079 "writer to block in backend append",
1080 );
1081
1082 let second = snapshot(2);
1083 assert!(
1084 writer
1085 .submit(MarkerMsg::Snapshot(second.clone()), second.marker_seq)
1086 .expect("submit second")
1087 );
1088
1089 for marker_seq in 3..=4 {
1090 let dropped = snapshot(marker_seq);
1091 assert!(
1092 !writer
1093 .submit(MarkerMsg::Snapshot(dropped), marker_seq)
1094 .expect("submit drop")
1095 );
1096 }
1097
1098 let close_thread = std::thread::spawn(move || writer.close());
1099 std::thread::sleep(Duration::from_millis(20));
1100 let (lock, cvar) = &*gate;
1101 *lock.lock() = true;
1102 cvar.notify_all();
1103 close_thread.join().expect("close thread");
1104
1105 let backend = inner.lock();
1106 assert_eq!(
1107 backend.scan_gaps().expect("scan gaps"),
1108 vec![MarkerGap {
1109 from_marker_seq: 3,
1110 to_marker_seq: 4,
1111 reason: MarkerGapReason::WriterClosed,
1112 }]
1113 );
1114 assert_eq!(
1115 backend
1116 .scan_snapshots()
1117 .expect("scan snapshots")
1118 .into_iter()
1119 .map(|snapshot| snapshot.marker_seq)
1120 .collect::<Vec<_>>(),
1121 vec![1, 2]
1122 );
1123 assert_eq!(
1124 backend.manifest().expect("manifest").status,
1125 RunStatus::Ended
1126 );
1127 }
1128
1129 #[rstest]
1130 fn drop_records_pending_drop_as_writer_closed_gap() {
1131 let inner = Arc::new(Mutex::new(MemoryMarkerBackend::new()));
1135 inner
1136 .lock()
1137 .open_run(manifest("run-drop-gap"))
1138 .expect("open marker run");
1139 let gate = Arc::new((Mutex::new(false), parking_lot::Condvar::new()));
1140 let appends_seen = Arc::new(AtomicUsize::new(0));
1141 let backend = BlockingMarkerBackend::new(
1142 Arc::clone(&inner),
1143 Arc::clone(&gate),
1144 Arc::clone(&appends_seen),
1145 );
1146
1147 let writer = MarkerWriter::spawn(
1148 Box::new(backend),
1149 get_atomic_clock_static(),
1150 MarkerWriterConfig {
1151 channel_capacity: 1,
1152 max_batch: 1,
1153 max_latency: Duration::from_secs(30),
1154 },
1155 )
1156 .expect("spawn marker writer");
1157
1158 let first = snapshot(1);
1159 assert!(
1160 writer
1161 .submit(MarkerMsg::Snapshot(first.clone()), first.marker_seq)
1162 .expect("submit first")
1163 );
1164 wait_until(
1165 || appends_seen.load(Ordering::SeqCst) == 1,
1166 "writer to block in backend append",
1167 );
1168
1169 let second = snapshot(2);
1170 assert!(
1171 writer
1172 .submit(MarkerMsg::Snapshot(second.clone()), second.marker_seq)
1173 .expect("submit second")
1174 );
1175
1176 let dropped = snapshot(3);
1177 assert!(
1178 !writer
1179 .submit(MarkerMsg::Snapshot(dropped), 3)
1180 .expect("submit drop")
1181 );
1182
1183 let drop_thread = std::thread::spawn(move || drop(writer));
1184 std::thread::sleep(Duration::from_millis(20));
1185 let (lock, cvar) = &*gate;
1186 *lock.lock() = true;
1187 cvar.notify_all();
1188 drop_thread.join().expect("drop thread");
1189
1190 let backend = inner.lock();
1191 assert_eq!(
1192 backend.scan_gaps().expect("scan gaps"),
1193 vec![MarkerGap {
1194 from_marker_seq: 3,
1195 to_marker_seq: 3,
1196 reason: MarkerGapReason::WriterClosed,
1197 }]
1198 );
1199 assert_eq!(
1200 backend
1201 .scan_snapshots()
1202 .expect("scan snapshots")
1203 .into_iter()
1204 .map(|snapshot| snapshot.marker_seq)
1205 .collect::<Vec<_>>(),
1206 vec![1, 2]
1207 );
1208 assert_eq!(
1209 backend.manifest().expect("manifest").status,
1210 RunStatus::Ended
1211 );
1212 }
1213
1214 #[rstest]
1215 fn close_drains_and_seals() {
1216 let (wrapper, shared) = SharedMemoryMarker::new();
1217 shared
1218 .lock()
1219 .open_run(manifest("run-close"))
1220 .expect("open marker run");
1221 let config = MarkerWriterConfig {
1222 channel_capacity: 16,
1223 max_batch: 10,
1224 max_latency: Duration::from_secs(30),
1225 };
1226
1227 let writer = MarkerWriter::spawn(Box::new(wrapper), get_atomic_clock_static(), config)
1228 .expect("spawn marker writer");
1229 let marker = hifi(1);
1230
1231 assert!(
1232 writer
1233 .submit(MarkerMsg::HiFi(marker.clone()), marker.marker_seq)
1234 .expect("submit hifi")
1235 );
1236 writer.close();
1237
1238 let backend = shared.lock();
1239 assert_eq!(backend.scan_hifi().expect("scan hifi"), vec![marker]);
1240 assert_eq!(
1241 backend.manifest().expect("manifest").status,
1242 RunStatus::Ended
1243 );
1244 }
1245}
1246
1247#[cfg(test)]
1248#[cfg(madsim)]
1249mod madsim_tests {
1250 use nautilus_core::{UnixNanos, time::get_atomic_clock_static};
1251 use rstest::rstest;
1252
1253 use super::{super::test_support::SharedMemoryMarker, *};
1254 use crate::{
1255 manifest::RunStatus,
1256 markers::{DataClass, MarkerManifest, StreamCursor},
1257 };
1258
1259 fn manifest(run_id: &str) -> MarkerManifest {
1260 MarkerManifest {
1261 run_id: run_id.to_string(),
1262 enabled_classes: vec![DataClass::Quote],
1263 high_fidelity: false,
1264 snapshot_count: 0,
1265 hifi_count: 0,
1266 gap_count: 0,
1267 dict_count: 0,
1268 status: RunStatus::Running,
1269 }
1270 }
1271
1272 fn snapshot(marker_seq: u64) -> DataCursorSnapshot {
1273 DataCursorSnapshot {
1274 marker_seq,
1275 event_seq_before: marker_seq.saturating_sub(1),
1276 ts_init: UnixNanos::from(1_700_000_000_000_000_000 + marker_seq),
1277 advanced: vec![StreamCursor {
1278 slot: 0,
1279 ts_init_hi: UnixNanos::from(1_700_000_000_000_000_000 + marker_seq),
1280 count: marker_seq,
1281 }],
1282 }
1283 }
1284
1285 #[rstest]
1286 fn submit_persists_synchronously_and_close_seals() {
1287 let (wrapper, shared) = SharedMemoryMarker::new();
1288 shared
1289 .lock()
1290 .open_run(manifest("run-madsim"))
1291 .expect("open marker run");
1292
1293 let writer = MarkerWriter::spawn(
1294 Box::new(wrapper),
1295 get_atomic_clock_static(),
1296 MarkerWriterConfig::default(),
1297 )
1298 .expect("spawn marker writer");
1299 let snap = snapshot(1);
1300
1301 assert!(
1302 writer
1303 .submit(MarkerMsg::Snapshot(snap.clone()), snap.marker_seq)
1304 .expect("submit")
1305 );
1306 assert_eq!(
1307 shared.lock().scan_snapshots().expect("scan snapshots"),
1308 vec![snap]
1309 );
1310
1311 writer.close();
1312
1313 assert_eq!(
1314 shared.lock().manifest().expect("manifest").status,
1315 RunStatus::Ended
1316 );
1317 }
1318}