Skip to main content

polars_time/windows/
window.rs

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
11/// Ensure that earliest datapoint (`t`) is in, or in front of, first window.
12///
13/// For example, if we have:
14///
15/// - first datapoint is `2020-01-01 01:00`
16/// - `every` is `'1d'`
17/// - `period` is `'2d'`
18/// - `offset` is `'6h'`
19///
20/// then truncating the earliest datapoint by `every` and adding `offset` results
21/// in the window `[2020-01-01 06:00, 2020-01-03 06:00)`. To give the earliest datapoint
22/// a chance of being included, we then shift the window back by `every` to
23/// `[2019-12-31 06:00, 2020-01-02 06:00)`.
24pub(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        // Ceil division
44        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/// Represents a window in time
55#[derive(Copy, Clone)]
56pub struct Window {
57    // The ith window start is expressed via this equation:
58    //   window_start_i = zero + every * i
59    //   window_stop_i = zero + every * i + period
60    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    /// Truncate the given timestamp in `tu` by the window boundary.
76    pub fn truncate(&self, tu: TimeUnit, t: i64, tz: Option<&Tz>) -> PolarsResult<i64> {
77        self.every.truncate(tu, t, tz)
78    }
79
80    /// Round the given timestamp in `tu` by the window boundary.
81    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    /// returns the bounds for the earliest window bounds
87    /// that contains the given time t.  For underlapping windows that
88    /// do not contain time t, the window directly after time t will be returned.
89    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    /// The start of the first window for data whose first value is `t0`, as placed by
119    /// `start_by`, `every` and `offset`.
120    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                // Find the beginning of the week in the time zone, then place the window
133                // start on the requested weekday plus `offset`.
134                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                // Make sure the first datapoint has a chance to be included.
148                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
163/// Iterates the windows `origin + every * k` for `k = 0, 1, ...`. Each start is computed from
164/// `origin` rather than from the previous start, so that calendar durations clamp the same way
165/// as in `datetime_range`, and skipping ahead gives the same windows as stepping one by one.
166pub struct BoundsIter<'a> {
167    window: Window,
168    // wrapping boundary
169    boundary: Bounds,
170    origin: i64,
171    // index of the window in `bi`
172    k: i64,
173    // boundary per window iterator
174    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        // TODO: find some way to propagate error instead of unwrapping?
206        // Issue is that `next` needs to return `Option`.
207        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    /// Number of iterations to advance, such that the bounds are on target; or, in
246    /// the case of non-constant duration, close to target.
247    /// Follows the `nth()` convention on Iterator indexing, i.e., a return value of 0
248    /// implies advancing 1 iteration.
249    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}