use aion_core::{WorkflowStatus, status_from_events};
use aion_store::{EventStore, OutboxStore, StoreError};
use tracing::info;
#[must_use]
pub fn is_settle_terminal(status: WorkflowStatus) -> bool {
matches!(
status,
WorkflowStatus::Completed
| WorkflowStatus::Failed
| WorkflowStatus::Cancelled
| WorkflowStatus::TimedOut
)
}
pub async fn settle_terminal_outbox_rows(
event_store: &dyn EventStore,
outbox_store: &dyn OutboxStore,
) -> Result<Vec<String>, StoreError> {
let workflow_ids = outbox_store.list_unsettled_outbox_workflow_ids().await?;
let mut settled = Vec::new();
for workflow_id in workflow_ids {
let history = event_store.read_history(&workflow_id).await?;
let status = status_from_events(&history);
if !is_settle_terminal(status) {
continue;
}
let keys = outbox_store
.cancel_outbox_rows_for_workflow(&workflow_id)
.await?;
if !keys.is_empty() {
info!(
workflow_id = %workflow_id,
projected_status = ?status,
settled = keys.len(),
dispatch_keys = ?keys,
"settled outbox rows for terminal workflow"
);
settled.extend(keys);
}
}
Ok(settled)
}
#[cfg(test)]
mod tests {
use super::is_settle_terminal;
use aion_core::WorkflowStatus;
#[test]
fn only_hard_terminals_are_settle_eligible() {
assert!(is_settle_terminal(WorkflowStatus::Completed));
assert!(is_settle_terminal(WorkflowStatus::Failed));
assert!(is_settle_terminal(WorkflowStatus::Cancelled));
assert!(is_settle_terminal(WorkflowStatus::TimedOut));
assert!(!is_settle_terminal(WorkflowStatus::Running));
assert!(!is_settle_terminal(WorkflowStatus::Paused));
assert!(!is_settle_terminal(WorkflowStatus::ContinuedAsNew));
}
}