use std::collections::{HashMap, HashSet};
#[derive(Debug, Clone)]
pub struct WorkspaceResult {
pub change_id: String,
pub workspace_name: String,
pub final_revision: Option<String>,
pub error: Option<String>,
pub rejected: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ResolveFailureClassification {
UnresolvedConflict,
ResolveAgentFailed,
EvidenceWithheld,
RetriesExhausted,
}
impl ResolveFailureClassification {
pub fn token(self) -> &'static str {
match self {
Self::UnresolvedConflict => "unresolved_conflict",
Self::ResolveAgentFailed => "resolve_agent_failed",
Self::EvidenceWithheld => "evidence_withheld",
Self::RetriesExhausted => "retries_exhausted",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AlreadyReportedFailureKind {
Push,
Hook,
RejectionReview,
}
impl AlreadyReportedFailureKind {
pub fn token(self) -> &'static str {
match self {
Self::Push => "push",
Self::Hook => "hook",
Self::RejectionReview => "rejection_review",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MergeTaskOutcome {
Merged,
Deferred {
reason: String,
auto_resumable: bool,
},
ResolveExhausted {
change_id: String,
attempts: u32,
classification: ResolveFailureClassification,
detail: String,
},
RecoverableAlreadyReported {
change_id: String,
kind: AlreadyReportedFailureKind,
detail: String,
},
RunFatal {
detail: String,
},
}
impl MergeTaskOutcome {
pub fn deferred(reason: impl Into<String>, auto_resumable: bool) -> Self {
Self::Deferred {
reason: reason.into(),
auto_resumable,
}
}
pub fn resolve_exhausted(
change_id: impl Into<String>,
attempts: u32,
classification: ResolveFailureClassification,
detail: impl Into<String>,
) -> Self {
Self::ResolveExhausted {
change_id: change_id.into(),
attempts,
classification,
detail: crate::events::sanitize_detail(&detail.into()),
}
}
pub fn already_reported(
change_id: impl Into<String>,
kind: AlreadyReportedFailureKind,
detail: impl Into<String>,
) -> Self {
Self::RecoverableAlreadyReported {
change_id: change_id.into(),
kind,
detail: crate::events::sanitize_detail(&detail.into()),
}
}
pub fn run_fatal(detail: impl Into<String>) -> Self {
Self::RunFatal {
detail: crate::events::sanitize_detail(&detail.into()),
}
}
pub fn disposition(&self) -> MergeResultDisposition {
match self {
Self::Merged => MergeResultDisposition::Merged,
Self::Deferred { .. } => MergeResultDisposition::Deferred,
Self::ResolveExhausted { .. } | Self::RecoverableAlreadyReported { .. } => {
MergeResultDisposition::ContinueWithErrors
}
Self::RunFatal { .. } => MergeResultDisposition::AbortRun,
}
}
pub fn scoped_change_id(&self) -> Option<&str> {
match self {
Self::ResolveExhausted { change_id, .. }
| Self::RecoverableAlreadyReported { change_id, .. } => Some(change_id.as_str()),
Self::Merged | Self::Deferred { .. } | Self::RunFatal { .. } => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MergeResultDisposition {
Merged,
Deferred,
ContinueWithErrors,
AbortRun,
}
impl MergeResultDisposition {
pub fn is_merged(self) -> bool {
matches!(self, Self::Merged)
}
}
pub fn resolve_failure_detail(
attempts: u32,
classification: ResolveFailureClassification,
summary: &str,
) -> String {
crate::events::sanitize_detail(&format!(
"resolve exhausted after {} attempt(s) [{}]: {}",
attempts,
classification.token(),
summary
))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MergeResultOrigin {
PostArchiveMerge,
ResolveWaitRetry,
RejectWaitRetry,
}
#[derive(Debug, Clone)]
pub struct MergeResult {
pub change_id: String,
pub workspace_name: String,
pub origin: MergeResultOrigin,
pub outcome: MergeTaskOutcome,
}
#[derive(Debug, Default)]
pub struct FailedChangeTracker {
failed_changes: HashSet<String>,
dependencies: HashMap<String, Vec<String>>,
notified_blocker_epochs: HashMap<String, Vec<String>>,
}
impl FailedChangeTracker {
pub fn new() -> Self {
Self::default()
}
pub fn set_dependencies(&mut self, dependencies: HashMap<String, Vec<String>>) {
self.dependencies = dependencies;
}
pub fn mark_failed(&mut self, change_id: &str) {
self.failed_changes.insert(change_id.to_string());
}
pub fn should_skip(&self, change_id: &str) -> Option<String> {
if let Some(deps) = self.dependencies.get(change_id) {
for dep in deps {
if self.failed_changes.contains(dep) {
return Some(dep.clone());
}
}
}
None
}
pub fn failed_blockers(&self, change_id: &str) -> Vec<String> {
let Some(deps) = self.dependencies.get(change_id) else {
return Vec::new();
};
let mut blockers: Vec<String> = deps
.iter()
.filter(|dep| self.failed_changes.contains(*dep))
.cloned()
.collect();
blockers.sort_unstable();
blockers.dedup();
blockers
}
pub fn clear_failed(&mut self, change_id: &str) -> bool {
let cleared = self.failed_changes.remove(change_id);
self.notified_blocker_epochs
.retain(|_, blockers| !blockers.iter().any(|blocker| blocker == change_id));
cleared
}
pub fn begin_blocker_epoch(&mut self, change_id: &str, blockers: &[String]) -> bool {
if self
.notified_blocker_epochs
.get(change_id)
.is_some_and(|recorded| recorded.as_slice() == blockers)
{
return false;
}
self.notified_blocker_epochs
.insert(change_id.to_string(), blockers.to_vec());
true
}
pub fn clear_blocker_epoch(&mut self, change_id: &str) -> bool {
self.notified_blocker_epochs.remove(change_id).is_some()
}
#[cfg(test)]
pub fn has_blocker_epoch(&self, change_id: &str) -> bool {
self.notified_blocker_epochs.contains_key(change_id)
}
#[allow(dead_code)] pub fn failed_changes(&self) -> &HashSet<String> {
&self.failed_changes
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn merged_outcome_disposes_as_merged() {
assert_eq!(
MergeTaskOutcome::Merged.disposition(),
MergeResultDisposition::Merged
);
}
#[test]
fn deferred_outcome_is_continuation_without_recorded_failure() {
for auto_resumable in [true, false] {
assert_eq!(
MergeTaskOutcome::deferred("base lane busy", auto_resumable).disposition(),
MergeResultDisposition::Deferred,
"a deferral is pending work, not a change failure"
);
}
}
#[test]
fn resolve_exhausted_outcome_continues_with_errors_and_keeps_change_scope() {
let outcome = MergeTaskOutcome::resolve_exhausted(
"alpha",
3,
ResolveFailureClassification::UnresolvedConflict,
"conflicts remain in a.rs",
);
assert_eq!(
outcome.disposition(),
MergeResultDisposition::ContinueWithErrors
);
assert_eq!(outcome.scoped_change_id(), Some("alpha"));
}
#[test]
fn already_reported_outcomes_continue_with_errors_for_every_kind() {
for kind in [
AlreadyReportedFailureKind::Push,
AlreadyReportedFailureKind::Hook,
AlreadyReportedFailureKind::RejectionReview,
] {
let outcome = MergeTaskOutcome::already_reported("alpha", kind, "already reported");
assert_eq!(
outcome.disposition(),
MergeResultDisposition::ContinueWithErrors,
"{:?} already has a typed owner and must never abort the run",
kind
);
assert_eq!(outcome.scoped_change_id(), Some("alpha"));
}
}
#[test]
fn run_fatal_outcome_aborts_the_run_and_carries_no_change_scope() {
let outcome = MergeTaskOutcome::run_fatal("base branch could not be identified");
assert_eq!(outcome.disposition(), MergeResultDisposition::AbortRun);
assert_eq!(outcome.scoped_change_id(), None);
}
#[test]
fn outcome_details_are_sanitized_and_bounded() {
let outcome = MergeTaskOutcome::run_fatal("line one\nline two\u{7}");
let MergeTaskOutcome::RunFatal { detail } = outcome else {
panic!("expected run-fatal outcome");
};
assert!(!detail.contains('\n'), "raw newlines must be escaped");
assert!(
!detail.contains('\u{7}'),
"control characters must be dropped"
);
}
#[test]
fn resolve_failure_detail_carries_attempts_classification_and_summary() {
let detail = resolve_failure_detail(
3,
ResolveFailureClassification::ResolveAgentFailed,
"resolve command exited 1",
);
assert!(
detail.contains("3 attempt(s)"),
"missing attempts: {detail}"
);
assert!(
detail.contains("resolve_agent_failed"),
"missing classification token: {detail}"
);
assert!(
detail.contains("resolve command exited 1"),
"missing summary: {detail}"
);
}
#[test]
fn classification_tokens_are_distinct_and_stable() {
let tokens = [
ResolveFailureClassification::UnresolvedConflict.token(),
ResolveFailureClassification::ResolveAgentFailed.token(),
ResolveFailureClassification::EvidenceWithheld.token(),
ResolveFailureClassification::RetriesExhausted.token(),
];
let unique: HashSet<&str> = tokens.iter().copied().collect();
assert_eq!(unique.len(), tokens.len(), "tokens must be distinguishable");
}
#[test]
fn test_failed_tracker_new() {
let tracker = FailedChangeTracker::new();
assert!(tracker.failed_changes.is_empty());
assert!(tracker.dependencies.is_empty());
}
#[test]
fn test_mark_failed() {
let mut tracker = FailedChangeTracker::new();
tracker.mark_failed("change-a");
assert!(tracker.failed_changes.contains("change-a"));
}
#[test]
fn test_should_skip_no_dependencies() {
let tracker = FailedChangeTracker::new();
assert!(tracker.should_skip("change-a").is_none());
}
#[test]
fn test_should_skip_with_failed_dependency() {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert("change-b".to_string(), vec!["change-a".to_string()]);
tracker.set_dependencies(deps);
tracker.mark_failed("change-a");
let result = tracker.should_skip("change-b");
assert_eq!(result, Some("change-a".to_string()));
}
#[test]
fn test_should_skip_no_failed_dependency() {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert("change-b".to_string(), vec!["change-a".to_string()]);
tracker.set_dependencies(deps);
assert!(tracker.should_skip("change-b").is_none());
}
fn failed_change_tracker_with_dependency(
dependent: &str,
dependency: &str,
) -> FailedChangeTracker {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert(dependent.to_string(), vec![dependency.to_string()]);
tracker.set_dependencies(deps);
tracker
}
#[test]
fn failed_change_tracker_reports_sorted_deduplicated_blockers() {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert(
"c".to_string(),
vec![
"z".to_string(),
"a".to_string(),
"z".to_string(),
"healthy".to_string(),
],
);
tracker.set_dependencies(deps);
tracker.mark_failed("z");
tracker.mark_failed("a");
assert_eq!(
tracker.failed_blockers("c"),
vec!["a".to_string(), "z".to_string()],
"blocker epochs must not depend on declaration order or duplicates"
);
assert!(
tracker.failed_blockers("unknown").is_empty(),
"an unknown change has no failed blockers"
);
}
#[test]
fn failed_change_tracker_epoch_is_announced_once_per_blocker_set() {
let mut tracker = failed_change_tracker_with_dependency("b", "a");
tracker.mark_failed("a");
let blockers = tracker.failed_blockers("b");
assert!(
tracker.begin_blocker_epoch("b", &blockers),
"the first blocked observation is a new epoch"
);
assert!(
!tracker.begin_blocker_epoch("b", &blockers),
"rediscovering the same blocker set must not open a new epoch"
);
let changed = vec!["a".to_string(), "other".to_string()];
assert!(
tracker.begin_blocker_epoch("b", &changed),
"a changed blocker set is a new epoch"
);
}
#[test]
fn failed_change_tracker_clear_failed_is_narrowly_scoped() {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert("b".to_string(), vec!["a".to_string()]);
deps.insert("d".to_string(), vec!["other".to_string()]);
tracker.set_dependencies(deps);
tracker.mark_failed("a");
tracker.mark_failed("other");
let a_blockers = tracker.failed_blockers("b");
let other_blockers = tracker.failed_blockers("d");
assert!(tracker.begin_blocker_epoch("b", &a_blockers));
assert!(tracker.begin_blocker_epoch("d", &other_blockers));
assert!(tracker.clear_failed("a"), "the failed marker was removed");
assert!(
!tracker.clear_failed("a"),
"clearing an unmarked change reports no change"
);
assert!(tracker.should_skip("b").is_none(), "b's gate is released");
assert!(
!tracker.has_blocker_epoch("b"),
"the epoch naming the retried change must be dropped"
);
assert_eq!(
tracker.should_skip("d"),
Some("other".to_string()),
"an unrelated failure must survive a narrowly scoped clear"
);
assert!(
tracker.has_blocker_epoch("d"),
"an unrelated epoch must survive a narrowly scoped clear"
);
tracker.mark_failed("a");
let refailed = tracker.failed_blockers("b");
assert!(
tracker.begin_blocker_epoch("b", &refailed),
"refailure must be announced as a new epoch"
);
}
#[test]
fn failed_change_tracker_clear_blocker_epoch_allows_one_new_announcement() {
let mut tracker = failed_change_tracker_with_dependency("b", "a");
tracker.mark_failed("a");
let blockers = tracker.failed_blockers("b");
assert!(tracker.begin_blocker_epoch("b", &blockers));
assert!(
tracker.clear_blocker_epoch("b"),
"a recorded epoch is reported as cleared"
);
assert!(
!tracker.clear_blocker_epoch("b"),
"clearing twice reports no recorded epoch"
);
assert!(
tracker.begin_blocker_epoch("b", &blockers),
"an explicit re-add after revocation may announce once more"
);
assert_eq!(
tracker.should_skip("b"),
Some("a".to_string()),
"clearing a notification epoch must not clear the failure itself"
);
}
#[test]
fn test_should_skip_with_multiple_dependencies() {
let mut tracker = FailedChangeTracker::new();
let mut deps = HashMap::new();
deps.insert(
"change-c".to_string(),
vec!["change-a".to_string(), "change-b".to_string()],
);
tracker.set_dependencies(deps);
tracker.mark_failed("change-b");
let result = tracker.should_skip("change-c");
assert_eq!(result, Some("change-b".to_string()));
}
}