1use std::time::{Duration, Instant};
4
5use super::Watermark;
6
7pub trait WatermarkGenerator: Send {
13 fn on_event(&mut self, timestamp: i64) -> Option<Watermark>;
18
19 fn on_periodic(&mut self) -> Option<Watermark>;
23
24 fn current_watermark(&self) -> i64;
26
27 fn advance_watermark(&mut self, timestamp: i64) -> Option<Watermark>;
33
34 fn restore_watermark_for_recovery(&mut self, timestamp: i64);
42
43 fn is_processing_time(&self) -> bool {
49 false
50 }
51}
52
53pub const DEFAULT_MAX_FUTURE_SKEW_MS: i64 = 5 * 60 * 1000;
55
56#[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
67pub struct BoundedOutOfOrdernessGenerator {
70 max_out_of_orderness: i64,
71 current_max_timestamp: i64,
72 current_watermark: i64,
73 max_future_skew_ms: i64,
75}
76
77impl BoundedOutOfOrdernessGenerator {
78 #[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 #[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 #[must_use]
102 #[allow(clippy::cast_possible_truncation)] pub fn from_duration(max_out_of_orderness: Duration) -> Self {
104 Self::new(max_out_of_orderness.as_millis() as i64)
105 }
106
107 #[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 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 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 self.current_max_timestamp = timestamp.saturating_add(self.max_out_of_orderness);
167 }
168}
169
170#[derive(Debug)]
173pub struct AscendingTimestampsGenerator {
174 current_watermark: i64,
175 max_future_skew_ms: i64,
177}
178
179impl Default for AscendingTimestampsGenerator {
180 fn default() -> Self {
181 Self::new()
182 }
183}
184
185impl AscendingTimestampsGenerator {
186 #[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 #[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
246pub 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 #[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 #[must_use]
274 pub fn inner(&self) -> &G {
275 &self.inner
276 }
277
278 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 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
332pub 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 #[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#[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#[derive(Debug)]
452pub struct WatermarkTracker {
453 source_watermarks: Vec<i64>,
455 combined_watermark: i64,
457 idle_sources: Vec<bool>,
459 last_activity: Vec<Instant>,
461 idle_timeout: Vec<Option<Duration>>,
464}
465
466impl WatermarkTracker {
467 #[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 #[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 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 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 self.idle_sources[source_id] = false;
510 self.last_activity[source_id] = Instant::now();
511
512 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 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 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; };
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 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 #[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 #[must_use]
613 pub fn source_watermark(&self, source_id: usize) -> Option<i64> {
614 self.source_watermarks.get(source_id).copied()
615 }
616
617 #[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 #[must_use]
625 pub fn num_sources(&self) -> usize {
626 self.source_watermarks.len()
627 }
628
629 #[must_use]
631 pub fn active_source_count(&self) -> usize {
632 self.idle_sources.iter().filter(|&&idle| !idle).count()
633 }
634
635 fn update_combined(&mut self) -> Option<Watermark> {
637 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 !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
667pub struct SourceProvidedGenerator {
672 source_watermark: i64,
674 fallback: BoundedOutOfOrdernessGenerator,
676 prefer_source: bool,
678}
679
680impl SourceProvidedGenerator {
681 #[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 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 if self.source_watermark > i64::MIN {
717 return None; }
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
746pub struct ProcessingTimeGenerator {
768 current_watermark: i64,
769}
770
771impl ProcessingTimeGenerator {
772 #[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 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 assert!(ProcessingTimeGenerator::new().is_processing_time());
848 assert!(!BoundedOutOfOrdernessGenerator::new(100).is_processing_time());
849 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 gen.on_event(1000);
864
865 let wm = gen.on_event(800);
867 assert_eq!(wm, None);
868 assert_eq!(gen.current_watermark(), 900); }
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 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); 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); 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 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))); }
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 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 let wm = tracker.mark_idle(1);
1048
1049 assert_eq!(wm, Some(Watermark::new(5000)));
1051 }
1052
1053 #[test]
1054 fn check_idle_sources_advances_then_reactivation_is_monotone() {
1055 let mut tracker = WatermarkTracker::with_idle_timeout(2, Duration::ZERO);
1060 tracker.update_source(0, 5_000); tracker.update_source(1, 1_000); assert_eq!(tracker.current_watermark(), Some(Watermark::new(1_000)));
1063
1064 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 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 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 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); }
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 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); 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))); }
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 #[test]
1206 fn test_advance_watermark_bounded_generator() {
1207 let mut gen = BoundedOutOfOrdernessGenerator::new(100);
1208
1209 let wm = gen.advance_watermark(500);
1211 assert_eq!(wm, Some(Watermark::new(500)));
1212 assert_eq!(gen.current_watermark(), 500);
1213
1214 let wm = gen.advance_watermark(800);
1216 assert_eq!(wm, Some(Watermark::new(800)));
1217 assert_eq!(gen.current_watermark(), 800);
1218
1219 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 gen.on_event(1000); gen.advance_watermark(1200);
1234 assert_eq!(gen.current_watermark(), 1200);
1235
1236 let wm = gen.on_event(1250);
1239 assert_eq!(wm, None);
1240 assert_eq!(gen.current_watermark(), 1200);
1241
1242 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 let wm = gen.advance_watermark(300);
1257 assert_eq!(wm, None);
1258 assert_eq!(gen.current_watermark(), 500);
1259
1260 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 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 let wm = gen.advance_watermark(500);
1291 assert_eq!(wm, Some(Watermark::new(500)));
1292 assert_eq!(gen.current_watermark(), 500);
1293
1294 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 let wm = gen.advance_watermark(300);
1309 assert_eq!(wm, None);
1310 }
1311
1312 #[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 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 let wm = gen.advance_watermark(300);
1342 assert_eq!(wm, None);
1343 assert_eq!(gen.current_watermark(), 500);
1344
1345 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 #[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 assert_eq!(gen.on_event(now + 2 * 60 * 60 * 1000), None);
1374 assert_eq!(gen.current_watermark(), i64::MIN);
1375 assert_eq!(gen.on_event(now), Some(Watermark::new(now)));
1377 }
1378}