use std::{collections::HashSet, ops::Range};
use chrono::{DateTime, Utc};
use tracing::warn;
use openleadr_wire::{
event::{EventContent, EventValuesMap, Priority},
interval::IntervalPeriod,
Program,
};
#[derive(Debug, Clone, PartialEq, Eq)]
struct InternalInterval {
id: u32,
priority: Priority,
randomize_start: Option<chrono::Duration>,
value_map: Vec<EventValuesMap>,
}
#[allow(unused)]
#[derive(Clone, Default, Debug)]
pub struct Timeline {
data: rangemap::RangeMap<DateTime<Utc>, InternalInterval>,
}
impl Timeline {
pub fn new() -> Self {
Self {
data: rangemap::RangeMap::new(),
}
}
pub fn from_events(program: &Program, mut events: Vec<&EventContent>) -> Option<Self> {
let mut data = Self::default();
events.sort_by_key(|e| e.priority);
for (id, event) in events.iter().enumerate() {
if event.program_id != program.id {
warn!(?event, %program.id, "skipping event that does not belong into the program; different program id");
continue;
}
let default_period = event.interval_period.as_ref();
let mut current_start = default_period.map(|p| p.start);
for event_interval in &event.intervals {
let (start, duration, randomize_start) =
match event_interval.interval_period.as_ref() {
Some(IntervalPeriod {
start,
duration,
randomize_start,
}) => (start, duration, randomize_start),
None => (
¤t_start?,
&default_period?.duration,
&default_period?.randomize_start,
),
};
let range = match duration {
Some(duration) => *start..*start + duration.to_chrono_at_datetime(*start),
None => *start..DateTime::<Utc>::MAX_UTC,
};
current_start = Some(range.end);
let interval = InternalInterval {
id: id as u32,
randomize_start: randomize_start
.as_ref()
.map(|d| d.to_chrono_at_datetime(*start)),
value_map: event_interval.payloads.clone(),
priority: event.priority,
};
for (existing_range, existing) in data.data.overlapping(&range) {
if existing.priority == event.priority {
warn!(?existing_range, ?existing, new_range = ?range, new = ?interval, "Overlapping ranges with equal priority");
}
}
data.data.insert(range, interval);
}
}
Some(data)
}
pub fn iter(&self) -> Iter<'_> {
Iter {
iter: self.data.iter(),
seen: HashSet::default(),
}
}
pub fn at_datetime(
&self,
datetime: &DateTime<Utc>,
) -> Option<(&Range<DateTime<Utc>>, Interval<'_>)> {
let (range, internal_interval) = self.data.get_key_value(datetime)?;
let interval = Interval {
randomize_start: internal_interval.randomize_start,
value_map: &internal_interval.value_map,
};
Some((range, interval))
}
pub fn next_update(&self, datetime: &DateTime<Utc>) -> Option<DateTime<Utc>> {
if let Some((k, _)) = self.at_datetime(datetime) {
return Some(k.end);
}
let (last_range, _) = self.data.last_range_value()?;
let (range, _) = self.data.overlapping(*datetime..last_range.end).next()?;
Some(range.start)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Interval<'a> {
randomize_start: Option<chrono::Duration>,
value_map: &'a [EventValuesMap],
}
impl Interval<'_> {
pub fn randomize_start(&self) -> Option<chrono::Duration> {
self.randomize_start
}
pub fn value_map(&self) -> &[EventValuesMap] {
self.value_map
}
}
pub struct Iter<'a> {
iter: rangemap::map::Iter<'a, DateTime<Utc>, InternalInterval>,
seen: HashSet<u32>,
}
impl<'a> Iterator for Iter<'a> {
type Item = (&'a Range<DateTime<Utc>>, Interval<'a>);
fn next(&mut self) -> Option<Self::Item> {
let (range, internal) = self.iter.next()?;
let interval = Interval {
randomize_start: match self.seen.insert(internal.id) {
true => internal.randomize_start,
false => None,
},
value_map: &internal.value_map,
};
Some((range, interval))
}
}
#[cfg(test)]
mod test {
use std::ops::Range;
use chrono::{DateTime, Duration, Utc};
use super::*;
use openleadr_wire::{
event::EventInterval,
program::{ProgramContent, ProgramId},
values_map::Value,
};
fn test_program_id() -> ProgramId {
ProgramId::new("test-program-id").unwrap()
}
fn test_event_content(range: Range<u32>, value: i64) -> EventContent {
EventContent::new(
test_program_id(),
vec![event_interval_with_value(range, value)],
)
}
fn test_program(name: &str) -> Program {
Program {
id: test_program_id(),
created_date_time: Default::default(),
modification_date_time: Default::default(),
content: ProgramContent::new(name),
}
}
fn event_interval_with_value(range: Range<u32>, value: i64) -> EventInterval {
EventInterval {
id: range.start as _,
interval_period: Some(IntervalPeriod {
start: DateTime::UNIX_EPOCH + Duration::hours(range.start.into()),
duration: Some(openleadr_wire::Duration::hours(
(range.end - range.start) as _,
)),
randomize_start: None,
}),
payloads: vec![EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(value)],
}],
}
}
fn interval_with_value(
id: u32,
range: Range<u32>,
value: i64,
priority: Priority,
) -> (Range<DateTime<Utc>>, InternalInterval) {
let start = DateTime::UNIX_EPOCH + Duration::hours(range.start.into());
let end = DateTime::UNIX_EPOCH + Duration::hours(range.end.into());
(
start..end,
InternalInterval {
id,
randomize_start: None,
value_map: vec![EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(value)],
}],
priority,
},
)
}
#[test]
fn overlap_same_priority() {
let program = test_program("p");
let event1 = test_event_content(0..10, 42);
let event2 = test_event_content(5..15, 43);
let tl1 = Timeline::from_events(&program, vec![&event1, &event2]).unwrap();
assert_eq!(
tl1.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(0, 0..5, 42, Priority::UNSPECIFIED),
interval_with_value(1, 5..15, 43, Priority::UNSPECIFIED),
]
);
let tl2 = Timeline::from_events(&program, vec![&event2, &event1]).unwrap();
assert_eq!(
tl2.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(1, 0..10, 42, Priority::UNSPECIFIED),
interval_with_value(0, 10..15, 43, Priority::UNSPECIFIED),
]
);
}
#[test]
fn overlap_lower_priority() {
let event1 = test_event_content(0..10, 42).with_priority(Priority::new(1));
let event2 = test_event_content(5..15, 43).with_priority(Priority::new(2));
let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
assert_eq!(
tl.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(1, 0..10, 42, Priority::new(1)),
interval_with_value(0, 10..15, 43, Priority::new(2)),
],
"a lower priority event MUST NOT overwrite a higher priority one",
);
let tl = Timeline::from_events(&test_program("p"), vec![&event2, &event1]).unwrap();
assert_eq!(
tl.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(1, 0..10, 42, Priority::new(1)),
interval_with_value(0, 10..15, 43, Priority::new(2)),
],
"a lower priority event MUST NOT overwrite a higher priority one",
);
}
#[test]
fn overlap_higher_priority() {
let event1 = test_event_content(0..10, 42).with_priority(Priority::new(2));
let event2 = test_event_content(5..15, 43).with_priority(Priority::new(1));
let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
assert_eq!(
tl.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(0, 0..5, 42, Priority::new(2)),
interval_with_value(1, 5..15, 43, Priority::new(1)),
],
"a higher priority event MUST overwrite a lower priority one",
);
let tl = Timeline::from_events(&test_program("p"), vec![&event2, &event1]).unwrap();
assert_eq!(
tl.data.into_iter().collect::<Vec<_>>(),
vec![
interval_with_value(0, 0..5, 42, Priority::new(2)),
interval_with_value(1, 5..15, 43, Priority::new(1)),
],
"a higher priority event MUST overwrite a lower priority one",
);
}
#[test]
fn default_interval() {
let program = test_program("p");
let event_intervals = vec![
EventInterval::new(
0,
vec![EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Number(1.23)],
}],
),
EventInterval::new(
1,
vec![EventValuesMap {
value_type: openleadr_wire::event::EventType::Simple,
values: vec![Value::Number(2.34)],
}],
),
];
let mut event = EventContent::new(program.id.clone(), event_intervals);
event.interval_period = Some(IntervalPeriod {
start: DateTime::UNIX_EPOCH,
duration: Some(openleadr_wire::Duration::hours(5.)),
randomize_start: None,
});
let timeline = Timeline::from_events(&program, vec![&event]).unwrap();
let interval = timeline
.at_datetime(&(DateTime::UNIX_EPOCH + Duration::hours(2)))
.unwrap();
assert_eq!(
interval.1.value_map[0].value_type,
openleadr_wire::event::EventType::Price
);
let interval = timeline
.at_datetime(&(DateTime::UNIX_EPOCH + Duration::hours(8)))
.unwrap();
assert_eq!(
interval.1.value_map[0].value_type,
openleadr_wire::event::EventType::Simple
);
}
#[test]
fn randomize_start_not_duplicated() {
let event1 = test_event_content(5..10, 42).with_priority(Priority::MAX);
let event2 = {
let range = 0..15;
let value = 43;
EventContent::new(
test_program_id(),
vec![EventInterval {
id: range.start as _,
interval_period: Some(IntervalPeriod {
start: DateTime::UNIX_EPOCH + Duration::hours(range.start.into()),
duration: Some(openleadr_wire::Duration::hours(
(range.end - range.start) as _,
)),
randomize_start: Some(openleadr_wire::Duration::hours(5.0)),
}),
payloads: vec![EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(value)],
}],
}],
)
};
let tl = Timeline::from_events(&test_program("p"), vec![&event1, &event2]).unwrap();
assert_eq!(
tl.iter().map(|(_, i)| i).collect::<Vec<_>>(),
vec![
Interval {
randomize_start: Some(Duration::hours(5)),
value_map: &[EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(43)],
}],
},
Interval {
randomize_start: None,
value_map: &[EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(42)],
}],
},
Interval {
randomize_start: None,
value_map: &[EventValuesMap {
value_type: openleadr_wire::event::EventType::Price,
values: vec![Value::Integer(43)],
}],
},
],
"when an event is split, only the first interval should retain `randomize_start`",
);
}
}