1#[cfg(feature = "timezones")]
2use chrono::TimeZone;
3use now::DateTimeNow;
4use polars_arrow::legacy::time_zone::Tz;
5use polars_core::prelude::*;
6use polars_defs::time::duration::Duration;
7use polars_defs::time::group_by::{ClosedWindow, StartBy};
8
9use crate::prelude::*;
10
11pub(crate) fn ensure_t_in_or_in_front_of_window(
25 mut every: Duration,
26 t: i64,
27 tu: TimeUnit,
28 period: Duration,
29 mut start: i64,
30 closed_window: ClosedWindow,
31 tz: Option<&Tz>,
32) -> PolarsResult<Bounds> {
33 every.negative = !every.negative;
34 let mut stop = period.add(tu, start, tz)?;
35
36 while Bounds::new(start, stop).is_past(t, closed_window) {
37 let mut gap = start - t;
38 if matches!(closed_window, ClosedWindow::Right | ClosedWindow::None) {
39 gap += 1;
40 }
41 debug_assert!(gap >= 1);
42
43 let stride = (gap + every.nte_duration(tu) - 1) / every.nte_duration(tu);
45 debug_assert!(stride >= 1);
46 let stride = std::cmp::max(stride, 1);
47
48 start = (every * stride).add(tu, start, tz)?;
49 stop = period.add(tu, start, tz)?;
50 }
51 Ok(Bounds::new_checked(start, stop))
52}
53
54#[derive(Copy, Clone)]
56pub struct Window {
57 pub(crate) every: Duration,
61 pub(crate) period: Duration,
62 pub offset: Duration,
63}
64
65impl Window {
66 pub fn new(every: Duration, period: Duration, offset: Duration) -> Self {
67 debug_assert!(!every.negative);
68 Self {
69 every,
70 period,
71 offset,
72 }
73 }
74
75 pub fn truncate(&self, tu: TimeUnit, t: i64, tz: Option<&Tz>) -> PolarsResult<i64> {
77 self.every.truncate(tu, t, tz)
78 }
79
80 pub fn round(&self, tu: TimeUnit, t: i64, tz: Option<&Tz>) -> PolarsResult<i64> {
82 let t = t + self.every.duration(tu) / 2;
83 self.truncate(tu, t, tz)
84 }
85
86 pub fn get_earliest_bounds(
90 &self,
91 tu: TimeUnit,
92 t: i64,
93 closed_window: ClosedWindow,
94 tz: Option<&Tz>,
95 ) -> PolarsResult<Bounds> {
96 let start = self.truncate(tu, t, tz)?;
97 let start = self.offset.add(tu, start, tz)?;
98 ensure_t_in_or_in_front_of_window(self.every, t, tu, self.period, start, closed_window, tz)
99 }
100
101 pub(crate) fn estimate_overlapping_bounds(&self, tu: TimeUnit, boundary: Bounds) -> usize {
102 (boundary.duration() / self.every.duration(tu)
103 + self.period.duration(tu) / self.every.duration(tu)) as usize
104 }
105
106 pub fn get_overlapping_bounds_iter<'a>(
107 &'a self,
108 boundary: Bounds,
109 closed_window: ClosedWindow,
110 tu: TimeUnit,
111 tz: Option<&'a Tz>,
112 start_by: StartBy,
113 origin: Option<i64>,
114 ) -> PolarsResult<BoundsIter<'a>> {
115 BoundsIter::new(*self, closed_window, boundary, tu, tz, start_by, origin)
116 }
117
118 pub fn first_window_start(
121 &self,
122 t0: i64,
123 closed_window: ClosedWindow,
124 tu: TimeUnit,
125 tz: Option<&Tz>,
126 start_by: StartBy,
127 ) -> PolarsResult<i64> {
128 match start_by {
129 StartBy::DataPoint => Ok(t0),
130 StartBy::WindowBound => Ok(self.get_earliest_bounds(tu, t0, closed_window, tz)?.start),
131 _ => {
132 let dt = tu.timestamp_to_datetime(t0);
135 let (week_start, tz) = match tz {
136 #[cfg(feature = "timezones")]
137 Some(tz) => (
138 tz.from_utc_datetime(&dt).beginning_of_week().naive_utc(),
139 Some(tz),
140 ),
141 _ => (dt.and_utc().beginning_of_week().naive_utc(), None),
142 };
143 let start = tu.datetime_to_timestamp(week_start);
144 let start = Duration::parse(&format!("{}d", start_by.weekday().unwrap()))
145 .add(tu, start, tz)?;
146 let start = self.offset.add(tu, start, tz)?;
147 let bounds = ensure_t_in_or_in_front_of_window(
149 self.every,
150 t0,
151 tu,
152 self.period,
153 start,
154 closed_window,
155 tz,
156 )?;
157 Ok(bounds.start)
158 },
159 }
160 }
161}
162
163pub struct BoundsIter<'a> {
167 window: Window,
168 boundary: Bounds,
170 origin: i64,
171 k: i64,
173 bi: Bounds,
175 tu: TimeUnit,
176 tz: Option<&'a Tz>,
177}
178impl<'a> BoundsIter<'a> {
179 fn new(
180 window: Window,
181 closed_window: ClosedWindow,
182 boundary: Bounds,
183 tu: TimeUnit,
184 tz: Option<&'a Tz>,
185 start_by: StartBy,
186 origin: Option<i64>,
187 ) -> PolarsResult<Self> {
188 let origin = match origin {
189 Some(origin) => origin,
190 None => window.first_window_start(boundary.start, closed_window, tu, tz, start_by)?,
191 };
192 let stop = window.period.add(tu, origin, tz)?;
193 Ok(Self {
194 window,
195 boundary,
196 origin,
197 k: 0,
198 bi: Bounds::new(origin, stop),
199 tu,
200 tz,
201 })
202 }
203
204 fn advance(&mut self, n: i64) {
205 self.k += n;
208 self.bi.start = (self.window.every * self.k)
209 .add(self.tu, self.origin, self.tz)
210 .unwrap();
211 self.bi.stop = self
212 .window
213 .period
214 .add(self.tu, self.bi.start, self.tz)
215 .unwrap();
216 }
217}
218
219impl Iterator for BoundsIter<'_> {
220 type Item = Bounds;
221
222 fn next(&mut self) -> Option<Self::Item> {
223 if self.bi.start < self.boundary.stop {
224 let out = self.bi;
225 self.advance(1);
226 Some(out)
227 } else {
228 None
229 }
230 }
231
232 fn nth(&mut self, n: usize) -> Option<Self::Item> {
233 if self.bi.start < self.boundary.stop {
234 if n > 0 {
235 self.advance(n.try_into().unwrap());
236 }
237 self.next()
238 } else {
239 None
240 }
241 }
242}
243
244impl<'a> BoundsIter<'a> {
245 pub fn get_stride(&self, target: i64) -> usize {
250 let mut stride = 0;
251 if self.bi.start < self.boundary.stop && target > self.bi.start {
252 let gap = target - self.bi.start;
253 let every = self.window.every.nte_duration(self.tu);
254 let period = self.window.period.nte_duration(self.tu);
255 if gap > every + period {
256 stride = ((gap - period) as usize) / (every as usize);
257 }
258 }
259 stride
260 }
261}