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: u32,
16 priority: Priority,
18 randomize_start: Option<chrono::Duration>,
20 value_map: Vec<EventValuesMap>,
22}
23
24#[allow(unused)]
29#[derive(Clone, Default, Debug)]
30pub struct Timeline {
31 data: rangemap::RangeMap<DateTime<Utc>, InternalInterval>,
32}
33
34impl Timeline {
35 pub fn new() -> Self {
37 Self {
38 data: rangemap::RangeMap::new(),
39 }
40 }
41
42 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 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 ¤t_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 pub fn iter(&self) -> Iter<'_> {
138 Iter {
139 iter: self.data.iter(),
140 seen: HashSet::default(),
141 }
142 }
143
144 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 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#[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 pub fn randomize_start(&self) -> Option<chrono::Duration> {
199 self.randomize_start
200 }
201
202 pub fn value_map(&self) -> &[EventValuesMap] {
204 self.value_map
205 }
206}
207
208pub 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 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 #[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 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 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}