use std::collections::{HashMap, HashSet};
use crate::error_history::{CircuitBreakerConfig, ErrorHistory};
#[derive(Debug, Clone, PartialEq, Default)]
pub enum QueueIntent {
#[default]
NotQueued,
Queued,
}
#[derive(Debug, Clone, PartialEq, Default)]
pub enum ActivityState {
#[default]
Idle,
Preparing,
Applying,
Accepting,
Rejecting,
Archiving,
Resolving,
}
#[derive(Debug, Clone, PartialEq, Default)]
pub enum WaitState {
#[default]
None,
MergeWait,
ResolveWait,
RejectWait,
DependencyBlocked,
ExternalBlocked,
Stalled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum BlockerKind {
#[default]
None,
Dependency,
External,
}
impl BlockerKind {
pub fn as_str(self) -> Option<&'static str> {
match self {
Self::None => None,
Self::Dependency => Some("dependency"),
Self::External => Some("external"),
}
}
}
fn dependency_blocker_detail(dependency_ids: &[String]) -> String {
if dependency_ids.is_empty() {
return "waiting for unresolved dependencies".to_string();
}
format!("waiting for dependencies: {}", dependency_ids.join(", "))
}
fn dependency_unblock_condition(dependency_ids: &[String]) -> String {
if dependency_ids.is_empty() {
return "every declared dependency is integrated into the effective dependency base"
.to_string();
}
format!(
"dependencies {} are integrated into the effective dependency base",
dependency_ids.join(", ")
)
}
#[derive(Debug, Clone, PartialEq, Default)]
pub struct BlockedMetadata {
pub blocker_reason: Option<String>,
pub unblock_metadata: Option<String>,
pub worktree_snapshot: Option<String>,
pub acceptance_stall: bool,
pub resumable: bool,
pub blocker_kind: BlockerKind,
pub unblock_condition: Option<String>,
pub prerequisite_owner: Option<String>,
pub blocker_origin: Option<String>,
pub dependency_ids: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BlockerView {
pub status: &'static str,
pub kind: BlockerKind,
pub category: Option<String>,
pub detail: Option<String>,
pub unblock_condition: Option<String>,
pub prerequisite_owner: Option<String>,
pub origin: Option<String>,
pub resumable: bool,
pub dependencies: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Default)]
pub enum TerminalState {
#[default]
None,
Merged,
Pushed,
Rejected(String),
Error(String),
Stopped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StaleResolveEvidence {
Proven,
Absent,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StaleResolveSettlement {
Merged,
MergeWaitRetained,
AlreadySettled,
}
#[derive(Debug, Clone, PartialEq, Default)]
pub enum WorkspaceObservation {
#[default]
None,
WorkspaceArchived,
WorktreeNotAhead,
}
#[derive(Debug, Clone, Default)]
pub struct ChangeRuntimeState {
pub queue_intent: QueueIntent,
pub activity: ActivityState,
pub wait_state: WaitState,
pub blocked_metadata: BlockedMetadata,
pub terminal: TerminalState,
pub observation: WorkspaceObservation,
pub dequeued: bool,
pub commit_phase_attempt: Option<u32>,
pub manual_resolve_retry: bool,
}
impl ChangeRuntimeState {
fn clear_blocked_metadata(&mut self) {
self.blocked_metadata = BlockedMetadata::default();
}
fn set_blocked_metadata(
&mut self,
blocker_reason: impl Into<String>,
unblock_metadata: impl Into<String>,
worktree_snapshot: impl Into<String>,
) {
self.blocked_metadata = BlockedMetadata {
blocker_reason: Some(blocker_reason.into()),
unblock_metadata: Some(unblock_metadata.into()),
worktree_snapshot: Some(worktree_snapshot.into()),
acceptance_stall: false,
resumable: false,
blocker_kind: BlockerKind::None,
unblock_condition: None,
prerequisite_owner: None,
blocker_origin: None,
dependency_ids: Vec::new(),
};
}
fn reconcile_dependency_blocker(&mut self, dependency_ids: &[String]) -> bool {
let unchanged = matches!(self.wait_state, WaitState::DependencyBlocked)
&& self.blocked_metadata.dependency_ids == dependency_ids;
if unchanged {
return false;
}
self.wait_state = WaitState::DependencyBlocked;
self.blocked_metadata = BlockedMetadata {
blocker_reason: Some("dependency_blocked".to_string()),
unblock_metadata: Some(dependency_blocker_detail(dependency_ids)),
worktree_snapshot: Some(
"worktree snapshot not required for dependency blocker".to_string(),
),
acceptance_stall: false,
resumable: true,
blocker_kind: BlockerKind::Dependency,
unblock_condition: Some(dependency_unblock_condition(dependency_ids)),
prerequisite_owner: None,
blocker_origin: Some("scheduler".to_string()),
dependency_ids: dependency_ids.to_vec(),
};
true
}
fn clear_activity_wait_and_blocker(&mut self) {
self.activity = ActivityState::Idle;
self.wait_state = WaitState::None;
self.manual_resolve_retry = false;
self.clear_blocked_metadata();
}
fn transition_to_terminal(&mut self, terminal: TerminalState) {
self.terminal = terminal;
self.activity = ActivityState::Idle;
self.wait_state = WaitState::None;
self.manual_resolve_retry = false;
self.clear_blocked_metadata();
}
fn transition_to_stalled(
&mut self,
blocker_reason: impl Into<String>,
unblock_metadata: impl Into<String>,
worktree_snapshot: impl Into<String>,
) {
self.activity = ActivityState::Idle;
self.wait_state = WaitState::Stalled;
self.terminal = TerminalState::None;
self.set_blocked_metadata(blocker_reason, unblock_metadata, worktree_snapshot);
}
fn transition_to_acceptance_stalled(
&mut self,
blocker_reason: impl Into<String>,
unblock_metadata: impl Into<String>,
worktree_snapshot: impl Into<String>,
resumable: bool,
) {
self.transition_to_stalled(blocker_reason, unblock_metadata, worktree_snapshot);
self.blocked_metadata.acceptance_stall = true;
self.blocked_metadata.resumable = resumable;
}
fn transition_to_external_blocked(
&mut self,
blocker: &crate::runtime::proposal::ExternalBlockerInfo,
worktree_snapshot: impl Into<String>,
) {
self.activity = ActivityState::Idle;
self.wait_state = WaitState::ExternalBlocked;
self.terminal = TerminalState::None;
self.blocked_metadata = BlockedMetadata {
blocker_reason: Some(format!("external-blocked:{}", blocker.category)),
unblock_metadata: Some(blocker.summary()),
worktree_snapshot: Some(worktree_snapshot.into()),
acceptance_stall: matches!(
blocker.origin,
crate::runtime::proposal::BlockerOrigin::Acceptance
),
resumable: blocker.resumable,
blocker_kind: BlockerKind::External,
unblock_condition: Some(blocker.unblock_condition.clone()),
prerequisite_owner: blocker.prerequisite_owner.clone(),
blocker_origin: Some(blocker.origin.as_str().to_string()),
dependency_ids: Vec::new(),
};
}
pub fn is_acceptance_stalled(&self) -> bool {
matches!(
self.wait_state,
WaitState::Stalled | WaitState::ExternalBlocked
) && self.blocked_metadata.acceptance_stall
}
pub fn is_resumable_acceptance_stall(&self) -> bool {
self.is_acceptance_stalled() && self.blocked_metadata.resumable
}
#[allow(dead_code)] pub fn is_external_blocked(&self) -> bool {
matches!(self.wait_state, WaitState::ExternalBlocked)
}
pub fn has_structured_blocker_hold(&self) -> bool {
match self.wait_state {
WaitState::ExternalBlocked => true,
WaitState::Stalled => self.blocked_metadata.acceptance_stall,
_ => false,
}
}
pub fn blocker_kind(&self) -> BlockerKind {
match self.wait_state {
WaitState::DependencyBlocked => BlockerKind::Dependency,
WaitState::ExternalBlocked => BlockerKind::External,
_ => BlockerKind::None,
}
}
pub fn blocker_detail(&self) -> Option<&str> {
self.blocked_metadata.unblock_metadata.as_deref()
}
pub fn blocker_view(&self) -> Option<BlockerView> {
let status = self.display_status();
if !matches!(status, "blocked" | "stalled") {
return None;
}
Some(BlockerView {
status,
kind: self.blocker_kind(),
category: self.blocked_metadata.blocker_reason.clone(),
detail: self.blocker_detail().map(str::to_string),
unblock_condition: self.blocked_metadata.unblock_condition.clone(),
prerequisite_owner: self.blocked_metadata.prerequisite_owner.clone(),
origin: self.blocked_metadata.blocker_origin.clone(),
resumable: self.blocked_metadata.resumable,
dependencies: match self.blocker_kind() {
BlockerKind::Dependency => self.blocked_metadata.dependency_ids.clone(),
_ => Vec::new(),
},
})
}
pub fn is_active(&self) -> bool {
!matches!(self.activity, ActivityState::Idle)
}
pub fn is_terminal(&self) -> bool {
!matches!(self.terminal, TerminalState::None)
}
fn is_fresh_idle(&self) -> bool {
matches!(self.activity, ActivityState::Idle)
&& matches!(self.wait_state, WaitState::None)
&& matches!(self.queue_intent, QueueIntent::NotQueued)
&& matches!(self.terminal, TerminalState::None)
&& !self.dequeued
}
fn can_success_supersede_terminal(&self) -> bool {
matches!(self.terminal, TerminalState::None | TerminalState::Error(_))
}
#[allow(dead_code)]
pub fn invariants_hold(&self) -> bool {
if self.is_terminal() && self.is_active() {
return false;
}
if matches!(self.wait_state, WaitState::ResolveWait)
&& matches!(self.activity, ActivityState::Resolving)
{
return false;
}
if matches!(self.wait_state, WaitState::RejectWait)
&& matches!(self.activity, ActivityState::Rejecting)
{
return false;
}
true
}
pub fn apply_operation_label(&self) -> &'static str {
match self.commit_phase_attempt {
Some(_) => "commit",
None => "apply",
}
}
pub fn display_status(&self) -> &'static str {
match &self.terminal {
TerminalState::Merged => return "merged",
TerminalState::Pushed => return "pushed",
TerminalState::Rejected(_) => return "rejected",
TerminalState::Error(_) => return "error",
TerminalState::Stopped => return "stopped",
TerminalState::None => {}
}
match self.activity {
ActivityState::Preparing => return "preparing",
ActivityState::Applying => return "applying",
ActivityState::Accepting => return "accepting",
ActivityState::Rejecting => return "rejecting",
ActivityState::Archiving => return "archiving",
ActivityState::Resolving => return "resolving",
ActivityState::Idle => {}
}
match self.wait_state {
WaitState::MergeWait => return "merge wait",
WaitState::ResolveWait => return "resolve pending",
WaitState::RejectWait => return "reject pending",
WaitState::DependencyBlocked | WaitState::ExternalBlocked => return "blocked",
WaitState::Stalled => return "stalled",
WaitState::None => {}
}
match self.queue_intent {
QueueIntent::Queued => "queued",
QueueIntent::NotQueued => "not queued",
}
}
#[allow(dead_code)]
pub fn display_color(&self) -> ratatui::style::Color {
match self.display_status() {
"not queued" => ratatui::style::Color::DarkGray,
"queued" => ratatui::style::Color::Yellow,
"blocked" => ratatui::style::Color::Gray,
"stalled" => ratatui::style::Color::LightYellow,
"preparing" => ratatui::style::Color::Green,
"applying" => ratatui::style::Color::Cyan,
"accepting" => ratatui::style::Color::LightGreen,
"rejecting" => ratatui::style::Color::LightYellow,
"archiving" => ratatui::style::Color::Magenta,
"merged" => ratatui::style::Color::LightBlue,
"rejected" => ratatui::style::Color::LightRed,
"merge wait" => ratatui::style::Color::LightMagenta,
"resolving" => ratatui::style::Color::LightCyan,
"resolve pending" => ratatui::style::Color::Magenta,
"reject pending" => ratatui::style::Color::LightMagenta,
"error" => ratatui::style::Color::Red,
"stopped" => ratatui::style::Color::DarkGray,
_ => ratatui::style::Color::DarkGray,
}
}
#[allow(dead_code)]
pub fn error_message(&self) -> Option<&str> {
match &self.terminal {
TerminalState::Error(message) => Some(message.as_str()),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub enum ReducerCommand {
AddToQueue(String),
RetryError(String),
RemoveFromQueue(String),
ResolveMerge(String),
DequeueChange(String),
StopChange(String),
ResumeStopped(String),
}
#[derive(Debug, Clone)]
#[allow(dead_code)]
pub enum ReduceOutcome {
Changed(ReducerEffect),
NoOp,
}
#[derive(Debug, Clone)]
#[allow(clippy::enum_variant_names, dead_code)]
pub enum ReducerEffect {
QueueIntentSet {
change_id: String,
intent: QueueIntent,
},
WaitStateSet { change_id: String, wait: WaitState },
TerminalStateSet {
change_id: String,
terminal: TerminalState,
},
}
#[derive(Debug, Clone)]
pub struct OrchestratorState {
initial_change_ids: HashSet<String>,
pending_changes: HashSet<String>,
archived_changes: HashSet<String>,
apply_counts: HashMap<String, u32>,
stalled_change_ids: HashSet<String>,
skipped_change_ids: HashSet<String>,
task_progress: HashMap<String, (u32, u32)>,
error_histories: HashMap<String, ErrorHistory>,
changes_processed: usize,
total_changes: usize,
max_iterations: u32,
iteration: u32,
current_change_id: Option<String>,
change_runtime: HashMap<String, ChangeRuntimeState>,
resolve_wait_queue: Vec<String>,
reject_wait_queue: Vec<String>,
apply_iteration_limits: Vec<ApplyIterationLimit>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApplyIterationLimit {
pub change_id: String,
pub attempts: u32,
pub max: u32,
}
#[allow(dead_code)] impl OrchestratorState {
pub fn new(change_ids: Vec<String>, max_iterations: u32) -> Self {
let initial_set: HashSet<String> = change_ids.iter().cloned().collect();
let pending_set = initial_set.clone();
let total = change_ids.len();
let change_runtime: HashMap<String, ChangeRuntimeState> = change_ids
.iter()
.map(|id| (id.clone(), ChangeRuntimeState::default()))
.collect();
Self {
initial_change_ids: initial_set,
pending_changes: pending_set,
archived_changes: HashSet::new(),
apply_counts: HashMap::new(),
stalled_change_ids: HashSet::new(),
skipped_change_ids: HashSet::new(),
task_progress: HashMap::new(),
error_histories: HashMap::new(),
changes_processed: 0,
total_changes: total,
max_iterations,
iteration: 0,
current_change_id: None,
change_runtime,
resolve_wait_queue: Vec::new(),
reject_wait_queue: Vec::new(),
apply_iteration_limits: Vec::new(),
}
}
pub fn initial_change_ids(&self) -> &HashSet<String> {
&self.initial_change_ids
}
pub fn pending_changes(&self) -> &HashSet<String> {
&self.pending_changes
}
pub fn archived_changes(&self) -> &HashSet<String> {
&self.archived_changes
}
pub fn stalled_change_ids(&self) -> &HashSet<String> {
&self.stalled_change_ids
}
pub fn skipped_change_ids(&self) -> &HashSet<String> {
&self.skipped_change_ids
}
pub fn externally_blocked_change_ids(&self) -> HashSet<String> {
self.change_runtime
.iter()
.filter(|(_, rt)| rt.is_external_blocked())
.map(|(id, _)| id.clone())
.collect()
}
pub fn acceptance_stalled_change_ids(&self) -> HashSet<String> {
self.change_runtime
.iter()
.filter(|(_, rt)| rt.is_acceptance_stalled())
.map(|(change_id, _)| change_id.clone())
.collect()
}
pub fn mark_stalled(&mut self, change_id: String) {
self.stalled_change_ids.insert(change_id);
}
pub fn record_apply_iteration_limit(&mut self, change_id: &str, attempts: u32, max: u32) {
if self
.apply_iteration_limits
.iter()
.any(|record| record.change_id == change_id)
{
return;
}
self.apply_iteration_limits.push(ApplyIterationLimit {
change_id: change_id.to_string(),
attempts,
max,
});
}
pub fn apply_iteration_limits(&self) -> &[ApplyIterationLimit] {
&self.apply_iteration_limits
}
pub fn clear_apply_iteration_limit(&mut self, change_id: &str) -> bool {
let before = self.apply_iteration_limits.len();
self.apply_iteration_limits
.retain(|record| record.change_id != change_id);
self.apply_iteration_limits.len() != before
}
pub fn apply_iteration_limit(&self, change_id: &str) -> Option<&ApplyIterationLimit> {
self.apply_iteration_limits
.iter()
.find(|record| record.change_id == change_id)
}
pub fn parallel_finish_report(&self) -> (&'static str, u32) {
match self.apply_iteration_limits.first() {
Some(record) => ("iteration_limit", record.attempts),
None => ("completed", 0),
}
}
pub fn mark_skipped(&mut self, change_id: String) -> bool {
self.skipped_change_ids.insert(change_id)
}
pub fn clear_stalled_change(&mut self, change_id: &str) {
self.stalled_change_ids.remove(change_id);
self.skipped_change_ids.remove(change_id);
self.apply_counts.remove(change_id);
if self.current_change_id.as_deref() == Some(change_id) {
self.current_change_id = None;
}
}
pub fn changes_processed(&self) -> usize {
self.changes_processed
}
pub fn total_changes(&self) -> usize {
self.total_changes
}
pub fn remaining_changes(&self) -> usize {
self.pending_changes.len()
}
pub fn iteration(&self) -> u32 {
self.iteration
}
pub fn max_iterations(&self) -> u32 {
self.max_iterations
}
pub fn current_change_id(&self) -> Option<&String> {
self.current_change_id.as_ref()
}
pub fn apply_count(&self, change_id: &str) -> u32 {
*self.apply_counts.get(change_id).unwrap_or(&0)
}
pub fn task_progress(&self, change_id: &str) -> (u32, u32) {
*self.task_progress.get(change_id).unwrap_or(&(0, 0))
}
pub fn set_task_progress(&mut self, change_id: String, completed: u32, total: u32) {
self.task_progress.insert(change_id, (completed, total));
}
pub fn is_in_snapshot(&self, change_id: &str) -> bool {
self.initial_change_ids.contains(change_id)
}
pub fn is_pending(&self, change_id: &str) -> bool {
self.pending_changes.contains(change_id)
}
pub fn is_archived(&self, change_id: &str) -> bool {
self.archived_changes.contains(change_id)
}
pub fn is_complete(&self) -> bool {
self.pending_changes.is_empty()
}
pub fn is_iteration_limit_reached(&self) -> bool {
self.max_iterations > 0 && self.iteration >= self.max_iterations
}
pub fn is_approaching_iteration_limit(&self) -> bool {
if self.max_iterations == 0 {
return false;
}
let threshold = (self.max_iterations as f32 * 0.8) as u32;
self.iteration == threshold
}
pub fn increment_iteration(&mut self) {
self.iteration += 1;
}
pub fn set_current_change(&mut self, change_id: Option<String>) {
self.current_change_id = change_id;
}
pub fn increment_apply_count(&mut self, change_id: &str) -> u32 {
let count = self.apply_counts.entry(change_id.to_string()).or_insert(0);
*count += 1;
*count
}
pub fn mark_archived(&mut self, change_id: &str) {
if self.pending_changes.remove(change_id) {
self.archived_changes.insert(change_id.to_string());
self.changes_processed += 1;
self.apply_counts.remove(change_id);
if self.current_change_id.as_deref() == Some(change_id) {
self.current_change_id = None;
}
}
}
pub fn add_dynamic_change(&mut self, change_id: String) {
if !self.initial_change_ids.contains(&change_id)
&& !self.pending_changes.contains(&change_id)
&& !self.archived_changes.contains(&change_id)
{
self.initial_change_ids.insert(change_id.clone());
self.pending_changes.insert(change_id.clone());
self.total_changes += 1;
self.change_runtime.entry(change_id).or_default();
}
}
fn runtime_entry(&mut self, change_id: &str) -> &mut ChangeRuntimeState {
self.change_runtime
.entry(change_id.to_string())
.or_default()
}
pub fn tracked_change_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self.change_runtime.keys().cloned().collect();
ids.sort();
ids
}
pub fn change_runtime(&self, change_id: &str) -> Option<&ChangeRuntimeState> {
self.change_runtime.get(change_id)
}
pub fn is_tracked_change(&self, change_id: &str) -> bool {
self.change_runtime.contains_key(change_id)
}
pub fn display_status(&self, change_id: &str) -> &'static str {
match self.change_runtime.get(change_id) {
Some(rt) => rt.display_status(),
None => "not queued",
}
}
pub fn is_active_change(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.map(|rt| rt.is_active())
.unwrap_or(false)
}
pub fn is_agent_execution_active(&self) -> bool {
self.change_runtime.values().any(|rt| {
matches!(
rt.activity,
ActivityState::Applying | ActivityState::Accepting | ActivityState::Archiving
) && !rt.is_terminal()
})
}
pub fn is_resolving_active(&self) -> bool {
self.change_runtime
.values()
.any(|rt| matches!(rt.activity, ActivityState::Resolving) && !rt.is_terminal())
}
pub fn is_rejecting_active(&self) -> bool {
self.change_runtime
.values()
.any(|rt| matches!(rt.activity, ActivityState::Rejecting) && !rt.is_terminal())
}
pub fn has_other_post_archive_lane_blocker(&self, change_id: &str) -> bool {
self.change_runtime.iter().any(|(id, rt)| {
id != change_id
&& matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
)
&& !rt.is_terminal()
})
}
pub fn base_mutating_lane_occupant(&self) -> Option<String> {
self.change_runtime.iter().find_map(|(id, rt)| {
if matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
) && !rt.is_terminal()
{
Some(id.clone())
} else {
None
}
})
}
pub fn is_base_mutating_lane_occupied(&self) -> bool {
self.base_mutating_lane_occupant().is_some()
}
pub fn global_invariants_hold(&self) -> bool {
let lane_occupants = self
.change_runtime
.values()
.filter(|rt| {
matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
) && !rt.is_terminal()
})
.count();
lane_occupants <= 1
&& self
.change_runtime
.values()
.all(ChangeRuntimeState::invariants_hold)
}
pub fn merge_wait_change_ids(&self) -> Vec<String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
if matches!(rt.wait_state, WaitState::MergeWait) && !rt.is_terminal() {
Some(id.clone())
} else {
None
}
})
.collect()
}
pub fn resolve_wait_change_ids(&self) -> Vec<String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
if matches!(rt.wait_state, WaitState::ResolveWait) && !rt.is_terminal() {
Some(id.clone())
} else {
None
}
})
.collect()
}
fn remove_from_resolve_wait_queue(&mut self, change_id: &str) {
self.resolve_wait_queue.retain(|id| id != change_id);
}
fn remove_from_reject_wait_queue(&mut self, change_id: &str) {
self.reject_wait_queue.retain(|id| id != change_id);
}
fn clear_base_mutating_wait_queues(&mut self, change_id: &str) {
self.remove_from_resolve_wait_queue(change_id);
self.remove_from_reject_wait_queue(change_id);
}
fn enqueue_unique_resolve_wait(&mut self, change_id: &str) {
if !self.resolve_wait_queue.iter().any(|id| id == change_id) {
self.resolve_wait_queue.push(change_id.to_string());
}
}
fn enqueue_unique_reject_wait(&mut self, change_id: &str) {
if !self.reject_wait_queue.iter().any(|id| id == change_id) {
self.reject_wait_queue.push(change_id.to_string());
}
}
pub fn clear_resolve_wait_intent(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if matches!(rt.wait_state, WaitState::ResolveWait) {
rt.wait_state = if rt.is_terminal()
|| rt.dequeued
|| rt.is_active()
|| matches!(rt.queue_intent, QueueIntent::Queued)
{
WaitState::None
} else {
WaitState::MergeWait
};
rt.clear_blocked_metadata();
}
self.remove_from_resolve_wait_queue(change_id);
}
pub fn settle_stale_resolve_retry(
&mut self,
change_id: &str,
evidence: StaleResolveEvidence,
) -> StaleResolveSettlement {
match evidence {
StaleResolveEvidence::Proven => {
if !self
.runtime_entry(change_id)
.can_success_supersede_terminal()
{
self.clear_base_mutating_wait_queues(change_id);
return StaleResolveSettlement::AlreadySettled;
}
self.transition_change_to_merged(change_id);
StaleResolveSettlement::Merged
}
StaleResolveEvidence::Absent | StaleResolveEvidence::Unknown => {
let rt = self.runtime_entry(change_id);
if rt.is_terminal() || rt.dequeued {
self.clear_base_mutating_wait_queues(change_id);
return StaleResolveSettlement::AlreadySettled;
}
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::MergeWait;
rt.queue_intent = QueueIntent::NotQueued;
rt.clear_blocked_metadata();
self.clear_base_mutating_wait_queues(change_id);
StaleResolveSettlement::MergeWaitRetained
}
}
}
pub fn reject_wait_change_ids(&self) -> Vec<String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
if matches!(rt.wait_state, WaitState::RejectWait) && !rt.is_terminal() {
Some(id.clone())
} else {
None
}
})
.collect()
}
pub fn mark_reject_wait(&mut self, change_id: &str) {
let lane_blocked = self.has_other_post_archive_lane_blocker(change_id);
let rt = self.runtime_entry(change_id);
if rt.is_terminal() || rt.dequeued {
return;
}
rt.activity = if lane_blocked {
ActivityState::Idle
} else {
ActivityState::Rejecting
};
rt.wait_state = if lane_blocked {
WaitState::RejectWait
} else {
WaitState::None
};
rt.clear_blocked_metadata();
self.remove_from_resolve_wait_queue(change_id);
if lane_blocked {
self.enqueue_unique_reject_wait(change_id);
} else {
self.remove_from_reject_wait_queue(change_id);
}
}
pub fn clear_reject_wait_intent(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if matches!(rt.wait_state, WaitState::RejectWait) {
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
self.remove_from_reject_wait_queue(change_id);
}
pub fn release_base_mutating_lane_after_retry(
&mut self,
change_id: &str,
wait_state: WaitState,
) -> bool {
if !matches!(wait_state, WaitState::ResolveWait | WaitState::RejectWait) {
return false;
}
let rt = self.runtime_entry(change_id);
if rt.is_terminal()
|| !matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
)
{
return false;
}
rt.activity = ActivityState::Idle;
rt.wait_state = wait_state.clone();
rt.clear_blocked_metadata();
self.clear_base_mutating_wait_queues(change_id);
match wait_state {
WaitState::ResolveWait => self.enqueue_unique_resolve_wait(change_id),
WaitState::RejectWait => self.enqueue_unique_reject_wait(change_id),
_ => unreachable!("wait_state was validated above"),
}
true
}
pub fn abandon_base_mutating_lane_occupant(&mut self, change_id: &str) -> bool {
let rt = self.runtime_entry(change_id);
if rt.is_terminal()
|| !matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
)
{
return false;
}
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
self.clear_base_mutating_wait_queues(change_id);
true
}
pub fn promote_next_base_mutating_lane_waiter(&mut self) -> Option<(String, WaitState)> {
if self.is_base_mutating_lane_occupied() {
return None;
}
while let Some(change_id) = self.resolve_wait_queue.first().cloned() {
self.resolve_wait_queue.remove(0);
let rt = self.runtime_entry(&change_id);
if !rt.is_terminal() && matches!(rt.wait_state, WaitState::ResolveWait) {
rt.wait_state = WaitState::None;
rt.activity = ActivityState::Resolving;
rt.clear_blocked_metadata();
return Some((change_id, WaitState::ResolveWait));
}
}
while let Some(change_id) = self.reject_wait_queue.first().cloned() {
self.reject_wait_queue.remove(0);
let rt = self.runtime_entry(&change_id);
if !rt.is_terminal() && matches!(rt.wait_state, WaitState::RejectWait) {
rt.wait_state = WaitState::None;
rt.activity = ActivityState::Rejecting;
rt.clear_blocked_metadata();
return Some((change_id, WaitState::RejectWait));
}
}
None
}
pub fn has_manual_resolve_retry(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.is_some_and(|rt| rt.manual_resolve_retry)
}
pub fn consume_manual_resolve_retry(&mut self, change_id: &str) -> bool {
match self.change_runtime.get_mut(change_id) {
Some(rt) => std::mem::take(&mut rt.manual_resolve_retry),
None => false,
}
}
pub fn queued_change_ids(&self) -> Vec<String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
if matches!(rt.queue_intent, QueueIntent::Queued) && !rt.is_terminal() {
Some(id.clone())
} else {
None
}
})
.collect()
}
pub fn active_change_ids(&self) -> Vec<String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
if rt.is_active() {
Some(id.clone())
} else {
None
}
})
.collect()
}
pub fn all_display_statuses(&self) -> HashMap<String, &'static str> {
self.change_runtime
.iter()
.map(|(id, rt)| (id.clone(), rt.display_status()))
.collect()
}
pub fn all_apply_operation_labels(&self) -> HashMap<String, &'static str> {
self.change_runtime
.iter()
.map(|(id, rt)| (id.clone(), rt.apply_operation_label()))
.collect()
}
pub fn all_error_details(&self) -> HashMap<String, String> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| {
rt.error_message()
.map(|message| (id.clone(), crate::events::sanitize_detail(message)))
})
.collect()
}
pub fn all_blocker_views(&self) -> HashMap<String, BlockerView> {
self.change_runtime
.iter()
.filter_map(|(id, rt)| rt.blocker_view().map(|view| (id.clone(), view)))
.collect()
}
pub fn blocker_view(&self, change_id: &str) -> Option<BlockerView> {
self.change_runtime
.get(change_id)
.and_then(ChangeRuntimeState::blocker_view)
}
pub fn reconcile_dependency_blocker(
&mut self,
change_id: &str,
dependency_ids: &[String],
) -> bool {
let rt = self.runtime_entry(change_id);
if rt.is_terminal() || rt.is_active() {
return false;
}
rt.reconcile_dependency_blocker(dependency_ids)
}
pub fn clear_dependency_blocker(&mut self, change_id: &str) -> bool {
let Some(rt) = self.change_runtime.get_mut(change_id) else {
return false;
};
if !matches!(rt.wait_state, WaitState::DependencyBlocked) {
return false;
}
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
true
}
pub fn is_terminal_change(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.map(|rt| rt.is_terminal())
.unwrap_or(false)
}
pub fn remove_from_pending(&mut self, change_id: &str) {
self.pending_changes.remove(change_id);
if self.current_change_id.as_deref() == Some(change_id) {
self.current_change_id = None;
}
}
pub fn drop_pending_change(&mut self, change_id: &str) -> bool {
let removed = self.pending_changes.remove(change_id);
if removed {
self.total_changes = self.total_changes.saturating_sub(1);
}
if self.current_change_id.as_deref() == Some(change_id) {
self.current_change_id = None;
}
removed
}
pub fn clear_pending_changes(&mut self) {
self.pending_changes.clear();
}
pub fn record_error_and_check_circuit_breaker(
&mut self,
change_id: &str,
error: &str,
config: CircuitBreakerConfig,
) -> bool {
let history = self
.error_histories
.entry(change_id.to_string())
.or_insert_with(|| ErrorHistory::new(config.clone()));
history.record_error(error);
history.detect_same_error()
}
pub fn last_error(&self, change_id: &str) -> Option<&str> {
self.error_histories
.get(change_id)
.and_then(ErrorHistory::last_error)
}
pub fn clear_error_history(&mut self, change_id: &str) {
self.error_histories.remove(change_id);
}
pub fn is_final_terminal_dispatch_stop(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.map(|rt| {
matches!(
rt.terminal,
TerminalState::Merged | TerminalState::Pushed | TerminalState::Rejected(_)
)
})
.unwrap_or(false)
}
pub fn is_ordinary_queue_eligible(&self, change_id: &str) -> bool {
match self.change_runtime.get(change_id) {
None => false,
Some(rt) => {
!rt.is_terminal()
&& !rt.dequeued
&& (rt.is_active() || matches!(rt.queue_intent, QueueIntent::Queued))
}
}
}
pub fn ordinary_queue_eligible_change_ids(&self) -> HashSet<String> {
self.change_runtime
.keys()
.filter(|id| self.is_ordinary_queue_eligible(id))
.cloned()
.collect()
}
pub fn final_terminal_dispatch_stop_change_ids(&self) -> HashSet<String> {
self.change_runtime
.keys()
.filter(|id| self.is_final_terminal_dispatch_stop(id))
.cloned()
.collect()
}
pub fn settled_lifecycle_change_ids(&self) -> HashSet<String> {
self.change_runtime
.iter()
.filter(|(_, rt)| rt.is_terminal())
.map(|(id, _)| id.clone())
.collect()
}
pub fn terminal_error_change_ids(&self) -> HashSet<String> {
self.change_runtime
.iter()
.filter(|(_, rt)| matches!(rt.terminal, TerminalState::Error(_)))
.map(|(id, _)| id.clone())
.collect()
}
pub fn resolving_change_ids(&self) -> HashSet<String> {
self.active_change_ids()
.into_iter()
.filter(|id| self.display_status(id) == "resolving")
.collect()
}
pub fn is_terminal_error_change(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.map(|rt| matches!(rt.terminal, TerminalState::Error(_)))
.unwrap_or(false)
}
pub fn is_resumable_stopped(&self, change_id: &str) -> bool {
self.change_runtime
.get(change_id)
.map(|rt| {
matches!(rt.terminal, TerminalState::Stopped)
&& matches!(rt.queue_intent, QueueIntent::NotQueued)
&& matches!(rt.activity, ActivityState::Idle)
&& matches!(rt.wait_state, WaitState::None)
})
.unwrap_or(false)
}
pub fn resume_stopped_change(&mut self, change_id: &str) -> ReduceOutcome {
if !self.is_resumable_stopped(change_id) {
return ReduceOutcome::NoOp;
}
{
let rt = self.runtime_entry(change_id);
rt.terminal = TerminalState::None;
rt.clear_activity_wait_and_blocker();
rt.queue_intent = QueueIntent::Queued;
rt.dequeued = false;
rt.observation = WorkspaceObservation::None;
rt.commit_phase_attempt = None;
}
self.clear_stalled_change(change_id);
self.clear_base_mutating_wait_queues(change_id);
self.add_dynamic_change(change_id.to_string());
ReduceOutcome::Changed(ReducerEffect::QueueIntentSet {
change_id: change_id.to_string(),
intent: QueueIntent::Queued,
})
}
pub fn retry_terminal_error(&mut self, change_id: &str) -> ReduceOutcome {
{
let rt = self.runtime_entry(change_id);
if !matches!(rt.terminal, TerminalState::Error(_)) {
return ReduceOutcome::NoOp;
}
rt.terminal = TerminalState::None;
rt.clear_activity_wait_and_blocker();
rt.queue_intent = QueueIntent::Queued;
rt.dequeued = false;
rt.observation = WorkspaceObservation::None;
}
self.clear_apply_iteration_limit(change_id);
self.clear_stalled_change(change_id);
self.clear_error_history(change_id);
self.clear_base_mutating_wait_queues(change_id);
self.add_dynamic_change(change_id.to_string());
ReduceOutcome::Changed(ReducerEffect::QueueIntentSet {
change_id: change_id.to_string(),
intent: QueueIntent::Queued,
})
}
fn transition_change_to_dequeued(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
rt.terminal = TerminalState::None;
rt.clear_activity_wait_and_blocker();
rt.queue_intent = QueueIntent::NotQueued;
rt.dequeued = true;
rt.commit_phase_attempt = None;
self.clear_base_mutating_wait_queues(change_id);
}
fn on_run_stopped(&mut self) {
let interrupted: Vec<String> = self
.change_runtime
.iter()
.filter(|(_, rt)| {
!rt.is_terminal()
&& (rt.is_active()
|| matches!(rt.queue_intent, QueueIntent::Queued)
|| !matches!(rt.wait_state, WaitState::None))
})
.map(|(change_id, _)| change_id.clone())
.collect();
for change_id in interrupted {
self.transition_change_to_dequeued(&change_id);
self.stalled_change_ids.remove(&change_id);
self.skipped_change_ids.remove(&change_id);
if self.current_change_id.as_deref() == Some(change_id.as_str()) {
self.current_change_id = None;
}
}
}
fn transition_change_to_stopped(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
rt.transition_to_terminal(TerminalState::Stopped);
rt.queue_intent = QueueIntent::NotQueued;
rt.commit_phase_attempt = None;
self.clear_base_mutating_wait_queues(change_id);
}
fn transition_change_to_error(&mut self, change_id: &str, error: String) {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() {
rt.transition_to_terminal(TerminalState::Error(error));
}
rt.commit_phase_attempt = None;
self.clear_base_mutating_wait_queues(change_id);
}
fn transition_change_to_merged(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
rt.transition_to_terminal(TerminalState::Merged);
rt.queue_intent = QueueIntent::NotQueued;
self.clear_base_mutating_wait_queues(change_id);
}
fn transition_change_to_stalled(
&mut self,
change_id: &str,
blocker_reason: impl Into<String>,
unblock_metadata: impl Into<String>,
worktree_snapshot: impl Into<String>,
) {
self.runtime_entry(change_id).transition_to_stalled(
blocker_reason,
unblock_metadata,
worktree_snapshot,
);
}
fn clear_recoverable_terminal_for_success(rt: &mut ChangeRuntimeState) -> bool {
if rt.can_success_supersede_terminal() {
rt.terminal = TerminalState::None;
rt.clear_activity_wait_and_blocker();
true
} else {
false
}
}
fn apply_add_to_queue_command(&mut self, change_id: String) -> ReduceOutcome {
self.clear_stalled_change(&change_id);
{
let rt = self.runtime_entry(&change_id);
if rt.is_terminal() {
return ReduceOutcome::NoOp;
}
if !rt.dequeued
&& !matches!(
rt.wait_state,
WaitState::Stalled | WaitState::ExternalBlocked
)
&& (rt.is_active() || rt.queue_intent == QueueIntent::Queued)
{
return ReduceOutcome::NoOp;
}
rt.dequeued = false;
rt.queue_intent = QueueIntent::Queued;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
self.clear_base_mutating_wait_queues(&change_id);
self.add_dynamic_change(change_id.clone());
ReduceOutcome::Changed(ReducerEffect::QueueIntentSet {
change_id,
intent: QueueIntent::Queued,
})
}
fn apply_remove_from_queue_command(&mut self, change_id: String) -> ReduceOutcome {
let rt = self.runtime_entry(&change_id);
if rt.queue_intent == QueueIntent::NotQueued {
return ReduceOutcome::NoOp;
}
rt.queue_intent = QueueIntent::NotQueued;
ReduceOutcome::Changed(ReducerEffect::QueueIntentSet {
change_id,
intent: QueueIntent::NotQueued,
})
}
fn apply_resolve_merge_command(&mut self, change_id: String) -> ReduceOutcome {
{
let rt = self.runtime_entry(&change_id);
if !matches!(rt.wait_state, WaitState::MergeWait | WaitState::ResolveWait) {
if matches!(rt.activity, ActivityState::Idle)
&& matches!(rt.terminal, TerminalState::None)
{
rt.terminal = TerminalState::None;
rt.wait_state = WaitState::MergeWait;
} else {
return ReduceOutcome::NoOp;
}
}
rt.terminal = TerminalState::None;
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::ResolveWait;
rt.queue_intent = QueueIntent::NotQueued;
rt.manual_resolve_retry = true;
rt.clear_blocked_metadata();
}
self.remove_from_reject_wait_queue(&change_id);
self.enqueue_unique_resolve_wait(&change_id);
ReduceOutcome::Changed(ReducerEffect::WaitStateSet {
change_id,
wait: WaitState::ResolveWait,
})
}
fn apply_dequeue_change_command(&mut self, change_id: String) -> ReduceOutcome {
let rt = self.runtime_entry(&change_id);
if matches!(
rt.terminal,
TerminalState::Merged | TerminalState::Pushed | TerminalState::Rejected(_)
) {
return ReduceOutcome::NoOp;
}
self.transition_change_to_dequeued(&change_id);
ReduceOutcome::Changed(ReducerEffect::QueueIntentSet {
change_id,
intent: QueueIntent::NotQueued,
})
}
fn apply_stop_change_command(&mut self, change_id: String) -> ReduceOutcome {
let rt = self.runtime_entry(&change_id);
if rt.is_terminal() {
return ReduceOutcome::NoOp;
}
self.transition_change_to_stopped(&change_id);
ReduceOutcome::Changed(ReducerEffect::TerminalStateSet {
change_id,
terminal: TerminalState::Stopped,
})
}
pub fn apply_command(&mut self, cmd: ReducerCommand) -> ReduceOutcome {
match cmd {
ReducerCommand::AddToQueue(change_id) => self.apply_add_to_queue_command(change_id),
ReducerCommand::RetryError(change_id) => self.retry_terminal_error(&change_id),
ReducerCommand::RemoveFromQueue(change_id) => {
self.apply_remove_from_queue_command(change_id)
}
ReducerCommand::ResolveMerge(change_id) => self.apply_resolve_merge_command(change_id),
ReducerCommand::DequeueChange(change_id) => {
self.apply_dequeue_change_command(change_id)
}
ReducerCommand::StopChange(change_id) => self.apply_stop_change_command(change_id),
ReducerCommand::ResumeStopped(change_id) => self.resume_stopped_change(&change_id),
}
}
pub fn apply_observation(&mut self, change_id: &str, obs: WorkspaceObservation) {
let scheduler_pending = self.resolve_wait_queue.iter().any(|id| id == change_id)
|| self.reject_wait_queue.iter().any(|id| id == change_id);
let rt = self.runtime_entry(change_id);
if rt.is_active() {
return;
}
if rt.is_terminal() {
return;
}
match obs {
WorkspaceObservation::WorkspaceArchived => {
if matches!(rt.wait_state, WaitState::MergeWait) {
rt.observation = WorkspaceObservation::WorkspaceArchived;
} else if rt.is_fresh_idle() && !scheduler_pending {
rt.wait_state = WaitState::MergeWait;
rt.clear_blocked_metadata();
rt.observation = WorkspaceObservation::WorkspaceArchived;
}
}
WorkspaceObservation::WorktreeNotAhead => {
if matches!(rt.wait_state, WaitState::MergeWait) {
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
rt.observation = WorkspaceObservation::WorktreeNotAhead;
}
WorkspaceObservation::None => {
rt.observation = WorkspaceObservation::None;
}
}
}
fn on_processing_started(&mut self, change_id: &str) {
if self.is_terminal_error_change(change_id) {
return;
}
self.set_current_change(Some(change_id.to_string()));
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.queue_intent = QueueIntent::Queued;
}
}
fn on_processing_error(&mut self, change_id: &str, error: String) {
self.remove_from_pending(change_id);
let rt = self.runtime_entry(change_id);
rt.commit_phase_attempt = None;
if !rt.is_terminal() && !rt.dequeued {
rt.transition_to_terminal(TerminalState::Error(error));
}
}
fn on_workspace_preparation_started(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if rt.is_terminal() || rt.dequeued || !matches!(rt.activity, ActivityState::Idle) {
return;
}
rt.activity = ActivityState::Preparing;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
fn on_workspace_preparation_ended(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if matches!(rt.activity, ActivityState::Preparing) {
rt.activity = ActivityState::Idle;
}
}
fn on_apply_started(&mut self, change_id: &str) {
let mut should_start = false;
{
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Applying;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
should_start = true;
}
rt.commit_phase_attempt = None;
}
if should_start {
self.set_current_change(Some(change_id.to_string()));
}
}
fn on_apply_completed(&mut self, change_id: &str) {
self.increment_apply_count(change_id);
let rt = self.runtime_entry(change_id);
rt.commit_phase_attempt = None;
if matches!(rt.activity, ActivityState::Applying) {
rt.activity = ActivityState::Idle;
}
}
fn on_apply_commit_phase(
&mut self,
change_id: &str,
phase: crate::events::ApplyCommitPhase,
attempt: u32,
) {
let rt = self.runtime_entry(change_id);
rt.commit_phase_attempt = match phase.is_active() && !rt.is_terminal() && !rt.dequeued {
true => Some(attempt),
false => None,
};
}
fn on_acceptance_started(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Accepting;
}
}
fn on_acceptance_completed(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if matches!(rt.activity, ActivityState::Accepting) {
rt.activity = ActivityState::Idle;
}
}
fn on_change_rejected(&mut self, change_id: &str, reason: String) {
self.remove_from_pending(change_id);
self.remove_from_reject_wait_queue(change_id);
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() {
rt.transition_to_terminal(TerminalState::Rejected(reason));
rt.queue_intent = QueueIntent::NotQueued;
}
}
fn on_rejection_review_completed(
&mut self,
change_id: &str,
outcome: crate::events::RejectionOutcome,
) {
if matches!(outcome, crate::events::RejectionOutcome::Confirm) {
self.remove_from_pending(change_id);
}
self.remove_from_reject_wait_queue(change_id);
let rt = self.runtime_entry(change_id);
if rt.is_terminal() || rt.dequeued {
return;
}
match outcome {
crate::events::RejectionOutcome::Confirm => {
rt.transition_to_terminal(TerminalState::Rejected(
"rejecting review confirmed rejection".to_string(),
));
rt.queue_intent = QueueIntent::NotQueued;
}
crate::events::RejectionOutcome::Resume => {
rt.activity = ActivityState::Applying;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
rt.terminal = TerminalState::None;
}
crate::events::RejectionOutcome::Block => {
rt.transition_to_stalled(
"rejection review returned block; unresolved blocker remains",
"resolve unresolved blocker tasks in openspec/changes/<change_id>/tasks.md, then trigger explicit resume",
"existing worktree and WIP context are preserved for stalled rejection review",
);
}
}
}
fn on_rejection_review_failed(&mut self, change_id: &str, error: String) {
self.remove_from_pending(change_id);
self.remove_from_reject_wait_queue(change_id);
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.transition_to_terminal(TerminalState::Error(error));
}
}
fn on_workspace_status_updated(
&mut self,
change_id: &str,
status: &crate::vcs::WorkspaceStatus,
) {
if !self
.change_runtime
.get(change_id)
.is_some_and(|rt| !rt.is_terminal() && !rt.dequeued)
{
return;
}
match status {
crate::vcs::WorkspaceStatus::Blocked
if self
.change_runtime
.get(change_id)
.is_some_and(ChangeRuntimeState::has_structured_blocker_hold) =>
{
if let Some(rt) = self.change_runtime.get(change_id) {
tracing::debug!(
change_id = %change_id,
current_status = rt.display_status(),
blocker_kind = ?rt.blocker_kind(),
"Ignoring generic blocked workspace status because a structured blocker hold owns this wait"
);
}
}
crate::vcs::WorkspaceStatus::Blocked if self.stalled_change_ids.contains(change_id) => {
self.transition_change_to_stalled(
change_id,
"change stalled with recoverable blocker",
"resolve blocker evidence and retry explicitly",
"workflow state is preserved in the workspace",
);
}
crate::vcs::WorkspaceStatus::Rejecting => {
self.mark_reject_wait(change_id);
}
crate::vcs::WorkspaceStatus::Applying => {
self.runtime_entry(change_id).activity = ActivityState::Applying;
}
crate::vcs::WorkspaceStatus::Accepting => {
self.runtime_entry(change_id).activity = ActivityState::Accepting;
}
crate::vcs::WorkspaceStatus::Blocked => {
self.transition_change_to_stalled(
change_id,
"apply reported recoverable blocker; workspace remains stalled",
"resolve implementation blocker section and pending unblock tasks before explicit retry",
"existing worktree and WIP context are preserved while stalled",
);
}
crate::vcs::WorkspaceStatus::Archiving => {
self.runtime_entry(change_id).activity = ActivityState::Archiving;
}
crate::vcs::WorkspaceStatus::Resolving => {
self.on_workspace_status_resolving(change_id);
}
crate::vcs::WorkspaceStatus::MergeWait => {
self.on_workspace_status_merge_wait(change_id);
}
_ => {}
}
}
fn on_workspace_status_resolving(&mut self, change_id: &str) {
let has_other_lane_blocker = self.has_other_post_archive_lane_blocker(change_id);
let rt = self.runtime_entry(change_id);
if has_other_lane_blocker {
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::ResolveWait;
self.enqueue_unique_resolve_wait(change_id);
} else {
rt.activity = ActivityState::Resolving;
rt.wait_state = WaitState::None;
self.remove_from_resolve_wait_queue(change_id);
}
}
fn on_workspace_status_merge_wait(&mut self, change_id: &str) {
let rt = self.runtime_entry(change_id);
if matches!(
rt.activity,
ActivityState::Resolving | ActivityState::Rejecting
) || matches!(
rt.wait_state,
WaitState::ResolveWait | WaitState::RejectWait
) {
tracing::debug!(
change_id = %change_id,
current_status = rt.display_status(),
"Ignoring workspace MergeWait status because reducer owns stronger post-archive state"
);
} else {
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::MergeWait;
rt.queue_intent = QueueIntent::NotQueued;
self.remove_from_resolve_wait_queue(change_id);
}
}
pub fn apply_execution_event(&mut self, event: &crate::events::ExecutionEvent) {
use crate::events::ExecutionEvent;
match event {
ExecutionEvent::ProcessingStarted(change_id) => self.on_processing_started(change_id),
ExecutionEvent::ProcessingError { id, error } => {
self.on_processing_error(id, error.clone());
}
ExecutionEvent::ApplyStarted { change_id, .. } => self.on_apply_started(change_id),
ExecutionEvent::ApplyCompleted { change_id, .. } => self.on_apply_completed(change_id),
ExecutionEvent::ApplyFailed { change_id, error } => {
self.runtime_entry(change_id).commit_phase_attempt = None;
self.remove_from_pending(change_id);
self.transition_change_to_error(change_id, error.clone());
}
ExecutionEvent::ApplyCommitPhase {
change_id,
phase,
attempt,
} => self.on_apply_commit_phase(change_id, *phase, *attempt),
ExecutionEvent::AcceptanceStarted { change_id, .. } => {
self.on_acceptance_started(change_id);
}
ExecutionEvent::AcceptanceCompleted { change_id } => {
self.on_acceptance_completed(change_id);
}
ExecutionEvent::AcceptanceFailed { change_id, error } => {
self.transition_change_to_error(change_id, error.clone());
}
ExecutionEvent::ChangeRejected { change_id, reason } => {
self.on_change_rejected(change_id, reason.clone());
}
ExecutionEvent::RejectionReviewCompleted { change_id, outcome } => {
self.on_rejection_review_completed(change_id, *outcome);
}
ExecutionEvent::RejectionReviewFailed { change_id, error } => {
self.on_rejection_review_failed(change_id, error.clone());
}
ExecutionEvent::WorkspacePreparationStarted { change_id } => {
self.on_workspace_preparation_started(change_id);
}
ExecutionEvent::WorkspacePreparationEnded { change_id } => {
self.on_workspace_preparation_ended(change_id);
}
ExecutionEvent::WorkspaceStatusUpdated {
change_id, status, ..
} => {
self.on_workspace_status_updated(change_id, status);
}
ExecutionEvent::ArchiveStarted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Archiving;
}
}
ExecutionEvent::ArchiveResumed { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if matches!(rt.terminal, TerminalState::Error(_)) {
rt.terminal = TerminalState::None;
}
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Archiving;
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
}
ExecutionEvent::ArchiveRetryScheduled { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() {
rt.activity = ActivityState::Archiving;
}
}
ExecutionEvent::ChangeArchived(change_id) => {
self.mark_archived(change_id);
let has_other_resolve_lane_blocker =
self.change_runtime.iter().any(|(id, runtime)| {
id != change_id
&& matches!(
runtime.activity,
ActivityState::Resolving | ActivityState::Rejecting
)
&& !runtime.is_terminal()
});
let rt = self.runtime_entry(change_id);
if Self::clear_recoverable_terminal_for_success(rt) {
rt.queue_intent = QueueIntent::NotQueued;
if has_other_resolve_lane_blocker {
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::ResolveWait;
self.enqueue_unique_resolve_wait(change_id);
} else {
rt.activity = ActivityState::Resolving;
rt.wait_state = WaitState::None;
self.remove_from_resolve_wait_queue(change_id);
}
}
}
ExecutionEvent::ArchiveFailed {
change_id, error, ..
} => {
self.transition_change_to_error(change_id, error.clone());
}
ExecutionEvent::MergeDeferred {
change_id,
auto_resumable,
..
} => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() {
if *auto_resumable {
if !rt.is_active() {
rt.wait_state = WaitState::ResolveWait;
}
self.enqueue_unique_resolve_wait(change_id);
} else {
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::MergeWait;
rt.queue_intent = QueueIntent::NotQueued;
self.remove_from_resolve_wait_queue(change_id);
}
}
}
ExecutionEvent::MergeCompleted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if rt.can_success_supersede_terminal() {
self.transition_change_to_merged(change_id);
}
}
ExecutionEvent::PushStarted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Resolving;
rt.wait_state = WaitState::None;
}
}
ExecutionEvent::PushCompleted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if rt.can_success_supersede_terminal() {
rt.transition_to_terminal(TerminalState::Pushed);
rt.queue_intent = QueueIntent::NotQueued;
self.clear_base_mutating_wait_queues(change_id);
}
}
ExecutionEvent::PushFailed {
change_id, error, ..
} => {
self.transition_change_to_error(change_id, error.clone());
}
ExecutionEvent::ResolveStarted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Resolving;
rt.wait_state = WaitState::None;
}
}
ExecutionEvent::ResolveCompleted { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if rt.can_success_supersede_terminal()
&& !rt.dequeued
&& (matches!(rt.terminal, TerminalState::Error(_))
|| matches!(rt.activity, ActivityState::Resolving)
|| matches!(rt.wait_state, WaitState::ResolveWait))
{
self.transition_change_to_merged(change_id);
} else {
self.clear_base_mutating_wait_queues(change_id);
}
}
ExecutionEvent::ResolveFailed { change_id, .. } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() {
if matches!(rt.activity, ActivityState::Resolving) {
rt.activity = ActivityState::Idle;
}
if matches!(rt.activity, ActivityState::Idle)
&& matches!(rt.wait_state, WaitState::ResolveWait | WaitState::None)
{
rt.wait_state = WaitState::MergeWait;
self.remove_from_resolve_wait_queue(change_id);
}
}
}
ExecutionEvent::DependencyBlocked {
change_id,
dependency_ids,
} => {
let ids = dependency_ids.clone();
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.is_active() {
rt.reconcile_dependency_blocker(&ids);
}
}
ExecutionEvent::DependencyResolved { change_id } => {
let rt = self.runtime_entry(change_id);
if matches!(rt.wait_state, WaitState::DependencyBlocked) {
rt.wait_state = WaitState::None;
rt.clear_blocked_metadata();
}
}
ExecutionEvent::AcceptanceGated { change_id, blocker } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
let classification =
crate::orchestration::blocker_classification::classify_reported_facts(
blocker,
);
tracing::debug!(
change_id = %change_id,
reported_category = %blocker.category,
classified_as = ?classification.display_status(),
"classified acceptance blocker facts"
);
match classification {
crate::orchestration::blocker_classification::LifecycleClassification::ExternalBlocked(info) => {
rt.transition_to_external_blocked(&info, blocker.worktree_snapshot());
}
crate::orchestration::blocker_classification::LifecycleClassification::Stalled { detail, .. } => {
rt.transition_to_acceptance_stalled(
format!("acceptance-gated:{}", blocker.category),
detail,
blocker.worktree_snapshot(),
blocker.resumable,
);
}
crate::orchestration::blocker_classification::LifecycleClassification::ProtocolCorrection { .. } => {}
}
}
}
ExecutionEvent::ExecutionBlocked { change_id, blocker } => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
use crate::orchestration::blocker_classification as classification;
let classified = if classification::is_permission_denial(blocker) {
classification::classify_execution_stop(
classification::ExecutionStopReason::PermissionDenial,
blocker.summary(),
)
} else {
classification::classify_reported_facts(blocker)
};
let blocker_reason = format!("execution-blocked:{}", blocker.category);
match classified {
crate::orchestration::blocker_classification::LifecycleClassification::ExternalBlocked(info) => {
rt.transition_to_external_blocked(&info, blocker.worktree_snapshot());
}
_ if blocker.phase == "acceptance" => {
rt.transition_to_acceptance_stalled(
blocker_reason,
blocker.summary(),
blocker.worktree_snapshot(),
blocker.resumable,
);
}
_ => {
rt.transition_to_stalled(
blocker_reason,
blocker.summary(),
blocker.worktree_snapshot(),
);
}
}
}
}
ExecutionEvent::HookFailed {
change_id,
hook_type,
..
} if hook_type == crate::hooks::HookType::OnMerged.config_key() => {
let rt = self.runtime_entry(change_id);
if !rt.is_terminal() && !rt.dequeued {
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::MergeWait;
rt.queue_intent = QueueIntent::NotQueued;
self.remove_from_resolve_wait_queue(change_id);
}
}
ExecutionEvent::ChangeDequeued { change_id }
| ExecutionEvent::ChangeStopped { change_id } => {
let rt = self.runtime_entry(change_id);
if matches!(
rt.terminal,
TerminalState::Merged
| TerminalState::Pushed
| TerminalState::Rejected(_)
| TerminalState::Stopped
) {
return;
}
self.transition_change_to_dequeued(change_id);
}
ExecutionEvent::Stopped => self.on_run_stopped(),
ExecutionEvent::ChangesRefreshed {
changes,
merge_wait_ids,
rejected_changes: _,
worktree_not_ahead_ids,
..
} => {
let new_ids: Vec<String> = changes
.iter()
.filter(|c| {
!self.initial_change_ids.contains(&c.id)
&& !self.archived_changes.contains(&c.id)
})
.map(|c| c.id.clone())
.collect();
for id in new_ids {
self.add_dynamic_change(id);
}
for change in changes {
let rt = self.runtime_entry(&change.id);
if matches!(rt.terminal, TerminalState::Rejected(_)) {
rt.terminal = TerminalState::None;
rt.clear_activity_wait_and_blocker();
rt.queue_intent = QueueIntent::NotQueued;
}
}
let mw: Vec<String> = merge_wait_ids.iter().cloned().collect();
let nah: Vec<String> = worktree_not_ahead_ids.iter().cloned().collect();
for id in mw {
self.apply_observation(&id.clone(), WorkspaceObservation::WorkspaceArchived);
}
for id in nah {
self.apply_observation(&id.clone(), WorkspaceObservation::WorktreeNotAhead);
}
}
_ => {}
}
}
}
impl Default for OrchestratorState {
fn default() -> Self {
Self::new(Vec::new(), 0)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_new_state() {
let state =
OrchestratorState::new(vec!["change-a".to_string(), "change-b".to_string()], 10);
assert_eq!(state.total_changes(), 2);
assert_eq!(state.remaining_changes(), 2);
assert_eq!(state.changes_processed(), 0);
assert_eq!(state.iteration(), 0);
assert_eq!(state.max_iterations(), 10);
assert!(state.current_change_id().is_none());
}
#[test]
fn test_is_in_snapshot() {
let state = OrchestratorState::new(vec!["change-a".to_string()], 0);
assert!(state.is_in_snapshot("change-a"));
assert!(!state.is_in_snapshot("change-b"));
}
#[test]
fn test_apply_count_increment() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
assert_eq!(state.apply_count("change-a"), 0);
assert_eq!(state.increment_apply_count("change-a"), 1);
assert_eq!(state.increment_apply_count("change-a"), 2);
assert_eq!(state.apply_count("change-a"), 2);
}
#[test]
fn test_mark_archived() {
let mut state =
OrchestratorState::new(vec!["change-a".to_string(), "change-b".to_string()], 0);
state.set_current_change(Some("change-a".to_string()));
state.increment_apply_count("change-a");
assert!(state.is_pending("change-a"));
assert!(!state.is_archived("change-a"));
assert_eq!(state.remaining_changes(), 2);
state.mark_archived("change-a");
assert!(!state.is_pending("change-a"));
assert!(state.is_archived("change-a"));
assert_eq!(state.remaining_changes(), 1);
assert_eq!(state.changes_processed(), 1);
assert!(state.current_change_id().is_none());
assert_eq!(state.apply_count("change-a"), 0); }
#[test]
fn test_is_complete() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
assert!(!state.is_complete());
state.mark_archived("change-a");
assert!(state.is_complete());
}
#[test]
fn test_iteration_limit() {
let mut state = OrchestratorState::new(vec![], 10);
assert!(!state.is_iteration_limit_reached());
for _ in 0..8 {
state.increment_iteration();
}
assert!(state.is_approaching_iteration_limit()); assert!(!state.is_iteration_limit_reached());
state.increment_iteration(); assert!(!state.is_iteration_limit_reached());
state.increment_iteration(); assert!(state.is_iteration_limit_reached());
}
#[test]
fn test_no_iteration_limit() {
let mut state = OrchestratorState::new(vec![], 0);
for _ in 0..100 {
state.increment_iteration();
}
assert!(!state.is_iteration_limit_reached());
assert!(!state.is_approaching_iteration_limit());
}
#[test]
fn test_add_dynamic_change() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
assert_eq!(state.total_changes(), 1);
assert!(!state.is_pending("change-b"));
state.add_dynamic_change("change-b".to_string());
assert_eq!(state.total_changes(), 2);
assert!(state.is_pending("change-b"));
assert!(state.is_in_snapshot("change-b"));
}
#[test]
fn test_add_dynamic_change_idempotent() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
state.add_dynamic_change("change-a".to_string()); state.add_dynamic_change("change-b".to_string());
state.add_dynamic_change("change-b".to_string());
assert_eq!(state.total_changes(), 2); }
#[test]
fn test_remove_from_pending() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
state.set_current_change(Some("change-a".to_string()));
assert!(state.is_pending("change-a"));
assert!(state.current_change_id().is_some());
state.remove_from_pending("change-a");
assert!(!state.is_pending("change-a"));
assert!(!state.is_archived("change-a")); assert!(state.current_change_id().is_none());
}
#[test]
fn test_error_history_management_and_circuit_breaker() {
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
let config = CircuitBreakerConfig {
enabled: true,
threshold: 2,
};
assert!(!state.record_error_and_check_circuit_breaker(
"change-a",
"same error",
config.clone()
));
assert!(state.record_error_and_check_circuit_breaker("change-a", "same error", config));
assert!(state.last_error("change-a").is_some());
state.clear_error_history("change-a");
assert!(state.last_error("change-a").is_none());
}
#[test]
fn test_orchestrator_state_initializes_change_runtime() {
let state = OrchestratorState::new(vec!["change-a".to_string(), "change-b".to_string()], 0);
let rt_a = state
.change_runtime("change-a")
.expect("runtime for change-a");
assert_eq!(rt_a.queue_intent, QueueIntent::NotQueued);
assert_eq!(rt_a.activity, ActivityState::Idle);
assert_eq!(rt_a.wait_state, WaitState::None);
assert!(matches!(rt_a.terminal, TerminalState::None));
let rt_b = state
.change_runtime("change-b")
.expect("runtime for change-b");
assert_eq!(rt_b.queue_intent, QueueIntent::NotQueued);
}
#[test]
fn test_change_runtime_invariants() {
let valid = ChangeRuntimeState::default();
assert!(valid.invariants_hold());
let invalid = ChangeRuntimeState {
terminal: TerminalState::Merged,
activity: ActivityState::Applying,
..Default::default()
};
assert!(!invalid.invariants_hold());
let invalid2 = ChangeRuntimeState {
wait_state: WaitState::ResolveWait,
activity: ActivityState::Resolving,
..Default::default()
};
assert!(!invalid2.invariants_hold());
let ok = ChangeRuntimeState {
wait_state: WaitState::MergeWait,
..Default::default()
};
assert!(ok.invariants_hold());
let in_flight_rejecting = ChangeRuntimeState {
activity: ActivityState::Rejecting,
terminal: TerminalState::None,
..Default::default()
};
assert!(in_flight_rejecting.invariants_hold());
let invalid3 = ChangeRuntimeState {
wait_state: WaitState::RejectWait,
activity: ActivityState::Rejecting,
..Default::default()
};
assert!(!invalid3.invariants_hold());
}
#[test]
fn test_lifecycle_display_distinguishes_rejected_stalled_blocked_and_error() {
let rejected = ChangeRuntimeState {
terminal: TerminalState::Rejected("terminal".to_string()),
..Default::default()
};
assert_eq!(rejected.display_status(), "rejected");
let stalled = ChangeRuntimeState {
wait_state: WaitState::Stalled,
..Default::default()
};
assert_eq!(stalled.display_status(), "stalled");
let dependency_blocked = ChangeRuntimeState {
wait_state: WaitState::DependencyBlocked,
..Default::default()
};
assert_eq!(dependency_blocked.display_status(), "blocked");
let error = ChangeRuntimeState {
terminal: TerminalState::Error("repo-fixable failure".to_string()),
..Default::default()
};
assert_eq!(error.display_status(), "error");
}
#[test]
fn pushed_terminal_status_is_distinct_from_merged() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 1);
state.apply_execution_event(&crate::events::ExecutionEvent::PushCompleted {
change_id: "c".to_string(),
remote: "origin".to_string(),
branch: "c".to_string(),
});
assert_eq!(state.display_status("c"), "pushed");
assert!(matches!(
state.runtime_entry("c").terminal,
TerminalState::Pushed
));
}
#[test]
fn test_display_status_derivation() {
let state = OrchestratorState::new(vec!["c".to_string()], 0);
assert_eq!(state.display_status("c"), "not queued");
assert_eq!(state.display_status("unknown"), "not queued");
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.display_status("c"), "queued");
let rt = state.runtime_entry("c");
rt.activity = ActivityState::Applying;
assert_eq!(state.display_status("c"), "applying");
let rt = state.runtime_entry("c");
rt.terminal = TerminalState::Merged;
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_display_color_derivation() {
let mut rt = ChangeRuntimeState::default();
assert_eq!(rt.display_color(), ratatui::style::Color::DarkGray);
rt.queue_intent = QueueIntent::Queued;
assert_eq!(rt.display_color(), ratatui::style::Color::Yellow);
rt.wait_state = WaitState::DependencyBlocked;
assert_eq!(rt.display_color(), ratatui::style::Color::Gray);
rt.activity = ActivityState::Preparing;
assert_eq!(rt.display_status(), "preparing");
assert_eq!(rt.display_color(), ratatui::style::Color::Green);
rt.activity = ActivityState::Applying;
assert_eq!(rt.display_color(), ratatui::style::Color::Cyan);
rt.activity = ActivityState::Accepting;
assert_eq!(rt.display_color(), ratatui::style::Color::LightGreen);
rt.activity = ActivityState::Archiving;
assert_eq!(rt.display_color(), ratatui::style::Color::Magenta);
rt.activity = ActivityState::Resolving;
assert_eq!(rt.display_color(), ratatui::style::Color::LightCyan);
rt.activity = ActivityState::Idle;
rt.wait_state = WaitState::MergeWait;
assert_eq!(rt.display_color(), ratatui::style::Color::LightMagenta);
rt.wait_state = WaitState::ResolveWait;
assert_eq!(rt.display_color(), ratatui::style::Color::Magenta);
rt.wait_state = WaitState::RejectWait;
assert_eq!(rt.display_status(), "reject pending");
assert_eq!(rt.display_color(), ratatui::style::Color::LightMagenta);
rt.wait_state = WaitState::None;
rt.terminal = TerminalState::Merged;
assert_eq!(rt.display_color(), ratatui::style::Color::LightBlue);
rt.terminal = TerminalState::Rejected("blocked".to_string());
assert_eq!(rt.display_color(), ratatui::style::Color::LightRed);
rt.terminal = TerminalState::Stopped;
assert_eq!(rt.display_color(), ratatui::style::Color::DarkGray);
rt.terminal = TerminalState::Error("boom".to_string());
assert_eq!(rt.display_color(), ratatui::style::Color::Red);
}
#[test]
fn test_error_message_derivation() {
let rt = ChangeRuntimeState {
terminal: TerminalState::Error("fatal".to_string()),
..Default::default()
};
assert_eq!(rt.error_message(), Some("fatal"));
let rt2 = ChangeRuntimeState {
terminal: TerminalState::Merged,
..Default::default()
};
assert_eq!(rt2.error_message(), None);
}
#[test]
fn test_runtime_state_active_classification() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
assert!(!state.is_active_change("c"));
state.runtime_entry("c").activity = ActivityState::Applying;
assert!(state.is_active_change("c"));
state.runtime_entry("c").terminal = TerminalState::Merged;
state.runtime_entry("c").activity = ActivityState::Idle;
assert!(!state.is_active_change("c"));
assert!(state.is_terminal_change("c"));
}
#[test]
fn final_terminal_dispatch_stop_excludes_recoverable_terminal_error() {
let mut state = OrchestratorState::new(
vec![
"merged".to_string(),
"archived".to_string(),
"rejected".to_string(),
"error".to_string(),
],
0,
);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "merged".to_string(),
revision: "rev".to_string(),
});
state.runtime_entry("archived").terminal = TerminalState::Pushed;
state.runtime_entry("rejected").terminal = TerminalState::Rejected("no".to_string());
state.apply_execution_event(&crate::events::ExecutionEvent::ProcessingError {
id: "error".to_string(),
error: "boom".to_string(),
});
assert!(state.is_final_terminal_dispatch_stop("merged"));
assert!(state.is_final_terminal_dispatch_stop("archived"));
assert!(state.is_final_terminal_dispatch_stop("rejected"));
assert!(!state.is_final_terminal_dispatch_stop("error"));
assert!(state.is_terminal_error_change("error"));
assert!(matches!(
state.apply_command(ReducerCommand::AddToQueue("error".to_string())),
ReduceOutcome::NoOp
));
assert!(matches!(
state.apply_command(ReducerCommand::RetryError("error".to_string())),
ReduceOutcome::Changed(_)
));
assert_eq!(state.display_status("error"), "queued");
}
#[test]
fn test_apply_command_queue_intent() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
let outcome = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
let outcome2 = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome2, ReduceOutcome::NoOp));
state.mark_stalled("c".to_string());
state.runtime_entry("c").wait_state = WaitState::Stalled;
let outcome_from_stalled = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome_from_stalled, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
assert!(!state.stalled_change_ids().contains("c"));
let outcome3 = state.apply_command(ReducerCommand::RemoveFromQueue("c".to_string()));
assert!(matches!(outcome3, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "not queued");
let outcome4 = state.apply_command(ReducerCommand::RemoveFromQueue("c".to_string()));
assert!(matches!(outcome4, ReduceOutcome::NoOp));
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
let outcome5 = state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
assert!(matches!(outcome5, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "not queued");
let outcome6 = state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
assert!(matches!(outcome6, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "not queued");
state.runtime_entry("c").terminal = TerminalState::Error("boom".to_string());
let outcome_error_queue = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome_error_queue, ReduceOutcome::NoOp));
assert_eq!(state.display_status("c"), "error");
let outcome_error_retry = state.apply_command(ReducerCommand::RetryError("c".to_string()));
assert!(matches!(outcome_error_retry, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
state.runtime_entry("c").terminal = TerminalState::Rejected("blocked".to_string());
let outcome7 = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome7, ReduceOutcome::NoOp));
let outcome_rejected_retry =
state.apply_command(ReducerCommand::RetryError("c".to_string()));
assert!(matches!(outcome_rejected_retry, ReduceOutcome::NoOp));
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use std::collections::{HashMap, HashSet};
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![Change {
id: "c".to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}],
rejected_changes: Vec::new(),
committed_change_ids: HashSet::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
});
assert_eq!(state.display_status("c"), "not queued");
let outcome8 = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome8, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
}
#[test]
fn retry_terminal_error_clears_error_gate_and_stale_retry_metadata() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
state.apply_execution_event(&crate::events::ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "apply".to_string(),
});
state.record_error_and_check_circuit_breaker(
"c",
"boom",
CircuitBreakerConfig {
enabled: true,
threshold: 2,
},
);
state.mark_stalled("c".to_string());
state.apply_execution_event(&crate::events::ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "boom".to_string(),
});
assert_eq!(state.display_status("c"), "error");
assert!(state.queued_change_ids().is_empty());
assert!(state.last_error("c").is_some());
let outcome = state.apply_command(ReducerCommand::RetryError("c".to_string()));
assert!(matches!(outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
assert_eq!(state.queued_change_ids(), vec!["c".to_string()]);
assert!(state.last_error("c").is_none());
assert!(!state.stalled_change_ids().contains("c"));
let runtime = state.change_runtime("c").expect("runtime should exist");
assert!(matches!(runtime.terminal, TerminalState::None));
assert_eq!(runtime.wait_state, WaitState::None);
assert_eq!(runtime.blocked_metadata, BlockedMetadata::default());
}
#[test]
fn late_success_supersedes_recoverable_error_without_requeueing_apply() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::ApplyFailed {
change_id: "c".to_string(),
error: "boom".to_string(),
});
assert_eq!(state.display_status("c"), "error");
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
assert!(state.queued_change_ids().is_empty());
let runtime = state.change_runtime("c").expect("runtime should exist");
assert_eq!(runtime.queue_intent, QueueIntent::NotQueued);
}
#[test]
fn archive_resumed_clears_recoverable_archive_error_and_restores_archiving_activity() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ArchiveFailed {
change_id: "c".to_string(),
error: "Archive commit finalization failed".to_string(),
reason: Some("verification_failed".to_string()),
summary: Some("commit incomplete".to_string()),
});
assert_eq!(state.display_status("c"), "error");
state.apply_execution_event(&ExecutionEvent::ArchiveResumed {
change_id: "c".to_string(),
reason: Some("archive_commit_incomplete".to_string()),
summary: Some("resuming archived dirty repair".to_string()),
});
assert_eq!(state.display_status("c"), "archiving");
let runtime = state.change_runtime("c").expect("runtime should exist");
assert!(matches!(runtime.terminal, TerminalState::None));
assert!(state.global_invariants_hold());
}
#[test]
fn test_base_mutating_lane_is_single_occupant_across_resolving_and_rejecting() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("a".to_string()));
assert_eq!(state.display_status("a"), "resolving");
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "b".to_string(),
workspace_name: "ws-b".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("b"), "reject pending");
assert_eq!(state.reject_wait_change_ids(), vec!["b".to_string()]);
assert!(state.has_other_post_archive_lane_blocker("b"));
assert!(state.global_invariants_hold());
}
#[test]
fn manual_resolve_permission_is_granted_only_by_explicit_operator_intent() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "merge lane busy".to_string(),
auto_resumable: true,
});
assert!(
!state.has_manual_resolve_retry("alpha"),
"an auto-resumable deferral is not operator intent"
);
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert!(state.has_manual_resolve_retry("alpha"));
}
#[test]
fn manual_resolve_permission_is_consumed_once_and_is_not_sticky() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "bounded resolve exhausted".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("alpha".to_string(), WaitState::ResolveWait))
);
assert!(state.has_manual_resolve_retry("alpha"));
assert!(state.consume_manual_resolve_retry("alpha"));
assert!(!state.has_manual_resolve_retry("alpha"));
assert!(
!state.consume_manual_resolve_retry("alpha"),
"a second dispatch gets the unchanged generic preflight"
);
assert!(!state.consume_manual_resolve_retry("unknown-change"));
}
#[test]
fn manual_resolve_permission_does_not_survive_dequeue() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "bounded resolve exhausted".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert!(state.has_manual_resolve_retry("alpha"));
state.apply_command(ReducerCommand::DequeueChange("alpha".to_string()));
assert!(!state.has_manual_resolve_retry("alpha"));
assert!(state.global_invariants_hold());
}
#[test]
fn release_base_mutating_lane_after_retry_restores_resolve_wait_uniquely() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "conflict".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("alpha".to_string(), WaitState::ResolveWait))
);
assert!(state.is_base_mutating_lane_occupied());
assert!(state.release_base_mutating_lane_after_retry("alpha", WaitState::ResolveWait));
assert!(!state.is_base_mutating_lane_occupied());
assert_eq!(state.display_status("alpha"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["alpha".to_string()]);
assert!(state.global_invariants_hold());
assert!(!state.release_base_mutating_lane_after_retry("alpha", WaitState::ResolveWait));
assert_eq!(state.resolve_wait_queue, vec!["alpha".to_string()]);
assert!(state.reject_wait_queue.is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn abandon_base_mutating_lane_occupant_releases_resolve_without_requeueing() {
let mut state = OrchestratorState::new(vec!["alpha".to_string(), "beta".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "conflict".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("alpha".to_string(), WaitState::ResolveWait))
);
state.resolve_wait_queue.push("alpha".to_string());
state.reject_wait_queue.push("alpha".to_string());
state
.runtime_entry("alpha")
.set_blocked_metadata("reason", "unblock", "snapshot");
assert!(state.abandon_base_mutating_lane_occupant("alpha"));
let runtime = state.change_runtime("alpha").expect("runtime for alpha");
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::None);
assert_eq!(runtime.blocked_metadata, BlockedMetadata::default());
assert!(!state
.resolve_wait_change_ids()
.contains(&"alpha".to_string()));
assert!(!state
.reject_wait_change_ids()
.contains(&"alpha".to_string()));
assert!(!state.resolve_wait_queue.contains(&"alpha".to_string()));
assert!(!state.reject_wait_queue.contains(&"alpha".to_string()));
assert!(!state.is_base_mutating_lane_occupied());
assert!(state.global_invariants_hold());
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "beta".to_string(),
reason: "conflict".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("beta".to_string()));
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("beta".to_string(), WaitState::ResolveWait))
);
}
#[test]
fn abandon_base_mutating_lane_occupant_releases_reject_and_noops_terminal_or_non_occupant() {
let mut state = OrchestratorState::new(
vec![
"lane".to_string(),
"reject".to_string(),
"terminal".to_string(),
],
0,
);
state.apply_execution_event(&crate::events::ExecutionEvent::ChangeArchived(
"lane".to_string(),
));
state.mark_reject_wait("reject");
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "lane".to_string(),
revision: "rev".to_string(),
});
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("reject".to_string(), WaitState::RejectWait))
);
assert!(state.abandon_base_mutating_lane_occupant("reject"));
let runtime = state.change_runtime("reject").expect("runtime for reject");
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::None);
assert!(state.reject_wait_change_ids().is_empty());
assert!(!state.is_base_mutating_lane_occupied());
assert!(state.global_invariants_hold());
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "terminal".to_string(),
revision: "rev".to_string(),
});
assert!(!state.abandon_base_mutating_lane_occupant("terminal"));
assert!(!state.abandon_base_mutating_lane_occupant("unknown"));
assert!(state.global_invariants_hold());
}
#[test]
fn release_base_mutating_lane_after_retry_restores_reject_wait_and_noops_terminal() {
let mut state = OrchestratorState::new(vec!["lane".to_string(), "reject".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::ChangeArchived(
"lane".to_string(),
));
state.mark_reject_wait("reject");
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "lane".to_string(),
revision: "rev".to_string(),
});
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("reject".to_string(), WaitState::RejectWait))
);
assert!(state.release_base_mutating_lane_after_retry("reject", WaitState::RejectWait));
assert!(!state.is_base_mutating_lane_occupied());
assert_eq!(state.display_status("reject"), "reject pending");
assert_eq!(state.reject_wait_change_ids(), vec!["reject".to_string()]);
assert!(state.global_invariants_hold());
state.apply_execution_event(&crate::events::ExecutionEvent::ChangeRejected {
change_id: "reject".to_string(),
reason: "confirmed".to_string(),
});
assert!(!state.release_base_mutating_lane_after_retry("reject", WaitState::RejectWait));
assert!(state.reject_wait_change_ids().is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn test_reject_wait_queue_membership_and_clear_on_start_completion() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("a".to_string()));
state.mark_reject_wait("b");
assert_eq!(state.display_status("b"), "reject pending");
assert_eq!(state.reject_wait_change_ids(), vec!["b".to_string()]);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "a".to_string(),
revision: "rev".to_string(),
});
let promoted = state.promote_next_base_mutating_lane_waiter();
assert_eq!(promoted, Some(("b".to_string(), WaitState::RejectWait)));
assert_eq!(state.display_status("b"), "rejecting");
assert!(state.reject_wait_change_ids().is_empty());
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "b".to_string(),
outcome: RejectionOutcome::Confirm,
});
assert_eq!(state.display_status("b"), "rejected");
assert!(state.reject_wait_change_ids().is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn command_side_effects_update_terminal_wait_and_base_lane_queues() {
let mut state = OrchestratorState::new(vec!["c".to_string(), "d".to_string()], 0);
state.resolve_wait_queue.push("c".to_string());
state.reject_wait_queue.push("c".to_string());
let add_outcome = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(add_outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
assert!(state.resolve_wait_queue.is_empty());
assert!(state.reject_wait_queue.is_empty());
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "manual conflict".to_string(),
auto_resumable: false,
});
let resolve_outcome = state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
assert!(matches!(
resolve_outcome,
ReduceOutcome::Changed(ReducerEffect::WaitStateSet {
wait: WaitState::ResolveWait,
..
})
));
assert_eq!(state.display_status("c"), "resolve pending");
assert_eq!(state.resolve_wait_queue, vec!["c".to_string()]);
assert!(state.reject_wait_queue.is_empty());
let dequeue_outcome = state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
assert!(matches!(dequeue_outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "not queued");
assert!(state.resolve_wait_queue.is_empty());
assert!(state.reject_wait_queue.is_empty());
let runtime = state.change_runtime("c").expect("runtime for c");
assert!(matches!(runtime.terminal, TerminalState::None));
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::None);
assert_eq!(runtime.queue_intent, QueueIntent::NotQueued);
assert!(runtime.dequeued);
state.apply_command(ReducerCommand::AddToQueue("d".to_string()));
let stop_outcome = state.apply_command(ReducerCommand::StopChange("d".to_string()));
assert!(matches!(stop_outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("d"), "stopped");
let runtime = state.change_runtime("d").expect("runtime for d");
assert!(matches!(runtime.terminal, TerminalState::Stopped));
assert_eq!(runtime.queue_intent, QueueIntent::NotQueued);
assert!(state.global_invariants_hold());
}
#[test]
fn execution_event_side_effects_update_wait_queues_and_blocked_metadata() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(
vec![
"resolving".to_string(),
"archived".to_string(),
"reject".to_string(),
],
0,
);
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "resolving".to_string(),
command: "resolve".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("archived".to_string()));
assert_eq!(state.display_status("archived"), "resolve pending");
assert_eq!(state.resolve_wait_queue, vec!["archived".to_string()]);
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "archived".to_string(),
error: "base dirty".to_string(),
});
assert_eq!(state.display_status("archived"), "merge wait");
assert!(state.resolve_wait_queue.is_empty());
state.mark_reject_wait("reject");
assert_eq!(state.reject_wait_queue, vec!["reject".to_string()]);
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "reject".to_string(),
outcome: RejectionOutcome::Block,
});
assert!(state.reject_wait_queue.is_empty());
let runtime = state.change_runtime("reject").expect("runtime for reject");
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert!(runtime.blocked_metadata.blocker_reason.is_some());
assert!(matches!(runtime.terminal, TerminalState::None));
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "reject".to_string(),
outcome: RejectionOutcome::Confirm,
});
assert_eq!(state.display_status("reject"), "rejected");
assert!(state.reject_wait_queue.is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn commit_presentation_never_changes_the_canonical_lifecycle() {
use crate::events::{ApplyCommitPhase, ExecutionEvent};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
let before = state.change_runtime("c").cloned().expect("runtime");
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: ApplyCommitPhase::Started,
attempt: 4,
});
let during = state.change_runtime("c").expect("runtime");
assert_eq!(
state.display_status("c"),
"applying",
"the canonical status must not become a commit state"
);
assert_eq!(during.apply_operation_label(), "commit");
assert_eq!(during.commit_phase_attempt, Some(4));
assert_eq!(during.activity, before.activity);
assert_eq!(during.wait_state, before.wait_state);
assert_eq!(during.queue_intent, before.queue_intent);
assert!(state.global_invariants_hold());
}
#[test]
fn every_terminal_commit_phase_clears_presentation() {
use crate::events::{ApplyCommitPhase, ExecutionEvent};
for closing in [ApplyCommitPhase::Completed, ApplyCommitPhase::Failed] {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: ApplyCommitPhase::Started,
attempt: 1,
});
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: closing,
attempt: 1,
});
let runtime = state.change_runtime("c").expect("runtime");
assert_eq!(
runtime.commit_phase_attempt, None,
"{closing:?} must clear commit presentation"
);
assert_eq!(runtime.apply_operation_label(), "apply");
assert_eq!(state.display_status("c"), "applying");
}
}
#[test]
fn cancelling_outside_the_commit_sequence_clears_commit_presentation() {
use crate::events::{ApplyCommitPhase, ExecutionEvent};
fn active_commit_then(cancel: impl Fn(&mut OrchestratorState)) -> Option<u32> {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: ApplyCommitPhase::Started,
attempt: 3,
});
assert_eq!(
state
.change_runtime("c")
.expect("runtime")
.commit_phase_attempt,
Some(3),
"precondition: the row is presenting an active commit phase"
);
cancel(&mut state);
assert!(state.global_invariants_hold());
state
.change_runtime("c")
.expect("runtime")
.commit_phase_attempt
}
assert_eq!(
active_commit_then(|state| {
state.apply_execution_event(&ExecutionEvent::ChangeDequeued {
change_id: "c".to_string(),
});
}),
None,
"ChangeDequeued must clear commit presentation"
);
assert_eq!(
active_commit_then(|state| {
state.apply_execution_event(&ExecutionEvent::ChangeStopped {
change_id: "c".to_string(),
});
}),
None,
"ChangeStopped must clear commit presentation"
);
assert_eq!(
active_commit_then(|state| {
state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
}),
None,
"the DequeueChange command must clear commit presentation"
);
assert_eq!(
active_commit_then(|state| {
state.apply_command(ReducerCommand::StopChange("c".to_string()));
}),
None,
"the StopChange command must clear commit presentation"
);
assert_eq!(
active_commit_then(|state| {
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "WIP snapshot failed".to_string(),
});
}),
None,
"a processing error must clear commit presentation"
);
}
#[test]
fn a_repair_iteration_restores_apply_presentation() {
use crate::events::{ApplyCommitPhase, ExecutionEvent};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: ApplyCommitPhase::Started,
attempt: 2,
});
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(
state
.change_runtime("c")
.expect("runtime")
.commit_phase_attempt,
None,
"ApplyStarted always clears commit presentation"
);
assert_eq!(
state.all_apply_operation_labels().get("c").copied(),
Some("apply")
);
}
#[test]
fn a_terminal_row_never_adopts_commit_presentation() {
use crate::events::{ApplyCommitPhase, ExecutionEvent};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyFailed {
change_id: "c".to_string(),
error: "boom".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ApplyCommitPhase {
change_id: "c".to_string(),
phase: ApplyCommitPhase::Started,
attempt: 1,
});
assert_eq!(state.display_status("c"), "error");
assert_eq!(
state
.change_runtime("c")
.expect("runtime")
.commit_phase_attempt,
None
);
}
#[test]
fn a_fresh_reducer_has_no_commit_presentation() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
assert_eq!(
state.all_apply_operation_labels().get("c").copied(),
Some("apply"),
"a fresh reducer presents the Apply lane, never a commit subphase"
);
state.apply_execution_event(&crate::events::ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(
state
.change_runtime("c")
.expect("runtime")
.commit_phase_attempt,
None,
"process-local commit presentation starts empty"
);
}
#[test]
fn test_apply_execution_event_transitions() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "applying");
state.apply_execution_event(&ExecutionEvent::ApplyCompleted {
change_id: "c".to_string(),
revision: "rev1".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "accepting");
state.apply_execution_event(&ExecutionEvent::AcceptanceCompleted {
change_id: "c".to_string(),
});
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("c"), "rejecting");
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "c".to_string(),
outcome: RejectionOutcome::Resume,
});
assert_eq!(state.display_status("c"), "applying");
assert_ne!(
state
.change_runtime("c")
.expect("runtime for c after rejecting resume")
.activity,
ActivityState::Rejecting
);
state.apply_execution_event(&ExecutionEvent::RejectionReviewFailed {
change_id: "c".to_string(),
error: "rejecting failed".to_string(),
});
assert_eq!(state.display_status("c"), "error");
assert_ne!(
state
.change_runtime("c")
.expect("runtime for c after rejecting failure")
.activity,
ActivityState::Rejecting
);
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ProcessingStarted("c".to_string()));
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("c"), "rejecting");
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "c".to_string(),
outcome: RejectionOutcome::Confirm,
});
assert_eq!(state.display_status("c"), "rejected");
assert_ne!(
state
.change_runtime("c")
.expect("runtime for c after rejecting confirm")
.activity,
ActivityState::Rejecting
);
}
#[test]
fn test_explicit_blocker_categories_reach_blocked_metadata_without_prose_inference() {
let cases = [
(
"credential",
"docker image pull failed: lookup registry-1.docker.io i/o timeout",
),
("infrastructure", "missing non-mockable external credential"),
(
"human_decision",
"agent-exec managed verification job still running",
),
("external_approval", "Cannot connect to the Docker daemon"),
(
"schema_incompatibility",
"package registry timeout fetching crate",
),
];
for (idx, (explicit_category, message)) in cases.iter().enumerate() {
let change_id = format!("c-{idx}");
let mut state = OrchestratorState::new(vec![change_id.clone()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::AcceptanceGated {
change_id: change_id.clone(),
blocker: crate::events::StalledBlocker::acceptance_external(
*explicit_category,
*message,
),
});
let runtime = state
.change_runtime(&change_id)
.expect("runtime after blocker classification");
assert_eq!(state.display_status(&change_id), "blocked");
assert_eq!(runtime.blocker_kind(), BlockerKind::External);
assert!(matches!(runtime.terminal, TerminalState::None));
assert_eq!(
runtime.blocked_metadata.blocker_reason.as_deref(),
Some(format!("external-blocked:{explicit_category}").as_str()),
"the explicit category must survive verbatim for {message}"
);
assert!(
runtime
.blocked_metadata
.unblock_metadata
.as_deref()
.unwrap_or_default()
.contains(message),
"metadata should preserve observed error summary for {message}"
);
}
}
#[test]
fn test_execution_blocked_permission_denial_transitions_to_stalled_with_operator_guidance() {
use crate::events::{ExecutionEvent, StalledBlocker};
let denial = crate::permission::classify_permission_denial(&[Some(
"Read permission denied for /private/secret.txt",
)])
.expect("denial should classify");
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "accept".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: "c".to_string(),
blocker: StalledBlocker::permission_denial("acceptance", &denial),
});
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(state.display_status("c"), "stalled");
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert!(matches!(runtime.terminal, TerminalState::None));
let metadata = runtime
.blocked_metadata
.unblock_metadata
.as_deref()
.expect("operator guidance metadata");
assert!(metadata.contains("permission/tool policy denial"));
assert!(metadata.contains("operator action"));
assert!(!metadata.to_ascii_lowercase().contains("dependency blocked"));
}
#[test]
fn test_acceptance_gated_transitions_to_external_blocked_with_structured_metadata() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "accept".to_string(),
});
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker::acceptance_external(
"pending_verification",
"docker image pull failed: lookup registry-1.docker.io i/o timeout",
),
});
let runtime = state
.change_runtime("c")
.expect("runtime for c after acceptance gated");
assert_eq!(state.display_status("c"), "blocked");
assert_eq!(runtime.blocker_kind(), BlockerKind::External);
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::ExternalBlocked);
assert!(matches!(runtime.terminal, TerminalState::None));
assert_eq!(
runtime.blocked_metadata.blocker_reason.as_deref(),
Some("external-blocked:pending_verification")
);
let unblock = runtime
.blocked_metadata
.unblock_metadata
.as_deref()
.expect("unblock metadata");
assert!(unblock.contains("external blocker (pending_verification)"));
assert!(unblock.contains("reported by acceptance"));
assert!(unblock.contains("unblock when"));
assert!(unblock.contains("next action"));
assert!(unblock.contains("docker image pull failed"));
assert!(runtime.blocked_metadata.unblock_condition.is_some());
assert_eq!(
runtime.blocked_metadata.blocker_origin.as_deref(),
Some("acceptance")
);
assert_eq!(
runtime.blocked_metadata.worktree_snapshot.as_deref(),
Some("existing worktree and WIP context are preserved while stalled")
);
}
#[test]
fn test_unvalidated_acceptance_blocker_facts_remain_stalled() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker {
category: "acceptance_finding".to_string(),
phase: "acceptance".to_string(),
gate: "acceptance".to_string(),
error_summary: "unresolved finding".to_string(),
evidence: vec!["src/lib.rs:1 missing coverage".to_string()],
unblock_condition: None,
prerequisite_owner: None,
next_action: "resolve the finding and retry".to_string(),
resumable: true,
worktree_preserved: true,
},
});
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(state.display_status("c"), "stalled");
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert_eq!(runtime.blocker_kind(), BlockerKind::None);
assert!(runtime.blocked_metadata.unblock_condition.is_none());
assert!(runtime
.blocked_metadata
.unblock_metadata
.as_deref()
.unwrap_or_default()
.contains("'acceptance_finding' is not one of"));
}
#[test]
fn test_dependency_and_external_waits_share_blocked_but_keep_their_kind() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(
vec![
"dependency-wait".to_string(),
"external-wait".to_string(),
"execution-stall".to_string(),
],
0,
);
state.apply_execution_event(&ExecutionEvent::DependencyBlocked {
change_id: "dependency-wait".to_string(),
dependency_ids: vec!["alpha".to_string()],
});
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "external-wait".to_string(),
blocker: StalledBlocker::acceptance_external("credential", "STAGING_API_KEY is unset"),
});
let denial = crate::permission::classify_permission_denial(&[Some(
"Read permission denied for /private/secret.txt",
)])
.expect("denial should classify");
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: "execution-stall".to_string(),
blocker: StalledBlocker::permission_denial("acceptance", &denial),
});
assert_eq!(state.display_status("dependency-wait"), "blocked");
assert_eq!(state.display_status("external-wait"), "blocked");
assert_eq!(state.display_status("execution-stall"), "stalled");
assert_eq!(
state
.change_runtime("dependency-wait")
.unwrap()
.blocker_kind(),
BlockerKind::Dependency
);
assert_eq!(
state
.change_runtime("external-wait")
.unwrap()
.blocker_kind(),
BlockerKind::External
);
assert_eq!(
state
.change_runtime("execution-stall")
.unwrap()
.blocker_kind(),
BlockerKind::None
);
let views = state.all_blocker_views();
assert_eq!(views["dependency-wait"].status, "blocked");
assert_eq!(views["dependency-wait"].kind, BlockerKind::Dependency);
assert_eq!(views["external-wait"].status, "blocked");
assert_eq!(views["external-wait"].kind, BlockerKind::External);
assert_eq!(views["external-wait"].origin.as_deref(), Some("acceptance"));
assert!(views["external-wait"].unblock_condition.is_some());
assert_eq!(views["execution-stall"].status, "stalled");
assert_eq!(views["execution-stall"].kind, BlockerKind::None);
assert!(views["execution-stall"].unblock_condition.is_none());
assert_eq!(
state.externally_blocked_change_ids(),
HashSet::from(["external-wait".to_string()])
);
assert!(state
.change_runtime("external-wait")
.unwrap()
.blocked_metadata
.blocker_reason
.as_deref()
.unwrap_or_default()
.starts_with("external-blocked:"));
}
#[test]
fn acceptance_hold_is_in_memory_only_and_clears_on_retry_or_restart() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker::acceptance_external("credential", "STAGING_API_KEY is unset"),
});
let runtime = state.change_runtime("c").expect("runtime for c");
assert!(runtime.is_external_blocked());
assert!(runtime.is_acceptance_stalled());
assert!(runtime.is_resumable_acceptance_stall());
assert_eq!(
state.acceptance_stalled_change_ids(),
HashSet::from(["c".to_string()]),
"the hold must be visible to queue classification"
);
assert_eq!(state.display_status("c"), "blocked");
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(state.acceptance_stalled_change_ids().is_empty());
let runtime = state.change_runtime("c").expect("runtime for c");
assert!(!runtime.is_acceptance_stalled());
assert!(!runtime.is_external_blocked());
assert_eq!(runtime.blocker_kind(), BlockerKind::None);
assert!(runtime.blocked_metadata.unblock_condition.is_none());
assert_eq!(state.display_status("c"), "queued");
let restarted = OrchestratorState::new(vec!["c".to_string()], 0);
assert!(restarted.acceptance_stalled_change_ids().is_empty());
assert!(!restarted
.change_runtime("c")
.map(ChangeRuntimeState::is_external_blocked)
.unwrap_or(false));
}
#[test]
fn apply_phase_blocker_stalls_without_creating_an_acceptance_hold() {
use crate::events::{ExecutionEvent, StalledBlocker};
let denial = crate::permission::classify_permission_denial(&[Some(
"Read permission denied for /private/secret.txt",
)])
.expect("denial should classify");
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: "c".to_string(),
blocker: StalledBlocker::permission_denial("apply", &denial),
});
assert_eq!(state.display_status("c"), "stalled");
assert!(state.acceptance_stalled_change_ids().is_empty());
}
#[test]
fn non_resumable_acceptance_hold_is_not_retry_eligible() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker {
resumable: false,
..StalledBlocker::acceptance_external("human_decision", "owner must decide")
},
});
let runtime = state.change_runtime("c").expect("runtime for c");
assert!(runtime.is_acceptance_stalled());
assert!(!runtime.is_resumable_acceptance_stall());
assert_eq!(
state.acceptance_stalled_change_ids(),
HashSet::from(["c".to_string()])
);
}
fn observe_generic_blocked_workspace(state: &mut OrchestratorState, change_id: &str) {
state.apply_execution_event(&crate::events::ExecutionEvent::WorkspaceStatusUpdated {
change_id: change_id.to_string(),
workspace_name: format!("ws-{change_id}"),
status: crate::vcs::WorkspaceStatus::Blocked,
});
}
fn apply_external_blocker() -> crate::events::StalledBlocker {
crate::events::StalledBlocker {
category: "infrastructure".to_string(),
phase: "apply".to_string(),
gate: "apply".to_string(),
error_summary: "the build cache service is unreachable".to_string(),
evidence: vec!["cargo test: connection refused to cache.internal".to_string()],
unblock_condition: Some("cache.internal accepts connections again".to_string()),
prerequisite_owner: Some("platform".to_string()),
next_action: "restore cache.internal then retry apply".to_string(),
resumable: true,
worktree_preserved: true,
}
}
#[test]
fn structured_blocker_metadata_survives_workspace_blocked_for_an_acceptance_external_wait() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "accept".to_string(),
});
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker {
prerequisite_owner: Some("platform".to_string()),
..StalledBlocker::acceptance_external("credential", "STAGING_API_KEY is unset")
},
});
let before = state
.blocker_view("c")
.expect("the structured event must establish a blocker view");
observe_generic_blocked_workspace(&mut state, "c");
assert_eq!(state.display_status("c"), "blocked");
let after = state
.blocker_view("c")
.expect("the hold must survive the generic observation");
assert_eq!(
after, before,
"a lower-fidelity observation must not change any projected blocker field"
);
assert_eq!(after.kind, BlockerKind::External);
assert_eq!(after.origin.as_deref(), Some("acceptance"));
assert_eq!(after.prerequisite_owner.as_deref(), Some("platform"));
assert_eq!(
after.category.as_deref(),
Some("external-blocked:credential")
);
assert!(after
.unblock_condition
.as_deref()
.unwrap_or_default()
.contains("STAGING_API_KEY is unset"));
assert!(after
.detail
.as_deref()
.unwrap_or_default()
.contains("next action"));
assert!(after.resumable, "resumability drives explicit retry");
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(runtime.wait_state, WaitState::ExternalBlocked);
assert_eq!(runtime.activity, ActivityState::Idle);
assert!(matches!(runtime.terminal, TerminalState::None));
assert!(runtime.blocked_metadata.acceptance_stall);
assert!(runtime.is_resumable_acceptance_stall());
assert_eq!(
state.externally_blocked_change_ids(),
HashSet::from(["c".to_string()]),
"dispatch suppression must survive the generic observation"
);
assert_eq!(
state.acceptance_stalled_change_ids(),
HashSet::from(["c".to_string()])
);
}
#[test]
fn structured_blocker_metadata_survives_workspace_blocked_for_an_apply_external_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: "c".to_string(),
blocker: apply_external_blocker(),
});
let before = state
.blocker_view("c")
.expect("a validated apply claim must establish a blocker view");
observe_generic_blocked_workspace(&mut state, "c");
assert_eq!(state.display_status("c"), "blocked");
let after = state
.blocker_view("c")
.expect("the hold must survive the generic observation");
assert_eq!(after, before);
assert_eq!(after.kind, BlockerKind::External);
assert_eq!(after.origin.as_deref(), Some("apply"));
assert_eq!(after.prerequisite_owner.as_deref(), Some("platform"));
assert_eq!(
after.category.as_deref(),
Some("external-blocked:infrastructure")
);
assert_eq!(
after.unblock_condition.as_deref(),
Some("cache.internal accepts connections again")
);
assert!(after.resumable);
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(runtime.wait_state, WaitState::ExternalBlocked);
assert!(
!runtime.blocked_metadata.acceptance_stall,
"an apply-origin external wait is not an Acceptance-owned hold"
);
assert_eq!(
state.externally_blocked_change_ids(),
HashSet::from(["c".to_string()]),
"queue suppression uses the external set for either origin"
);
assert!(state.acceptance_stalled_change_ids().is_empty());
}
#[test]
fn structured_blocker_metadata_survives_workspace_blocked_for_an_acceptance_owned_stall() {
use crate::events::{ExecutionEvent, StalledBlocker};
let denial = crate::permission::classify_permission_denial(&[Some(
"Read permission denied for /private/secret.txt",
)])
.expect("denial should classify");
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: "c".to_string(),
blocker: StalledBlocker::permission_denial("acceptance", &denial),
});
let before = state
.blocker_view("c")
.expect("an acceptance-owned hold must establish a blocker view");
observe_generic_blocked_workspace(&mut state, "c");
assert_eq!(state.display_status("c"), "stalled");
let after = state
.blocker_view("c")
.expect("the hold must survive the generic observation");
assert_eq!(after, before);
assert_eq!(after.kind, BlockerKind::None);
assert!(after.resumable);
assert_eq!(
after.category.as_deref(),
Some("execution-blocked:permission:file_read")
);
assert!(after
.detail
.as_deref()
.unwrap_or_default()
.contains("permission"));
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert!(runtime.blocked_metadata.acceptance_stall);
assert!(runtime.is_resumable_acceptance_stall());
assert_eq!(
state.acceptance_stalled_change_ids(),
HashSet::from(["c".to_string()]),
"ordinary dispatch stays suppressed until explicit retry or restart"
);
assert!(state.externally_blocked_change_ids().is_empty());
}
#[test]
fn structured_blocker_metadata_survives_workspace_blocked_fallback_stays_generic_without_a_hold(
) {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "apply".to_string(),
});
observe_generic_blocked_workspace(&mut state, "c");
assert_eq!(state.display_status("c"), "stalled");
let view = state
.blocker_view("c")
.expect("the fallback still records a blocker view");
assert_eq!(view.kind, BlockerKind::None);
assert!(view.unblock_condition.is_none());
assert!(view.prerequisite_owner.is_none());
assert!(view.origin.is_none());
assert!(!view.resumable);
assert_eq!(
view.category.as_deref(),
Some("apply reported recoverable blocker; workspace remains stalled")
);
let runtime = state.change_runtime("c").expect("runtime for c");
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert!(!runtime.blocked_metadata.acceptance_stall);
assert!(state.acceptance_stalled_change_ids().is_empty());
assert!(state.externally_blocked_change_ids().is_empty());
let mut marked = OrchestratorState::new(vec!["d".to_string()], 0);
marked.mark_stalled("d".to_string());
observe_generic_blocked_workspace(&mut marked, "d");
assert_eq!(marked.display_status("d"), "stalled");
assert_eq!(
marked
.change_runtime("d")
.unwrap()
.blocked_metadata
.blocker_reason
.as_deref(),
Some("change stalled with recoverable blocker")
);
}
#[test]
fn structured_blocker_metadata_survives_workspace_blocked_under_duplicate_delivery() {
use crate::events::{ExecutionEvent, StalledBlocker};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: "c".to_string(),
blocker: StalledBlocker::acceptance_external("credential", "STAGING_API_KEY is unset"),
});
let before = state.blocker_view("c").expect("blocker view");
for _ in 0..3 {
observe_generic_blocked_workspace(&mut state, "c");
}
assert_eq!(state.blocker_view("c").as_ref(), Some(&before));
assert_eq!(state.display_status("c"), "blocked");
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.display_status("c"), "queued");
assert!(state.blocker_view("c").is_none());
assert!(state.externally_blocked_change_ids().is_empty());
}
#[test]
fn test_rejection_review_block_transitions_to_blocked_with_metadata() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ProcessingStarted("c".to_string()));
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("c"), "rejecting");
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "c".to_string(),
outcome: RejectionOutcome::Block,
});
let runtime = state
.change_runtime("c")
.expect("runtime for c after rejecting block");
assert_eq!(state.display_status("c"), "stalled");
assert_eq!(runtime.activity, ActivityState::Idle);
assert_eq!(runtime.wait_state, WaitState::Stalled);
assert!(matches!(runtime.terminal, TerminalState::None));
assert_eq!(
runtime.blocked_metadata.blocker_reason.as_deref(),
Some("rejection review returned block; unresolved blocker remains")
);
assert_eq!(
runtime.blocked_metadata.unblock_metadata.as_deref(),
Some(
"resolve unresolved blocker tasks in openspec/changes/<change_id>/tasks.md, then trigger explicit resume"
)
);
assert_eq!(
runtime.blocked_metadata.worktree_snapshot.as_deref(),
Some("existing worktree and WIP context are preserved for stalled rejection review")
);
}
#[test]
fn test_rejection_review_resume_from_archived_workspace_context_sets_applying() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "merge wait");
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("c"), "rejecting");
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "c".to_string(),
outcome: RejectionOutcome::Resume,
});
assert_eq!(state.display_status("c"), "applying");
}
#[test]
fn test_rejection_review_block_from_archived_workspace_context_sets_stalled() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "merge wait");
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("c"), "rejecting");
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "c".to_string(),
outcome: RejectionOutcome::Block,
});
assert_eq!(state.display_status("c"), "stalled");
}
#[test]
fn test_apply_observation_reconcile_merge_wait() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "merge wait");
state.apply_observation("c", WorkspaceObservation::WorktreeNotAhead);
assert_eq!(state.display_status("c"), "not queued");
state.runtime_entry("c").activity = ActivityState::Applying;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "applying");
}
#[test]
fn test_changes_refreshed_uses_reducer_observation_path() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
let mut merge_wait_ids = HashSet::new();
merge_wait_ids.insert("c".to_string());
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids,
});
assert_eq!(state.display_status("c"), "merge wait");
}
#[test]
fn test_change_rejected_clears_only_target_queue_intent() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("a".to_string()));
state.apply_command(ReducerCommand::AddToQueue("b".to_string()));
assert_eq!(state.display_status("a"), "queued");
assert_eq!(state.display_status("b"), "queued");
state.apply_execution_event(&ExecutionEvent::ChangeRejected {
change_id: "a".to_string(),
reason: "blocked".to_string(),
});
assert_eq!(state.display_status("a"), "rejected");
assert_eq!(state.display_status("b"), "queued");
}
#[test]
fn test_changes_refreshed_reactivates_rejected_change() {
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
let rt = state.runtime_entry("c");
rt.terminal = TerminalState::Rejected("blocked".to_string());
rt.activity = ActivityState::Rejecting;
rt.wait_state = WaitState::MergeWait;
rt.queue_intent = QueueIntent::Queued;
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![Change {
id: "c".to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}],
rejected_changes: Vec::new(),
committed_change_ids: HashSet::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
});
let rt = state.change_runtime("c").expect("runtime for c");
assert!(matches!(rt.terminal, TerminalState::None));
assert!(matches!(rt.activity, ActivityState::Idle));
assert!(matches!(rt.wait_state, WaitState::None));
assert!(matches!(rt.queue_intent, QueueIntent::NotQueued));
assert_eq!(state.display_status("c"), "not queued");
}
#[test]
fn test_merge_wait_release_after_external_merge() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "merge wait");
let mut not_ahead = HashSet::new();
not_ahead.insert("c".to_string());
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: not_ahead,
merge_wait_ids: HashSet::new(),
});
assert_eq!(state.display_status("c"), "not queued");
}
fn startup_refresh(merge_wait_ids: &[&str]) -> crate::events::ExecutionEvent {
use std::collections::{HashMap, HashSet};
crate::events::ExecutionEvent::ChangesRefreshed {
changes: vec![],
rejected_changes: Vec::new(),
committed_change_ids: HashSet::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: merge_wait_ids.iter().map(|id| id.to_string()).collect(),
}
}
#[test]
fn startup_refresh_restores_merge_wait_for_a_fresh_idle_change() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
assert_eq!(state.display_status("alpha"), "not queued");
state.apply_execution_event(&startup_refresh(&["alpha"]));
assert_eq!(state.display_status("alpha"), "merge wait");
assert!(matches!(
state.change_runtime("alpha").unwrap().wait_state,
WaitState::MergeWait
));
assert!(
state.resolve_wait_change_ids().is_empty(),
"restoring the wait must not enqueue scheduler-owned resolve work"
);
assert!(state.global_invariants_hold());
}
#[test]
fn startup_refresh_merge_wait_restoration_is_idempotent() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&startup_refresh(&["alpha"]));
state.apply_execution_event(&startup_refresh(&["alpha"]));
assert_eq!(state.display_status("alpha"), "merge wait");
assert!(state.resolve_wait_change_ids().is_empty());
}
#[test]
fn startup_refresh_does_not_regress_stronger_reducer_state_to_merge_wait() {
use crate::events::ExecutionEvent;
type CaseSetup = Box<dyn Fn(&mut OrchestratorState)>;
let cases: Vec<(&str, &str, CaseSetup)> = vec![
(
"resolving",
"resolving",
Box::new(|state: &mut OrchestratorState| {
state.apply_execution_event(&ExecutionEvent::ChangeArchived("x".to_string()));
}),
),
(
"resolve pending",
"resolve pending",
Box::new(|state: &mut OrchestratorState| {
state.runtime_entry("x").wait_state = WaitState::MergeWait;
state.apply_command(ReducerCommand::ResolveMerge("x".to_string()));
}),
),
(
"rejecting",
"rejecting",
Box::new(|state: &mut OrchestratorState| {
state.runtime_entry("x").activity = ActivityState::Rejecting;
}),
),
(
"reject pending",
"reject pending",
Box::new(|state: &mut OrchestratorState| {
state.runtime_entry("lane").activity = ActivityState::Resolving;
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "x".to_string(),
workspace_name: "ws-x".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
}),
),
(
"queued",
"queued",
Box::new(|state: &mut OrchestratorState| {
state.apply_command(ReducerCommand::AddToQueue("x".to_string()));
}),
),
(
"merged",
"merged",
Box::new(|state: &mut OrchestratorState| {
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "x".to_string(),
revision: "rev-x".to_string(),
});
}),
),
(
"rejected",
"rejected",
Box::new(|state: &mut OrchestratorState| {
state.apply_execution_event(&ExecutionEvent::ChangeRejected {
change_id: "x".to_string(),
reason: "blocked".to_string(),
});
}),
),
(
"error",
"error",
Box::new(|state: &mut OrchestratorState| {
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "x".to_string(),
error: "boom".to_string(),
});
}),
),
(
"stopped",
"stopped",
Box::new(|state: &mut OrchestratorState| {
state.apply_command(ReducerCommand::StopChange("x".to_string()));
}),
),
(
"blocked",
"blocked",
Box::new(|state: &mut OrchestratorState| {
state.apply_execution_event(&ExecutionEvent::DependencyBlocked {
change_id: "x".to_string(),
dependency_ids: vec!["dep".to_string()],
});
}),
),
(
"dequeued",
"not queued",
Box::new(|state: &mut OrchestratorState| {
state.apply_command(ReducerCommand::AddToQueue("x".to_string()));
state.apply_command(ReducerCommand::DequeueChange("x".to_string()));
}),
),
];
for (label, expected, setup) in cases {
let mut state = OrchestratorState::new(vec!["x".to_string(), "lane".to_string()], 0);
setup(&mut state);
assert_eq!(state.display_status("x"), expected, "setup for {}", label);
let resolve_wait_before = state.resolve_wait_change_ids();
let reject_wait_before = state.reject_wait_change_ids();
state.apply_execution_event(&startup_refresh(&["x"]));
assert_eq!(
state.display_status("x"),
expected,
"refresh evidence must not regress {} to merge wait",
label
);
assert_eq!(
state.resolve_wait_change_ids(),
resolve_wait_before,
"refresh must not create a duplicate manual resolve reservation for {}",
label
);
assert_eq!(
state.reject_wait_change_ids(),
reject_wait_before,
"refresh must not disturb reject-wait membership for {}",
label
);
assert!(state.global_invariants_hold(), "invariants for {}", label);
}
}
#[test]
fn test_workspace_archived_recovers_merge_wait() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.runtime_entry("c").wait_state = WaitState::MergeWait;
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "merge wait");
assert!(
matches!(
state.change_runtime("c").unwrap().wait_state,
WaitState::MergeWait
),
"observation should set MergeWait, not ResolveWait"
);
}
#[test]
fn test_queue_add_not_overwritten_by_merge_wait_refresh() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.display_status("c"), "queued");
state.apply_observation("c", WorkspaceObservation::WorkspaceArchived);
assert_eq!(state.display_status("c"), "queued");
let mut not_ahead = HashSet::new();
not_ahead.insert("c".to_string());
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: not_ahead,
merge_wait_ids: HashSet::new(),
});
assert_eq!(state.display_status("c"), "queued");
}
#[test]
fn test_base_mutating_lane_exclusivity_and_wait_membership() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(
vec![
"resolving-a".to_string(),
"rejecting-b".to_string(),
"archive-c".to_string(),
],
0,
);
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "resolving-a".to_string(),
workspace_name: "ws-a".to_string(),
status: crate::vcs::WorkspaceStatus::Resolving,
});
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "rejecting-b".to_string(),
workspace_name: "ws-b".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("resolving-a"), "resolving");
assert_eq!(state.display_status("rejecting-b"), "reject pending");
assert_eq!(
state.reject_wait_change_ids(),
vec!["rejecting-b".to_string()]
);
assert!(state.resolve_wait_change_ids().is_empty());
assert!(state.has_other_post_archive_lane_blocker("archive-c"));
assert!(state.global_invariants_hold());
state.apply_execution_event(&ExecutionEvent::ChangeArchived("archive-c".to_string()));
assert_eq!(state.display_status("archive-c"), "resolve pending");
assert_eq!(
state.resolve_wait_change_ids(),
vec!["archive-c".to_string()]
);
assert!(state.global_invariants_hold());
}
#[test]
fn test_reject_wait_clear_and_deterministic_single_promotion() {
use crate::events::{ExecutionEvent, RejectionOutcome};
let mut state = OrchestratorState::new(
vec![
"lane-a".to_string(),
"resolve-b".to_string(),
"reject-c".to_string(),
],
0,
);
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "lane-a".to_string(),
workspace_name: "ws-a".to_string(),
status: crate::vcs::WorkspaceStatus::Resolving,
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("resolve-b".to_string()));
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "reject-c".to_string(),
workspace_name: "ws-c".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("resolve-b"), "resolve pending");
assert_eq!(state.display_status("reject-c"), "reject pending");
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "lane-a".to_string(),
error: "manual blocker".to_string(),
});
let promoted = state.promote_next_base_mutating_lane_waiter();
assert_eq!(
promoted,
Some(("resolve-b".to_string(), WaitState::ResolveWait))
);
assert_eq!(state.display_status("resolve-b"), "resolving");
assert_eq!(state.display_status("reject-c"), "reject pending");
assert!(state.global_invariants_hold());
assert_eq!(state.promote_next_base_mutating_lane_waiter(), None);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "resolve-b".to_string(),
revision: "rev-b".to_string(),
});
let promoted = state.promote_next_base_mutating_lane_waiter();
assert_eq!(
promoted,
Some(("reject-c".to_string(), WaitState::RejectWait))
);
assert_eq!(state.display_status("reject-c"), "rejecting");
assert!(state.reject_wait_change_ids().is_empty());
assert!(state.global_invariants_hold());
state.apply_execution_event(&ExecutionEvent::RejectionReviewCompleted {
change_id: "reject-c".to_string(),
outcome: RejectionOutcome::Resume,
});
assert!(state.reject_wait_change_ids().is_empty());
assert_eq!(state.display_status("reject-c"), "applying");
}
#[test]
fn test_parallel_change_archived_no_blocker_enters_resolving_not_merge_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(state.display_status("c"), "resolving");
assert!(
state.resolve_wait_change_ids().is_empty(),
"no-blocker archive completion must enter immediate merge handling, not resolve wait"
);
assert!(
state.queued_change_ids().is_empty(),
"no-blocker archive completion must not be reintroduced as ordinary queued work"
);
}
#[test]
fn test_no_blocker_merge_wait_to_merged_vibration_regression() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(state.display_status("c"), "resolving");
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::from(["c".to_string()]),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::from(["c".to_string()]),
});
assert_eq!(
state.display_status("c"),
"resolving",
"workspace refresh must not turn active no-blocker merge handling into merge wait"
);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_merged_archived_vibration_regression() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::from(["c".to_string()]),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::from(["c".to_string()]),
});
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_manual_merge_deferred_clears_normal_queue_intent() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.queued_change_ids(), vec!["c".to_string()]);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
assert!(
state.queued_change_ids().is_empty(),
"manual merge deferral must consume ordinary queue intent"
);
assert!(
state.resolve_wait_change_ids().is_empty(),
"manual merge deferral must not remain scheduler-owned resolve retry intent"
);
}
#[test]
fn test_manual_merge_deferred_clears_existing_resolve_wait_queue_membership() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "another merge is in progress".to_string(),
auto_resumable: true,
});
assert_eq!(state.display_status("c"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["c".to_string()]);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
assert!(state.queued_change_ids().is_empty());
assert!(state.resolve_wait_change_ids().is_empty());
}
#[test]
fn test_manual_merge_deferred_resolve_merge_explicit_retry_sets_resolve_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
assert!(state.queued_change_ids().is_empty());
let outcome = state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
assert!(matches!(outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["c".to_string()]);
assert!(
state.queued_change_ids().is_empty(),
"explicit merge retry must use resolve-wait intent, not normal queue intent"
);
}
#[test]
fn test_archived_change_enters_active_merge_handling_before_manual_retry() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::ChangeArchived(
"alpha".to_string(),
));
assert_eq!(
state.display_status("alpha"),
"resolving",
"archive alone must never be terminal"
);
assert!(
matches!(
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string())),
ReduceOutcome::NoOp
),
"a row already in active merge handling needs no manual retry"
);
assert_eq!(state.display_status("alpha"), "resolving");
state.apply_execution_event(&crate::events::ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "manual resolution required".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("alpha"), "merge wait");
let outcome = state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert!(
matches!(outcome, ReduceOutcome::Changed(_)),
"explicit manual retry for a merge-wait row must not be dropped"
);
assert_eq!(state.display_status("alpha"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["alpha".to_string()]);
assert!(state.queued_change_ids().is_empty());
}
#[test]
fn test_merged_manual_retry_remains_noop() {
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&crate::events::ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "rev-alpha".to_string(),
});
let outcome = state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert!(matches!(outcome, ReduceOutcome::NoOp));
assert_eq!(state.display_status("alpha"), "merged");
assert!(state.resolve_wait_change_ids().is_empty());
}
#[test]
fn test_auto_resumable_merge_deferred_keeps_scheduler_retry_intent() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "another merge is in progress".to_string(),
auto_resumable: true,
});
assert_eq!(state.display_status("c"), "resolve pending");
assert_eq!(state.queued_change_ids(), vec!["c".to_string()]);
assert_eq!(state.resolve_wait_change_ids(), vec!["c".to_string()]);
}
#[test]
fn test_archive_merge_defers_to_resolve_pending_when_rejecting_active() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "a".to_string(),
workspace_name: "ws-a".to_string(),
status: crate::vcs::WorkspaceStatus::Rejecting,
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("b".to_string()));
assert_eq!(state.display_status("a"), "rejecting");
assert_eq!(state.display_status("b"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["b".to_string()]);
assert!(state.reject_wait_change_ids().is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn test_active_applying_does_not_create_resolve_pending_on_archive() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "a".to_string(),
command: "apply".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("b".to_string()));
assert_eq!(state.display_status("a"), "applying");
assert_eq!(state.display_status("b"), "resolving");
assert!(state.resolve_wait_change_ids().is_empty());
}
#[test]
fn test_parallel_merge_events_drive_reducer_wait_states() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "c".to_string(),
command: "resolve".to_string(),
});
assert_eq!(state.display_status("c"), "resolving");
state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "c".to_string(),
worktree_change_ids: None,
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "late".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_merge_completed_ignores_later_stale_manual_merge_deferred() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "rev-alpha".to_string(),
});
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "Archive incomplete for 'alpha': worktree may be dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("alpha"), "merged");
assert!(state.queued_change_ids().is_empty());
assert!(state.resolve_wait_change_ids().is_empty());
}
#[test]
fn test_merge_completed_clears_resolve_wait_intent() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "Resolve in progress for another change".to_string(),
auto_resumable: true,
});
assert_eq!(state.resolve_wait_change_ids(), vec!["alpha".to_string()]);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "rev-alpha".to_string(),
});
assert_eq!(state.display_status("alpha"), "merged");
assert!(
state.resolve_wait_change_ids().is_empty(),
"MergeCompleted must clear reducer-owned ResolveWait retry intent"
);
}
#[test]
fn test_clear_resolve_wait_intent_removes_retry_without_terminal_transition() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "Resolve in progress for another change".to_string(),
auto_resumable: true,
});
assert_eq!(state.resolve_wait_change_ids(), vec!["alpha".to_string()]);
state.clear_resolve_wait_intent("alpha");
assert!(state.resolve_wait_change_ids().is_empty());
assert_eq!(state.display_status("alpha"), "merge wait");
assert!(!state.is_terminal_change("alpha"));
}
fn manual_resolve_retry_state(change_id: &str) -> OrchestratorState {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec![change_id.to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "manual conflict".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status(change_id), "merge wait");
state.apply_command(ReducerCommand::ResolveMerge(change_id.to_string()));
assert_eq!(state.display_status(change_id), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec![change_id.to_string()]);
state
}
#[test]
fn stale_deferred_merge_retry_with_proven_base_integration_settles_merged() {
let mut state = manual_resolve_retry_state("alpha");
let settlement = state.settle_stale_resolve_retry("alpha", StaleResolveEvidence::Proven);
assert_eq!(settlement, StaleResolveSettlement::Merged);
assert_eq!(state.display_status("alpha"), "merged");
assert_ne!(state.display_status("alpha"), "not queued");
assert!(state.is_terminal_change("alpha"));
assert!(state.resolve_wait_change_ids().is_empty());
assert!(state.resolve_wait_queue.is_empty());
assert!(!state.is_base_mutating_lane_occupied());
assert!(state.global_invariants_hold());
}
#[test]
fn stale_deferred_merge_retry_without_proven_integration_retains_merge_wait() {
for evidence in [StaleResolveEvidence::Absent, StaleResolveEvidence::Unknown] {
let mut state = manual_resolve_retry_state("alpha");
let settlement = state.settle_stale_resolve_retry("alpha", evidence);
assert_eq!(
settlement,
StaleResolveSettlement::MergeWaitRetained,
"{evidence:?} proves nothing and must stay retryable"
);
assert_eq!(
state.display_status("alpha"),
"merge wait",
"{evidence:?} must not expose not queued or merged"
);
assert!(!state.is_terminal_change("alpha"));
assert!(state.resolve_wait_change_ids().is_empty());
assert!(state.resolve_wait_queue.is_empty());
assert!(state.global_invariants_hold());
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string()));
assert_eq!(state.display_status("alpha"), "resolve pending");
}
}
#[test]
fn stale_deferred_merge_retry_settles_promoted_lane_occupant() {
let mut state = manual_resolve_retry_state("alpha");
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("alpha".to_string(), WaitState::ResolveWait))
);
assert!(state.is_base_mutating_lane_occupied());
let settlement = state.settle_stale_resolve_retry("alpha", StaleResolveEvidence::Absent);
assert_eq!(settlement, StaleResolveSettlement::MergeWaitRetained);
assert_eq!(state.display_status("alpha"), "merge wait");
assert!(
!state.is_base_mutating_lane_occupied(),
"settlement releases base-lane ownership once the reducer state agrees"
);
assert!(state.global_invariants_hold());
}
#[test]
fn stale_deferred_merge_retry_preserves_immutable_terminal_outcomes() {
let mut state = manual_resolve_retry_state("alpha");
state.apply_command(ReducerCommand::StopChange("alpha".to_string()));
assert_eq!(state.display_status("alpha"), "stopped");
let settlement = state.settle_stale_resolve_retry("alpha", StaleResolveEvidence::Proven);
assert_eq!(settlement, StaleResolveSettlement::AlreadySettled);
assert_eq!(state.display_status("alpha"), "stopped");
assert!(state.resolve_wait_change_ids().is_empty());
assert!(state.resolve_wait_queue.is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn stale_deferred_merge_retry_keeps_bounded_resolve_failure_in_merge_wait() {
use crate::events::ExecutionEvent;
let mut state = manual_resolve_retry_state("alpha");
assert_eq!(
state.promote_next_base_mutating_lane_waiter(),
Some(("alpha".to_string(), WaitState::ResolveWait))
);
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "alpha".to_string(),
error: "bounded resolve exhausted".to_string(),
});
assert_eq!(state.display_status("alpha"), "merge wait");
assert!(!state.is_terminal_change("alpha"));
assert!(state.resolve_wait_change_ids().is_empty());
assert!(!state.is_base_mutating_lane_occupied());
assert!(state.global_invariants_hold());
}
#[test]
fn test_resolve_completed_clears_resolve_wait_and_survives_refresh() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
assert_eq!(state.display_status("c"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["c".to_string()]);
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: std::collections::HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: std::collections::HashSet::new(),
worktree_change_ids: std::collections::HashSet::from(["c".to_string()]),
worktree_paths: std::collections::HashMap::new(),
worktree_not_ahead_ids: std::collections::HashSet::new(),
merge_wait_ids: std::collections::HashSet::from(["c".to_string()]),
});
assert_eq!(state.display_status("c"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "c".to_string(),
worktree_change_ids: None,
});
assert_eq!(state.display_status("c"), "merged");
assert!(
state.resolve_wait_change_ids().is_empty(),
"resolve completion must clear queued resolve intent"
);
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: std::collections::HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: std::collections::HashSet::new(),
worktree_change_ids: std::collections::HashSet::from(["c".to_string()]),
worktree_paths: std::collections::HashMap::new(),
worktree_not_ahead_ids: std::collections::HashSet::new(),
merge_wait_ids: std::collections::HashSet::from(["c".to_string()]),
});
assert_eq!(state.display_status("c"), "merged");
assert_ne!(
state.display_status("c"),
"resolve pending",
"row must not regress to resolve pending after successful resolve + refresh"
);
}
#[test]
fn test_resolve_wait_manual_merge_deferred_demotes_to_merge_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["change-a".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "change-a".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("change-a".to_string()));
assert_eq!(state.display_status("change-a"), "resolve pending");
assert_eq!(
state.resolve_wait_change_ids(),
vec!["change-a".to_string()]
);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "change-a".to_string(),
reason: "Working tree has uncommitted changes".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("change-a"), "merge wait");
assert!(
state.resolve_wait_change_ids().is_empty(),
"manual retry deferral must remove reducer-owned ResolveWait membership"
);
assert!(state.queued_change_ids().is_empty());
assert!(state.global_invariants_hold());
}
#[test]
fn test_resolve_wait_clean_lane_promotion_promotes_exactly_one_waiter() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["alpha".to_string(), "beta".to_string()], 0);
for change_id in ["alpha", "beta"] {
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge(change_id.to_string()));
}
let promoted = state.promote_next_base_mutating_lane_waiter();
assert_eq!(
promoted,
Some(("alpha".to_string(), WaitState::ResolveWait))
);
assert_eq!(state.display_status("alpha"), "resolving");
assert_eq!(state.display_status("beta"), "resolve pending");
assert_eq!(state.resolve_wait_change_ids(), vec!["beta".to_string()]);
assert_eq!(state.promote_next_base_mutating_lane_waiter(), None);
assert!(state.global_invariants_hold());
}
#[test]
fn test_resolve_failed_restores_merge_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "c".to_string(),
command: "resolve".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "conflict".to_string(),
});
assert_eq!(state.display_status("c"), "merge wait");
let mut state2 = OrchestratorState::new(vec!["c".to_string()], 0);
state2.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
state2.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "late".to_string(),
});
assert_eq!(state2.display_status("c"), "merged");
}
#[test]
fn test_workspace_status_update_targets_explicit_change_id() {
use crate::events::ExecutionEvent;
use crate::vcs::WorkspaceStatus;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "a".to_string(),
command: "apply-a".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "b".to_string(),
command: "apply-b".to_string(),
});
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "b".to_string(),
workspace_name: "ws-b".to_string(),
status: WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("a"), "applying");
assert_eq!(state.display_status("b"), "rejecting");
}
#[test]
fn test_late_events_after_stop_do_not_regress_state() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeDequeued {
change_id: "c".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "late error".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
}
#[test]
fn test_reducer_idempotency_and_precedence() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(state.display_status("c"), "resolving");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "late".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_reducer_runtime_and_legacy_aggregates_stay_consistent() {
use crate::events::ExecutionEvent;
let mut state =
OrchestratorState::new(vec!["a".to_string(), "b".to_string(), "c".to_string()], 5);
assert_eq!(state.display_status("a"), "not queued");
assert_eq!(state.display_status("b"), "not queued");
assert_eq!(state.display_status("c"), "not queued");
assert!(state.is_pending("a"));
assert!(state.is_pending("b"));
assert!(state.is_pending("c"));
state.apply_command(ReducerCommand::AddToQueue("a".to_string()));
state.apply_command(ReducerCommand::AddToQueue("b".to_string()));
assert_eq!(state.display_status("a"), "queued");
assert_eq!(state.display_status("b"), "queued");
assert!(state.is_pending("a"));
assert!(state.is_pending("b"));
state.set_current_change(Some("a".to_string()));
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "a".to_string(),
command: "cmd".to_string(),
});
state.increment_apply_count("a");
assert_eq!(state.display_status("a"), "applying");
assert_eq!(state.current_change_id(), Some(&"a".to_string()));
assert_eq!(state.apply_count("a"), 1);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("a".to_string()));
state.mark_archived("a");
assert_eq!(state.display_status("a"), "resolving");
assert!(state.is_archived("a"));
assert!(!state.is_pending("a"));
state.apply_command(ReducerCommand::DequeueChange("b".to_string()));
assert_eq!(state.display_status("b"), "not queued");
assert_eq!(state.display_status("c"), "not queued");
assert!(state.is_pending("c")); assert!(!state.is_archived("c"));
}
#[test]
fn test_parallel_mode_change_archived_transitions_to_resolving_without_blocker() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ApplyCompleted {
change_id: "c".to_string(),
revision: "rev1".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ArchiveStarted {
change_id: "c".to_string(),
command: "archive".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(
state.display_status("c"),
"resolving",
"Parallel: no-blocker ChangeArchived must transition to active merge handling, not merge wait"
);
assert!(
!state.is_terminal_change("c"),
"Parallel: change must not be terminal after archive"
);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "merge-rev".to_string(),
});
assert_eq!(
state.display_status("c"),
"merged",
"Parallel: MergeCompleted must transition to merged terminal"
);
assert!(state.is_terminal_change("c"));
}
#[test]
fn test_parallel_mode_archive_success_clears_prior_acceptance_error() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceFailed {
change_id: "alpha".to_string(),
error: "transient acceptance failure".to_string(),
});
assert_eq!(state.display_status("alpha"), "error");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("alpha".to_string()));
assert_eq!(
state.display_status("alpha"),
"resolving",
"same-change archive success must supersede a recoverable acceptance error and enter active no-blocker merge handling"
);
assert!(!state.is_terminal_change("alpha"));
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "merge-rev".to_string(),
});
assert_eq!(state.display_status("alpha"), "merged");
assert!(state.is_terminal_change("alpha"));
}
#[test]
fn test_parallel_mode_acceptance_error_archive_then_merge_finishes_merged() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["add-skill-secret-ingestion".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceFailed {
change_id: "add-skill-secret-ingestion".to_string(),
error: "acceptance failed before fix".to_string(),
});
assert_eq!(state.display_status("add-skill-secret-ingestion"), "error");
state.apply_execution_event(&ExecutionEvent::ChangeArchived(
"add-skill-secret-ingestion".to_string(),
));
assert_ne!(state.display_status("add-skill-secret-ingestion"), "error");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "add-skill-secret-ingestion".to_string(),
revision: "merge-rev".to_string(),
});
assert_eq!(state.display_status("add-skill-secret-ingestion"), "merged");
assert!(state.is_terminal_change("add-skill-secret-ingestion"));
}
#[test]
fn test_merge_and_resolve_success_clear_prior_processing_error_but_not_rejected() {
use crate::events::ExecutionEvent;
let mut merge_state = OrchestratorState::new(vec!["alpha".to_string()], 0);
merge_state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "alpha".to_string(),
error: "recoverable process failure".to_string(),
});
assert_eq!(merge_state.display_status("alpha"), "error");
merge_state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "merge-rev".to_string(),
});
assert_eq!(merge_state.display_status("alpha"), "merged");
let mut resolve_state = OrchestratorState::new(vec!["alpha".to_string()], 0);
resolve_state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "alpha".to_string(),
error: "recoverable process failure".to_string(),
});
assert_eq!(resolve_state.display_status("alpha"), "error");
resolve_state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "alpha".to_string(),
worktree_change_ids: None,
});
assert_eq!(resolve_state.display_status("alpha"), "merged");
let mut rejected_state = OrchestratorState::new(vec!["alpha".to_string()], 0);
rejected_state.apply_execution_event(&ExecutionEvent::ChangeRejected {
change_id: "alpha".to_string(),
reason: "final rejection".to_string(),
});
rejected_state.apply_execution_event(&ExecutionEvent::ChangeArchived("alpha".to_string()));
rejected_state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "stale-merge-rev".to_string(),
});
rejected_state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "alpha".to_string(),
worktree_change_ids: None,
});
assert_eq!(rejected_state.display_status("alpha"), "rejected");
}
#[test]
fn test_parallel_mode_change_archived_uses_resolve_pending_when_other_change_is_resolving() {
use crate::events::ExecutionEvent;
let mut state =
OrchestratorState::new(vec!["resolving".to_string(), "archived".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "resolving".to_string(),
command: "resolve resolving".to_string(),
});
assert_eq!(state.display_status("resolving"), "resolving");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("archived".to_string()));
assert_eq!(
state.display_status("archived"),
"resolve pending",
"Parallel: ChangeArchived must transition to resolve pending while another change is resolving"
);
assert!(!state.is_terminal_change("archived"));
}
#[test]
fn test_parallel_mode_change_archived_uses_resolve_pending_when_other_change_is_rejecting() {
use crate::events::ExecutionEvent;
use crate::vcs::WorkspaceStatus;
let mut state =
OrchestratorState::new(vec!["rejecting".to_string(), "archived".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: "rejecting".to_string(),
workspace_name: "ws-rejecting".to_string(),
status: WorkspaceStatus::Rejecting,
});
assert_eq!(state.display_status("rejecting"), "rejecting");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("archived".to_string()));
assert_eq!(
state.display_status("archived"),
"resolve pending",
"Parallel: ChangeArchived must transition to resolve pending while another change is rejecting"
);
assert!(!state.is_terminal_change("archived"));
}
#[test]
fn test_parallel_mode_change_archived_enters_resolving_when_other_change_is_accepting() {
use crate::events::ExecutionEvent;
let mut state =
OrchestratorState::new(vec!["accepting".to_string(), "archived".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "accepting".to_string(),
command: "accept accepting".to_string(),
});
assert_eq!(state.display_status("accepting"), "accepting");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("archived".to_string()));
assert_eq!(
state.display_status("archived"),
"resolving",
"Parallel: accepting activity is not a merge/resolve lane blocker, so no-blocker archive handling must enter resolving instead of merge wait or resolve pending"
);
}
#[test]
fn test_parallel_mode_full_lifecycle() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("a".to_string()));
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "a".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("a"), "applying");
state.apply_execution_event(&ExecutionEvent::ChangeArchived("a".to_string()));
assert_eq!(
state.display_status("a"),
"resolving",
"Parallel: no-blocker archive completion must enter active resolving before merged"
);
assert!(!state.is_terminal_change("a"));
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "a".to_string(),
revision: "rev-a".to_string(),
});
assert_eq!(state.display_status("a"), "merged");
assert!(state.is_terminal_change("a"));
assert_eq!(state.display_status("b"), "not queued");
assert!(!state.is_terminal_change("b"));
}
#[test]
fn test_parallel_mode_merge_deferred_then_completed() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(
state.display_status("c"),
"resolving",
"Parallel: ChangeArchived alone must not imply merge wait before manual blocker evidence exists"
);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
assert!(state.is_terminal_change("c"));
}
#[test]
fn test_parallel_mode_late_events_do_not_regress_merged() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "late".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "cmd".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_auto_resumable_merge_deferred_sets_resolve_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "b".to_string(),
reason: "Merge in progress (MERGE_HEAD exists)".to_string(),
auto_resumable: true,
});
assert_eq!(
state.display_status("b"),
"resolve pending",
"auto-resumable deferred change must enter ResolveWait, not MergeWait"
);
}
#[test]
fn test_auto_resumable_deferred_survives_workspace_refresh() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "b".to_string(),
reason: "Working tree has uncommitted changes".to_string(),
auto_resumable: true,
});
assert_eq!(state.display_status("b"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: Default::default(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: Default::default(),
worktree_change_ids: Default::default(),
worktree_paths: Default::default(),
worktree_not_ahead_ids: Default::default(),
merge_wait_ids: ["b".to_string()].into_iter().collect(),
});
assert_eq!(
state.display_status("b"),
"resolve pending",
"workspace refresh must not regress auto-resumable deferred change to merge wait"
);
}
#[test]
fn test_auto_resumable_deferred_then_merge_completed() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "b".to_string(),
reason: "Merge in progress (MERGE_HEAD exists)".to_string(),
auto_resumable: true,
});
assert_eq!(state.display_status("b"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "b".to_string(),
revision: "abc123".to_string(),
});
assert_eq!(state.display_status("b"), "merged");
assert!(state.is_terminal_change("b"));
}
#[test]
fn test_manual_resolve_blocked_by_merge_in_progress_becomes_resolve_pending() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
assert_eq!(state.display_status("c"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "Merge in progress (MERGE_HEAD exists)".to_string(),
auto_resumable: true,
});
assert_eq!(
state.display_status("c"),
"resolve pending",
"auto-resumable dirty base must stay as resolve pending"
);
}
#[test]
fn test_manual_resolve_blocked_by_uncommitted_changes_stays_merge_wait() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "c".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("c"), "merge wait");
state.apply_command(ReducerCommand::ResolveMerge("c".to_string()));
assert_eq!(state.display_status("c"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::ResolveFailed {
change_id: "c".to_string(),
error: "Base is dirty: Working tree has uncommitted changes".to_string(),
});
assert_eq!(
state.display_status("c"),
"merge wait",
"uncommitted-changes dirty base must revert to merge wait"
);
}
#[test]
fn test_auto_resumable_deferred_resolve_auto_retries_after_preceding_completes() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["a".to_string(), "b".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "a".to_string(),
command: "resolve a".to_string(),
});
assert_eq!(state.display_status("a"), "resolving");
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "b".to_string(),
reason: "Merge in progress (MERGE_HEAD exists)".to_string(),
auto_resumable: true,
});
assert_eq!(state.display_status("b"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "a".to_string(),
worktree_change_ids: None,
});
assert_eq!(state.display_status("a"), "merged");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "b".to_string(),
revision: "rev-b".to_string(),
});
assert_eq!(state.display_status("b"), "merged");
assert!(state.is_terminal_change("b"));
}
#[test]
fn test_explicit_retry_retries_error_terminal() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "apply failed".to_string(),
});
assert_eq!(state.display_status("c"), "error");
assert!(state.is_terminal_change("c"));
let ordinary_outcome = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(ordinary_outcome, ReduceOutcome::NoOp));
assert_eq!(state.display_status("c"), "error");
let outcome = state.apply_command(ReducerCommand::RetryError("c".to_string()));
assert!(
matches!(outcome, ReduceOutcome::Changed(_)),
"RetryError on error change must be Changed, not NoOp"
);
assert_eq!(
state.display_status("c"),
"queued",
"after retry, change must display as queued"
);
assert!(
!state.is_terminal_change("c"),
"error terminal must be cleared by AddToQueue"
);
}
#[test]
fn test_dequeue_change_resets_to_not_queued_idle() {
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
let outcome = state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
assert!(matches!(outcome, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "not queued");
assert!(!state.is_terminal_change("c"));
let outcome2 = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(matches!(outcome2, ReduceOutcome::Changed(_)));
assert_eq!(state.display_status("c"), "queued");
}
#[test]
fn test_add_to_queue_noop_on_merged() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived("c".to_string()));
assert_eq!(
state.display_status("c"),
"resolving",
"archive enters post-archive handling instead of terminating"
);
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "c".to_string(),
revision: "rev".to_string(),
});
assert_eq!(state.display_status("c"), "merged");
let outcome = state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert!(
matches!(outcome, ReduceOutcome::NoOp),
"AddToQueue on a merged change must be NoOp"
);
assert_eq!(state.display_status("c"), "merged");
}
#[test]
fn test_dependency_resolved_restores_queued_after_block() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.display_status("c"), "queued");
state.apply_execution_event(&ExecutionEvent::DependencyBlocked {
change_id: "c".to_string(),
dependency_ids: vec!["dep".to_string()],
});
assert_eq!(state.display_status("c"), "blocked");
state.apply_execution_event(&ExecutionEvent::DependencyResolved {
change_id: "c".to_string(),
});
assert_eq!(
state.display_status("c"),
"queued",
"DependencyResolved must restore queued (not not-queued)"
);
}
#[test]
fn test_dependency_blocked_and_resolved_preserve_queue_intent_until_user_dequeue() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
state.apply_execution_event(&ExecutionEvent::DependencyBlocked {
change_id: "c".to_string(),
dependency_ids: vec!["dep-a".to_string()],
});
assert_eq!(state.display_status("c"), "blocked");
state.apply_execution_event(&ExecutionEvent::DependencyResolved {
change_id: "c".to_string(),
});
assert_eq!(state.display_status("c"), "queued");
state.apply_command(ReducerCommand::RemoveFromQueue("c".to_string()));
assert_eq!(
state.display_status("c"),
"not queued",
"queue intent should only clear on explicit dequeue command"
);
}
#[test]
fn test_changes_refreshed_preserves_queue_intent() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["c".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("c".to_string()));
assert_eq!(state.display_status("c"), "queued");
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: HashSet::new(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
});
assert_eq!(
state.display_status("c"),
"queued",
"ChangesRefreshed must not overwrite queue_intent = Queued"
);
}
#[test]
fn test_fast_forward_merged_survives_archived_observation() {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec!["ff".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "ff".to_string(),
reason: "base dirty".to_string(),
auto_resumable: false,
});
assert_eq!(state.display_status("ff"), "merge wait");
state.apply_command(ReducerCommand::ResolveMerge("ff".to_string()));
assert_eq!(state.display_status("ff"), "resolve pending");
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "ff".to_string(),
command: "resolve-cmd".to_string(),
});
assert_eq!(state.display_status("ff"), "resolving");
state.apply_execution_event(&ExecutionEvent::ResolveCompleted {
change_id: "ff".to_string(),
worktree_change_ids: None,
});
assert_eq!(state.display_status("ff"), "merged");
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![],
committed_change_ids: Default::default(),
rejected_changes: Vec::new(),
uncommitted_file_change_ids: Default::default(),
worktree_change_ids: Default::default(),
worktree_paths: Default::default(),
worktree_not_ahead_ids: Default::default(),
merge_wait_ids: ["ff".to_string()].into_iter().collect(),
});
assert_eq!(
state.display_status("ff"),
"merged",
"Terminal Merged from fast-forward resolve must not regress to merge wait"
);
}
fn per_change_upstream_state(change_id: &str) -> OrchestratorState {
use crate::events::ExecutionEvent;
let mut state = OrchestratorState::new(vec![change_id.to_string()], 0);
state.apply_execution_event(&ExecutionEvent::ChangeArchived(change_id.to_string()));
state
}
fn per_change_upstream_push_started(change_id: &str) -> crate::events::ExecutionEvent {
crate::events::ExecutionEvent::PushStarted {
change_id: change_id.to_string(),
remote: "origin".to_string(),
branch: "main".to_string(),
}
}
fn per_change_upstream_push_completed(change_id: &str) -> crate::events::ExecutionEvent {
crate::events::ExecutionEvent::PushCompleted {
change_id: change_id.to_string(),
remote: "origin".to_string(),
branch: "main".to_string(),
}
}
fn per_change_upstream_push_failed(change_id: &str) -> crate::events::ExecutionEvent {
crate::events::ExecutionEvent::PushFailed {
change_id: change_id.to_string(),
remote: "origin".to_string(),
branch: "main".to_string(),
error: "upstream publication incomplete: verification failed".to_string(),
}
}
#[test]
fn per_change_upstream_disabled_cumulative_merge_stays_merged() {
use crate::events::ExecutionEvent;
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "alpha".to_string(),
revision: "merge-rev".to_string(),
});
assert_eq!(state.display_status("alpha"), "merged");
assert!(state.is_terminal_change("alpha"));
}
#[test]
fn per_change_upstream_local_integration_is_not_terminal_merged() {
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
assert_ne!(state.display_status("alpha"), "merged");
assert!(
!state.is_terminal_change("alpha"),
"publication progress must remain non-terminal"
);
assert!(
!state.queued_change_ids().contains(&"alpha".to_string()),
"a publishing change must not be ordinary queued apply work"
);
}
fn per_change_upstream_apply_fresh_integration(state: &mut OrchestratorState, change_id: &str) {
use crate::events::ExecutionEvent;
state.apply_execution_event(&ExecutionEvent::MergeStarted {
revisions: vec!["ws-alpha".to_string()],
});
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: change_id.to_string(),
command: "merge archived change into base branch (1 revision(s))".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ConflictResolutionStarted);
state.apply_execution_event(&ExecutionEvent::ConflictResolutionCompleted);
}
#[test]
fn per_change_upstream_fresh_integration_never_displays_merged() {
let mut state = per_change_upstream_state("alpha");
per_change_upstream_apply_fresh_integration(&mut state, "alpha");
assert_ne!(
state.display_status("alpha"),
"merged",
"local integration must not finalize an opted-in change"
);
assert!(
!state.is_terminal_change("alpha"),
"publication is still owed, so nothing is terminal yet"
);
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
assert_ne!(state.display_status("alpha"), "merged");
assert!(!state.is_terminal_change("alpha"));
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert_eq!(
state.display_status("alpha"),
"pushed",
"remote confirmation is the opted-in terminal state"
);
assert!(state.is_terminal_change("alpha"));
}
#[test]
fn per_change_upstream_fresh_integration_failure_stays_recoverable() {
let mut state = per_change_upstream_state("alpha");
per_change_upstream_apply_fresh_integration(&mut state, "alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_failed("alpha"));
assert_eq!(
state.display_status("alpha"),
"error",
"a failed publication after a fresh integration must surface as recoverable error"
);
assert!(
state.is_terminal_error_change("alpha"),
"explicit F5 retry needs a recoverable error to act on"
);
assert!(matches!(
state.apply_command(ReducerCommand::RetryError("alpha".to_string())),
ReduceOutcome::Changed(_)
));
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert_eq!(state.display_status("alpha"), "pushed");
}
#[test]
fn per_change_upstream_remote_confirmation_becomes_pushed_terminal() {
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert_eq!(state.display_status("alpha"), "pushed");
assert!(state.is_terminal_change("alpha"));
assert_ne!(state.display_status("alpha"), "merged");
}
#[test]
fn per_change_upstream_publication_failure_is_recoverable_not_merged() {
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_failed("alpha"));
assert_eq!(
state.display_status("alpha"),
"error",
"a failed publication projects into the existing recoverable error flow"
);
assert!(
state.is_terminal_error_change("alpha"),
"the recoverable error must gate ordinary apply dispatch"
);
let outcome = state.apply_command(ReducerCommand::RetryError("alpha".to_string()));
assert!(matches!(outcome, ReduceOutcome::Changed(_)));
}
#[test]
fn per_change_upstream_late_confirmation_supersedes_publication_failure() {
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_failed("alpha"));
assert_eq!(state.display_status("alpha"), "error");
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert_eq!(state.display_status("alpha"), "pushed");
assert!(
!state.queued_change_ids().contains(&"alpha".to_string()),
"confirmed publication must not leave ordinary apply dispatch behind"
);
}
#[test]
fn per_change_upstream_pushed_terminal_is_not_retryable() {
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&per_change_upstream_push_started("alpha"));
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert!(matches!(
state.apply_command(ReducerCommand::RetryError("alpha".to_string())),
ReduceOutcome::NoOp
));
assert_eq!(state.display_status("alpha"), "pushed");
assert!(matches!(
state.apply_command(ReducerCommand::ResolveMerge("alpha".to_string())),
ReduceOutcome::NoOp
));
assert_eq!(state.display_status("alpha"), "pushed");
}
#[test]
fn per_change_upstream_confirmation_clears_base_lane_retry_intent() {
use crate::events::ExecutionEvent;
let mut state = per_change_upstream_state("alpha");
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "base lane busy".to_string(),
auto_resumable: false,
});
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: "alpha".to_string(),
reason: "waiting for publication of beta".to_string(),
auto_resumable: true,
});
assert!(state
.resolve_wait_change_ids()
.contains(&"alpha".to_string()));
state.apply_execution_event(&per_change_upstream_push_completed("alpha"));
assert_eq!(state.display_status("alpha"), "pushed");
assert!(
!state
.resolve_wait_change_ids()
.contains(&"alpha".to_string()),
"confirmed publication must clear base-lane retry intent"
);
}
#[test]
fn changes_refreshed_registers_without_creating_queue_or_lane_eligibility() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["fresh".to_string()], 1);
state.apply_command(ReducerCommand::AddToQueue("fresh".to_string()));
let change = |id: &str| crate::openspec::Change {
id: id.to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: crate::openspec::ProposalMetadata::default(),
};
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![change("fresh"), change("stale")],
rejected_changes: Vec::new(),
committed_change_ids: HashSet::from(["fresh".to_string(), "stale".to_string()]),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::from(["stale".to_string()]),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
});
assert!(
state.is_in_snapshot("stale"),
"refresh may register a newly observed change"
);
assert_eq!(state.display_status("stale"), "not queued");
assert_eq!(
state
.change_runtime("stale")
.expect("refresh registers runtime state")
.queue_intent,
QueueIntent::NotQueued
);
assert_eq!(state.queued_change_ids(), vec!["fresh".to_string()]);
assert!(state.merge_wait_change_ids().is_empty());
assert!(state.resolve_wait_change_ids().is_empty());
assert!(state.reject_wait_change_ids().is_empty());
assert!(state.active_change_ids().is_empty());
assert!(
!state.is_ordinary_queue_eligible("stale"),
"a registered but unqueued change must not be dispatchable"
);
}
#[test]
fn removal_and_dequeue_revoke_ordinary_eligibility_until_explicit_requeue() {
use crate::events::ExecutionEvent;
use std::collections::{HashMap, HashSet};
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 1);
state.apply_command(ReducerCommand::AddToQueue("alpha".to_string()));
assert!(state.is_ordinary_queue_eligible("alpha"));
state.apply_command(ReducerCommand::RemoveFromQueue("alpha".to_string()));
assert!(!state.is_ordinary_queue_eligible("alpha"));
assert!(state.queued_change_ids().is_empty());
state.apply_execution_event(&ExecutionEvent::ChangesRefreshed {
changes: vec![crate::openspec::Change {
id: "alpha".to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: "now".to_string(),
dependencies: Vec::new(),
metadata: crate::openspec::ProposalMetadata::default(),
}],
rejected_changes: Vec::new(),
committed_change_ids: HashSet::from(["alpha".to_string()]),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::from(["alpha".to_string()]),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::new(),
});
assert!(!state.is_ordinary_queue_eligible("alpha"));
state.apply_command(ReducerCommand::AddToQueue("alpha".to_string()));
state.apply_command(ReducerCommand::DequeueChange("alpha".to_string()));
assert!(
!state.is_ordinary_queue_eligible("alpha"),
"stop-and-dequeue revokes ordinary eligibility"
);
state.apply_command(ReducerCommand::AddToQueue("alpha".to_string()));
assert!(
state.is_ordinary_queue_eligible("alpha"),
"explicit requeue restores eligibility"
);
assert_eq!(state.queued_change_ids(), vec!["alpha".to_string()]);
}
fn queued_state(change_id: &str) -> OrchestratorState {
let mut state = OrchestratorState::new(vec![change_id.to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
assert_eq!(state.display_status(change_id), "queued");
state
}
fn prepare(state: &mut OrchestratorState, change_id: &str) {
state.apply_execution_event(
&crate::events::ExecutionEvent::WorkspacePreparationStarted {
change_id: change_id.to_string(),
},
);
}
fn end_preparation(state: &mut OrchestratorState, change_id: &str) {
state.apply_execution_event(&crate::events::ExecutionEvent::WorkspacePreparationEnded {
change_id: change_id.to_string(),
});
}
#[test]
fn preparing_replaces_queued_and_then_yields_to_the_repository_derived_phase() {
use crate::events::ExecutionEvent;
let mut state = queued_state("c");
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "preparing");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "apply".to_string(),
});
assert_eq!(state.display_status("c"), "applying");
let mut state = queued_state("c");
prepare(&mut state, "c");
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "c".to_string(),
command: "accept".to_string(),
});
assert_eq!(state.display_status("c"), "accepting");
assert_eq!(state.apply_count("c"), 0);
let mut state = queued_state("c");
prepare(&mut state, "c");
state.apply_execution_event(&ExecutionEvent::ArchiveStarted {
change_id: "c".to_string(),
command: "archive".to_string(),
});
assert_eq!(state.display_status("c"), "archiving");
}
#[test]
fn preparing_failure_becomes_error_with_the_preparation_diagnostic() {
use crate::events::ExecutionEvent;
let mut state = queued_state("c");
prepare(&mut state, "c");
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "Worktree setup failed after 1.0s: .wt/setup exited with code 3".to_string(),
});
assert_eq!(state.display_status("c"), "error");
let rt = state.change_runtime("c").expect("runtime entry");
assert!(matches!(rt.activity, ActivityState::Idle));
assert!(rt
.error_message()
.is_some_and(|message| message.contains(".wt/setup")));
}
#[test]
fn preparing_is_cleared_by_a_pre_operation_exit() {
let mut state = queued_state("c");
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "preparing");
end_preparation(&mut state, "c");
assert_eq!(
state.display_status("c"),
"queued",
"preparation must not outlive the dispatch that announced it"
);
}
#[test]
fn preparing_leaves_through_a_stop_before_an_operation_agent_starts() {
use crate::events::ExecutionEvent;
let mut state = queued_state("c");
prepare(&mut state, "c");
state.apply_execution_event(&ExecutionEvent::ChangeDequeued {
change_id: "c".to_string(),
});
assert_eq!(state.display_status("c"), "not queued");
let rt = state.change_runtime("c").expect("runtime entry");
assert!(matches!(rt.activity, ActivityState::Idle));
end_preparation(&mut state, "c");
assert_eq!(state.display_status("c"), "not queued");
}
#[test]
fn preparing_never_overwrites_a_running_operation_or_a_terminal_change() {
use crate::events::ExecutionEvent;
let mut state = queued_state("c");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "apply".to_string(),
});
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "applying");
let mut state = queued_state("c");
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "c".to_string(),
error: "boom".to_string(),
});
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "error");
let mut state = queued_state("c");
state.apply_command(ReducerCommand::DequeueChange("c".to_string()));
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "not queued");
}
#[test]
fn preparing_clearing_is_narrow_and_idempotent() {
use crate::events::ExecutionEvent;
let mut state = queued_state("c");
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: "c".to_string(),
command: "apply".to_string(),
});
end_preparation(&mut state, "c");
assert_eq!(state.display_status("c"), "applying");
let mut state = queued_state("c");
prepare(&mut state, "c");
end_preparation(&mut state, "c");
end_preparation(&mut state, "c");
assert_eq!(state.display_status("c"), "queued");
}
#[test]
fn preparing_counts_as_active_execution_but_not_as_a_running_agent() {
let mut state = queued_state("c");
prepare(&mut state, "c");
assert!(
state.is_active_change("c"),
"an admitted change mutating its worktree is active execution"
);
assert!(
!state.is_agent_execution_active(),
"preparation must never justify a force-stopped-process claim"
);
assert!(
crate::orchestration::operator_command::is_active_status(state.display_status("c")),
"every operator surface must treat preparing as active"
);
}
#[test]
fn preparing_is_not_durable_routing_state() {
let mut state = queued_state("c");
prepare(&mut state, "c");
assert_eq!(state.display_status("c"), "preparing");
let restarted = OrchestratorState::new(vec!["c".to_string()], 0);
assert_eq!(restarted.display_status("c"), "not queued");
assert!(!restarted.is_active_change("c"));
}
fn enter_runtime_family(state: &mut OrchestratorState, change_id: &str, family: &str) {
use crate::events::ExecutionEvent;
use crate::vcs::WorkspaceStatus;
match family {
"queued" => {
state.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
}
"preparing" => {
state.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
state.apply_execution_event(&ExecutionEvent::WorkspacePreparationStarted {
change_id: change_id.to_string(),
});
}
"applying" => {
state.apply_execution_event(&ExecutionEvent::ApplyStarted {
change_id: change_id.to_string(),
command: "apply".to_string(),
});
}
"accepting" => {
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: change_id.to_string(),
command: "accept".to_string(),
});
}
"archiving" => {
state.apply_execution_event(&ExecutionEvent::ArchiveStarted {
change_id: change_id.to_string(),
command: "archive".to_string(),
});
}
"resolving" => {
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: change_id.to_string(),
command: "resolve".to_string(),
});
}
"rejecting" => {
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: change_id.to_string(),
workspace_name: format!("ws-{change_id}"),
status: WorkspaceStatus::Rejecting,
});
}
"merge wait" => {
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "base branch is dirty".to_string(),
auto_resumable: false,
});
}
"resolve pending" => {
state.apply_execution_event(&ExecutionEvent::MergeDeferred {
change_id: change_id.to_string(),
reason: "another merge is in progress".to_string(),
auto_resumable: true,
});
}
"reject pending" => {
state.apply_execution_event(&ExecutionEvent::ResolveStarted {
change_id: "lane-occupant".to_string(),
command: "resolve".to_string(),
});
state.apply_execution_event(&ExecutionEvent::WorkspaceStatusUpdated {
change_id: change_id.to_string(),
workspace_name: format!("ws-{change_id}"),
status: WorkspaceStatus::Rejecting,
});
}
"blocked (dependency)" => {
state.apply_execution_event(&ExecutionEvent::DependencyBlocked {
change_id: change_id.to_string(),
dependency_ids: vec!["dep".to_string()],
});
}
"blocked (external)" => {
state.apply_execution_event(&ExecutionEvent::AcceptanceGated {
change_id: change_id.to_string(),
blocker: crate::events::StalledBlocker::acceptance_external(
"credential",
"STAGING_API_KEY is unset",
),
});
}
"stalled" => {
state.mark_stalled(change_id.to_string());
state.apply_execution_event(&ExecutionEvent::ExecutionBlocked {
change_id: change_id.to_string(),
blocker: crate::events::StalledBlocker {
category: "no_progress".to_string(),
phase: "apply".to_string(),
gate: "apply".to_string(),
error_summary: "no semantic progress".to_string(),
evidence: vec!["tasks.md unchanged".to_string()],
unblock_condition: None,
prerequisite_owner: None,
next_action: "operator review".to_string(),
resumable: true,
worktree_preserved: true,
},
});
}
other => panic!("unknown runtime family: {other}"),
}
}
#[test]
fn global_stopped_reconciles_interrupted_runtime() {
use crate::events::ExecutionEvent;
let families = [
"queued",
"preparing",
"applying",
"accepting",
"archiving",
"resolving",
"rejecting",
"merge wait",
"resolve pending",
"reject pending",
"blocked (dependency)",
"blocked (external)",
"stalled",
];
for family in families {
let mut state =
OrchestratorState::new(vec!["target".to_string(), "lane-occupant".to_string()], 0);
enter_runtime_family(&mut state, "target", family);
assert_ne!(
state.display_status("target"),
"not queued",
"{family} setup did not produce interrupted runtime state"
);
state.apply_execution_event(&ExecutionEvent::Stopped);
assert_eq!(
state.display_status("target"),
"not queued",
"{family} was not reconciled by the run-boundary stop"
);
let rt = state.change_runtime.get("target").expect("runtime entry");
assert_eq!(rt.activity, ActivityState::Idle, "{family} stayed active");
assert_eq!(rt.wait_state, WaitState::None, "{family} stayed waiting");
assert_eq!(
rt.queue_intent,
QueueIntent::NotQueued,
"{family} kept queue intent"
);
assert_eq!(
rt.terminal,
TerminalState::None,
"{family} was given a terminal outcome by a process stop"
);
assert_eq!(rt.blocked_metadata, BlockedMetadata::default());
assert!(rt.commit_phase_attempt.is_none());
assert!(rt.dequeued, "{family} did not get the reactivation guard");
assert!(
state.resolve_wait_change_ids().is_empty(),
"{family} left scheduler resolve membership behind"
);
assert!(
state.reject_wait_change_ids().is_empty(),
"{family} left scheduler reject membership behind"
);
assert!(
!state.stalled_change_ids().contains("target"),
"{family} left stall membership behind"
);
assert!(
!state.is_ordinary_queue_eligible("target"),
"{family} remained dispatchable after the run ended"
);
}
let mut state = OrchestratorState::new(
vec![
"err".to_string(),
"merged".to_string(),
"pushed".to_string(),
"rejected".to_string(),
"fresh".to_string(),
],
0,
);
state.apply_execution_event(&ExecutionEvent::ProcessingError {
id: "err".to_string(),
error: "apply exited 1".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ChangeArchived("merged".to_string()));
state.apply_execution_event(&ExecutionEvent::MergeCompleted {
change_id: "merged".to_string(),
revision: "rev".to_string(),
});
state.apply_execution_event(&ExecutionEvent::PushCompleted {
change_id: "pushed".to_string(),
remote: "origin".to_string(),
branch: "main".to_string(),
});
state.apply_execution_event(&ExecutionEvent::ChangeRejected {
change_id: "rejected".to_string(),
reason: "acceptance confirmed rejection".to_string(),
});
state.apply_execution_event(&ExecutionEvent::Stopped);
assert_eq!(state.display_status("err"), "error");
assert_eq!(state.display_status("merged"), "merged");
assert_eq!(state.display_status("pushed"), "pushed");
assert_eq!(state.display_status("rejected"), "rejected");
assert_eq!(state.display_status("fresh"), "not queued");
assert!(
!state
.change_runtime
.get("fresh")
.expect("fresh runtime entry")
.dequeued,
"an unrelated idle row must not be claimed by the stopped run"
);
let mut state = OrchestratorState::new(vec!["alpha".to_string()], 0);
state.apply_execution_event(&ExecutionEvent::AcceptanceStarted {
change_id: "alpha".to_string(),
command: "accept".to_string(),
});
assert_eq!(state.display_status("alpha"), "accepting");
state.apply_execution_event(&ExecutionEvent::Stopped);
let after_first_stop = format!("{:?}", state.change_runtime.get("alpha"));
state.apply_execution_event(&ExecutionEvent::Stopped);
assert_eq!(
format!("{:?}", state.change_runtime.get("alpha")),
after_first_stop,
"a duplicate stop must not change reconciled state"
);
assert_eq!(state.display_status("alpha"), "not queued");
for late in [
ExecutionEvent::AcceptanceStarted {
change_id: "alpha".to_string(),
command: "accept".to_string(),
},
ExecutionEvent::ArchiveStarted {
change_id: "alpha".to_string(),
command: "archive".to_string(),
},
ExecutionEvent::ResolveStarted {
change_id: "alpha".to_string(),
command: "resolve".to_string(),
},
ExecutionEvent::WorkspacePreparationStarted {
change_id: "alpha".to_string(),
},
] {
state.apply_execution_event(&late);
assert_eq!(
state.display_status("alpha"),
"not queued",
"a late event from the stopped run reactivated the row"
);
}
state.apply_execution_event(&stop_refresh_event("alpha"));
assert_eq!(
state.display_status("alpha"),
"not queued",
"a same-process workspace observation resurrected stopped work"
);
state.apply_command(ReducerCommand::AddToQueue("alpha".to_string()));
assert_eq!(state.display_status("alpha"), "queued");
assert!(
!state
.change_runtime
.get("alpha")
.expect("runtime entry")
.dequeued,
"an explicit requeue must release the reactivation guard"
);
assert!(state.is_ordinary_queue_eligible("alpha"));
let mut restarted = OrchestratorState::new(vec!["alpha".to_string()], 0);
restarted.apply_execution_event(&stop_refresh_event("alpha"));
assert_eq!(
restarted.display_status("alpha"),
"merge wait",
"a restarted process must re-derive routing from workspace evidence alone"
);
}
fn stop_refresh_event(change_id: &str) -> crate::events::ExecutionEvent {
crate::events::ExecutionEvent::ChangesRefreshed {
changes: Vec::new(),
rejected_changes: Vec::new(),
committed_change_ids: HashSet::new(),
uncommitted_file_change_ids: HashSet::new(),
worktree_change_ids: HashSet::new(),
worktree_paths: HashMap::new(),
worktree_not_ahead_ids: HashSet::new(),
merge_wait_ids: HashSet::from([change_id.to_string()]),
}
}
#[test]
fn dependency_blocker_projection_initial() {
let mut state = OrchestratorState::new(vec!["beta".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("beta".to_string()));
assert!(
state.reconcile_dependency_blocker("beta", &["alpha".to_string()]),
"the first classification must publish the wait"
);
let rt = state.change_runtime("beta").expect("runtime entry");
assert_eq!(state.display_status("beta"), "blocked");
assert_eq!(
rt.queue_intent,
QueueIntent::Queued,
"a dependency wait excludes dispatch and never revokes admitted work"
);
assert_eq!(rt.wait_state, WaitState::DependencyBlocked);
assert!(!rt.is_active(), "no execution episode has started");
let view = state.blocker_view("beta").expect("structured blocker");
assert_eq!(view.status, "blocked");
assert_eq!(view.kind, BlockerKind::Dependency);
assert_eq!(view.dependencies, vec!["alpha".to_string()]);
assert_eq!(view.category.as_deref(), Some("dependency_blocked"));
assert!(view
.detail
.as_deref()
.is_some_and(|detail| detail.contains("alpha")));
assert_eq!(
crate::orchestration::execution_facts::project_execution_state(
rt,
crate::orchestration::execution_facts::ExecutionPhase::None,
false,
),
crate::orchestration::execution_facts::ChangeExecutionState::Queued
);
}
#[test]
fn dependency_blocker_projection_rebuild() {
let mut state = OrchestratorState::new(vec!["beta".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("beta".to_string()));
state.reconcile_dependency_blocker("beta", &["alpha".to_string()]);
assert!(
!state.reconcile_dependency_blocker("beta", &["alpha".to_string()]),
"an unchanged fingerprint must settle as a no-op, not a rewrite"
);
assert_eq!(state.display_status("beta"), "blocked");
assert!(
state.reconcile_dependency_blocker("beta", &["alpha".to_string(), "gamma".to_string()])
);
assert_eq!(
state.blocker_view("beta").expect("blocker").dependencies,
vec!["alpha".to_string(), "gamma".to_string()]
);
let mut rebuilt = OrchestratorState::new(vec!["beta".to_string()], 0);
rebuilt.apply_command(ReducerCommand::AddToQueue("beta".to_string()));
assert_eq!(
rebuilt.display_status("beta"),
"queued",
"a fresh reducer starts from retained queue intent alone"
);
assert!(rebuilt.blocker_view("beta").is_none());
rebuilt.reconcile_dependency_blocker("beta", &["alpha".to_string()]);
assert_eq!(rebuilt.display_status("beta"), "blocked");
assert_eq!(
rebuilt.blocker_view("beta").expect("blocker").dependencies,
vec!["alpha".to_string()]
);
}
#[test]
fn dependency_blocker_projection_resolution() {
let mut state = OrchestratorState::new(vec!["beta".to_string(), "held".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("beta".to_string()));
state.apply_command(ReducerCommand::AddToQueue("held".to_string()));
state.reconcile_dependency_blocker("beta", &["alpha".to_string()]);
assert!(state.clear_dependency_blocker("beta"));
assert_eq!(
state.display_status("beta"),
"queued",
"the retained queue intent decides the display again"
);
assert!(state.blocker_view("beta").is_none());
assert!(
!state.clear_dependency_blocker("beta"),
"clearing an already-cleared wait changes nothing"
);
state.runtime_entry("held").transition_to_stalled(
"acceptance_finding",
"no progress",
"snapshot",
);
assert!(!state.clear_dependency_blocker("held"));
assert_eq!(state.display_status("held"), "stalled");
}
#[test]
fn dependency_blocker_projection_capacity_only() {
let mut state = OrchestratorState::new(vec!["ready".to_string()], 0);
state.apply_command(ReducerCommand::AddToQueue("ready".to_string()));
assert!(!state.clear_dependency_blocker("ready"));
assert_eq!(state.display_status("ready"), "queued");
assert!(
state.blocker_view("ready").is_none(),
"an occupied execution slot is not a blocker"
);
let rt = state.change_runtime("ready").expect("runtime entry");
assert_eq!(rt.wait_state, WaitState::None);
assert_eq!(rt.blocker_kind(), BlockerKind::None);
}
#[test]
fn dependency_blocker_projection_refuses_terminal_and_active_rows() {
let mut state = OrchestratorState::new(vec!["done".to_string(), "running".to_string()], 0);
state.runtime_entry("done").terminal = TerminalState::Merged;
state.runtime_entry("running").activity = ActivityState::Applying;
assert!(!state.reconcile_dependency_blocker("done", &["alpha".to_string()]));
assert!(!state.reconcile_dependency_blocker("running", &["alpha".to_string()]));
assert_eq!(state.display_status("done"), "merged");
assert_eq!(state.display_status("running"), "applying");
}
}