use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{CorrelationKey, EffectKey, RunId, Timestamp};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct InboundEvent {
pub source: String,
pub id: String,
pub kind: String,
pub correlation: Vec<CorrelationKey>,
pub payload: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub by: Option<crate::core::Operator>,
}
impl InboundEvent {
pub fn new(
source: impl Into<String>,
id: impl Into<String>,
kind: impl Into<String>,
payload: Value,
) -> Self {
Self {
source: source.into(),
id: id.into(),
kind: kind.into(),
correlation: Vec::new(),
payload,
by: None,
}
}
#[must_use]
pub fn correlate(mut self, key: CorrelationKey) -> Self {
self.correlation.push(key);
self
}
#[must_use]
pub fn minted_by(mut self, by: crate::core::Operator) -> Self {
self.by = Some(by);
self
}
#[must_use]
pub fn dedup_key(&self) -> String {
origin_key(&self.source, &self.id)
}
}
pub(crate) fn is_reserved_kind(kind: &str) -> bool {
kind.starts_with("agentplane.")
}
#[must_use]
pub fn origin_key(source: &str, id: &str) -> String {
format!("{}\u{1f}{source}\u{1f}{id}", source.len())
}
#[must_use]
pub fn origin_source(key: &str) -> Option<&str> {
let (len, rest) = key.split_once('\u{1f}')?;
if len.is_empty() || !len.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
let len: usize = len.parse().ok()?;
let source = rest.get(..len)?;
rest[len..].starts_with('\u{1f}').then_some(source)
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AwaitSpec {
pub kind: String,
pub correlation: Vec<CorrelationKey>,
pub deadline: String,
}
impl AwaitSpec {
pub fn new(kind: impl Into<String>, deadline: impl Into<String>) -> Self {
Self {
kind: kind.into(),
correlation: Vec::new(),
deadline: deadline.into(),
}
}
#[must_use]
pub fn correlate(mut self, key: CorrelationKey) -> Self {
self.correlation.push(key);
self
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Subscription {
pub run: RunId,
pub case: Option<crate::core::CaseId>,
pub effect: EffectKey,
pub step: crate::core::StepId,
pub phase: crate::core::Phase,
pub kind: String,
pub correlation: Vec<CorrelationKey>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Timer {
pub run: crate::core::RunId,
pub case: Option<crate::core::CaseId>,
pub effect: EffectKey,
pub step: crate::core::StepId,
pub phase: crate::core::Phase,
pub fire_at: Timestamp,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "reason", rename_all = "snake_case", deny_unknown_fields)]
#[non_exhaustive]
pub enum SuspendReason {
AwaitingEvent {
kind: String,
correlation: Vec<CorrelationKey>,
#[serde(with = "time::serde::rfc3339")]
until: Timestamp,
},
AwaitingTime {
#[serde(with = "time::serde::rfc3339")]
until: Timestamp,
},
}
impl SuspendReason {
#[must_use]
pub const fn until(&self) -> Timestamp {
match self {
Self::AwaitingEvent { until, .. } | Self::AwaitingTime { until } => *until,
}
}
}
impl std::fmt::Display for SuspendReason {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::AwaitingEvent {
kind,
correlation,
until,
..
} => {
write!(f, "awaiting '{kind}'")?;
if let Some(k) = correlation.first() {
write!(f, " for {k}")?;
}
write!(f, " until {until}")
}
Self::AwaitingTime { until } => write!(f, "sleeping until {until}"),
}
}
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum Delivery {
Resumed { run: RunId },
Buffered,
Duplicate,
}
impl Delivery {
#[must_use]
pub fn resumed_run(&self) -> Option<RunId> {
match self {
Self::Resumed { run } => Some(*run),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DeadLetter {
pub event: InboundEvent,
#[serde(with = "time::serde::rfc3339")]
pub received_at: Timestamp,
pub reason: String,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn only_a_resumed_delivery_names_a_run() {
let run = RunId::generate();
assert_eq!(Delivery::Resumed { run }.resumed_run(), Some(run));
assert_eq!(Delivery::Buffered.resumed_run(), None);
assert_eq!(Delivery::Duplicate.resumed_run(), None);
}
#[test]
fn an_origin_keys_source_reads_back_exactly() {
for (source, id) in [
("peer:a", "m1"),
("a\u{1f}b", "c"),
("a", "b\u{1f}c"),
("", "x"),
] {
assert_eq!(origin_source(&origin_key(source, id)), Some(source));
}
for key in [
"",
"m1",
"x\u{1f}peer:a\u{1f}m1",
"9\u{1f}peer:a\u{1f}m1",
"6\u{1f}peer:a",
] {
assert_eq!(origin_source(key), None, "{key:?} is not an origin key");
}
}
}