use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::{Mutex, MutexGuard};
use chrono::{DateTime, Utc};
use crate::events::ExecutionEvent;
use crate::orchestration::state::{
ActivityState, ChangeRuntimeState, OrchestratorState, QueueIntent, TerminalState, WaitState,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum ExecutionPhase {
Preparing,
Apply,
Acceptance,
RejectionReview,
Archive,
Resolve,
Push,
Merge,
#[default]
None,
Unknown,
}
impl ExecutionPhase {
pub fn as_str(self) -> &'static str {
match self {
Self::Preparing => "preparing",
Self::Apply => "apply",
Self::Acceptance => "acceptance",
Self::RejectionReview => "rejection_review",
Self::Archive => "archive",
Self::Resolve => "resolve",
Self::Push => "push",
Self::Merge => "merge",
Self::None => "none",
Self::Unknown => "unknown",
}
}
pub fn is_active(self) -> bool {
!matches!(self, Self::None | Self::Unknown)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ChangeExecutionState {
Queued,
Active,
Waiting,
Stopping,
Stopped,
Failed,
Completed,
#[default]
Unknown,
}
impl ChangeExecutionState {
#[allow(dead_code)] pub fn as_str(self) -> &'static str {
match self {
Self::Queued => "queued",
Self::Active => "active",
Self::Waiting => "waiting",
Self::Stopping => "stopping",
Self::Stopped => "stopped",
Self::Failed => "failed",
Self::Completed => "completed",
Self::Unknown => "unknown",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ProcessActivity {
DependencyAnalysis,
BaseBranchMerge,
ConflictResolution,
BranchMerge,
WorkspaceCleanup,
}
impl ProcessActivity {
pub fn as_str(self) -> &'static str {
match self {
Self::DependencyAnalysis => "dependency_analysis",
Self::BaseBranchMerge => "base_branch_merge",
Self::ConflictResolution => "conflict_resolution",
Self::BranchMerge => "branch_merge",
Self::WorkspaceCleanup => "workspace_cleanup",
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ChangeExecutionFacts {
pub execution_id: Option<String>,
pub execution_state: ChangeExecutionState,
pub current_phase: ExecutionPhase,
pub phase_started_at: Option<DateTime<Utc>>,
pub last_completed_phase: Option<ExecutionPhase>,
pub last_completed_at: Option<DateTime<Utc>>,
pub apply_commit_oid: Option<String>,
}
impl ChangeExecutionFacts {
pub fn unknown() -> Self {
Self {
execution_id: None,
execution_state: ChangeExecutionState::Unknown,
current_phase: ExecutionPhase::Unknown,
phase_started_at: None,
last_completed_phase: None,
last_completed_at: None,
apply_commit_oid: None,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ExecutionFactsSnapshot {
pub changes: HashMap<String, ChangeExecutionFacts>,
pub activities: Vec<ProcessActivity>,
}
impl ExecutionFactsSnapshot {
pub fn change(&self, change_id: &str) -> ChangeExecutionFacts {
self.changes
.get(change_id)
.cloned()
.unwrap_or_else(ChangeExecutionFacts::unknown)
}
pub fn has_active_work(&self) -> bool {
!self.activities.is_empty()
|| self
.changes
.values()
.any(|facts| facts.current_phase.is_active())
}
}
pub fn is_admitted_execution_state(state: ChangeExecutionState) -> bool {
matches!(
state,
ChangeExecutionState::Queued
| ChangeExecutionState::Active
| ChangeExecutionState::Waiting
| ChangeExecutionState::Stopping
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EpisodeTerminal {
Completed,
Failed,
Stopped,
}
impl EpisodeTerminal {
#[cfg_attr(not(test), allow(dead_code))]
pub fn as_str(self) -> &'static str {
match self {
Self::Completed => "completed",
Self::Failed => "failed",
Self::Stopped => "stopped",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EpisodeTransitionKind {
Started,
BlockedEntered,
BlockedLeft,
Terminal(EpisodeTerminal),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EpisodeTransition {
pub change_id: String,
pub execution_id: String,
pub kind: EpisodeTransitionKind,
}
pub trait EpisodeObserver: std::fmt::Debug + Send + Sync {
fn observe_episode(&self, transition: &EpisodeTransition);
}
#[derive(Debug, Clone, Default)]
struct ChangeFactsState {
execution_id: Option<String>,
episode_open: bool,
blocked_edge: bool,
execution_state: ChangeExecutionState,
current_phase: ExecutionPhase,
phase_started_at: Option<DateTime<Utc>>,
last_completed_phase: Option<ExecutionPhase>,
last_completed_at: Option<DateTime<Utc>>,
apply_commit_oid: Option<String>,
push_open: bool,
}
#[derive(Debug, Default)]
struct AbsorbedDispatches {
order: VecDeque<u64>,
seen: HashSet<u64>,
}
impl AbsorbedDispatches {
const CAPACITY: usize = 1024;
fn admit(&mut self, id: u64) -> bool {
if !self.seen.insert(id) {
return false;
}
self.order.push_back(id);
while self.order.len() > Self::CAPACITY {
if let Some(evicted) = self.order.pop_front() {
self.seen.remove(&evicted);
}
}
true
}
}
#[derive(Debug, Default)]
struct Inner {
changes: HashMap<String, ChangeFactsState>,
activities: HashSet<ProcessActivity>,
stop_requested: bool,
absorbed: AbsorbedDispatches,
pending: Vec<EpisodeTransition>,
}
#[derive(Debug, Default)]
pub struct ExecutionFactsStore {
inner: Mutex<Inner>,
observer: std::sync::RwLock<Option<std::sync::Arc<dyn EpisodeObserver>>>,
}
impl ExecutionFactsStore {
pub fn new() -> Self {
Self::default()
}
pub fn bind_episode_observer(&self, observer: std::sync::Arc<dyn EpisodeObserver>) {
*self
.observer
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
}
fn publish(&self, transitions: Vec<EpisodeTransition>) {
if transitions.is_empty() {
return;
}
let observer = self
.observer
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.clone();
let Some(observer) = observer else {
return;
};
for transition in &transitions {
observer.observe_episode(transition);
}
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn execution_id(&self, change_id: &str) -> Option<String> {
self.lock()
.changes
.get(change_id)
.and_then(|facts| facts.execution_id.clone())
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn change_of_execution(&self, execution_id: &str) -> Option<String> {
self.lock()
.changes
.iter()
.find(|(_, facts)| facts.execution_id.as_deref() == Some(execution_id))
.map(|(id, _)| id.clone())
}
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub fn observe(
&self,
dispatch_id: u64,
event: &ExecutionEvent,
state: Option<&OrchestratorState>,
now: DateTime<Utc>,
) -> bool {
let transitions = {
let mut inner = self.lock();
if !inner.absorbed.admit(dispatch_id) {
return false;
}
Self::observe_completion(&mut inner, event, now);
Self::observe_push(&mut inner, event);
Self::observe_process(&mut inner, event);
if let Some(state) = state {
Self::refresh_from_reducer(&mut inner, state, now);
}
std::mem::take(&mut inner.pending)
};
self.publish(transitions);
true
}
fn observe_completion(inner: &mut Inner, event: &ExecutionEvent, now: DateTime<Utc>) {
use ExecutionEvent as E;
let (change_id, phase) = match event {
E::WorkspacePreparationEnded { change_id } => (change_id, ExecutionPhase::Preparing),
E::ApplyCompleted {
change_id,
revision,
} => {
if !revision.trim().is_empty() {
inner
.changes
.entry(change_id.clone())
.or_default()
.apply_commit_oid = Some(revision.trim().to_string());
}
(change_id, ExecutionPhase::Apply)
}
E::AcceptanceCompleted { change_id } => (change_id, ExecutionPhase::Acceptance),
E::RejectionReviewCompleted { change_id, .. } => {
(change_id, ExecutionPhase::RejectionReview)
}
E::ChangeArchived(change_id) => (change_id, ExecutionPhase::Archive),
E::ResolveCompleted { change_id, .. } => (change_id, ExecutionPhase::Resolve),
E::MergeCompleted { change_id, .. } => (change_id, ExecutionPhase::Merge),
E::PushCompleted { change_id, .. } => (change_id, ExecutionPhase::Push),
_ => return,
};
let facts = inner.changes.entry(change_id.clone()).or_default();
facts.last_completed_phase = Some(phase);
facts.last_completed_at = Some(now);
}
fn observe_push(inner: &mut Inner, event: &ExecutionEvent) {
use ExecutionEvent as E;
let (change_id, open) = match event {
E::PushStarted { change_id, .. } => (change_id, true),
E::PushCompleted { change_id, .. } | E::PushFailed { change_id, .. } => {
(change_id, false)
}
_ => return,
};
inner
.changes
.entry(change_id.clone())
.or_default()
.push_open = open;
}
fn observe_process(inner: &mut Inner, event: &ExecutionEvent) {
use ExecutionEvent as E;
match event {
E::AnalysisStarted { .. } => {
inner.activities.insert(ProcessActivity::DependencyAnalysis);
}
E::AnalysisCompleted { .. } => {
inner
.activities
.remove(&ProcessActivity::DependencyAnalysis);
}
E::MergeStarted { .. } => {
inner.activities.insert(ProcessActivity::BaseBranchMerge);
}
E::MergeCompleted { .. } | E::MergeDeferred { .. } | E::ResolveFailed { .. } => {
inner.activities.remove(&ProcessActivity::BaseBranchMerge);
}
E::ConflictResolutionStarted => {
inner.activities.insert(ProcessActivity::ConflictResolution);
}
E::ConflictResolutionCompleted | E::ConflictResolutionFailed { .. } => {
inner
.activities
.remove(&ProcessActivity::ConflictResolution);
}
E::BranchMergeStarted { .. } => {
inner.activities.insert(ProcessActivity::BranchMerge);
}
E::BranchMergeCompleted { .. } | E::BranchMergeFailed { .. } => {
inner.activities.remove(&ProcessActivity::BranchMerge);
}
E::CleanupStarted { .. } => {
inner.activities.insert(ProcessActivity::WorkspaceCleanup);
}
E::CleanupCompleted { .. } => {
inner.activities.remove(&ProcessActivity::WorkspaceCleanup);
}
E::Stopping => inner.stop_requested = true,
E::Stopped | E::Error { .. } | E::AllCompleted => {
inner.activities.clear();
inner.stop_requested = false;
}
E::ProcessingStarted(_) => inner.stop_requested = false,
_ => {}
}
}
fn refresh_from_reducer(inner: &mut Inner, state: &OrchestratorState, now: DateTime<Utc>) {
let stop_requested = inner.stop_requested;
let Inner {
changes, pending, ..
} = inner;
for change_id in state.tracked_change_ids() {
let Some(runtime) = state.change_runtime(&change_id) else {
continue;
};
let change_key = change_id.clone();
let facts = changes.entry(change_id).or_default();
let phase = project_phase(runtime, facts.push_open);
if phase != facts.current_phase {
facts.current_phase = phase;
facts.phase_started_at = phase.is_active().then_some(now);
}
let next = project_execution_state(runtime, phase, stop_requested);
facts.execution_state = next;
Self::advance_episode(pending, &change_key, facts, next);
}
}
fn advance_episode(
pending: &mut Vec<EpisodeTransition>,
change_id: &str,
facts: &mut ChangeFactsState,
next: ChangeExecutionState,
) {
let admitted = is_admitted_execution_state(next);
if admitted && !facts.episode_open {
let execution_id = crate::ids::new_hex_id();
facts.execution_id = Some(execution_id.clone());
facts.episode_open = true;
facts.blocked_edge = false;
pending.push(EpisodeTransition {
change_id: change_id.to_string(),
execution_id,
kind: EpisodeTransitionKind::Started,
});
}
let Some(execution_id) = facts.execution_id.clone() else {
return;
};
if !facts.episode_open {
return;
}
if admitted {
let blocked = matches!(next, ChangeExecutionState::Waiting);
if blocked != facts.blocked_edge {
facts.blocked_edge = blocked;
pending.push(EpisodeTransition {
change_id: change_id.to_string(),
execution_id,
kind: if blocked {
EpisodeTransitionKind::BlockedEntered
} else {
EpisodeTransitionKind::BlockedLeft
},
});
}
return;
}
let terminal = match next {
ChangeExecutionState::Completed => EpisodeTerminal::Completed,
ChangeExecutionState::Failed => EpisodeTerminal::Failed,
ChangeExecutionState::Stopped | ChangeExecutionState::Unknown => {
EpisodeTerminal::Stopped
}
_ => return,
};
facts.episode_open = false;
facts.blocked_edge = false;
pending.push(EpisodeTransition {
change_id: change_id.to_string(),
execution_id,
kind: EpisodeTransitionKind::Terminal(terminal),
});
}
pub fn snapshot(&self) -> ExecutionFactsSnapshot {
let inner = self.lock();
let mut activities: Vec<ProcessActivity> = inner.activities.iter().copied().collect();
activities.sort();
ExecutionFactsSnapshot {
changes: inner
.changes
.iter()
.map(|(id, facts)| {
(
id.clone(),
ChangeExecutionFacts {
execution_id: facts.execution_id.clone(),
execution_state: facts.execution_state,
current_phase: facts.current_phase,
phase_started_at: facts.phase_started_at,
last_completed_phase: facts.last_completed_phase,
last_completed_at: facts.last_completed_at,
apply_commit_oid: facts.apply_commit_oid.clone(),
},
)
})
.collect(),
activities,
}
}
pub fn change(&self, change_id: &str) -> ChangeExecutionFacts {
let inner = self.lock();
inner
.changes
.get(change_id)
.map(|facts| ChangeExecutionFacts {
execution_id: facts.execution_id.clone(),
execution_state: facts.execution_state,
current_phase: facts.current_phase,
phase_started_at: facts.phase_started_at,
last_completed_phase: facts.last_completed_phase,
last_completed_at: facts.last_completed_at,
apply_commit_oid: facts.apply_commit_oid.clone(),
})
.unwrap_or_else(ChangeExecutionFacts::unknown)
}
pub fn apply_commit_oid(&self, change_id: &str) -> Option<String> {
self.lock()
.changes
.get(change_id)
.and_then(|facts| facts.apply_commit_oid.clone())
}
}
pub fn project_phase(runtime: &ChangeRuntimeState, push_open: bool) -> ExecutionPhase {
if push_open
&& matches!(
runtime.activity,
ActivityState::Idle | ActivityState::Resolving
)
{
return ExecutionPhase::Push;
}
match runtime.activity {
ActivityState::Preparing => ExecutionPhase::Preparing,
ActivityState::Applying => ExecutionPhase::Apply,
ActivityState::Accepting => ExecutionPhase::Acceptance,
ActivityState::Rejecting => ExecutionPhase::RejectionReview,
ActivityState::Archiving => ExecutionPhase::Archive,
ActivityState::Resolving => ExecutionPhase::Resolve,
ActivityState::Idle => ExecutionPhase::None,
}
}
pub fn project_execution_state(
runtime: &ChangeRuntimeState,
phase: ExecutionPhase,
stop_requested: bool,
) -> ChangeExecutionState {
match &runtime.terminal {
TerminalState::Merged | TerminalState::Pushed => return ChangeExecutionState::Completed,
TerminalState::Rejected(_) | TerminalState::Error(_) => {
return ChangeExecutionState::Failed
}
TerminalState::Stopped => return ChangeExecutionState::Stopped,
TerminalState::None => {}
}
if runtime.dequeued {
return ChangeExecutionState::Stopped;
}
if phase.is_active() {
return if stop_requested {
ChangeExecutionState::Stopping
} else {
ChangeExecutionState::Active
};
}
if matches!(runtime.wait_state, WaitState::DependencyBlocked)
&& matches!(runtime.queue_intent, QueueIntent::Queued)
{
return ChangeExecutionState::Queued;
}
if !matches!(runtime.wait_state, WaitState::None) {
return ChangeExecutionState::Waiting;
}
if matches!(runtime.queue_intent, QueueIntent::Queued) {
return ChangeExecutionState::Queued;
}
ChangeExecutionState::Unknown
}
#[cfg(test)]
mod tests;