use chrono::{DateTime, Utc};
use crate::{ActivityId, Event, describe_live::OpenStep, describe_live::StepState};
#[derive(Clone, Debug)]
struct Open {
step: OpenStep,
last_advanced_at: usize,
}
#[must_use]
pub fn open_steps(history: &[Event]) -> Vec<OpenStep> {
fold(history).into_iter().map(|open| open.step).collect()
}
#[must_use]
pub fn current_step(history: &[Event]) -> Option<OpenStep> {
fold(history)
.into_iter()
.max_by_key(|open| open.last_advanced_at)
.map(|open| open.step)
}
fn fold(history: &[Event]) -> Vec<Open> {
let segment_start = history
.iter()
.rposition(|event| matches!(event, Event::WorkflowStarted { .. }))
.unwrap_or(0);
let mut open: Vec<Open> = Vec::new();
for (position, event) in history[segment_start..].iter().enumerate() {
match event {
Event::ActivityScheduled {
envelope,
activity_id,
activity_type,
task_queue,
node,
..
} => {
let entry = Open {
step: OpenStep {
activity_id: activity_id.clone(),
activity_type: activity_type.clone(),
task_queue: task_queue.clone(),
node: node.clone(),
scheduled_at: envelope.recorded_at,
state: StepState::Scheduled,
},
last_advanced_at: position,
};
replace_or_push(&mut open, entry);
}
Event::ActivityStarted {
envelope,
activity_id,
attempt,
} => {
if let Some(held) = find_mut(&mut open, activity_id) {
held.step.state = StepState::Dispatched {
attempt: *attempt,
dispatched_at: envelope.recorded_at,
};
held.last_advanced_at = position;
}
}
Event::ActivityCompleted { activity_id, .. }
| Event::ActivityFailed { activity_id, .. }
| Event::ActivityCancelled { activity_id, .. } => {
open.retain(|held| &held.step.activity_id != activity_id);
}
Event::WorkflowReopened {
envelope, reopened, ..
} => reopen(
&mut open,
history,
segment_start,
reopened,
envelope.recorded_at,
position,
),
_ => {}
}
}
open
}
fn reopen(
open: &mut Vec<Open>,
history: &[Event],
segment_start: usize,
reopened: &[ActivityId],
recorded_at: DateTime<Utc>,
position: usize,
) {
for activity_id in reopened {
let Some(scheduled) = last_scheduling(&history[segment_start..], activity_id) else {
continue;
};
let entry = Open {
step: OpenStep {
state: StepState::Reopened {
reopened_at: recorded_at,
},
..scheduled
},
last_advanced_at: position,
};
replace_or_push(open, entry);
}
}
fn last_scheduling(segment: &[Event], activity_id: &ActivityId) -> Option<OpenStep> {
segment.iter().rev().find_map(|event| match event {
Event::ActivityScheduled {
envelope,
activity_id: scheduled_id,
activity_type,
task_queue,
node,
..
} if scheduled_id == activity_id => Some(OpenStep {
activity_id: scheduled_id.clone(),
activity_type: activity_type.clone(),
task_queue: task_queue.clone(),
node: node.clone(),
scheduled_at: envelope.recorded_at,
state: StepState::Scheduled,
}),
_ => None,
})
}
fn replace_or_push(open: &mut Vec<Open>, entry: Open) {
match find_mut(open, &entry.step.activity_id) {
Some(held) => *held = entry,
None => open.push(entry),
}
}
fn find_mut<'a>(open: &'a mut [Open], activity_id: &ActivityId) -> Option<&'a mut Open> {
open.iter_mut()
.find(|held| &held.step.activity_id == activity_id)
}
#[cfg(test)]
#[path = "current_step_tests.rs"]
mod tests;