use std::sync::Arc;
use aion_core::{Event, RunId, WorkflowId, WorkflowSummary};
use aion_store::EventStore;
use aion_store::visibility::{
StreamHead, VisibilityRecord, VisibilityStore,
head::{self, RowVerdict},
};
use crate::EngineError;
#[must_use]
pub fn project_visibility(history: &[Event], run_id: &RunId) -> Option<VisibilityRecord> {
let window = run_window(history, run_id)?;
let summary = WorkflowSummary::from_history(window)?;
let search_attributes = aion_core::search_attributes_from_events(window);
let namespace = aion_core::namespace_from_attributes(&search_attributes);
let continued = summary.status == aion_core::WorkflowStatus::ContinuedAsNew;
let status = if continued {
aion_core::WorkflowStatus::Running
} else {
summary.status
};
let ended_at = if continued { None } else { summary.ended_at };
Some(VisibilityRecord {
namespace,
workflow_id: summary.workflow_id,
run_id: run_id.clone(),
workflow_type: summary.workflow_type,
status,
started_at: summary.started_at,
updated_at: summary.updated_at,
ended_at,
parent: summary.parent,
display_name: summary.display_name,
kind: summary.kind,
failed_step: summary.failed_step,
failure_reason: summary.failure_reason,
search_attributes,
outstanding_leases: aion_core::outstanding_leases(window),
package_version: summary.package_version,
head_seq: window.last().map_or(0, Event::seq),
})
}
#[must_use]
pub fn run_window<'history>(
history: &'history [Event],
run_id: &RunId,
) -> Option<&'history [Event]> {
let start = history.iter().position(|event| {
matches!(event, Event::WorkflowStarted { run_id: started, .. } if started == run_id)
})?;
let end = start.saturating_add(aion_core::run_segment(history, run_id).len());
Some(&history[..end])
}
#[must_use]
pub fn generation_run_ids(history: &[Event]) -> Vec<RunId> {
history
.iter()
.filter_map(|event| match event {
Event::WorkflowStarted { run_id, .. } => Some(run_id.clone()),
_ => None,
})
.collect()
}
#[must_use]
pub fn current_run_id(history: &[Event]) -> Option<RunId> {
history.iter().rev().find_map(|event| match event {
Event::WorkflowStarted { run_id, .. } => Some(run_id.clone()),
_ => None,
})
}
pub async fn upsert_workflow_visibility(
event_store: Arc<dyn EventStore>,
visibility_store: Arc<dyn VisibilityStore>,
workflow_id: &WorkflowId,
run_id: &RunId,
) -> Result<(), EngineError> {
let history = event_store.read_history(workflow_id).await?;
let record = project_visibility(&history, run_id).ok_or_else(|| EngineError::Load {
reason: format!(
"workflow `{workflow_id}` history has no WorkflowStarted event for visibility projection"
),
})?;
visibility_store.record_visibility(record).await?;
Ok(())
}
pub async fn reconcile_visibility(
event_store: Arc<dyn EventStore>,
visibility_store: Arc<dyn VisibilityStore>,
trigger: ReconcileTrigger,
) -> Result<(), EngineError> {
let started = std::time::Instant::now();
let mut streams = 0_usize;
let mut settled_by_row = 0_usize;
let mut histories_read = 0_usize;
let mut rows_written = 0_usize;
for StreamHead {
workflow_id,
head_seq,
} in event_store.stream_heads().await?
{
streams += 1;
let stored = visibility_store.get_visibility(&workflow_id).await?;
if head::verdict(stored.as_ref(), head_seq) != RowVerdict::Unsettled {
settled_by_row += 1;
continue;
}
histories_read += 1;
let history = event_store.read_history(&workflow_id).await?;
let generations = generation_run_ids(&history);
let Some(current) = generations.last().cloned() else {
return Err(EngineError::Load {
reason: format!(
"workflow `{workflow_id}` history has no WorkflowStarted event for \
visibility projection"
),
});
};
let Some(projected) = project_visibility(&history, ¤t) else {
return Err(EngineError::Load {
reason: format!(
"workflow `{workflow_id}` run `{current}` was started in its history but \
could not be projected"
),
});
};
if stored.as_ref() != Some(&projected) {
visibility_store.record_visibility(projected).await?;
rows_written += 1;
}
for run_id in generations {
if run_id != current {
visibility_store
.remove_visibility(&workflow_id, &run_id)
.await?;
}
}
}
let elapsed_ms = crate::engine::startup_telemetry::elapsed_ms(started);
match trigger {
ReconcileTrigger::Boot => tracing::info!(
streams,
settled_by_row,
histories_read,
rows_written,
elapsed_ms,
"visibility reconciled against history: rows at their stream head were trusted unread"
),
ReconcileTrigger::Periodic => tracing::debug!(
streams,
settled_by_row,
histories_read,
rows_written,
elapsed_ms,
"visibility reconciled against history: rows at their stream head were trusted unread"
),
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReconcileTrigger {
Boot,
Periodic,
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::error::Error;
use std::sync::Arc;
use aion_core::{
DISPLAY_NAME_ATTRIBUTE, Event, EventEnvelope, NAMESPACE_ATTRIBUTE, PackageVersion, Payload,
RunId, SearchAttributeValue, WorkflowError, WorkflowId, WorkflowStatus,
};
use aion_store::visibility::VisibilityStore;
use aion_store::{EventStore, InMemoryStore, WritableEventStore, WriteToken};
use chrono::{TimeZone, Utc};
use super::{
ReconcileTrigger, current_run_id, generation_run_ids, project_visibility,
reconcile_visibility,
};
type TestResult = Result<(), Box<dyn Error>>;
fn envelope(workflow_id: &WorkflowId, seq: u64) -> Result<EventEnvelope, Box<dyn Error>> {
let base = Utc
.with_ymd_and_hms(2026, 1, 1, 0, 0, 0)
.single()
.ok_or("test timestamp should be unambiguous")?;
Ok(EventEnvelope {
seq,
recorded_at: base + chrono::Duration::seconds(i64::try_from(seq)?),
workflow_id: workflow_id.clone(),
})
}
fn envelope_at(workflow_id: &WorkflowId, seq: u64) -> Result<EventEnvelope, Box<dyn Error>> {
envelope(workflow_id, seq)
}
fn payload() -> Result<Payload, Box<dyn Error>> {
Ok(Payload::from_json(&serde_json::json!({}))?)
}
fn workflow_started(
workflow_id: &WorkflowId,
run_id: &RunId,
parent: Option<WorkflowId>,
) -> Result<Event, Box<dyn Error>> {
Ok(Event::WorkflowStarted {
envelope: envelope(workflow_id, 1)?,
workflow_type: String::from("order_processing"),
input: payload()?,
run_id: run_id.clone(),
parent_run_id: None,
parent_workflow_id: parent,
package_version: PackageVersion::new("a".repeat(64)),
})
}
fn attributes(
workflow_id: &WorkflowId,
seq: u64,
pairs: &[(&str, &str)],
) -> Result<Event, Box<dyn Error>> {
Ok(Event::SearchAttributesUpdated {
envelope: envelope(workflow_id, seq)?,
workflow_id: workflow_id.clone(),
attributes: pairs
.iter()
.map(|(key, value)| {
(
(*key).to_owned(),
SearchAttributeValue::String((*value).to_owned()),
)
})
.collect::<HashMap<_, _>>(),
})
}
#[test]
fn a_history_without_a_start_projects_nothing() -> TestResult {
let wf_id = WorkflowId::new_v4();
let orphan = vec![Event::WorkflowCompleted {
envelope: envelope(&wf_id, 1)?,
result: payload()?,
}];
assert!(project_visibility(&[], &RunId::new_v4()).is_none());
assert!(project_visibility(&orphan, &RunId::new_v4()).is_none());
assert!(current_run_id(&orphan).is_none());
Ok(())
}
#[test]
fn a_running_history_projects_every_field() -> TestResult {
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
let parent = WorkflowId::new_v4();
let history = vec![
workflow_started(&wf_id, &run_id, Some(parent.clone()))?,
attributes(
&wf_id,
2,
&[
(NAMESPACE_ATTRIBUTE, "tenant-a"),
(DISPLAY_NAME_ATTRIBUTE, "Nightly close"),
("region", "eu-west-1"),
],
)?,
Event::SignalReceived {
envelope: envelope(&wf_id, 3)?,
name: String::from("wake"),
payload: payload()?,
},
];
let record = project_visibility(&history, &run_id).ok_or("a started history projects")?;
assert_eq!(record.namespace, "tenant-a");
assert_eq!(record.workflow_id, wf_id);
assert_eq!(record.run_id, run_id);
assert_eq!(record.workflow_type, "order_processing");
assert_eq!(record.status, WorkflowStatus::Running);
assert_eq!(record.started_at, envelope(&wf_id, 1)?.recorded_at);
assert_eq!(
record.updated_at,
envelope(&wf_id, 3)?.recorded_at,
"updated_at is the LAST event of any kind, not the last lifecycle event"
);
assert_eq!(record.ended_at, None);
assert_eq!(record.parent, Some(parent));
assert_eq!(record.display_name.as_deref(), Some("Nightly close"));
assert_eq!(record.kind, None);
assert_eq!(record.failed_step, None);
assert_eq!(record.failure_reason, None);
assert_eq!(
record.package_version,
Some(PackageVersion::new("a".repeat(64))),
"the row carries the hash of the package the run started under"
);
assert_eq!(
record.search_attributes.get("region"),
Some(&SearchAttributeValue::String(String::from("eu-west-1")))
);
assert_eq!(current_run_id(&history), Some(run_id));
Ok(())
}
#[test]
fn a_failed_history_projects_the_terminal_and_an_unplaced_run_is_in_the_default_namespace()
-> TestResult {
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
let terminal = envelope(&wf_id, 2)?;
let ended = terminal.recorded_at;
let history = vec![
workflow_started(&wf_id, &run_id, None)?,
Event::WorkflowFailed {
envelope: terminal,
error: WorkflowError {
message: String::from("boom"),
details: None,
},
},
];
let record = project_visibility(&history, &run_id).ok_or("a started history projects")?;
assert_eq!(record.namespace, aion_core::DEFAULT_NAMESPACE);
assert_eq!(record.status, WorkflowStatus::Failed);
assert_eq!(record.ended_at, Some(ended));
assert_eq!(record.updated_at, ended);
assert_eq!(record.failure_reason.as_deref(), Some("boom"));
Ok(())
}
#[test]
fn a_reopened_history_has_no_end_and_the_current_run_is_the_latest_start() -> TestResult {
let wf_id = WorkflowId::new_v4();
let first = RunId::new_v4();
let history = vec![
workflow_started(&wf_id, &first, None)?,
Event::WorkflowCompleted {
envelope: envelope(&wf_id, 2)?,
result: payload()?,
},
Event::WorkflowReopened {
envelope: envelope(&wf_id, 3)?,
run_id: first.clone(),
reopened: Vec::new(),
},
];
let record = project_visibility(&history, &first).ok_or("a started history projects")?;
assert_eq!(record.status, WorkflowStatus::Running);
assert_eq!(record.ended_at, None);
assert_eq!(record.updated_at, envelope(&wf_id, 3)?.recorded_at);
assert_eq!(current_run_id(&history), Some(first));
Ok(())
}
#[test]
fn a_continued_history_projects_each_generation_from_its_own_window() -> TestResult {
let wf_id = WorkflowId::new_v4();
let first = RunId::new_v4();
let second = RunId::new_v4();
let mut successor = workflow_started(&wf_id, &second, None)?;
if let Event::WorkflowStarted { envelope, .. } = &mut successor {
*envelope = envelope_at(&wf_id, 4)?;
}
let history = vec![
workflow_started(&wf_id, &first, None)?,
attributes(
&wf_id,
2,
&[
(NAMESPACE_ATTRIBUTE, "team-a"),
(DISPLAY_NAME_ATTRIBUTE, "Disk reaper"),
],
)?,
Event::WorkflowContinuedAsNew {
envelope: envelope(&wf_id, 3)?,
input: payload()?,
workflow_type: None,
parent_run_id: first.clone(),
},
successor,
];
let closed = project_visibility(&history, &first).ok_or("the predecessor projects")?;
assert_eq!(closed.run_id, first);
assert_eq!(closed.status, WorkflowStatus::Running);
assert_eq!(closed.ended_at, None);
assert_eq!(closed.updated_at, envelope(&wf_id, 3)?.recorded_at);
assert_eq!(closed.started_at, envelope(&wf_id, 1)?.recorded_at);
let live = project_visibility(&history, &second).ok_or("the successor projects")?;
assert_eq!(live.run_id, second);
assert_eq!(live.status, WorkflowStatus::Running);
assert_eq!(live.ended_at, None);
assert_eq!(live.started_at, envelope(&wf_id, 4)?.recorded_at);
assert_eq!(
live.display_name.as_deref(),
Some("Disk reaper"),
"an attribute the first generation set still names the second"
);
assert_eq!(live.namespace, closed.namespace);
assert!(
project_visibility(&history, &RunId::new_v4()).is_none(),
"a run this history never started projects nothing"
);
assert_eq!(generation_run_ids(&history), vec![first, second]);
Ok(())
}
#[tokio::test]
async fn reconcile_rights_every_generation_of_a_continued_history() -> TestResult {
let events = Arc::new(InMemoryStore::default());
let visibility: Arc<dyn VisibilityStore> = Arc::clone(&events) as Arc<dyn VisibilityStore>;
let wf_id = WorkflowId::new_v4();
let first = RunId::new_v4();
let second = RunId::new_v4();
let mut successor = workflow_started(&wf_id, &second, None)?;
if let Event::WorkflowStarted { envelope, .. } = &mut successor {
*envelope = envelope_at(&wf_id, 3)?;
}
let history = vec![
workflow_started(&wf_id, &first, None)?,
Event::WorkflowContinuedAsNew {
envelope: envelope(&wf_id, 2)?,
input: payload()?,
workflow_type: None,
parent_run_id: first.clone(),
},
successor,
];
events
.append(WriteToken::recorder(), &wf_id, &history, 0)
.await?;
let mut phantom = project_visibility(&history, &first).ok_or("projects")?;
phantom.status = WorkflowStatus::Running;
phantom.ended_at = None;
visibility.record_visibility(phantom).await?;
reconcile_visibility(
Arc::clone(&events) as Arc<dyn EventStore>,
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
let live = visibility
.get_visibility(&wf_id)
.await?
.ok_or("the workflow keeps its one row")?;
assert_eq!(live.run_id, second);
assert_eq!(live.status, WorkflowStatus::Running);
Ok(())
}
#[test]
fn the_projection_stamps_the_row_at_the_stream_head() -> TestResult {
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
let history = vec![
workflow_started(&wf_id, &run_id, None)?,
Event::WorkflowCompleted {
envelope: envelope(&wf_id, 2)?,
result: payload()?,
},
];
let row = project_visibility(&history, &run_id).ok_or("the run projects")?;
assert_eq!(
row.head_seq, 2,
"head_seq is the seq of the last event folded"
);
Ok(())
}
#[test]
fn a_superseded_generation_projects_behind_the_stream_head() -> TestResult {
let wf_id = WorkflowId::new_v4();
let first = RunId::new_v4();
let second = RunId::new_v4();
let mut successor = workflow_started(&wf_id, &second, None)?;
if let Event::WorkflowStarted { envelope, .. } = &mut successor {
*envelope = envelope_at(&wf_id, 3)?;
}
let history = vec![
workflow_started(&wf_id, &first, None)?,
Event::WorkflowContinuedAsNew {
envelope: envelope(&wf_id, 2)?,
input: payload()?,
workflow_type: None,
parent_run_id: first.clone(),
},
successor,
];
let old = project_visibility(&history, &first).ok_or("the old run projects")?;
let live = project_visibility(&history, &second).ok_or("the live run projects")?;
assert_eq!(
old.head_seq, 2,
"an old generation's row ends at its own window"
);
assert_eq!(
live.head_seq, 3,
"the live generation's row stands at the stream head"
);
Ok(())
}
#[tokio::test]
async fn reconcile_believes_a_row_at_its_stream_head_without_reading_history() -> TestResult {
let backing = Arc::new(InMemoryStore::default());
let events: Arc<dyn EventStore> = Arc::clone(&backing) as Arc<dyn EventStore>;
let visibility: Arc<dyn VisibilityStore> = Arc::clone(&backing) as Arc<dyn VisibilityStore>;
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
backing
.append(
WriteToken::recorder(),
&wf_id,
&[workflow_started(&wf_id, &run_id, None)?],
0,
)
.await?;
let history = events.read_history(&wf_id).await?;
let mut lying = project_visibility(&history, &run_id).ok_or("the run projects")?;
lying.status = WorkflowStatus::Completed;
assert_eq!(lying.head_seq, 1);
visibility.record_visibility(lying.clone()).await?;
reconcile_visibility(
Arc::clone(&events),
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
assert_eq!(
visibility.get_visibility(&wf_id).await?,
Some(lying.clone()),
"a row at its stream head is trusted; the history behind it is not opened"
);
backing
.append(
WriteToken::recorder(),
&wf_id,
&[Event::SearchAttributesUpdated {
envelope: envelope(&wf_id, 2)?,
workflow_id: wf_id.clone(),
attributes: HashMap::new(),
}],
1,
)
.await?;
reconcile_visibility(
Arc::clone(&events),
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
let healed = visibility
.get_visibility(&wf_id)
.await?
.ok_or("the row survives reconcile")?;
assert_eq!(
healed.status,
WorkflowStatus::Running,
"the fold corrects the lie"
);
assert_eq!(healed.head_seq, 2, "the healed row stands at the new head");
Ok(())
}
#[tokio::test]
async fn reconcile_re_derives_an_unstamped_row() -> TestResult {
let backing = Arc::new(InMemoryStore::default());
let events: Arc<dyn EventStore> = Arc::clone(&backing) as Arc<dyn EventStore>;
let visibility: Arc<dyn VisibilityStore> = Arc::clone(&backing) as Arc<dyn VisibilityStore>;
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
backing
.append(
WriteToken::recorder(),
&wf_id,
&[workflow_started(&wf_id, &run_id, None)?],
0,
)
.await?;
let history = events.read_history(&wf_id).await?;
let mut legacy = project_visibility(&history, &run_id).ok_or("the run projects")?;
legacy.head_seq = 0;
legacy.status = WorkflowStatus::Completed;
visibility.record_visibility(legacy).await?;
reconcile_visibility(events, Arc::clone(&visibility), ReconcileTrigger::Boot).await?;
let stamped = visibility
.get_visibility(&wf_id)
.await?
.ok_or("the row survives reconcile")?;
assert_eq!(
stamped.status,
WorkflowStatus::Running,
"an unstamped row is re-derived, never trusted"
);
assert_eq!(
stamped.head_seq, 1,
"the first reconcile after upgrade stamps the row"
);
Ok(())
}
#[tokio::test]
async fn reconcile_writes_only_rows_that_differ_from_history() -> TestResult {
let backing = Arc::new(InMemoryStore::default());
let events: Arc<dyn EventStore> = Arc::clone(&backing) as Arc<dyn EventStore>;
let visibility: Arc<dyn VisibilityStore> = Arc::clone(&backing) as Arc<dyn VisibilityStore>;
let wf_id = WorkflowId::new_v4();
let run_id = RunId::new_v4();
backing
.append(
WriteToken::recorder(),
&wf_id,
&[workflow_started(&wf_id, &run_id, None)?],
0,
)
.await?;
reconcile_visibility(
Arc::clone(&events),
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
let row = visibility
.get_visibility(&wf_id)
.await?
.ok_or("reconcile writes the missing row")?;
assert_eq!(row.status, WorkflowStatus::Running);
backing
.append(
WriteToken::recorder(),
&wf_id,
&[Event::WorkflowCompleted {
envelope: envelope(&wf_id, 2)?,
result: payload()?,
}],
1,
)
.await?;
reconcile_visibility(
Arc::clone(&events),
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
let row = visibility
.get_visibility(&wf_id)
.await?
.ok_or("the row survives reconcile")?;
assert_eq!(row.status, WorkflowStatus::Completed);
assert_eq!(row.ended_at, Some(envelope(&wf_id, 2)?.recorded_at));
assert_eq!(row.updated_at, envelope(&wf_id, 2)?.recorded_at);
let before = row.clone();
reconcile_visibility(
Arc::clone(&events),
Arc::clone(&visibility),
ReconcileTrigger::Boot,
)
.await?;
assert_eq!(
visibility.get_visibility(&wf_id).await?,
Some(before.clone())
);
let mut unstamped = before.clone();
unstamped.head_seq = 0;
visibility.record_visibility(unstamped).await?;
reconcile_visibility(events, Arc::clone(&visibility), ReconcileTrigger::Boot).await?;
assert_eq!(
visibility.get_visibility(&wf_id).await?,
Some(before),
"the comparison sees the stamp and rewrites the row at its head"
);
Ok(())
}
}