use std::time::Duration;
use serde::Deserialize;
#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum QueueServicePolicy {
#[default]
Strict,
DurablePending,
}
impl QueueServicePolicy {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Strict => "strict",
Self::DurablePending => "durable_pending",
}
}
}
impl std::fmt::Display for QueueServicePolicy {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct QueueServiceOverride {
#[serde(default)]
pub namespace: Option<String>,
pub task_queue: String,
pub policy: QueueServicePolicy,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq)]
#[serde(default, deny_unknown_fields)]
pub struct QueueServiceConfig {
pub default_policy: QueueServicePolicy,
pub overrides: Vec<QueueServiceOverride>,
#[serde(default, with = "optional_duration_millis")]
pub service_availability_deadline: Option<Duration>,
#[serde(default, with = "optional_duration_millis")]
pub schedule_to_start_timeout: Option<Duration>,
}
impl QueueServiceConfig {
#[must_use]
pub fn policy_for(&self, namespace: &str, task_queue: &str) -> QueueServicePolicy {
self.overrides
.iter()
.find(|entry| {
entry.task_queue == task_queue
&& entry
.namespace
.as_ref()
.is_none_or(|scoped| scoped == namespace)
})
.map_or(self.default_policy, |entry| entry.policy)
}
#[must_use]
pub const fn availability_deadline_for(&self, policy: QueueServicePolicy) -> Option<Duration> {
match policy {
QueueServicePolicy::Strict => self.service_availability_deadline,
QueueServicePolicy::DurablePending => None,
}
}
}
mod optional_duration_millis {
use std::time::Duration;
use serde::{Deserialize, Deserializer};
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<Option<Duration>, D::Error>
where
D: Deserializer<'de>,
{
let millis = Option::<u64>::deserialize(deserializer)?;
Ok(millis.map(Duration::from_millis))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn strict_is_the_default_policy_and_no_clock_is_invented() {
let config = QueueServiceConfig::default();
assert_eq!(config.default_policy, QueueServicePolicy::Strict);
assert_eq!(
config.policy_for("default", "general"),
QueueServicePolicy::Strict
);
assert_eq!(config.service_availability_deadline, None);
assert_eq!(config.schedule_to_start_timeout, None);
}
#[test]
fn a_written_override_opts_one_queue_into_durable_pending() {
let config = QueueServiceConfig {
overrides: vec![QueueServiceOverride {
namespace: None,
task_queue: "general".to_owned(),
policy: QueueServicePolicy::DurablePending,
}],
..QueueServiceConfig::default()
};
assert_eq!(
config.policy_for("default", "general"),
QueueServicePolicy::DurablePending
);
assert_eq!(
config.policy_for("default", "billing"),
QueueServicePolicy::Strict
);
}
#[test]
fn an_override_can_be_scoped_to_one_namespace() {
let config = QueueServiceConfig {
overrides: vec![QueueServiceOverride {
namespace: Some("lab".to_owned()),
task_queue: "general".to_owned(),
policy: QueueServicePolicy::DurablePending,
}],
..QueueServiceConfig::default()
};
assert_eq!(
config.policy_for("lab", "general"),
QueueServicePolicy::DurablePending
);
assert_eq!(
config.policy_for("prod", "general"),
QueueServicePolicy::Strict
);
}
#[test]
fn the_availability_deadline_bounds_strict_only() {
let config = QueueServiceConfig {
service_availability_deadline: Some(Duration::from_secs(30)),
..QueueServiceConfig::default()
};
assert_eq!(
config.availability_deadline_for(QueueServicePolicy::Strict),
Some(Duration::from_secs(30))
);
assert_eq!(
config.availability_deadline_for(QueueServicePolicy::DurablePending),
None
);
}
#[test]
fn queue_service_settings_deserialize_from_the_written_form() -> Result<(), toml::de::Error> {
let config: QueueServiceConfig = toml::from_str(
r#"
default_policy = "strict"
service_availability_deadline = 45000
schedule_to_start_timeout = 5000
[[overrides]]
task_queue = "general"
policy = "durable_pending"
"#,
)?;
assert_eq!(config.default_policy, QueueServicePolicy::Strict);
assert_eq!(
config.service_availability_deadline,
Some(Duration::from_secs(45))
);
assert_eq!(
config.schedule_to_start_timeout,
Some(Duration::from_secs(5))
);
assert_eq!(
config.policy_for("default", "general"),
QueueServicePolicy::DurablePending
);
Ok(())
}
}