use aion_core::{WorkflowId, WorkflowStatus};
use super::VisibilityRecord;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamHead {
pub workflow_id: WorkflowId,
pub head_seq: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RowVerdict {
Finished,
Paused,
InFlight,
Unsettled,
}
#[must_use]
pub fn verdict(row: Option<&VisibilityRecord>, head_seq: u64) -> RowVerdict {
let Some(row) = row else {
return RowVerdict::Unsettled;
};
if row.head_seq == 0 || row.head_seq != head_seq {
return RowVerdict::Unsettled;
}
match row.status {
WorkflowStatus::Completed
| WorkflowStatus::Failed
| WorkflowStatus::Cancelled
| WorkflowStatus::TimedOut
| WorkflowStatus::ContinuedAsNew => RowVerdict::Finished,
WorkflowStatus::Paused => RowVerdict::Paused,
WorkflowStatus::Running => RowVerdict::InFlight,
}
}
#[must_use]
pub const fn list_for(status: WorkflowStatus) -> Option<ListMembership> {
match status {
WorkflowStatus::Running => Some(ListMembership::Active),
WorkflowStatus::Paused => Some(ListMembership::Paused),
WorkflowStatus::Completed
| WorkflowStatus::Failed
| WorkflowStatus::Cancelled
| WorkflowStatus::TimedOut
| WorkflowStatus::ContinuedAsNew => None,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ListMembership {
Active,
Paused,
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use aion_core::{RunId, WorkflowId, WorkflowStatus};
use chrono::Utc;
use super::{ListMembership, RowVerdict, list_for, verdict};
use crate::visibility::VisibilityRecord;
fn row(status: WorkflowStatus, head_seq: u64) -> VisibilityRecord {
VisibilityRecord {
namespace: String::from("default"),
workflow_id: WorkflowId::new(uuid::Uuid::from_u128(7)),
run_id: RunId::new(uuid::Uuid::from_u128(8)),
workflow_type: String::from("checkout"),
status,
started_at: Utc::now(),
updated_at: Utc::now(),
ended_at: None,
parent: None,
display_name: None,
kind: None,
failed_step: None,
failure_reason: None,
search_attributes: HashMap::new(),
outstanding_leases: Vec::new(),
package_version: None,
head_seq,
}
}
#[test]
fn a_terminal_row_at_the_head_is_finished_and_unread() {
for status in [
WorkflowStatus::Completed,
WorkflowStatus::Failed,
WorkflowStatus::Cancelled,
WorkflowStatus::TimedOut,
] {
assert_eq!(verdict(Some(&row(status, 5)), 5), RowVerdict::Finished);
}
}
#[test]
fn a_paused_row_at_the_head_is_paused_and_unread() {
assert_eq!(
verdict(Some(&row(WorkflowStatus::Paused, 3)), 3),
RowVerdict::Paused
);
}
#[test]
fn a_running_row_at_the_head_is_in_flight_and_the_fold_decides() {
assert_eq!(
verdict(Some(&row(WorkflowStatus::Running, 3)), 3),
RowVerdict::InFlight
);
}
#[test]
fn a_row_behind_or_ahead_of_the_head_settles_nothing() {
assert_eq!(
verdict(Some(&row(WorkflowStatus::Completed, 4)), 5),
RowVerdict::Unsettled
);
assert_eq!(
verdict(Some(&row(WorkflowStatus::Completed, 6)), 5),
RowVerdict::Unsettled
);
}
#[test]
fn an_unstamped_row_is_never_at_the_head() {
assert_eq!(
verdict(Some(&row(WorkflowStatus::Completed, 0)), 0),
RowVerdict::Unsettled
);
}
#[test]
fn no_row_settles_nothing() {
assert_eq!(verdict(None, 5), RowVerdict::Unsettled);
}
#[test]
fn the_two_lists_partition_the_fold_exactly() {
assert_eq!(
list_for(WorkflowStatus::Running),
Some(ListMembership::Active)
);
assert_eq!(
list_for(WorkflowStatus::Paused),
Some(ListMembership::Paused)
);
for status in [
WorkflowStatus::Completed,
WorkflowStatus::Failed,
WorkflowStatus::Cancelled,
WorkflowStatus::TimedOut,
WorkflowStatus::ContinuedAsNew,
] {
assert_eq!(list_for(status), None);
}
}
}