Skip to main content

nautilus_event_store/markers/
writer.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Non-blocking writer lane for the data marker sidecar.
17
18use std::time::Duration;
19
20use crate::{
21    error::EventStoreError,
22    markers::{DataCursorSnapshot, HiFiMarker, MarkerBackend, MarkerGap, StreamDictEntry},
23};
24
25/// Default channel capacity for markers pending the writer thread.
26pub const DEFAULT_MARKER_CHANNEL_CAPACITY: usize = 10_000;
27/// Default maximum number of markers collected before forcing a flush.
28pub const DEFAULT_MARKER_MAX_BATCH: usize = 100;
29/// Default maximum time a marker batch may accumulate before forcing a flush.
30pub const DEFAULT_MARKER_MAX_LATENCY: Duration = Duration::from_millis(5);
31
32/// Configuration knobs for the data marker writer.
33#[derive(Clone, Debug)]
34pub struct MarkerWriterConfig {
35    /// Capacity of the bounded `sync_channel` between submit and the writer thread.
36    pub channel_capacity: usize,
37    /// Maximum marker messages collected before a flush is forced.
38    pub max_batch: usize,
39    /// Maximum time a marker batch may accumulate before a flush is forced.
40    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/// A message sent to the data marker writer.
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub enum MarkerMsg {
56    /// Cursor snapshot marker.
57    Snapshot(DataCursorSnapshot),
58    /// High-fidelity per-record marker.
59    HiFi(HiFiMarker),
60    /// Stream dictionary entry.
61    Dict(StreamDictEntry),
62    /// Closes and seals the marker run after draining older messages.
63    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    /// Dedicated marker writer thread.
97    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        /// Spawns the marker writer thread and takes ownership of `backend`.
153        ///
154        /// # Errors
155        ///
156        /// Returns [`EventStoreError::Backend`] when the backend has no open run or when the
157        /// writer thread cannot be spawned.
158        pub fn spawn(
159            backend: Box<dyn MarkerBackend + Send>,
160            _clock: &'static AtomicTime,
161            config: MarkerWriterConfig,
162        ) -> Result<Self, EventStoreError> {
163            backend.manifest()?;
164            // A zero capacity would create a rendezvous channel: try_send only succeeds
165            // while the writer thread is parked in recv, so most submits would overflow
166            // and put_dict would block the engine thread on every new stream.
167            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        /// Submits a marker without blocking the caller.
187        ///
188        /// Returns `Ok(false)` when the bounded channel is full; the dropped marker
189        /// sequence is recorded into the next overflow gap.
190        ///
191        /// # Errors
192        ///
193        /// Returns [`EventStoreError::Closed`] when the writer is closed, and
194        /// [`EventStoreError::Backend`] when `msg` is not a `Snapshot` or `HiFi` marker.
195        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        /// Submits a stream dictionary entry to the writer.
238        ///
239        /// Dictionary entries are one-time stream metadata, so this path waits for channel
240        /// capacity instead of dropping and gap-accounting them.
241        ///
242        /// # Errors
243        ///
244        /// Returns [`EventStoreError::Closed`] when the writer is closed.
245        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        /// Drains pending markers and seals the marker run.
260        pub fn close(mut self) {
261            self.close_inner();
262        }
263
264        // Shared by close() and Drop: an implicit drop must still convert the pending
265        // overflow range into its WriterClosed gap record and drain buffered markers,
266        // or the verifier sees an unexplained marker_seq hole.
267        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    /// Synchronous marker writer used under simulation.
422    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        /// Constructs a synchronous marker writer over `backend`.
445        ///
446        /// # Errors
447        ///
448        /// Returns [`EventStoreError::Backend`] when the backend has no open run.
449        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        /// Commits the marker synchronously.
465        ///
466        /// # Errors
467        ///
468        /// Returns [`EventStoreError::Closed`] when the writer is closed, and
469        /// [`EventStoreError::Backend`] when `msg` is not a `Snapshot` or `HiFi` marker.
470        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        /// Commits the stream dictionary entry synchronously.
496        ///
497        /// # Errors
498        ///
499        /// Returns [`EventStoreError::Closed`] when the writer is closed.
500        #[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        /// Seals the marker run.
516        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        // An implicit drop (panic unwind, owner teardown that skips close) must still
1132        // convert the pending overflow range into its WriterClosed gap record; losing
1133        // it leaves the verifier an unexplained marker_seq hole.
1134        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}