Skip to main content

openleadr_client/
timeline.rs

1use std::{collections::HashSet, ops::Range};
2
3use chrono::{DateTime, Utc};
4use tracing::warn;
5
6use openleadr_wire::{
7    Program,
8    event::{EventRequest, EventValuesMap, Priority},
9    interval::IntervalPeriod,
10};
11
12#[derive(Debug, Clone, PartialEq, Eq)]
13struct InternalInterval {
14    /// Id so that split intervals with a randomized start don't start randomly twice
15    id: u32,
16    /// Relative priority of event
17    priority: Priority,
18    /// Indicates a randomization time that may be applied to start.
19    randomize_start: Option<chrono::Duration>,
20    /// The actual values that are active during this interval
21    value_map: Vec<EventValuesMap>,
22}
23
24/// A sequence of ordered, non-overlapping intervals and associated values.
25///
26/// Intervals are sorted by their timestamp. The intervals will not overlap, but there may be gaps
27/// between intervals.
28#[allow(unused)]
29#[derive(Clone, Default, Debug)]
30pub struct Timeline {
31    data: rangemap::RangeMap<DateTime<Utc>, InternalInterval>,
32}
33
34impl Timeline {
35    /// Create an empty [`Timeline`]
36    pub fn new() -> Self {
37        Self {
38            data: rangemap::RangeMap::new(),
39        }
40    }
41
42    /// Creates a [`Timeline`] from a [`Program`], and the [`Event`](crate::Event)s that belong to it.
43    ///
44    /// It sorts events according to their priority and builds the timeline accordingly.
45    /// The timeline can have gaps if the intervals in the events contain gaps.
46    /// The event with the highest priority always takes presence at a specific time point.
47    /// Therefore, a long-lasting, low-priority event can be interrupted by a short, high-priority event,
48    /// for example.
49    /// In this case, the long-lasting event will be split into two parts,
50    /// such that the high-priority event fits in between.
51    ///
52    /// ```text
53    /// Input:
54    /// |------------------------long, low prio--------------------------|    |--another-interval--|
55    ///                    |-----short, high prio---|
56    /// Result:
57    /// |--long, low prio--|-----short, high prio---|---long, low prio---|    |--another-interval--|
58    /// ```
59    ///
60    /// This function logs at `warn` level if provided with an [`Event`](crate::Event)
61    /// those [`program_id`](openleadr_wire::event::EventRequest::program_id) does not match with the [`Program::id`].
62    /// The corresponding event will be ignored then building the timeline.
63    ///
64    /// This function also logs at `warn`
65    /// level if there are two overlapping events at the same priority.
66    /// There is no guarantee which event will take precedence, though in the current implementation
67    /// the event stored later in the `events` param will take precedence.
68    ///
69    /// There must be an [`IntervalPeriod`] present on the event level,
70    /// or the individual intervals.
71    /// If both are specified, the individual period takes precedence over the one specified in the event.
72    /// If for an interval, there is no period present, and none specified in the event,
73    /// then this function will return [`None`]
74    pub fn from_events(program: &Program, mut events: Vec<&EventRequest>) -> Option<Self> {
75        let mut data = Self::default();
76
77        events.sort_by_key(|e| e.priority);
78
79        for (id, event) in events.iter().enumerate() {
80            if event.program_id != program.id {
81                warn!(?event, %program.id, "skipping event that does not belong into the program; different program id");
82                continue;
83            }
84
85            let default_period = event.interval_period.as_ref();
86
87            let mut current_start = default_period.map(|p| p.start);
88
89            if let Some(event_intervals) = &event.intervals {
90                for event_interval in event_intervals {
91                    // use the event interval period when the interval doesn't specify one
92                    let (start, duration, randomize_start) =
93                        match event_interval.interval_period.as_ref() {
94                            Some(IntervalPeriod {
95                                start,
96                                duration,
97                                randomize_start,
98                            }) => (start, duration, randomize_start),
99                            None => (
100                                &current_start?,
101                                &default_period?.duration,
102                                &default_period?.randomize_start,
103                            ),
104                        };
105
106                    let range = match duration {
107                        Some(duration) => *start..*start + duration.to_chrono_at_datetime(*start),
108                        None => *start..DateTime::<Utc>::MAX_UTC,
109                    };
110
111                    current_start = Some(range.end);
112
113                    let interval = InternalInterval {
114                        id: id as u32,
115                        randomize_start: randomize_start
116                            .as_ref()
117                            .map(|d| d.to_chrono_at_datetime(*start)),
118                        value_map: event_interval.payloads.clone(),
119                        priority: event.priority,
120                    };
121
122                    for (existing_range, existing) in data.data.overlapping(&range) {
123                        if existing.priority == event.priority {
124                            warn!(?existing_range, ?existing, new_range = ?range, new = ?interval, "Overlapping ranges with equal priority");
125                        }
126                    }
127
128                    data.data.insert(range, interval);
129                }
130            }
131        }
132
133        Some(data)
134    }
135
136    /// Get an iterator over the [`Interval`]s in this [`Timeline`]
137    pub fn iter(&self) -> Iter<'_> {
138        Iter {
139            iter: self.data.iter(),
140            seen: HashSet::default(),
141        }
142    }
143
144    /// Returns the [`Interval`] applicable at the requested time point and the range it is valid for.
145    pub fn at_datetime(
146        &self,
147        datetime: &DateTime<Utc>,
148    ) -> Option<(&Range<DateTime<Utc>>, Interval<'_>)> {
149        let (range, internal_interval) = self.data.get_key_value(datetime)?;
150
151        let interval = Interval {
152            randomize_start: internal_interval.randomize_start,
153            value_map: &internal_interval.value_map,
154        };
155
156        Some((range, interval))
157    }
158
159    /// Returns the time when to next change takes effect.
160    ///
161    /// **Example:**
162    /// ```text
163    ///   |--interval 1--|---interval 2---|       |---interval 3---|
164    ///   ↑              ↑                ↑       ↑                ↑
165    /// 08:03          09:56            10:59   11:01            12:00
166    /// ```
167    /// For the timeline illustrated above,
168    /// * `next_update(09:56)` would return `Some(10:59)`
169    /// * `next_update(10:00)` would return `Some(10:59)`
170    /// * `next_update(11:00)` would return `Some(11:01)`
171    /// * `next_update(12:00)` would return `None`
172    /// * `next_update(12:01)` would return `None`
173    pub fn next_update(&self, datetime: &DateTime<Utc>) -> Option<DateTime<Utc>> {
174        if let Some((k, _)) = self.at_datetime(datetime) {
175            return Some(k.end);
176        }
177
178        let (last_range, _) = self.data.last_range_value()?;
179
180        let (range, _) = self.data.overlapping(*datetime..last_range.end).next()?;
181
182        Some(range.start)
183    }
184}
185
186/// Holds the data stored in a [`Timeline`].
187///
188/// This data type is returned by the Iterator over the [`Timeline`] and by the [`at_datetime`](Timeline::at_datetime)
189/// method.
190#[derive(Debug, Clone, PartialEq, Eq)]
191pub struct Interval<'a> {
192    randomize_start: Option<chrono::Duration>,
193    value_map: &'a [EventValuesMap],
194}
195
196impl Interval<'_> {
197    /// Indicates a randomization time that may be applied to start.
198    pub fn randomize_start(&self) -> Option<chrono::Duration> {
199        self.randomize_start
200    }
201
202    /// The actual values that are active during this interval
203    pub fn value_map(&self) -> &[EventValuesMap] {
204        self.value_map
205    }
206}
207
208/// Iterator over [`Timeline`].
209///
210/// **Important:** The specification does not specify how to handle the `randomize_start`
211/// for overlapping intervals.
212/// This implementation sets the [`randomize_start`](Interval::randomize_start)
213/// value to [`None`] if it's not the first part of that interval.
214/// A single interval can be split in multiple parts
215/// if it is partly covered by an event with higher priority.
216/// See [`Timeline::from_events`].
217///
218/// **Example:**
219/// ```text
220/// |----interval 1----|-----interval 2---|----interval 1----|
221/// ```
222/// If `interval 1` has a `randomize_start` specified,
223/// this iterator will only set this in the first part, but `None` in the second.
224pub struct Iter<'a> {
225    iter: rangemap::map::Iter<'a, DateTime<Utc>, InternalInterval>,
226    seen: HashSet<u32>,
227}
228
229impl<'a> Iterator for Iter<'a> {
230    type Item = (&'a Range<DateTime<Utc>>, Interval<'a>);
231
232    fn next(&mut self) -> Option<Self::Item> {
233        let (range, internal) = self.iter.next()?;
234
235        let interval = Interval {
236            // only the first occurrence of an id should randomize its start
237            randomize_start: match self.seen.insert(internal.id) {
238                true => internal.randomize_start,
239                false => None,
240            },
241            value_map: &internal.value_map,
242        };
243
244        Some((range, interval))
245    }
246}
247
248#[cfg(test)]
249mod test {
250    use std::ops::Range;
251
252    use chrono::{DateTime, Duration, Utc};
253
254    use super::*;
255    use openleadr_wire::{
256        event::EventInterval,
257        program::{ProgramId, ProgramRequest},
258        values_map::Value,
259    };
260
261    fn test_program_id() -> ProgramId {
262        ProgramId::new("test-program-id").unwrap()
263    }
264
265    fn test_event_content(range: Range<u32>, value: i64) -> EventRequest {
266        EventRequest::new(test_program_id())
267            .with_intervals(vec![event_interval_with_value(range, value)])
268    }
269
270    fn test_program(name: &str) -> Program {
271        Program {
272            id: test_program_id(),
273            created_date_time: Default::default(),
274            modification_date_time: Default::default(),
275            content: ProgramRequest::new(name),
276        }
277    }
278
279    fn event_interval_with_value(range: Range<u32>, value: i64) -> EventInterval {
280        EventInterval {
281            id: range.start as _,
282            interval_period: Some(IntervalPeriod {
283                start: DateTime::UNIX_EPOCH + Duration::hours(range.start.into()),
284                duration: Some(openleadr_wire::Duration::hours(
285                    (range.end - range.start) as _,
286                )),
287                randomize_start: None,
288            }),
289            payloads: vec![EventValuesMap {
290                value_type: openleadr_wire::event::EventType::Price,
291                values: vec![Value::Integer(value)],
292            }],
293        }
294    }
295
296    fn interval_with_value(
297        id: u32,
298        range: Range<u32>,
299        value: i64,
300        priority: Priority,
301    ) -> (Range<DateTime<Utc>>, InternalInterval) {
302        let start = DateTime::UNIX_EPOCH + Duration::hours(range.start.into());
303        let end = DateTime::UNIX_EPOCH + Duration::hours(range.end.into());
304
305        (
306            start..end,
307            InternalInterval {
308                id,
309                randomize_start: None,
310                value_map: vec![EventValuesMap {
311                    value_type: openleadr_wire::event::EventType::Price,
312                    values: vec![Value::Integer(value)],
313                }],
314                priority,
315            },
316        )
317    }
318
319    // The spec does not specify the behavior when two intervals with the same priority overlap.
320    // Our current implementation uses `RangeMap`, and its behavior is to overwrite the existing
321    // range with a new one.
322    // In other words: the event which is inserted last wins.
323    #[test]
324    fn overlap_same_priority() {
325        let program = test_program("p");
326
327        let event1 = test_event_content(0..10, 42);
328        let event2 = test_event_content(5..15, 43);
329
330        // first come, last serve
331        let tl1 = Timeline::from_events(&program, vec![&event1, &event2]).unwrap();
332        assert_eq!(
333            tl1.data.into_iter().collect::<Vec<_>>(),
334            vec![
335                interval_with_value(0, 0..5, 42, Priority::UNSPECIFIED),
336                interval_with_value(1, 5..15, 43, Priority::UNSPECIFIED),
337            ]
338        );
339
340        // first come, last serve
341        let tl2 = Timeline::from_events(&program, vec![&event2, &event1]).unwrap();
342        assert_eq!(
343            tl2.data.into_iter().collect::<Vec<_>>(),
344            vec![
345                interval_with_value(1, 0..10, 42, Priority::UNSPECIFIED),
346                interval_with_value(0, 10..15, 43, Priority::UNSPECIFIED),
347            ]
348        );
349    }
350
351    #[test]
352    fn overlap_lower_priority() {
353        let event1 = test_event_content(0..10, 42).with_priority(Priority::new(1));
354        let event2 = test_event_content(5..15, 43).with_priority(Priority::new(2));
355
356        let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
357        assert_eq!(
358            tl.data.into_iter().collect::<Vec<_>>(),
359            vec![
360                interval_with_value(1, 0..10, 42, Priority::new(1)),
361                interval_with_value(0, 10..15, 43, Priority::new(2)),
362            ],
363            "a lower priority event MUST NOT overwrite a higher priority one",
364        );
365
366        let tl = Timeline::from_events(&test_program("p"), vec![&event2, &event1]).unwrap();
367        assert_eq!(
368            tl.data.into_iter().collect::<Vec<_>>(),
369            vec![
370                interval_with_value(1, 0..10, 42, Priority::new(1)),
371                interval_with_value(0, 10..15, 43, Priority::new(2)),
372            ],
373            "a lower priority event MUST NOT overwrite a higher priority one",
374        );
375    }
376
377    #[test]
378    fn overlap_higher_priority() {
379        let event1 = test_event_content(0..10, 42).with_priority(Priority::new(2));
380        let event2 = test_event_content(5..15, 43).with_priority(Priority::new(1));
381
382        let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
383        assert_eq!(
384            tl.data.into_iter().collect::<Vec<_>>(),
385            vec![
386                interval_with_value(0, 0..5, 42, Priority::new(2)),
387                interval_with_value(1, 5..15, 43, Priority::new(1)),
388            ],
389            "a higher priority event MUST overwrite a lower priority one",
390        );
391
392        let tl = Timeline::from_events(&test_program("p"), vec![&event2, &event1]).unwrap();
393        assert_eq!(
394            tl.data.into_iter().collect::<Vec<_>>(),
395            vec![
396                interval_with_value(0, 0..5, 42, Priority::new(2)),
397                interval_with_value(1, 5..15, 43, Priority::new(1)),
398            ],
399            "a higher priority event MUST overwrite a lower priority one",
400        );
401    }
402
403    #[test]
404    fn default_interval() {
405        let program = test_program("p");
406
407        let event_intervals = vec![
408            EventInterval::new(
409                0,
410                vec![EventValuesMap {
411                    value_type: openleadr_wire::event::EventType::Price,
412                    values: vec![Value::Number(1.23)],
413                }],
414            ),
415            EventInterval::new(
416                1,
417                vec![EventValuesMap {
418                    value_type: openleadr_wire::event::EventType::Simple,
419                    values: vec![Value::Number(2.34)],
420                }],
421            ),
422        ];
423
424        let mut event = EventRequest::new(program.id.clone()).with_intervals(event_intervals);
425
426        event.interval_period = Some(IntervalPeriod {
427            start: DateTime::UNIX_EPOCH,
428            duration: Some(openleadr_wire::Duration::hours(5.)),
429            randomize_start: None,
430        });
431
432        let timeline = Timeline::from_events(&program, vec![&event]).unwrap();
433
434        let interval = timeline
435            .at_datetime(&(DateTime::UNIX_EPOCH + Duration::hours(2)))
436            .unwrap();
437        assert_eq!(
438            interval.1.value_map[0].value_type,
439            openleadr_wire::event::EventType::Price
440        );
441
442        let interval = timeline
443            .at_datetime(&(DateTime::UNIX_EPOCH + Duration::hours(8)))
444            .unwrap();
445        assert_eq!(
446            interval.1.value_map[0].value_type,
447            openleadr_wire::event::EventType::Simple
448        );
449    }
450
451    #[test]
452    fn randomize_start_not_duplicated() {
453        let event1 = test_event_content(5..10, 42).with_priority(Priority::MAX);
454
455        let event2 = {
456            let range = 0..15;
457            let value = 43;
458            EventRequest::new(test_program_id()).with_intervals(vec![EventInterval {
459                id: range.start as _,
460                interval_period: Some(IntervalPeriod {
461                    start: DateTime::UNIX_EPOCH + Duration::hours(range.start.into()),
462                    duration: Some(openleadr_wire::Duration::hours(
463                        (range.end - range.start) as _,
464                    )),
465                    randomize_start: Some(openleadr_wire::Duration::hours(5.0)),
466                }),
467                payloads: vec![EventValuesMap {
468                    value_type: openleadr_wire::event::EventType::Price,
469                    values: vec![Value::Integer(value)],
470                }],
471            }])
472        };
473
474        let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
475        assert_eq!(
476            tl.iter().map(|(_, i)| i).collect::<Vec<_>>(),
477            vec![
478                Interval {
479                    randomize_start: Some(Duration::hours(5)),
480                    value_map: &[EventValuesMap {
481                        value_type: openleadr_wire::event::EventType::Price,
482                        values: vec![Value::Integer(43)],
483                    }],
484                },
485                Interval {
486                    randomize_start: None,
487                    value_map: &[EventValuesMap {
488                        value_type: openleadr_wire::event::EventType::Price,
489                        values: vec![Value::Integer(42)],
490                    }],
491                },
492                Interval {
493                    randomize_start: None,
494                    value_map: &[EventValuesMap {
495                        value_type: openleadr_wire::event::EventType::Price,
496                        values: vec![Value::Integer(43)],
497                    }],
498                },
499            ],
500            "when an event is split, only the first interval should retain `randomize_start`",
501        );
502    }
503}