From 39e08d602b9a0df4dfba686fdeb6ccafaaa013ad Mon Sep 17 00:00:00 2001 From: Wez Furlong Date: Tue, 28 Feb 2023 23:16:11 -0700 Subject: [PATCH] add scheduling restriction datatype This will be used in the future to adjust the message sending schedule on a per-message basis. --- Cargo.lock | 77 ++++++ crates/message/Cargo.toml | 1 + crates/message/src/lib.rs | 1 + crates/message/src/message.rs | 29 ++- crates/message/src/scheduling.rs | 391 +++++++++++++++++++++++++++++++ 5 files changed, 498 insertions(+), 1 deletion(-) create mode 100644 crates/message/src/scheduling.rs diff --git a/Cargo.lock b/Cargo.lock index 9dccfc15..33824648 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -381,6 +381,29 @@ dependencies = [ "winapi", ] +[[package]] +name = "chrono-tz" +version = "0.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fa48fa079165080f11d7753fd0bc175b7d391f276b965fe4b55bfad67856e463" +dependencies = [ + "chrono", + "chrono-tz-build", + "phf", + "serde", +] + +[[package]] +name = "chrono-tz-build" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9998fb9f7e9b2111641485bf8beb32f92945f97f92a3d061f744cfef335f751" +dependencies = [ + "parse-zoneinfo", + "phf", + "phf_codegen", +] + [[package]] name = "cidr" version = "0.2.1" @@ -1560,6 +1583,7 @@ dependencies = [ "anyhow", "bitflags", "chrono", + "chrono-tz", "config", "futures", "k9", @@ -1898,6 +1922,15 @@ dependencies = [ "windows-sys 0.45.0", ] +[[package]] +name = "parse-zoneinfo" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c705f256449c60da65e11ff6626e0c16a0a0b96aaa348de61376b249bc340f41" +dependencies = [ + "regex", +] + [[package]] name = "password-hash" version = "0.4.2" @@ -1986,6 +2019,44 @@ dependencies = [ "sha2", ] +[[package]] +name = "phf" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "928c6535de93548188ef63bb7c4036bd415cd8f36ad25af44b9789b2ee72a48c" +dependencies = [ + "phf_shared", +] + +[[package]] +name = "phf_codegen" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a56ac890c5e3ca598bbdeaa99964edb5b0258a583a9eb6ef4e89fc85d9224770" +dependencies = [ + "phf_generator", + "phf_shared", +] + +[[package]] +name = "phf_generator" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1181c94580fa345f50f19d738aaa39c0ed30a600d95cb2d3e23f94266f14fbf" +dependencies = [ + "phf_shared", + "rand", +] + +[[package]] +name = "phf_shared" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e1fb5f6f826b772a8d4c0394209441e7d37cbbb967ae9c7e0e8134365c9ee676" +dependencies = [ + "siphasher", +] + [[package]] name = "pin-project" version = "1.0.12" @@ -2513,6 +2584,12 @@ dependencies = [ "libc", ] +[[package]] +name = "siphasher" +version = "0.3.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bd3e3206899af3f8b12af284fafc038cc1dc2b41d1b89dd17297221c5d225de" + [[package]] name = "sketches-ddsketch" version = "0.2.0" diff --git a/crates/message/Cargo.toml b/crates/message/Cargo.toml index d8c5b4a9..b3d63180 100644 --- a/crates/message/Cargo.toml +++ b/crates/message/Cargo.toml @@ -10,6 +10,7 @@ anyhow = "1.0" bitflags = "1.3" config = {path="../config"} chrono = {version="0.4", features=["serde"]} +chrono-tz = {version="0.8", features=["serde"]} futures = "0.3" lazy_static = "1.4" lruttl = {path="../lruttl"} diff --git a/crates/message/src/lib.rs b/crates/message/src/lib.rs index 5741fd62..3b42c5a2 100644 --- a/crates/message/src/lib.rs +++ b/crates/message/src/lib.rs @@ -1,6 +1,7 @@ pub mod address; pub mod dkim; pub mod message; +pub mod scheduling; pub use crate::address::EnvelopeAddress; pub use crate::message::Message; diff --git a/crates/message/src/message.rs b/crates/message/src/message.rs index 6ef40ca7..2b91aea5 100644 --- a/crates/message/src/message.rs +++ b/crates/message/src/message.rs @@ -1,4 +1,5 @@ use crate::dkim::Signer; +use crate::scheduling::Scheduling; use crate::EnvelopeAddress; use chrono::{DateTime, Utc}; use futures::FutureExt; @@ -18,6 +19,8 @@ bitflags::bitflags! { const META_DIRTY = 1; /// true if Data needs to be saved const DATA_DIRTY = 2; + /// true if scheduling restrictions are present in the metadata + const SCHEDULED = 4; } } @@ -56,6 +59,8 @@ struct MetaData { sender: EnvelopeAddress, recipient: EnvelopeAddress, meta: serde_json::Value, + #[serde(default)] + schedule: Option, } impl Drop for MessageInner { @@ -91,6 +96,7 @@ impl Message { sender, recipient, meta, + schedule: None, }), data, flags: MessageFlags::META_DIRTY | MessageFlags::DATA_DIRTY, @@ -108,12 +114,18 @@ impl Message { MESSAGE_COUNT.inc(); META_COUNT.inc(); + let flags = if metadata.schedule.is_some() { + MessageFlags::SCHEDULED + } else { + MessageFlags::empty() + }; + Ok(Self { id, inner: Arc::new(Mutex::new(MessageInner { metadata: Some(metadata), data: Arc::new(vec![].into_boxed_slice()), - flags: MessageFlags::empty(), + flags, num_attempts: 0, due: None, })), @@ -135,6 +147,20 @@ impl Message { inner.num_attempts += 1; } + pub fn set_scheduling(&self, scheduling: Option) -> anyhow::Result<()> { + let mut inner = self.inner.lock().unwrap(); + match &mut inner.metadata { + None => anyhow::bail!("metadata must be loaded first"), + Some(meta) => { + meta.schedule = scheduling; + inner + .flags + .set(MessageFlags::SCHEDULED, scheduling.is_some()); + Ok(()) + } + } + } + pub fn get_due(&self) -> Option> { let inner = self.inner.lock().unwrap(); inner.due @@ -737,6 +763,7 @@ impl<'a> QueueNameComponents<'a> { } } } + #[cfg(test)] mod test { use super::*; diff --git a/crates/message/src/scheduling.rs b/crates/message/src/scheduling.rs new file mode 100644 index 00000000..74462f72 --- /dev/null +++ b/crates/message/src/scheduling.rs @@ -0,0 +1,391 @@ +use chrono::naive::NaiveTime; +use chrono::{DateTime, Datelike, FixedOffset, LocalResult, TimeZone, Timelike, Utc, Weekday}; +use chrono_tz::Tz; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; + +bitflags::bitflags! { + pub struct DaysOfWeek: u8 { + const MON = 1; + const TUE = 2; + const WED = 4; + const THU = 8; + const FRI = 16; + const SAT = 32; + const SUN = 64; + } +} + +impl From for DaysOfWeek { + fn from(day: Weekday) -> DaysOfWeek { + match day { + Weekday::Mon => DaysOfWeek::MON, + Weekday::Tue => DaysOfWeek::TUE, + Weekday::Wed => DaysOfWeek::WED, + Weekday::Thu => DaysOfWeek::THU, + Weekday::Fri => DaysOfWeek::FRI, + Weekday::Sat => DaysOfWeek::SAT, + Weekday::Sun => DaysOfWeek::SUN, + } + } +} + +/// Represents a restriction on when the message can be sent. +/// This encodes the permitted times. +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Copy)] +pub struct ScheduleRestriction { + #[serde(rename = "dow")] + pub days_of_week: DaysOfWeek, + #[serde(rename = "tz")] + pub timezone: Tz, + pub start: NaiveTime, + pub end: NaiveTime, +} + +impl ScheduleRestriction { + fn start_end_on_day(&self, dt: DateTime) -> Option<(DateTime, DateTime)> { + let y = dt.year(); + let m = dt.month(); + let d = dt.day(); + + let start = match dbg!(self.timezone.with_ymd_and_hms( + y, + m, + d, + self.start.hour(), + self.start.minute(), + self.start.second(), + )) { + LocalResult::Single(t) => t, + _ => return None, + }; + + let end = match dbg!(self.timezone.with_ymd_and_hms( + y, + m, + d, + self.end.hour(), + self.end.minute(), + self.end.second(), + )) { + LocalResult::Single(t) => t, + _ => return None, + }; + Some((start, end)) + } +} + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Copy)] +pub struct Scheduling { + #[serde(flatten, skip_serializing_if = "Option::is_none")] + pub restriction: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub not_before: Option>, +} + +impl Scheduling { + pub fn adjust_for_schedule(&self, mut dt: DateTime) -> DateTime { + if let Some(start) = &self.not_before { + if dt < *start { + dt = (*start).into(); + } + } + + if let Some(restrict) = &self.restriction { + let mut dt = dt.with_timezone(&restrict.timezone); + println!("start with {dt:?}"); + + let one_day = chrono::Duration::days(1); + + // Worst case is 1 week off the current time; if we + // can't find a time in a reasonable number of iterations, + // something is wrong! + for iter in 0..8 { + let weekday = dt.weekday(); + println!("iter {iter} {weekday:?}"); + let dow: DaysOfWeek = weekday.into(); + + let (start, end) = match restrict.start_end_on_day(dt) { + Some(result) => result, + None => { + // Wonky date/time, try the next day + dt = dt + one_day; + println!("WONKY! using {dt:?}"); + continue; + } + }; + + if restrict.days_of_week.contains(dow) { + if dt < start { + // Delay until the start time + println!("round up to start"); + dt = start; + break; + } + + if dt < end { + // We're within the permitted range + println!("we are within range"); + break; + } + } + + // Try the same start time the next day + dt = start + one_day; + println!("try {start:?} + 1 day -> {dt:?}"); + } + println!("selected {dt:?}"); + dbg!(dt.with_timezone(&Utc)) + } else { + dt + } + } + + pub fn is_within_schedule(&self, dt: DateTime) -> bool { + if let Some(start) = &self.not_before { + if dt < *start { + return false; + } + } + + if let Some(restrict) = &self.restriction { + let dt = dt.with_timezone(&restrict.timezone); + + let weekday: DaysOfWeek = dt.weekday().into(); + + if !restrict.days_of_week.contains(weekday) { + return false; + } + + let (start, end) = match restrict.start_end_on_day(dt) { + Some(result) => result, + None => return false, + }; + + if dt < start { + return false; + } + if dt >= end { + return false; + } + } + + true + } +} + +const DAYS: &[(&str, DaysOfWeek)] = &[ + ("Monday", DaysOfWeek::MON), + ("Tuesday", DaysOfWeek::TUE), + ("Wednesday", DaysOfWeek::WED), + ("Thursday", DaysOfWeek::THU), + ("Friday", DaysOfWeek::FRI), + ("Saturday", DaysOfWeek::SAT), + ("Sunday", DaysOfWeek::SUN), +]; + +impl FromStr for DaysOfWeek { + type Err = String; + fn from_str(s: &str) -> Result { + let mut days = DaysOfWeek::empty(); + 'next: for dow in s.split(',') { + for (label, value) in DAYS { + if dow.eq_ignore_ascii_case(label) || dow.eq_ignore_ascii_case(&label[0..3]) { + days.set(*value, true); + continue 'next; + } + } + return Err(format!("invalid day '{dow}'")); + } + + Ok(days) + } +} + +impl Serialize for DaysOfWeek { + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + let mut result = String::new(); + for (label, value) in DAYS { + if self.contains(*value) { + if !result.is_empty() { + result.push(','); + } + result.push_str(&label[0..3]); + } + } + serializer.serialize_str(&result) + } +} + +impl<'de> Deserialize<'de> for DaysOfWeek { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + struct Visit {} + impl<'de> serde::de::Visitor<'de> for Visit { + type Value = DaysOfWeek; + + fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result { + formatter.write_str("a comma separated list of days of the week like 'Mon,Tue'") + } + + fn visit_str(self, value: &str) -> Result + where + E: serde::de::Error, + { + value.parse::().map_err(|err| E::custom(err)) + } + } + + deserializer.deserialize_str(Visit {}) + } +} + +#[cfg(test)] +mod test { + use super::*; + + #[test] + fn days_of_week() { + let all = "Mon,Tue,Wed,Thu,Fri,Sat,Sun".parse::().unwrap(); + k9::assert_equal!( + all, + DaysOfWeek::MON + | DaysOfWeek::TUE + | DaysOfWeek::WED + | DaysOfWeek::THU + | DaysOfWeek::FRI + | DaysOfWeek::SAT + | DaysOfWeek::SUN + ); + + let middle = "Wed,Tue,Thursday".parse::().unwrap(); + k9::assert_equal!(middle, DaysOfWeek::TUE | DaysOfWeek::WED | DaysOfWeek::THU); + + k9::assert_equal!( + "Wed,Sumday".parse::().unwrap_err(), + "invalid day 'Sumday'" + ); + } + + #[test] + fn schedule_parse_restriction() { + let sched = Scheduling { + restriction: Some(ScheduleRestriction { + days_of_week: DaysOfWeek::MON | DaysOfWeek::WED, + timezone: "America/Phoenix".parse().unwrap(), + start: NaiveTime::from_hms_opt(9, 0, 0).unwrap(), + end: NaiveTime::from_hms_opt(17, 0, 0).unwrap(), + }), + not_before: None, + }; + + let serialized = serde_json::to_string(&sched).unwrap(); + k9::snapshot!( + &serialized, + r#"{"dow":"Mon,Wed","tz":"America/Phoenix","start":"09:00:00","end":"17:00:00"}"# + ); + + let round_trip: Scheduling = serde_json::from_str(&serialized).unwrap(); + k9::assert_equal!(sched, round_trip); + } + + #[test] + fn schedule_parse_restriction_and_start() { + let sched = Scheduling { + restriction: Some(ScheduleRestriction { + days_of_week: DaysOfWeek::MON | DaysOfWeek::WED, + timezone: "America/Phoenix".parse().unwrap(), + start: NaiveTime::from_hms_opt(9, 0, 0).unwrap(), + end: NaiveTime::from_hms_opt(17, 0, 0).unwrap(), + }), + not_before: DateTime::parse_from_rfc3339("1996-12-19T16:39:57-08:00").ok(), + }; + + let serialized = serde_json::to_string(&sched).unwrap(); + k9::snapshot!( + &serialized, + r#"{"dow":"Mon,Wed","tz":"America/Phoenix","start":"09:00:00","end":"17:00:00","not_before":"1996-12-19T16:39:57-08:00"}"# + ); + + let round_trip: Scheduling = serde_json::from_str(&serialized).unwrap(); + k9::assert_equal!(sched, round_trip); + } + + #[test] + fn schedule_parse_no_restriction_and_start() { + let sched = Scheduling { + restriction: None, + not_before: DateTime::parse_from_rfc3339("1996-12-19T16:39:57-08:00").ok(), + }; + + let serialized = serde_json::to_string(&sched).unwrap(); + k9::snapshot!(&serialized, r#"{"not_before":"1996-12-19T16:39:57-08:00"}"#); + + let round_trip: Scheduling = serde_json::from_str(&serialized).unwrap(); + k9::assert_equal!(sched, round_trip); + } + + #[test] + fn schedule_adjust_start() { + let sched = Scheduling { + restriction: None, + not_before: DateTime::parse_from_rfc3339("2023-03-20T16:39:57-08:00").ok(), + }; + + let now: DateTime = DateTime::parse_from_rfc3339("2023-03-20T08:00:00-08:00") + .unwrap() + .into(); + k9::assert_equal!(sched.adjust_for_schedule(now), sched.not_before.unwrap()); + } + + #[test] + fn schedule_adjust_dow() { + let phoenix: Tz = "America/Phoenix".parse().unwrap(); + let sched = Scheduling { + restriction: Some(ScheduleRestriction { + days_of_week: DaysOfWeek::MON | DaysOfWeek::WED, + timezone: phoenix.clone(), + start: NaiveTime::from_hms_opt(9, 0, 0).unwrap(), + end: NaiveTime::from_hms_opt(17, 0, 0).unwrap(), + }), + not_before: None, + }; + + // This is a Tuesday + let now: DateTime = DateTime::parse_from_rfc3339("2023-03-28T08:00:00-08:00") + .unwrap() + .into(); + + let adjusted = sched.adjust_for_schedule(now).with_timezone(&phoenix); + // Expected to round into wednesday, the next day + k9::assert_equal!(adjusted.to_string(), "2023-03-29 09:00:00 MST"); + } + + #[test] + fn schedule_adjust_dow_2() { + let phoenix: Tz = "America/Phoenix".parse().unwrap(); + let sched = Scheduling { + restriction: Some(ScheduleRestriction { + days_of_week: DaysOfWeek::MON | DaysOfWeek::FRI, + timezone: phoenix.clone(), + start: NaiveTime::from_hms_opt(9, 0, 0).unwrap(), + end: NaiveTime::from_hms_opt(17, 0, 0).unwrap(), + }), + not_before: None, + }; + + // This is a Monday, but after hours + let now: DateTime = DateTime::parse_from_rfc3339("2023-03-27T18:00:00-08:00") + .unwrap() + .into(); + + let adjusted = sched.adjust_for_schedule(now).with_timezone(&phoenix); + // Expected to round into Friday, later that week + k9::assert_equal!(adjusted.to_string(), "2023-03-31 09:00:00 MST"); + } +}