use std::sync::Arc;
use crate::events::ExecutionEvent;
use crate::orchestration::operator_command::{
parallel_cleanup_targets, ExecutionMarkStore, ParallelCleanupRow, ParallelEligibility,
ParallelRuntime,
};
use crate::orchestration::state::{OrchestratorState, TerminalState, WaitState};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct MarkEvidence {
pub error: bool,
pub rejected: bool,
pub dequeued: bool,
pub merge_wait: bool,
}
impl MarkEvidence {
pub fn observe(state: &OrchestratorState, change_id: &str) -> Self {
let Some(runtime) = state.change_runtime(change_id) else {
return Self::default();
};
Self {
error: matches!(runtime.terminal, TerminalState::Error(_)),
rejected: matches!(runtime.terminal, TerminalState::Rejected(_)),
dequeued: runtime.dequeued,
merge_wait: matches!(runtime.wait_state, WaitState::MergeWait),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RevokingEdge {
ChangeLevelFailure,
Rejection,
Dequeue,
MergedHookRecovery,
}
impl RevokingEdge {
pub fn revokes(self, pre: MarkEvidence, post: MarkEvidence) -> bool {
match self {
Self::ChangeLevelFailure => post.error && !pre.error,
Self::Rejection => post.rejected && !pre.rejected,
Self::Dequeue => post.dequeued && !pre.dequeued,
Self::MergedHookRecovery => post.merge_wait && !pre.merge_wait,
}
}
}
pub fn revoking_edge_target(event: &ExecutionEvent) -> Option<(&str, RevokingEdge)> {
use ExecutionEvent as E;
match event {
E::ProcessingError { id, .. } => Some((id, RevokingEdge::ChangeLevelFailure)),
E::ApplyFailed { change_id, .. }
| E::AcceptanceFailed { change_id, .. }
| E::ArchiveFailed { change_id, .. }
| E::PushFailed { change_id, .. }
| E::RejectionReviewFailed { change_id, .. } => {
Some((change_id, RevokingEdge::ChangeLevelFailure))
}
E::ChangeRejected { change_id, .. } => Some((change_id, RevokingEdge::Rejection)),
E::ChangeDequeued { change_id } | E::ChangeStopped { change_id } => {
Some((change_id, RevokingEdge::Dequeue))
}
E::HookFailed {
change_id,
hook_type,
..
} if hook_type == crate::hooks::HookType::OnMerged.config_key() => {
Some((change_id, RevokingEdge::MergedHookRecovery))
}
_ => None,
}
}
#[derive(Debug, Clone)]
pub enum MarkPreState {
Target {
change_id: String,
edge: RevokingEdge,
evidence: MarkEvidence,
},
Refresh,
}
pub fn capture_pre_state(
event: &ExecutionEvent,
state: &OrchestratorState,
) -> Option<MarkPreState> {
if matches!(event, ExecutionEvent::ChangesRefreshed { .. }) {
return Some(MarkPreState::Refresh);
}
let (change_id, edge) = revoking_edge_target(event)?;
Some(MarkPreState::Target {
change_id: change_id.to_string(),
edge,
evidence: MarkEvidence::observe(state, change_id),
})
}
pub fn refresh_revocation_targets(
event: &ExecutionEvent,
marks: &ExecutionMarkStore,
post: &OrchestratorState,
) -> Vec<String> {
let ExecutionEvent::ChangesRefreshed {
changes,
rejected_changes,
committed_change_ids,
uncommitted_file_change_ids,
..
} = event
else {
return Vec::new();
};
let mut targets: Vec<String> = Vec::new();
for change in rejected_changes {
if marks.is_marked(&change.id) && !targets.contains(&change.id) {
targets.push(change.id.clone());
}
}
let rows: Vec<ParallelCleanupRow<'_>> = changes
.iter()
.map(|change| ParallelCleanupRow {
change_id: &change.id,
parallel_eligible: ParallelEligibility::observe(
&change.id,
committed_change_ids,
uncommitted_file_change_ids,
)
.is_eligible(),
marked: marks.is_marked(&change.id),
queued: post.display_status(&change.id) == "queued",
})
.collect();
for change_id in parallel_cleanup_targets(&rows) {
if marks.is_marked(&change_id) && !targets.contains(&change_id) {
targets.push(change_id);
}
}
targets
}
#[derive(Clone)]
pub struct ExecutionMarkReconciler {
marks: Arc<ExecutionMarkStore>,
guard: Arc<ParallelRuntime>,
}
impl std::fmt::Debug for ExecutionMarkReconciler {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ExecutionMarkReconciler")
.field("marked", &self.marks.marked_ids())
.finish()
}
}
impl ExecutionMarkReconciler {
pub fn new(marks: Arc<ExecutionMarkStore>, guard: Arc<ParallelRuntime>) -> Self {
Self { marks, guard }
}
pub async fn lock_mutations(&self) -> tokio::sync::MutexGuard<'_, ()> {
self.guard.lock_mutations().await
}
pub fn capture(
&self,
event: &ExecutionEvent,
state: &OrchestratorState,
) -> Option<MarkPreState> {
capture_pre_state(event, state)
}
pub fn reconcile(
&self,
event: &ExecutionEvent,
pre: &MarkPreState,
post: &OrchestratorState,
) -> Vec<String> {
match pre {
MarkPreState::Target {
change_id,
edge,
evidence,
} => {
let after = MarkEvidence::observe(post, change_id);
if edge.revokes(*evidence, after) && self.marks.set(change_id, false) {
vec![change_id.clone()]
} else {
Vec::new()
}
}
MarkPreState::Refresh => {
let targets = refresh_revocation_targets(event, &self.marks, post);
targets
.into_iter()
.filter(|change_id| self.marks.set(change_id, false))
.collect()
}
}
}
}
#[cfg(test)]
mod tests;