Skip to main content

laminar_core/time/
watermark.rs

1//! Watermark generators and multi-source tracking.
2
3use std::time::{Duration, Instant};
4
5use super::Watermark;
6
7/// Trait for generating watermarks from event timestamps.
8///
9/// Implementations track observed timestamps and produce watermarks that indicate
10/// event-time progress. The watermark is an assertion that no events with timestamps
11/// earlier than the watermark are expected.
12pub trait WatermarkGenerator: Send {
13    /// Process an event timestamp and potentially emit a new watermark.
14    ///
15    /// Called for each event processed. Returns `Some(watermark)` if the watermark
16    /// should advance, `None` otherwise.
17    fn on_event(&mut self, timestamp: i64) -> Option<Watermark>;
18
19    /// Called periodically to emit watermarks based on wall-clock time.
20    ///
21    /// Useful for generating watermarks even when no events are arriving.
22    fn on_periodic(&mut self) -> Option<Watermark>;
23
24    /// Returns the current watermark value without advancing it.
25    fn current_watermark(&self) -> i64;
26
27    /// Advances the watermark to at least the given timestamp from an external source.
28    ///
29    /// Called when the source provides an explicit watermark (e.g., `Source::watermark()`).
30    /// Returns `Some(Watermark)` if the watermark advanced, `None` if the timestamp
31    /// was not higher than the current watermark.
32    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark>;
33
34    /// Replaces the generator's watermark with an exact recovered value.
35    ///
36    /// Unlike [`Self::advance_watermark`], this may move the watermark
37    /// backwards. It is exclusively for restoring a committed checkpoint
38    /// while source intake and computation are fenced. Calling it while the
39    /// pipeline is running can make already-finalized event-time state visible
40    /// again and is therefore incorrect.
41    fn restore_watermark_for_recovery(&mut self, timestamp: i64);
42
43    /// Whether the watermark is processing-time based (wall clock), rather than
44    /// derived from the event-time column. Such a watermark lives in a different
45    /// time domain than the event timestamps, so comparing the two to drop "late"
46    /// rows would discard every event — callers skip source-side late-filtering
47    /// when this is `true`. Defaults to `false` (event-time generators).
48    fn is_processing_time(&self) -> bool {
49        false
50    }
51}
52
53/// Default `max_future_skew_ms`: 5 min.
54pub const DEFAULT_MAX_FUTURE_SKEW_MS: i64 = 5 * 60 * 1000;
55
56/// `true` if `timestamp` is more than `skew_ms` beyond wall-clock now.
57/// `skew_ms <= 0` or an unset system clock disables the guard.
58#[inline]
59fn is_grossly_future(timestamp: i64, skew_ms: i64) -> bool {
60    if skew_ms <= 0 {
61        return false;
62    }
63    let now = super::now_unix_millis();
64    now > 0 && timestamp > now.saturating_add(skew_ms)
65}
66
67/// Watermark = `max_timestamp_seen - max_out_of_orderness`. `on_event`
68/// ignores timestamps far beyond wall clock for advancement.
69pub struct BoundedOutOfOrdernessGenerator {
70    max_out_of_orderness: i64,
71    current_max_timestamp: i64,
72    current_watermark: i64,
73    /// `0` (or negative) disables the guard (unbounded — legacy behaviour).
74    max_future_skew_ms: i64,
75}
76
77impl BoundedOutOfOrdernessGenerator {
78    /// Creates a new generator with the specified maximum out-of-orderness.
79    ///
80    /// # Arguments
81    ///
82    /// * `max_out_of_orderness` - Maximum allowed lateness in milliseconds
83    #[must_use]
84    pub fn new(max_out_of_orderness: i64) -> Self {
85        Self {
86            max_out_of_orderness,
87            current_max_timestamp: i64::MIN,
88            current_watermark: i64::MIN,
89            max_future_skew_ms: DEFAULT_MAX_FUTURE_SKEW_MS,
90        }
91    }
92
93    /// Sets the future-skew ceiling; `<= 0` disables the guard.
94    #[must_use]
95    pub fn with_max_future_skew(mut self, skew_ms: i64) -> Self {
96        self.max_future_skew_ms = skew_ms;
97        self
98    }
99
100    /// Creates a new generator from a `Duration`.
101    #[must_use]
102    #[allow(clippy::cast_possible_truncation)] // Duration.as_millis() fits i64 for practical values
103    pub fn from_duration(max_out_of_orderness: Duration) -> Self {
104        Self::new(max_out_of_orderness.as_millis() as i64)
105    }
106
107    /// Returns the maximum out-of-orderness in milliseconds.
108    #[must_use]
109    pub fn max_out_of_orderness(&self) -> i64 {
110        self.max_out_of_orderness
111    }
112}
113
114impl WatermarkGenerator for BoundedOutOfOrdernessGenerator {
115    #[inline]
116    fn on_event(&mut self, timestamp: i64) -> Option<Watermark> {
117        if is_grossly_future(timestamp, self.max_future_skew_ms) {
118            return None;
119        }
120        if timestamp > self.current_max_timestamp {
121            self.current_max_timestamp = timestamp;
122            let new_watermark = timestamp.saturating_sub(self.max_out_of_orderness);
123            if new_watermark > self.current_watermark {
124                self.current_watermark = new_watermark;
125                return Some(Watermark::new(new_watermark));
126            }
127        }
128        None
129    }
130
131    #[inline]
132    fn on_periodic(&mut self) -> Option<Watermark> {
133        // Bounded out-of-orderness doesn't emit periodic watermarks
134        None
135    }
136
137    #[inline]
138    fn current_watermark(&self) -> i64 {
139        self.current_watermark
140    }
141
142    #[inline]
143    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
144        if is_grossly_future(timestamp, self.max_future_skew_ms) {
145            return None;
146        }
147        if timestamp > self.current_watermark {
148            self.current_watermark = timestamp;
149            // Maintain invariant: current_max_timestamp >= current_watermark + max_out_of_orderness
150            let min_max = timestamp.saturating_add(self.max_out_of_orderness);
151            if min_max > self.current_max_timestamp {
152                self.current_max_timestamp = min_max;
153            }
154            Some(Watermark::new(timestamp))
155        } else {
156            None
157        }
158    }
159
160    #[inline]
161    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
162        self.current_watermark = timestamp;
163        // `on_event` only considers timestamps above this baseline. Restoring
164        // it along with the watermark prevents pre-recovery observations from
165        // suppressing valid replay after a backwards restore.
166        self.current_max_timestamp = timestamp.saturating_add(self.max_out_of_orderness);
167    }
168}
169
170/// Watermark generator for strictly ascending timestamps; the watermark
171/// equals the current timestamp.
172#[derive(Debug)]
173pub struct AscendingTimestampsGenerator {
174    current_watermark: i64,
175    /// `0` ⇒ disabled.
176    max_future_skew_ms: i64,
177}
178
179impl Default for AscendingTimestampsGenerator {
180    fn default() -> Self {
181        Self::new()
182    }
183}
184
185impl AscendingTimestampsGenerator {
186    /// Creates a new ascending timestamps generator.
187    #[must_use]
188    pub fn new() -> Self {
189        Self {
190            current_watermark: i64::MIN,
191            max_future_skew_ms: DEFAULT_MAX_FUTURE_SKEW_MS,
192        }
193    }
194
195    /// Override the future-skew ceiling (`0` disables).
196    #[must_use]
197    pub fn with_max_future_skew(mut self, skew_ms: i64) -> Self {
198        self.max_future_skew_ms = skew_ms;
199        self
200    }
201}
202
203impl WatermarkGenerator for AscendingTimestampsGenerator {
204    #[inline]
205    fn on_event(&mut self, timestamp: i64) -> Option<Watermark> {
206        if is_grossly_future(timestamp, self.max_future_skew_ms) {
207            return None;
208        }
209        if timestamp > self.current_watermark {
210            self.current_watermark = timestamp;
211            Some(Watermark::new(timestamp))
212        } else {
213            None
214        }
215    }
216
217    #[inline]
218    fn on_periodic(&mut self) -> Option<Watermark> {
219        None
220    }
221
222    #[inline]
223    fn current_watermark(&self) -> i64 {
224        self.current_watermark
225    }
226
227    #[inline]
228    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
229        if is_grossly_future(timestamp, self.max_future_skew_ms) {
230            return None;
231        }
232        if timestamp > self.current_watermark {
233            self.current_watermark = timestamp;
234            Some(Watermark::new(timestamp))
235        } else {
236            None
237        }
238    }
239
240    #[inline]
241    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
242        self.current_watermark = timestamp;
243    }
244}
245
246/// Wraps another generator and emits watermarks at fixed wall-clock
247/// intervals so idle sources don't stall time-based windows.
248pub struct PeriodicGenerator<G: WatermarkGenerator> {
249    inner: G,
250    period: Duration,
251    last_emit_time: Instant,
252    last_emitted_watermark: i64,
253}
254
255impl<G: WatermarkGenerator> PeriodicGenerator<G> {
256    /// Creates a new periodic generator wrapping another generator.
257    ///
258    /// # Arguments
259    ///
260    /// * `inner` - The underlying watermark generator
261    /// * `period` - How often to emit watermarks (wall-clock time)
262    #[must_use]
263    pub fn new(inner: G, period: Duration) -> Self {
264        Self {
265            inner,
266            period,
267            last_emit_time: Instant::now(),
268            last_emitted_watermark: i64::MIN,
269        }
270    }
271
272    /// Returns a reference to the inner generator.
273    #[must_use]
274    pub fn inner(&self) -> &G {
275        &self.inner
276    }
277
278    /// Returns a mutable reference to the inner generator.
279    pub fn inner_mut(&mut self) -> &mut G {
280        &mut self.inner
281    }
282}
283
284impl<G: WatermarkGenerator> WatermarkGenerator for PeriodicGenerator<G> {
285    fn on_event(&mut self, timestamp: i64) -> Option<Watermark> {
286        let wm = self.inner.on_event(timestamp);
287        if let Some(ref w) = wm {
288            self.last_emitted_watermark = w.timestamp();
289            self.last_emit_time = Instant::now();
290        }
291        wm
292    }
293
294    fn on_periodic(&mut self) -> Option<Watermark> {
295        // Check if enough wall-clock time has passed
296        if self.last_emit_time.elapsed() >= self.period {
297            let current = self.inner.current_watermark();
298            if current > self.last_emitted_watermark {
299                self.last_emitted_watermark = current;
300                self.last_emit_time = Instant::now();
301                return Some(Watermark::new(current));
302            }
303            self.last_emit_time = Instant::now();
304        }
305        None
306    }
307
308    fn current_watermark(&self) -> i64 {
309        self.inner.current_watermark()
310    }
311
312    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
313        let wm = self.inner.advance_watermark(timestamp);
314        if let Some(ref w) = wm {
315            self.last_emitted_watermark = w.timestamp();
316            self.last_emit_time = Instant::now();
317        }
318        wm
319    }
320
321    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
322        self.inner.restore_watermark_for_recovery(timestamp);
323        self.last_emitted_watermark = timestamp;
324        self.last_emit_time = Instant::now();
325    }
326
327    fn is_processing_time(&self) -> bool {
328        self.inner.is_processing_time()
329    }
330}
331
332/// Punctuated watermark generator that emits based on special events.
333///
334/// Uses a predicate function to identify watermark-carrying events. When the
335/// predicate returns `Some(watermark)`, that watermark is emitted.
336///
337/// # Example
338///
339/// ```rust
340/// use laminar_core::time::{PunctuatedGenerator, WatermarkGenerator, Watermark};
341///
342/// // Emit watermark on every 1000ms boundary
343/// let mut gen = PunctuatedGenerator::new(|ts| {
344///     if ts % 1000 == 0 {
345///         Some(Watermark::new(ts))
346///     } else {
347///         None
348///     }
349/// });
350///
351/// assert_eq!(gen.on_event(999), None);
352/// assert_eq!(gen.on_event(1000), Some(Watermark::new(1000)));
353/// ```
354pub struct PunctuatedGenerator<F>
355where
356    F: Fn(i64) -> Option<Watermark> + Send,
357{
358    predicate: F,
359    current_watermark: i64,
360}
361
362impl<F> PunctuatedGenerator<F>
363where
364    F: Fn(i64) -> Option<Watermark> + Send,
365{
366    /// Creates a new punctuated generator with the given predicate.
367    ///
368    /// # Arguments
369    ///
370    /// * `predicate` - Function that returns `Some(Watermark)` for watermark events
371    #[must_use]
372    pub fn new(predicate: F) -> Self {
373        Self {
374            predicate,
375            current_watermark: i64::MIN,
376        }
377    }
378}
379
380impl<F> WatermarkGenerator for PunctuatedGenerator<F>
381where
382    F: Fn(i64) -> Option<Watermark> + Send,
383{
384    fn on_event(&mut self, timestamp: i64) -> Option<Watermark> {
385        if let Some(wm) = (self.predicate)(timestamp) {
386            if wm.timestamp() > self.current_watermark {
387                self.current_watermark = wm.timestamp();
388                return Some(wm);
389            }
390        }
391        None
392    }
393
394    fn on_periodic(&mut self) -> Option<Watermark> {
395        None
396    }
397
398    fn current_watermark(&self) -> i64 {
399        self.current_watermark
400    }
401
402    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
403        if timestamp > self.current_watermark {
404            self.current_watermark = timestamp;
405            Some(Watermark::new(timestamp))
406        } else {
407            None
408        }
409    }
410
411    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
412        self.current_watermark = timestamp;
413    }
414}
415
416/// A recovered tracker snapshot did not match the tracker's source topology.
417#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
418#[error(
419    "watermark recovery source count mismatch: expected {expected}, got {watermarks} watermarks and {idle_statuses} idle statuses"
420)]
421pub struct WatermarkRestoreError {
422    expected: usize,
423    watermarks: usize,
424    idle_statuses: usize,
425}
426
427/// Tracks watermarks across multiple input sources.
428///
429/// For operators with multiple inputs (e.g., joins, unions), the combined
430/// watermark is the minimum across all sources. This ensures no late events
431/// from any source are missed.
432///
433/// # Example
434///
435/// ```rust
436/// use laminar_core::time::{WatermarkTracker, Watermark};
437///
438/// let mut tracker = WatermarkTracker::new(3); // 3 sources
439///
440/// // Source 0 advances to 1000
441/// let wm = tracker.update_source(0, 1000);
442/// assert_eq!(wm, None); // Other sources still at MIN
443///
444/// // Source 1 advances to 2000
445/// tracker.update_source(1, 2000);
446///
447/// // Source 2 advances to 500
448/// let wm = tracker.update_source(2, 500);
449/// assert_eq!(wm, Some(Watermark::new(500))); // Min of all sources
450/// ```
451#[derive(Debug)]
452pub struct WatermarkTracker {
453    /// Watermark for each source
454    source_watermarks: Vec<i64>,
455    /// Combined minimum watermark
456    combined_watermark: i64,
457    /// Idle status for each source
458    idle_sources: Vec<bool>,
459    /// Last activity time for each source
460    last_activity: Vec<Instant>,
461    /// Per-source idle timeout; `None` ⇒ that source never auto-idles
462    /// (the default — idleness is opt-in, like Flink `withIdleness`).
463    idle_timeout: Vec<Option<Duration>>,
464}
465
466impl WatermarkTracker {
467    /// Creates a tracker with idleness **disabled** for all sources
468    /// (strict correctness; configure per source via
469    /// [`Self::set_idle_timeout`]).
470    #[must_use]
471    pub fn new(num_sources: usize) -> Self {
472        Self {
473            source_watermarks: vec![i64::MIN; num_sources],
474            combined_watermark: i64::MIN,
475            idle_sources: vec![false; num_sources],
476            last_activity: vec![Instant::now(); num_sources],
477            idle_timeout: vec![None; num_sources],
478        }
479    }
480
481    /// Creates a tracker with the same idle timeout for every source.
482    #[must_use]
483    pub fn with_idle_timeout(num_sources: usize, idle_timeout: Duration) -> Self {
484        Self {
485            source_watermarks: vec![i64::MIN; num_sources],
486            combined_watermark: i64::MIN,
487            idle_sources: vec![false; num_sources],
488            last_activity: vec![Instant::now(); num_sources],
489            idle_timeout: vec![Some(idle_timeout); num_sources],
490        }
491    }
492
493    /// Sets (or clears, with `None`) the idle timeout for one source.
494    pub fn set_idle_timeout(&mut self, source_id: usize, timeout: Option<Duration>) {
495        if let Some(slot) = self.idle_timeout.get_mut(source_id) {
496            *slot = timeout;
497        }
498    }
499
500    /// Updates the watermark for a specific source.
501    ///
502    /// Returns `Some(Watermark)` if the combined watermark advances.
503    pub fn update_source(&mut self, source_id: usize, watermark: i64) -> Option<Watermark> {
504        if source_id >= self.source_watermarks.len() {
505            return None;
506        }
507
508        // Mark source as active
509        self.idle_sources[source_id] = false;
510        self.last_activity[source_id] = Instant::now();
511
512        // Update source watermark
513        if watermark > self.source_watermarks[source_id] {
514            self.source_watermarks[source_id] = watermark;
515            self.update_combined()
516        } else {
517            None
518        }
519    }
520
521    /// Marks a source as idle, excluding it from watermark calculation.
522    ///
523    /// Idle sources don't hold back the combined watermark.
524    pub fn mark_idle(&mut self, source_id: usize) -> Option<Watermark> {
525        if source_id >= self.idle_sources.len() {
526            return None;
527        }
528
529        self.idle_sources[source_id] = true;
530        self.update_combined()
531    }
532
533    /// Checks for sources that have been idle longer than the timeout.
534    ///
535    /// Should be called periodically to detect stalled sources.
536    pub fn check_idle_sources(&mut self) -> Option<Watermark> {
537        let mut any_marked = false;
538        for i in 0..self.idle_sources.len() {
539            let Some(timeout) = self.idle_timeout[i] else {
540                continue; // idleness disabled for this source
541            };
542            if !self.idle_sources[i] && self.last_activity[i].elapsed() >= timeout {
543                self.idle_sources[i] = true;
544                any_marked = true;
545            }
546        }
547        if any_marked {
548            self.update_combined()
549        } else {
550            None
551        }
552    }
553
554    /// Replaces all watermark state from a committed recovery snapshot.
555    ///
556    /// This is the only tracker operation allowed to lower source or combined
557    /// watermarks. The caller must hold the intake/compute recovery fence so no
558    /// events, idle transitions, or watermark reads race with the replacement.
559    /// Configured idle timeouts are preserved and activity timers restart at
560    /// the instant of restoration.
561    ///
562    /// `None` represents a source or combined frontier that has not been
563    /// initialized. `combined_watermark` is installed exactly rather than
564    /// recomputed because a committed cluster frontier may intentionally lag
565    /// this tracker's local source frontiers.
566    ///
567    /// # Errors
568    ///
569    /// Returns [`WatermarkRestoreError`] when either input does not contain
570    /// exactly one entry for every tracked source. No state is changed on
571    /// error.
572    pub fn restore_for_recovery(
573        &mut self,
574        source_watermarks: &[Option<i64>],
575        idle_sources: &[bool],
576        combined_watermark: Option<i64>,
577    ) -> Result<(), WatermarkRestoreError> {
578        let expected = self.source_watermarks.len();
579        if source_watermarks.len() != expected || idle_sources.len() != expected {
580            return Err(WatermarkRestoreError {
581                expected,
582                watermarks: source_watermarks.len(),
583                idle_statuses: idle_sources.len(),
584            });
585        }
586
587        for (target, recovered) in self
588            .source_watermarks
589            .iter_mut()
590            .zip(source_watermarks.iter())
591        {
592            *target = recovered.unwrap_or(i64::MIN);
593        }
594        self.idle_sources.copy_from_slice(idle_sources);
595        self.combined_watermark = combined_watermark.unwrap_or(i64::MIN);
596        let restored_at = Instant::now();
597        self.last_activity.fill(restored_at);
598        Ok(())
599    }
600
601    /// Returns the current combined watermark.
602    #[must_use]
603    pub fn current_watermark(&self) -> Option<Watermark> {
604        if self.combined_watermark == i64::MIN {
605            None
606        } else {
607            Some(Watermark::new(self.combined_watermark))
608        }
609    }
610
611    /// Returns the watermark for a specific source.
612    #[must_use]
613    pub fn source_watermark(&self, source_id: usize) -> Option<i64> {
614        self.source_watermarks.get(source_id).copied()
615    }
616
617    /// Returns whether a source is marked as idle.
618    #[must_use]
619    pub fn is_idle(&self, source_id: usize) -> bool {
620        self.idle_sources.get(source_id).copied().unwrap_or(false)
621    }
622
623    /// Returns the number of sources being tracked.
624    #[must_use]
625    pub fn num_sources(&self) -> usize {
626        self.source_watermarks.len()
627    }
628
629    /// Returns the number of active (non-idle) sources.
630    #[must_use]
631    pub fn active_source_count(&self) -> usize {
632        self.idle_sources.iter().filter(|&&idle| !idle).count()
633    }
634
635    /// Updates the combined watermark based on all active sources.
636    fn update_combined(&mut self) -> Option<Watermark> {
637        // Calculate minimum across active sources only
638        let mut min_watermark = i64::MAX;
639        let mut has_active = false;
640
641        for (i, &wm) in self.source_watermarks.iter().enumerate() {
642            if !self.idle_sources[i] {
643                has_active = true;
644                min_watermark = min_watermark.min(wm);
645            }
646        }
647
648        // If all sources are idle, use the max watermark
649        if !has_active {
650            min_watermark = self
651                .source_watermarks
652                .iter()
653                .copied()
654                .max()
655                .unwrap_or(i64::MIN);
656        }
657
658        if min_watermark > self.combined_watermark && min_watermark != i64::MAX {
659            self.combined_watermark = min_watermark;
660            Some(Watermark::new(min_watermark))
661        } else {
662            None
663        }
664    }
665}
666
667/// Watermark generator for sources with embedded watermarks.
668///
669/// Some sources provide a durable upstream event-time frontier directly.
670/// This generator tracks both event timestamps and explicit watermarks.
671pub struct SourceProvidedGenerator {
672    /// Last watermark from the source
673    source_watermark: i64,
674    /// Fallback generator for when source doesn't provide watermarks
675    fallback: BoundedOutOfOrdernessGenerator,
676    /// Whether to use source watermarks when available
677    prefer_source: bool,
678}
679
680impl SourceProvidedGenerator {
681    /// Creates a new source-provided generator.
682    ///
683    /// # Arguments
684    ///
685    /// * `fallback_lateness` - Lateness for fallback bounded generator
686    /// * `prefer_source` - If true, source watermarks take precedence
687    #[must_use]
688    pub fn new(fallback_lateness: i64, prefer_source: bool) -> Self {
689        Self {
690            source_watermark: i64::MIN,
691            fallback: BoundedOutOfOrdernessGenerator::new(fallback_lateness),
692            prefer_source,
693        }
694    }
695
696    /// Updates the watermark from the source.
697    ///
698    /// Call this when the source provides an explicit watermark.
699    pub fn on_source_watermark(&mut self, watermark: i64) -> Option<Watermark> {
700        if watermark > self.source_watermark {
701            self.source_watermark = watermark;
702            if self.prefer_source || watermark > self.fallback.current_watermark() {
703                return Some(Watermark::new(watermark));
704            }
705        }
706        None
707    }
708}
709
710impl WatermarkGenerator for SourceProvidedGenerator {
711    fn on_event(&mut self, timestamp: i64) -> Option<Watermark> {
712        let fallback_wm = self.fallback.on_event(timestamp);
713
714        if self.prefer_source {
715            // Only emit if source watermark allows and it's an advancement
716            if self.source_watermark > i64::MIN {
717                return None; // Wait for source watermark
718            }
719        }
720
721        fallback_wm
722    }
723
724    fn on_periodic(&mut self) -> Option<Watermark> {
725        None
726    }
727
728    fn current_watermark(&self) -> i64 {
729        if self.prefer_source && self.source_watermark > i64::MIN {
730            self.source_watermark
731        } else {
732            self.fallback.current_watermark().max(self.source_watermark)
733        }
734    }
735
736    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
737        self.on_source_watermark(timestamp)
738    }
739
740    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
741        self.source_watermark = timestamp;
742        self.fallback.restore_watermark_for_recovery(timestamp);
743    }
744}
745
746/// Processing-time watermark generator.
747///
748/// Ignores event timestamps entirely and advances the watermark based on
749/// wall-clock time. Use with `PROCTIME()` sources where events are processed
750/// in arrival order without event-time semantics.
751///
752/// - [`on_event`](WatermarkGenerator::on_event) returns `None` (event timestamps are ignored)
753/// - [`on_periodic`](WatermarkGenerator::on_periodic) returns `Some(Watermark::new(now_millis()))`
754///
755/// # Example
756///
757/// ```rust
758/// use laminar_core::time::{ProcessingTimeGenerator, WatermarkGenerator};
759///
760/// let mut gen = ProcessingTimeGenerator::new();
761/// // on_event ignores the timestamp
762/// assert_eq!(gen.on_event(1000), None);
763/// // on_periodic returns wall-clock time
764/// let wm = gen.on_periodic();
765/// assert!(wm.is_some());
766/// ```
767pub struct ProcessingTimeGenerator {
768    current_watermark: i64,
769}
770
771impl ProcessingTimeGenerator {
772    /// Creates a new processing-time watermark generator.
773    #[must_use]
774    pub fn new() -> Self {
775        Self {
776            current_watermark: i64::MIN,
777        }
778    }
779}
780
781impl Default for ProcessingTimeGenerator {
782    fn default() -> Self {
783        Self::new()
784    }
785}
786
787impl WatermarkGenerator for ProcessingTimeGenerator {
788    #[inline]
789    fn on_event(&mut self, _timestamp: i64) -> Option<Watermark> {
790        // Processing-time mode ignores event timestamps
791        None
792    }
793
794    #[inline]
795    fn on_periodic(&mut self) -> Option<Watermark> {
796        let now = super::now_unix_millis();
797        if now > self.current_watermark {
798            self.current_watermark = now;
799            Some(Watermark::new(now))
800        } else {
801            None
802        }
803    }
804
805    #[inline]
806    fn current_watermark(&self) -> i64 {
807        self.current_watermark
808    }
809
810    #[inline]
811    fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark> {
812        if timestamp > self.current_watermark {
813            self.current_watermark = timestamp;
814            Some(Watermark::new(timestamp))
815        } else {
816            None
817        }
818    }
819
820    #[inline]
821    fn restore_watermark_for_recovery(&mut self, timestamp: i64) {
822        self.current_watermark = timestamp;
823    }
824
825    #[inline]
826    fn is_processing_time(&self) -> bool {
827        true
828    }
829}
830
831#[cfg(test)]
832mod tests {
833    use super::*;
834
835    #[test]
836    fn test_bounded_generator_first_event() {
837        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
838        let wm = gen.on_event(1000);
839        assert_eq!(wm, Some(Watermark::new(900)));
840        assert_eq!(gen.current_watermark(), 900);
841    }
842
843    #[test]
844    fn processing_time_domain_is_reported_for_late_filter_skip() {
845        // Source-side late-filtering keys off this: a wall-clock watermark must
846        // not be compared against event-time timestamps, or it drops every row.
847        assert!(ProcessingTimeGenerator::new().is_processing_time());
848        assert!(!BoundedOutOfOrdernessGenerator::new(100).is_processing_time());
849        // The periodic wrapper reports its inner generator's time domain.
850        let p = Duration::from_millis(1);
851        assert!(PeriodicGenerator::new(ProcessingTimeGenerator::new(), p).is_processing_time());
852        assert!(
853            !PeriodicGenerator::new(BoundedOutOfOrdernessGenerator::new(100), p)
854                .is_processing_time()
855        );
856    }
857
858    #[test]
859    fn test_bounded_generator_out_of_order() {
860        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
861
862        // First event
863        gen.on_event(1000);
864
865        // Out of order - should not emit new watermark
866        let wm = gen.on_event(800);
867        assert_eq!(wm, None);
868        assert_eq!(gen.current_watermark(), 900); // Still at 1000 - 100
869    }
870
871    #[test]
872    fn test_bounded_generator_advancement() {
873        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
874
875        gen.on_event(1000);
876        let wm = gen.on_event(1200);
877
878        assert_eq!(wm, Some(Watermark::new(1100)));
879    }
880
881    #[test]
882    fn bounded_recovery_restore_lowers_timestamp_baseline() {
883        let mut gen = BoundedOutOfOrdernessGenerator::new(100).with_max_future_skew(0);
884        assert_eq!(gen.on_event(2_000), Some(Watermark::new(1_900)));
885
886        gen.restore_watermark_for_recovery(500);
887        assert_eq!(gen.current_watermark(), 500);
888        // This event is below the pre-recovery maximum but above the restored
889        // baseline, so it proves stale generator state was not retained.
890        assert_eq!(gen.on_event(650), Some(Watermark::new(550)));
891    }
892
893    #[test]
894    fn test_bounded_generator_from_duration() {
895        let gen = BoundedOutOfOrdernessGenerator::from_duration(Duration::from_secs(5));
896        assert_eq!(gen.max_out_of_orderness(), 5000);
897    }
898
899    #[test]
900    fn test_bounded_generator_no_periodic() {
901        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
902        assert_eq!(gen.on_periodic(), None);
903    }
904
905    #[test]
906    fn test_ascending_generator_advances_on_each_event() {
907        let mut gen = AscendingTimestampsGenerator::new();
908
909        let wm1 = gen.on_event(1000);
910        assert_eq!(wm1, Some(Watermark::new(1000)));
911
912        let wm2 = gen.on_event(2000);
913        assert_eq!(wm2, Some(Watermark::new(2000)));
914    }
915
916    #[test]
917    fn test_ascending_generator_ignores_backwards() {
918        let mut gen = AscendingTimestampsGenerator::new();
919
920        gen.on_event(2000);
921        let wm = gen.on_event(1000); // Earlier timestamp
922
923        assert_eq!(wm, None);
924        assert_eq!(gen.current_watermark(), 2000);
925    }
926
927    #[test]
928    fn ascending_recovery_restore_lowers_then_advances() {
929        let mut gen = AscendingTimestampsGenerator::new().with_max_future_skew(0);
930        gen.on_event(2_000);
931
932        gen.restore_watermark_for_recovery(500);
933        assert_eq!(gen.current_watermark(), 500);
934        assert_eq!(gen.on_event(600), Some(Watermark::new(600)));
935    }
936
937    #[test]
938    fn test_periodic_generator_passes_through() {
939        let inner = BoundedOutOfOrdernessGenerator::new(100);
940        let mut gen = PeriodicGenerator::new(inner, Duration::from_millis(100));
941
942        let wm = gen.on_event(1000);
943        assert_eq!(wm, Some(Watermark::new(900)));
944    }
945
946    #[test]
947    fn test_periodic_generator_inner_access() {
948        let inner = BoundedOutOfOrdernessGenerator::new(100);
949        let gen = PeriodicGenerator::new(inner, Duration::from_millis(100));
950
951        assert_eq!(gen.inner().max_out_of_orderness(), 100);
952    }
953
954    #[test]
955    fn periodic_recovery_restore_resets_inner_and_emission_frontier() {
956        let inner = AscendingTimestampsGenerator::new().with_max_future_skew(0);
957        let mut gen = PeriodicGenerator::new(inner, Duration::from_millis(100));
958        gen.on_event(2_000);
959
960        gen.restore_watermark_for_recovery(500);
961        assert_eq!(gen.current_watermark(), 500);
962        assert_eq!(gen.last_emitted_watermark, 500);
963        assert_eq!(gen.on_event(600), Some(Watermark::new(600)));
964    }
965
966    #[test]
967    fn test_punctuated_generator_predicate() {
968        let mut gen = PunctuatedGenerator::new(|ts| {
969            if ts % 1000 == 0 {
970                Some(Watermark::new(ts))
971            } else {
972                None
973            }
974        });
975
976        assert_eq!(gen.on_event(500), None);
977        assert_eq!(gen.on_event(999), None);
978        assert_eq!(gen.on_event(1000), Some(Watermark::new(1000)));
979        assert_eq!(gen.on_event(1500), None);
980        assert_eq!(gen.on_event(2000), Some(Watermark::new(2000)));
981    }
982
983    #[test]
984    fn test_punctuated_generator_no_regression() {
985        let mut gen = PunctuatedGenerator::new(|ts| Some(Watermark::new(ts)));
986
987        gen.on_event(2000);
988        let wm = gen.on_event(1000); // Lower watermark
989
990        assert_eq!(wm, None);
991        assert_eq!(gen.current_watermark(), 2000);
992    }
993
994    #[test]
995    fn punctuated_recovery_restore_lowers_then_advances() {
996        let mut gen = PunctuatedGenerator::new(|ts| Some(Watermark::new(ts)));
997        gen.on_event(2_000);
998
999        gen.restore_watermark_for_recovery(500);
1000        assert_eq!(gen.current_watermark(), 500);
1001        assert_eq!(gen.on_event(600), Some(Watermark::new(600)));
1002    }
1003
1004    #[test]
1005    fn test_tracker_single_source() {
1006        let mut tracker = WatermarkTracker::new(1);
1007
1008        let wm = tracker.update_source(0, 1000);
1009        assert_eq!(wm, Some(Watermark::new(1000)));
1010        assert_eq!(tracker.current_watermark(), Some(Watermark::new(1000)));
1011    }
1012
1013    #[test]
1014    fn test_tracker_multiple_sources() {
1015        let mut tracker = WatermarkTracker::new(3);
1016
1017        // All sources need to report before watermark advances
1018        tracker.update_source(0, 1000);
1019        tracker.update_source(1, 2000);
1020        let wm = tracker.update_source(2, 500);
1021
1022        assert_eq!(wm, Some(Watermark::new(500))); // Minimum
1023    }
1024
1025    #[test]
1026    fn test_tracker_min_watermark() {
1027        let mut tracker = WatermarkTracker::new(2);
1028
1029        tracker.update_source(0, 5000);
1030        tracker.update_source(1, 3000);
1031
1032        assert_eq!(tracker.current_watermark(), Some(Watermark::new(3000)));
1033
1034        // Source 1 advances
1035        tracker.update_source(1, 4000);
1036        assert_eq!(tracker.current_watermark(), Some(Watermark::new(4000)));
1037    }
1038
1039    #[test]
1040    fn test_tracker_idle_source() {
1041        let mut tracker = WatermarkTracker::new(2);
1042
1043        tracker.update_source(0, 5000);
1044        tracker.update_source(1, 1000);
1045
1046        // Source 1 is slow, mark it idle
1047        let wm = tracker.mark_idle(1);
1048
1049        // Now only source 0's watermark counts
1050        assert_eq!(wm, Some(Watermark::new(5000)));
1051    }
1052
1053    #[test]
1054    fn check_idle_sources_advances_then_reactivation_is_monotone() {
1055        // A quiet watermarked source must not pin the combined-min, and a
1056        // later-reactivating source with an OLD watermark must not regress
1057        // it. `idle_timeout = 0` makes any source immediately eligible for
1058        // `check_idle_sources` without sleeping.
1059        let mut tracker = WatermarkTracker::with_idle_timeout(2, Duration::ZERO);
1060        tracker.update_source(0, 5_000); // active, fast
1061        tracker.update_source(1, 1_000); // will go quiet
1062        assert_eq!(tracker.current_watermark(), Some(Watermark::new(1_000)));
1063
1064        // Source 1 idle past timeout → excluded → combined jumps to s0.
1065        let advanced = tracker.check_idle_sources();
1066        assert_eq!(advanced, Some(Watermark::new(5_000)));
1067        assert_eq!(tracker.current_watermark(), Some(Watermark::new(5_000)));
1068
1069        // Source 1 reactivates with a STALE watermark: must not regress.
1070        let res = tracker.update_source(1, 1_500);
1071        assert_eq!(res, None, "stale reactivation must not emit a regress");
1072        assert_eq!(tracker.current_watermark(), Some(Watermark::new(5_000)));
1073
1074        // Once it catches up past the combined, progress resumes from min.
1075        tracker.update_source(0, 9_000);
1076        tracker.update_source(1, 8_000);
1077        assert_eq!(tracker.current_watermark(), Some(Watermark::new(8_000)));
1078    }
1079
1080    #[test]
1081    fn test_tracker_all_idle() {
1082        let mut tracker = WatermarkTracker::new(2);
1083
1084        tracker.update_source(0, 5000);
1085        tracker.update_source(1, 3000);
1086
1087        tracker.mark_idle(0);
1088        let wm = tracker.mark_idle(1);
1089
1090        // Use max when all idle
1091        assert_eq!(wm, Some(Watermark::new(5000)));
1092    }
1093
1094    #[test]
1095    fn test_tracker_source_watermark() {
1096        let mut tracker = WatermarkTracker::new(2);
1097
1098        tracker.update_source(0, 1000);
1099        tracker.update_source(1, 2000);
1100
1101        assert_eq!(tracker.source_watermark(0), Some(1000));
1102        assert_eq!(tracker.source_watermark(1), Some(2000));
1103        assert_eq!(tracker.source_watermark(5), None); // Out of bounds
1104    }
1105
1106    #[test]
1107    fn test_tracker_active_source_count() {
1108        let mut tracker = WatermarkTracker::new(3);
1109
1110        assert_eq!(tracker.active_source_count(), 3);
1111
1112        tracker.mark_idle(0);
1113        assert_eq!(tracker.active_source_count(), 2);
1114
1115        tracker.mark_idle(2);
1116        assert_eq!(tracker.active_source_count(), 1);
1117
1118        // Reactivate by updating
1119        tracker.update_source(0, 1000);
1120        assert_eq!(tracker.active_source_count(), 2);
1121    }
1122
1123    #[test]
1124    fn test_tracker_invalid_source() {
1125        let mut tracker = WatermarkTracker::new(2);
1126
1127        let wm = tracker.update_source(5, 1000); // Invalid source ID
1128        assert_eq!(wm, None);
1129
1130        let wm = tracker.mark_idle(5);
1131        assert_eq!(wm, None);
1132    }
1133
1134    #[test]
1135    fn tracker_recovery_restore_is_exact_and_runtime_progress_resumes() {
1136        let mut tracker = WatermarkTracker::new(2);
1137        tracker.update_source(0, 9_000);
1138        tracker.update_source(1, 8_000);
1139        assert_eq!(tracker.current_watermark(), Some(Watermark::new(8_000)));
1140
1141        tracker
1142            .restore_for_recovery(&[Some(1_000), None], &[false, true], Some(750))
1143            .unwrap();
1144
1145        assert_eq!(tracker.source_watermark(0), Some(1_000));
1146        assert_eq!(tracker.source_watermark(1), Some(i64::MIN));
1147        assert!(!tracker.is_idle(0));
1148        assert!(tracker.is_idle(1));
1149        assert_eq!(tracker.current_watermark(), Some(Watermark::new(750)));
1150
1151        assert_eq!(tracker.update_source(0, 1_200), Some(Watermark::new(1_200)));
1152        assert_eq!(tracker.current_watermark(), Some(Watermark::new(1_200)));
1153    }
1154
1155    #[test]
1156    fn tracker_recovery_restore_rejects_topology_mismatch_without_mutation() {
1157        let mut tracker = WatermarkTracker::new(2);
1158        tracker.update_source(0, 2_000);
1159        tracker.update_source(1, 1_000);
1160
1161        let error = tracker
1162            .restore_for_recovery(&[Some(500)], &[false, true], Some(500))
1163            .unwrap_err();
1164
1165        assert_eq!(error.expected, 2);
1166        assert_eq!(error.watermarks, 1);
1167        assert_eq!(error.idle_statuses, 2);
1168        assert_eq!(tracker.source_watermark(0), Some(2_000));
1169        assert_eq!(tracker.source_watermark(1), Some(1_000));
1170        assert_eq!(tracker.current_watermark(), Some(Watermark::new(1_000)));
1171    }
1172
1173    #[test]
1174    fn test_source_provided_fallback() {
1175        let mut gen = SourceProvidedGenerator::new(100, false);
1176
1177        let wm = gen.on_event(1000);
1178        assert_eq!(wm, Some(Watermark::new(900))); // Fallback behavior
1179    }
1180
1181    #[test]
1182    fn test_source_provided_explicit_watermark() {
1183        let mut gen = SourceProvidedGenerator::new(100, true);
1184
1185        let wm = gen.on_source_watermark(500);
1186        assert_eq!(wm, Some(Watermark::new(500)));
1187        assert_eq!(gen.current_watermark(), 500);
1188    }
1189
1190    #[test]
1191    fn source_provided_recovery_restore_resets_source_and_fallback() {
1192        let mut gen = SourceProvidedGenerator::new(100, false);
1193        gen.on_source_watermark(2_000);
1194        gen.on_event(2_000);
1195
1196        gen.restore_watermark_for_recovery(500);
1197        assert_eq!(gen.source_watermark, 500);
1198        assert_eq!(gen.fallback.current_watermark(), 500);
1199        assert_eq!(gen.on_event(700), Some(Watermark::new(600)));
1200        assert_eq!(gen.current_watermark(), 600);
1201    }
1202
1203    // --- advance_watermark() tests ---
1204
1205    #[test]
1206    fn test_advance_watermark_bounded_generator() {
1207        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
1208
1209        // Advance from initial state
1210        let wm = gen.advance_watermark(500);
1211        assert_eq!(wm, Some(Watermark::new(500)));
1212        assert_eq!(gen.current_watermark(), 500);
1213
1214        // Advance further
1215        let wm = gen.advance_watermark(800);
1216        assert_eq!(wm, Some(Watermark::new(800)));
1217        assert_eq!(gen.current_watermark(), 800);
1218
1219        // No regression
1220        let wm = gen.advance_watermark(600);
1221        assert_eq!(wm, None);
1222        assert_eq!(gen.current_watermark(), 800);
1223    }
1224
1225    #[test]
1226    fn test_advance_watermark_maintains_invariant() {
1227        let mut gen = BoundedOutOfOrdernessGenerator::new(100);
1228
1229        // Process an event to set initial state
1230        gen.on_event(1000); // wm=900, max_ts=1000
1231
1232        // Advance watermark beyond current
1233        gen.advance_watermark(1200);
1234        assert_eq!(gen.current_watermark(), 1200);
1235
1236        // Now on_event at 1250 should work correctly: max_ts should be >= 1300
1237        // wm = 1250 - 100 = 1150 which is < 1200, so no new watermark from on_event
1238        let wm = gen.on_event(1250);
1239        assert_eq!(wm, None);
1240        assert_eq!(gen.current_watermark(), 1200);
1241
1242        // But event at 1400: max_ts = 1400, wm = 1300 > 1200
1243        let wm = gen.on_event(1400);
1244        assert_eq!(wm, Some(Watermark::new(1300)));
1245    }
1246
1247    #[test]
1248    fn test_advance_watermark_ascending_generator() {
1249        let mut gen = AscendingTimestampsGenerator::new();
1250
1251        let wm = gen.advance_watermark(500);
1252        assert_eq!(wm, Some(Watermark::new(500)));
1253        assert_eq!(gen.current_watermark(), 500);
1254
1255        // No regression
1256        let wm = gen.advance_watermark(300);
1257        assert_eq!(wm, None);
1258        assert_eq!(gen.current_watermark(), 500);
1259
1260        // Further advance
1261        let wm = gen.advance_watermark(1000);
1262        assert_eq!(wm, Some(Watermark::new(1000)));
1263    }
1264
1265    #[test]
1266    fn test_advance_watermark_periodic_generator() {
1267        let inner = BoundedOutOfOrdernessGenerator::new(100);
1268        let mut gen = PeriodicGenerator::new(inner, Duration::from_millis(100));
1269
1270        let wm = gen.advance_watermark(500);
1271        assert_eq!(wm, Some(Watermark::new(500)));
1272        assert_eq!(gen.current_watermark(), 500);
1273
1274        // No regression
1275        let wm = gen.advance_watermark(300);
1276        assert_eq!(wm, None);
1277    }
1278
1279    #[test]
1280    fn test_advance_watermark_punctuated_generator() {
1281        let mut gen = PunctuatedGenerator::new(|ts| {
1282            if ts % 1000 == 0 {
1283                Some(Watermark::new(ts))
1284            } else {
1285                None
1286            }
1287        });
1288
1289        // External advance (does not invoke predicate)
1290        let wm = gen.advance_watermark(500);
1291        assert_eq!(wm, Some(Watermark::new(500)));
1292        assert_eq!(gen.current_watermark(), 500);
1293
1294        // No regression
1295        let wm = gen.advance_watermark(200);
1296        assert_eq!(wm, None);
1297    }
1298
1299    #[test]
1300    fn test_advance_watermark_source_provided_generator() {
1301        let mut gen = SourceProvidedGenerator::new(100, true);
1302
1303        let wm = gen.advance_watermark(500);
1304        assert_eq!(wm, Some(Watermark::new(500)));
1305        assert_eq!(gen.current_watermark(), 500);
1306
1307        // No regression
1308        let wm = gen.advance_watermark(300);
1309        assert_eq!(wm, None);
1310    }
1311
1312    // --- ProcessingTimeGenerator tests ---
1313
1314    #[test]
1315    fn test_processing_time_generator_ignores_events() {
1316        let mut gen = ProcessingTimeGenerator::new();
1317        assert_eq!(gen.on_event(1000), None);
1318        assert_eq!(gen.on_event(2000), None);
1319        assert_eq!(gen.current_watermark(), i64::MIN);
1320    }
1321
1322    #[test]
1323    fn test_processing_time_generator_periodic() {
1324        let mut gen = ProcessingTimeGenerator::new();
1325        let wm = gen.on_periodic();
1326        assert!(wm.is_some());
1327        let ts = wm.unwrap().timestamp();
1328        // Should be a reasonable timestamp (after 2020-01-01)
1329        assert!(ts > 1_577_836_800_000, "timestamp too old: {ts}");
1330    }
1331
1332    #[test]
1333    fn test_processing_time_generator_advance_watermark() {
1334        let mut gen = ProcessingTimeGenerator::new();
1335
1336        let wm = gen.advance_watermark(500);
1337        assert_eq!(wm, Some(Watermark::new(500)));
1338        assert_eq!(gen.current_watermark(), 500);
1339
1340        // No regression
1341        let wm = gen.advance_watermark(300);
1342        assert_eq!(wm, None);
1343        assert_eq!(gen.current_watermark(), 500);
1344
1345        // Further advance
1346        let wm = gen.advance_watermark(1000);
1347        assert_eq!(wm, Some(Watermark::new(1000)));
1348    }
1349
1350    #[test]
1351    fn processing_time_recovery_restore_lowers_then_advances() {
1352        let mut gen = ProcessingTimeGenerator::new();
1353        gen.advance_watermark(2_000);
1354
1355        gen.restore_watermark_for_recovery(500);
1356        assert_eq!(gen.current_watermark(), 500);
1357        assert_eq!(gen.advance_watermark(600), Some(Watermark::new(600)));
1358    }
1359
1360    #[test]
1361    fn test_processing_time_generator_default() {
1362        let gen = ProcessingTimeGenerator::default();
1363        assert_eq!(gen.current_watermark(), i64::MIN);
1364    }
1365
1366    // --- future-skew guard ---
1367
1368    #[test]
1369    fn future_skew_event_does_not_advance_watermark() {
1370        let mut gen = BoundedOutOfOrdernessGenerator::new(0);
1371        let now = crate::time::now_unix_millis();
1372        // A ~2h-future event must not poison the watermark...
1373        assert_eq!(gen.on_event(now + 2 * 60 * 60 * 1000), None);
1374        assert_eq!(gen.current_watermark(), i64::MIN);
1375        // ...but a normal event still advances it.
1376        assert_eq!(gen.on_event(now), Some(Watermark::new(now)));
1377    }
1378}