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 {
format!("{}\u{1f}{}", self.source, self.id)
}
}
#[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")]
#[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);
}
}