Skip to main content

laminar_core/operator/sliding_window/
mod.rs

1//! Sliding (hopping) window assigner.
2
3use super::window::{WindowAssigner, WindowAssignmentError, WindowId, WindowIdVec};
4use std::time::Duration;
5
6/// Lazy sliding-window assignments for one timestamp.
7#[derive(Debug, Clone)]
8pub struct SlidingWindowIter {
9    next_start: i64,
10    size_ms: i64,
11    slide_ms: i64,
12    remaining: usize,
13}
14
15impl Iterator for SlidingWindowIter {
16    type Item = WindowId;
17
18    #[inline]
19    fn next(&mut self) -> Option<Self::Item> {
20        if self.remaining == 0 {
21            return None;
22        }
23
24        let start = self.next_start;
25        self.remaining -= 1;
26        if self.remaining != 0 {
27            self.next_start += self.slide_ms;
28        }
29        Some(WindowId::new(start, start + self.size_ms))
30    }
31
32    fn size_hint(&self) -> (usize, Option<usize>) {
33        (self.remaining, Some(self.remaining))
34    }
35}
36
37impl ExactSizeIterator for SlidingWindowIter {}
38
39/// Sliding window assigner.
40///
41/// Each event is assigned to one or more overlapping windows.
42/// The maximum number of windows per event is `ceil(size / slide)`.
43#[derive(Debug, Clone)]
44pub struct SlidingWindowAssigner {
45    /// Window size in milliseconds
46    size_ms: i64,
47    /// Slide interval in milliseconds
48    slide_ms: i64,
49    /// Maximum windows per event (cached for admission)
50    windows_per_event: usize,
51    /// Offset in milliseconds for timezone-aligned windows
52    offset_ms: i64,
53}
54
55impl SlidingWindowAssigner {
56    /// Creates a new sliding window assigner.
57    ///
58    /// # Panics
59    ///
60    /// Panics if size or slide is zero/negative, or if slide > size.
61    #[must_use]
62    pub fn new(size: Duration, slide: Duration) -> Self {
63        let size_ms = i64::try_from(size.as_millis()).expect("Window size must fit in i64");
64        let slide_ms = i64::try_from(slide.as_millis()).expect("Slide interval must fit in i64");
65
66        assert!(size_ms > 0, "Window size must be positive");
67        assert!(slide_ms > 0, "Slide interval must be positive");
68        assert!(
69            slide_ms <= size_ms,
70            "Slide must not exceed size (use tumbling windows for non-overlapping)"
71        );
72
73        let windows_per_event = usize::try_from(1 + (size_ms - 1) / slide_ms)
74            .expect("Windows per event should fit in usize");
75
76        Self {
77            size_ms,
78            slide_ms,
79            windows_per_event,
80            offset_ms: 0,
81        }
82    }
83
84    /// Creates a new sliding window assigner with sizes in milliseconds.
85    ///
86    /// # Panics
87    ///
88    /// Panics if size or slide is zero/negative, or if slide > size.
89    #[must_use]
90    pub fn from_millis(size_ms: i64, slide_ms: i64) -> Self {
91        assert!(size_ms > 0, "Window size must be positive");
92        assert!(slide_ms > 0, "Slide interval must be positive");
93        assert!(
94            slide_ms <= size_ms,
95            "Slide must not exceed size (use tumbling windows for non-overlapping)"
96        );
97
98        let windows_per_event = usize::try_from(1 + (size_ms - 1) / slide_ms)
99            .expect("Windows per event should fit in usize");
100
101        Self {
102            size_ms,
103            slide_ms,
104            windows_per_event,
105            offset_ms: 0,
106        }
107    }
108
109    /// Set window offset in milliseconds.
110    #[must_use]
111    pub fn with_offset_ms(mut self, offset_ms: i64) -> Self {
112        self.offset_ms = offset_ms;
113        self
114    }
115
116    /// Returns the window size in milliseconds.
117    #[must_use]
118    pub fn size_ms(&self) -> i64 {
119        self.size_ms
120    }
121
122    /// Returns the slide interval in milliseconds.
123    #[must_use]
124    pub fn slide_ms(&self) -> i64 {
125        self.slide_ms
126    }
127
128    /// Returns the maximum number of windows an event can belong to.
129    #[must_use]
130    pub fn windows_per_event(&self) -> usize {
131        self.windows_per_event
132    }
133
134    /// Returns the window offset in milliseconds.
135    #[must_use]
136    pub fn offset_ms(&self) -> i64 {
137        self.offset_ms
138    }
139
140    /// Iterates over containing windows in ascending start-time order.
141    ///
142    /// # Errors
143    ///
144    /// Returns an error when any required boundary is outside the `i64` timestamp range.
145    #[inline]
146    pub fn try_iter_windows(
147        &self,
148        timestamp: i64,
149    ) -> Result<SlidingWindowIter, WindowAssignmentError> {
150        if let Some(adjusted) = timestamp.checked_sub(self.offset_ms) {
151            let quotient = adjusted.div_euclid(self.slide_ms);
152            if let Some(last_start) = quotient
153                .checked_mul(self.slide_ms)
154                .and_then(|start| start.checked_add(self.offset_ms))
155            {
156                let since_last_start = timestamp - last_start;
157                let preceding_windows = (self.size_ms - since_last_start - 1) / self.slide_ms;
158                let size_remainder = self.size_ms % self.slide_ms;
159                let remaining = self.windows_per_event
160                    - usize::from(size_remainder != 0 && since_last_start >= size_remainder);
161                if let (Some(first_start), Some(_)) = (
162                    preceding_windows
163                        .checked_mul(self.slide_ms)
164                        .and_then(|delta| last_start.checked_sub(delta)),
165                    last_start.checked_add(self.size_ms),
166                ) {
167                    return Ok(SlidingWindowIter {
168                        next_start: first_start,
169                        size_ms: self.size_ms,
170                        slide_ms: self.slide_ms,
171                        remaining,
172                    });
173                }
174            }
175        }
176
177        self.iter_windows_wide(timestamp)
178    }
179
180    #[cold]
181    fn iter_windows_wide(
182        &self,
183        timestamp: i64,
184    ) -> Result<SlidingWindowIter, WindowAssignmentError> {
185        let timestamp_wide = i128::from(timestamp);
186        let size_ms = i128::from(self.size_ms);
187        let slide_ms = i128::from(self.slide_ms);
188        let offset_ms = i128::from(self.offset_ms);
189        let adjusted = timestamp_wide - offset_ms;
190        let last_start = adjusted.div_euclid(slide_ms) * slide_ms + offset_ms;
191        let since_last_start = timestamp_wide - last_start;
192        let preceding_windows = (size_ms - since_last_start - 1) / slide_ms;
193        let first_start = last_start - preceding_windows * slide_ms;
194        let last_end = last_start + size_ms;
195        let size_remainder = size_ms % slide_ms;
196        let remaining = self.windows_per_event
197            - usize::from(size_remainder != 0 && since_last_start >= size_remainder);
198
199        let first_start = i64::try_from(first_start).map_err(|_| {
200            WindowAssignmentError::new(timestamp, first_start, first_start + size_ms)
201        })?;
202        i64::try_from(last_end)
203            .map_err(|_| WindowAssignmentError::new(timestamp, last_start, last_end))?;
204
205        Ok(SlidingWindowIter {
206            next_start: first_start,
207            size_ms: self.size_ms,
208            slide_ms: self.slide_ms,
209            remaining,
210        })
211    }
212}
213
214impl WindowAssigner for SlidingWindowAssigner {
215    /// Assigns a timestamp to all overlapping windows.
216    ///
217    /// Returns windows in order from earliest to latest start time.
218    #[inline]
219    fn assign_windows(&self, timestamp: i64) -> WindowIdVec {
220        self.try_iter_windows(timestamp)
221            .expect("sliding window boundaries must fit in i64")
222            .collect()
223    }
224
225    fn max_timestamp(&self, window_end: i64) -> i64 {
226        window_end - 1
227    }
228}
229
230#[cfg(test)]
231mod tests;