laminar_core/operator/sliding_window/
mod.rs1use super::window::{WindowAssigner, WindowAssignmentError, WindowId, WindowIdVec};
4use std::time::Duration;
5
6#[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#[derive(Debug, Clone)]
44pub struct SlidingWindowAssigner {
45 size_ms: i64,
47 slide_ms: i64,
49 windows_per_event: usize,
51 offset_ms: i64,
53}
54
55impl SlidingWindowAssigner {
56 #[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 #[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 #[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 #[must_use]
118 pub fn size_ms(&self) -> i64 {
119 self.size_ms
120 }
121
122 #[must_use]
124 pub fn slide_ms(&self) -> i64 {
125 self.slide_ms
126 }
127
128 #[must_use]
130 pub fn windows_per_event(&self) -> usize {
131 self.windows_per_event
132 }
133
134 #[must_use]
136 pub fn offset_ms(&self) -> i64 {
137 self.offset_ms
138 }
139
140 #[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 #[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;