diff --git a/server/src/agents/pool/auto_assign/merge_failure_block_subscriber.rs b/server/src/agents/pool/auto_assign/merge_failure_block_subscriber.rs index 47534810..212787f8 100644 --- a/server/src/agents/pool/auto_assign/merge_failure_block_subscriber.rs +++ b/server/src/agents/pool/auto_assign/merge_failure_block_subscriber.rs @@ -1,18 +1,43 @@ -//! TransitionFired subscriber that auto-blocks stories after N consecutive MergeFailure transitions. +//! TransitionFired subscriber that owns the consecutive-MergeFailure budget: +//! auto-blocks stories at the threshold and auto-retries `GatesFailed` +//! failures below it. //! //! Listens on the pipeline transition broadcast channel and, for each story, //! counts how many times it has entered [`Stage::MergeFailure`] consecutively. -//! When the count reaches the configurable threshold (default 3), the story is -//! transitioned to [`Stage::Blocked`] with a reason that names the failure kind. +//! One counter drives two policies sharing the `merge_failure_block_threshold` +//! budget (default 3): //! -//! The counter for a story resets whenever a non-`MergeFailure` transition fires -//! for that story (e.g. after a successful merge or a `FixupRequested` demotion -//! back to coding). +//! - **Below the threshold**, a `GatesFailed` failure schedules a delayed +//! re-trigger of the deterministic server-side merge (story 1185) — gates +//! failures are dominated by transients (flaky tests, stale base) that a +//! plain re-run fixes. Other kinds still count toward the budget but are +//! not retried: `ConflictDetected` has its own mergemaster recovery path via +//! [`super::merge_failure_subscriber`]; `EmptyDiff`/`NoCommits`/`Other` +//! require human intervention. +//! - **At the threshold**, the story is transitioned to [`Stage::Blocked`] +//! with a reason naming the failure kind. +//! +//! The counter resets when the story leaves `MergeFailure` for a real reason +//! (successful merge, `FixupRequested`, `Block`), but **not** on +//! [`PipelineEvent::MergeRetryStarted`] — that is the `MergeFailure → Merge` +//! bounce a retry itself causes. Treating it as a reset made the budget +//! unreachable and let a deterministic gates failure retry forever (1185 +//! review finding 1); counting across the bounce is what makes the budget +//! real. +//! +//! Bug 1025: while a mergemaster is actively running on the story, its +//! iteration loop (squash → fail → fix → retry) generates multiple +//! MergeFailure transitions. Those are NOT consecutive give-ups — they are +//! recovery iterations in progress. We neither count nor schedule retries +//! while a mergemaster is in the pool for the story. use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::sync::Mutex; +use std::time::Duration; +use crate::io::watcher::WatcherEvent; use crate::pipeline_state::{MergeFailureKind, PipelineEvent, Stage, Status, StoryId}; use crate::slog; use crate::slog_warn; @@ -20,42 +45,102 @@ use crate::slog_warn; use super::super::super::PipelineStage; use super::super::AgentPool; use super::scan::is_story_assigned_for_stage; +use super::story_checks::{ + has_review_hold, has_unmet_dependencies, is_story_blocked, is_story_frozen, +}; -/// Reconcile: no-op for the merge-failure block subscriber. +/// Delay before an auto-retry re-triggers the server-side merge for a +/// `GatesFailed` failure. Gives transient conditions (a concurrently landing +/// master merge, an exhausted runner) a moment to clear; retrying instantly +/// would just replay the same failure. +const AUTO_RETRY_DELAY: Duration = Duration::from_secs(30); + +/// Per-story scheduling generation, shared between the subscriber loop and the +/// delayed retry tasks it spawns. /// -/// The block subscriber maintains an in-memory per-story consecutive-failure counter -/// that cannot be reconstructed from CRDT state alone (only the current stage is -/// stored, not the history of how many times each story failed). Eventual consistency -/// is guaranteed by the live subscriber reacting to each new `MergeFailure` event; -/// the periodic reconciler cannot add value here without risking spurious blocks. +/// Every scheduled retry captures the generation current at schedule time; the +/// timer only acts if that generation is still current when it fires. The +/// subscriber bumps the generation on every (re)schedule and clears the entry +/// on counter reset, so stale timers left over from an earlier failure cycle +/// become no-ops instead of firing unaccounted retries (1185 review finding 6). +type Generations = Arc>>; + +/// What the subscriber decided to do about one transition. Split out from the +/// event loop so the counter/budget policy is synchronous and unit-testable. +#[derive(Debug, PartialEq, Eq)] +enum Decision { + /// Nothing to do (not a MergeFailure, recovery in progress, retry bounce, + /// budget disabled, or a non-retryable kind below the threshold). + Nothing, + /// Schedule a delayed auto-retry: this is consecutive failure `attempt` of + /// a `budget`-sized budget, and the kind is `GatesFailed`. + ScheduleRetry { attempt: u32, budget: u32 }, + /// The budget is exhausted: block the story. + Block { count: u32 }, +} + +/// Reconcile: no-op for the periodic pass. +/// +/// The consecutive-failure counter is in-memory and cannot be reconstructed +/// from CRDT state (only the current stage is stored, not the failure +/// history). Restart catch-up for stories already parked in +/// `MergeFailure{GatesFailed}` is handled once, at subscriber startup, by +/// [`reconcile_stranded_gates_failed`] — running it from the periodic +/// reconciler instead would re-schedule retries for budget-exhausted stories +/// on every tick, reintroducing the unbounded-retry bug the startup-only scan +/// avoids. pub(crate) fn reconcile_merge_failure_block() {} -/// Spawn a background task that blocks stories after N consecutive `MergeFailure` transitions. -/// -/// Subscribes to the pipeline transition broadcast channel and tracks a per-story -/// consecutive-failure counter. When a story's count reaches the threshold configured -/// in `project.toml` (`merge_failure_block_threshold`, default 3), the story is -/// transitioned to `Stage::Blocked` with a reason that names the failure kind. -/// -/// The counter resets when the story leaves `MergeFailure` (e.g. on `FixupRequested`, -/// `ReQueuedForQa`, or a successful merge via `Unblock → Merge → Done`). -/// -/// Bug 1025: while a mergemaster is actively running on the story, its -/// iteration loop (squash → fail → fix → retry) generates multiple -/// MergeFailure transitions. Those are NOT consecutive give-ups — they are -/// recovery iterations in progress. We skip counter increments while a -/// mergemaster is in the pool for the story; the counter only increments on -/// transitions that happen with no recovery agent attached. +/// Spawn the background task that owns the consecutive-MergeFailure budget: +/// auto-retry for `GatesFailed` below the threshold, auto-block at it. pub(crate) fn spawn_merge_failure_block_subscriber(pool: Arc, project_root: PathBuf) { let mut rx = crate::pipeline_state::subscribe_transitions(); tokio::spawn(async move { let mut counters: HashMap = HashMap::new(); + let generations: Generations = Generations::default(); + + // One-shot restart catch-up: stories already sitting in GatesFailed + // when the process starts will never fire another transition on their + // own, so without this they'd silently lose auto-retry coverage + // (1185 review finding 3). + reconcile_stranded_gates_failed(&pool, &project_root, &mut counters, &generations).await; + loop { match rx.recv().await { Ok(fired) => { let recovery_running = is_mergemaster_running(&pool, &project_root, &fired.story_id.0).await; - on_transition(&project_root, &fired, &mut counters, recovery_running); + let threshold = load_threshold(&project_root).await; + match decide(&fired, &mut counters, recovery_running, threshold) { + Decision::Nothing => { + // A real departure from MergeFailure also + // invalidates any pending retry timer. + if fired.after.status() != Status::MergeFailure + && !matches!(fired.event, PipelineEvent::MergeRetryStarted) + { + invalidate_generation(&generations, &fired.story_id.0); + } + } + Decision::ScheduleRetry { attempt, budget } => { + schedule_auto_retry( + Arc::clone(&pool), + project_root.clone(), + fired.story_id.0.clone(), + attempt, + budget, + Arc::clone(&generations), + ); + } + Decision::Block { count } => { + let kind = counters + .get(&fired.story_id) + .map(|(_, k)| k.clone()) + .unwrap_or(MergeFailureKind::Other(String::new())); + apply_block(&fired.story_id, count, &kind); + counters.remove(&fired.story_id); + invalidate_generation(&generations, &fired.story_id.0); + } + } } Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => { slog_warn!( @@ -69,91 +154,294 @@ pub(crate) fn spawn_merge_failure_block_subscriber(pool: Arc, project }); } +/// Pure budget policy: given a fired transition and the per-story counter, +/// decide whether to do nothing, schedule a `GatesFailed` auto-retry, or +/// block the story. +/// +/// `threshold == 0` disables both policies (feature off), matching the +/// pre-1185 block-subscriber behaviour. +fn decide( + fired: &crate::pipeline_state::TransitionFired, + counters: &mut HashMap, + recovery_running: bool, + threshold: u32, +) -> Decision { + // Story 1086: gate on the typed `Status` projection — `Status::MergeFailure` + // is precisely the set of stages we count toward the budget. + if fired.after.status() != Status::MergeFailure { + // The MergeFailure → Merge bounce caused by a retry (auto or manual) + // is not a recovery: the budget must survive it, or a deterministic + // failure retries forever (1185 review finding 1). + if !matches!(fired.event, PipelineEvent::MergeRetryStarted) { + counters.remove(&fired.story_id); + } + return Decision::Nothing; + } + let Stage::MergeFailure { kind, .. } = &fired.after else { + counters.remove(&fired.story_id); + return Decision::Nothing; + }; + + if recovery_running { + slog!( + "[merge-block-sub] Story '{}' MergeFailure while mergemaster is running; \ + not counting toward budget (recovery in progress).", + fired.story_id.0 + ); + return Decision::Nothing; + } + + let entry = counters + .entry(fired.story_id.clone()) + .or_insert_with(|| (0, kind.clone())); + entry.0 += 1; + entry.1 = kind.clone(); + let count = entry.0; + + if threshold == 0 { + return Decision::Nothing; + } + if count >= threshold { + return Decision::Block { count }; + } + if matches!(kind, MergeFailureKind::GatesFailed(_)) { + return Decision::ScheduleRetry { + attempt: count, + budget: threshold, + }; + } + Decision::Nothing +} + +/// Transition `story_id` to `Blocked` with a reason naming the failure kind. +fn apply_block(story_id: &StoryId, count: u32, kind: &MergeFailureKind) { + let kind_str = failure_kind_label(kind); + let reason = + format!("Auto-blocked after {count} consecutive MergeFailure ({kind_str}) transitions."); + let story_id = story_id.0.as_str(); + slog!( + "[merge-block-sub] Story '{story_id}' reached {count} consecutive \ + MergeFailure ({kind_str}); blocking." + ); + if let Err(e) = + crate::pipeline_state::apply_transition(story_id, PipelineEvent::Block { reason }, None) + { + slog_warn!("[merge-block-sub] Failed to block '{story_id}': {e}"); + } +} + +/// Spawn the delayed auto-retry task for one `GatesFailed` failure. +/// +/// Bumps the story's scheduling generation so any earlier still-sleeping timer +/// for the story becomes stale and skips itself. +fn schedule_auto_retry( + pool: Arc, + project_root: PathBuf, + story_id: String, + attempt: u32, + budget: u32, + generations: Generations, +) { + let generation = next_generation(&generations, &story_id); + slog!( + "[merge-block-sub] Story '{story_id}' GatesFailed (attempt {attempt}/{budget}); \ + scheduling auto-retry in {AUTO_RETRY_DELAY:?}." + ); + tokio::spawn(async move { + tokio::time::sleep(AUTO_RETRY_DELAY).await; + perform_auto_retry( + &pool, + &project_root, + &story_id, + attempt, + budget, + &generations, + generation, + ) + .await; + }); +} + +/// Fire one auto-retry: re-check that acting is still correct, start the +/// server-side merge, and notify chat only when the merge actually started +/// (1185 review finding 7). +/// +/// Guards, in order: +/// - the scheduling generation is still current (no newer schedule/reset); +/// - the story is still parked in `MergeFailure{GatesFailed}`; +/// - the story is not frozen/blocked/on hold/dep-blocked — the same +/// eligibility gates `assign_merge_stage` applies (1185 review finding 4); +/// - no mergemaster is actively recovering the story (1185 review finding 2). +async fn perform_auto_retry( + pool: &Arc, + project_root: &Path, + story_id: &str, + attempt: u32, + budget: u32, + generations: &Generations, + generation: u64, +) { + if !is_generation_current(generations, story_id, generation) { + slog!( + "[merge-block-sub] Story '{story_id}' auto-retry ({attempt}/{budget}) is stale \ + (superseded or reset); skipping." + ); + return; + } + + let still_gates_failed = matches!( + crate::pipeline_state::read_typed(story_id), + Ok(Some(item)) if matches!( + item.stage, + Stage::MergeFailure { + kind: MergeFailureKind::GatesFailed(_), + .. + } + ) + ); + if !still_gates_failed { + slog!( + "[merge-block-sub] Story '{story_id}' left GatesFailed before auto-retry \ + ({attempt}/{budget}) fired; skipping." + ); + return; + } + + if has_review_hold(story_id) + || is_story_frozen(story_id) + || is_story_blocked(story_id) + || has_unmet_dependencies(story_id) + { + slog!( + "[merge-block-sub] Story '{story_id}' is held/frozen/blocked/dep-blocked; \ + skipping auto-retry ({attempt}/{budget})." + ); + return; + } + + if is_mergemaster_running(pool, project_root, story_id).await { + slog!( + "[merge-block-sub] Story '{story_id}' has an active mergemaster; \ + skipping auto-retry ({attempt}/{budget}) — recovery owns the story." + ); + return; + } + + match pool.start_merge_agent_work(project_root, story_id) { + Ok(()) => { + slog!( + "[merge-block-sub] Auto-retrying merge for '{story_id}' \ + (attempt {attempt}/{budget})." + ); + let _ = pool.watcher_tx.send(WatcherEvent::MergeAutoRetry { + story_id: story_id.to_string(), + attempt, + budget, + }); + } + Err(e) => { + slog_warn!( + "[merge-block-sub] Auto-retry for '{story_id}' ({attempt}/{budget}) \ + could not start: {e}; not notifying." + ); + } + } +} + +/// One-shot startup scan: schedule a first auto-retry for every story already +/// parked in `MergeFailure{GatesFailed}`. +/// +/// The pre-restart attempt count is unrecoverable, so the counter restarts at +/// 1 — worst case a story gets up to `threshold - 1` extra retries across a +/// restart, still bounded per process lifetime. +async fn reconcile_stranded_gates_failed( + pool: &Arc, + project_root: &Path, + counters: &mut HashMap, + generations: &Generations, +) { + let threshold = load_threshold(project_root).await; + if threshold == 0 { + return; + } + for item in crate::pipeline_state::read_all_typed() { + let Stage::MergeFailure { kind, .. } = &item.stage else { + continue; + }; + if !matches!(kind, MergeFailureKind::GatesFailed(_)) { + continue; + } + counters.insert(item.story_id.clone(), (1, kind.clone())); + slog!( + "[merge-block-sub] Story '{}' found parked in GatesFailed at startup; \ + scheduling catch-up auto-retry (attempt 1/{threshold}).", + item.story_id.0 + ); + schedule_auto_retry( + Arc::clone(pool), + project_root.to_path_buf(), + item.story_id.0.clone(), + 1, + threshold, + Arc::clone(generations), + ); + } +} + +/// Bump and return the scheduling generation for `story_id`. +fn next_generation(generations: &Generations, story_id: &str) -> u64 { + let mut map = generations.lock().unwrap_or_else(|p| p.into_inner()); + let entry = map.entry(story_id.to_string()).or_insert(0); + *entry += 1; + *entry +} + +/// Drop the generation entry for `story_id`, making every pending timer stale. +fn invalidate_generation(generations: &Generations, story_id: &str) { + generations + .lock() + .unwrap_or_else(|p| p.into_inner()) + .remove(story_id); +} + +/// True when `expected` is still the current scheduling generation. +fn is_generation_current(generations: &Generations, story_id: &str, expected: u64) -> bool { + generations + .lock() + .unwrap_or_else(|p| p.into_inner()) + .get(story_id) + == Some(&expected) +} + /// Return true if a mergemaster agent is currently in the pool for `story_id`. -/// Used to suppress counter increments while recovery is actively iterating +/// Used to suppress counting and retries while recovery is actively iterating /// (bug 1025). async fn is_mergemaster_running(pool: &AgentPool, project_root: &Path, story_id: &str) -> bool { - let config = match crate::config::ProjectConfig::load(project_root) { - Ok(c) => c, - Err(_) => return false, + let root = project_root.to_path_buf(); + let config = match tokio::task::spawn_blocking(move || { + crate::config::ProjectConfig::load(&root) + }) + .await + { + Ok(Ok(c)) => c, + _ => return false, }; let agents = pool.agents.lock().await; is_story_assigned_for_stage(&config, &agents, story_id, &PipelineStage::Mergemaster) } -/// Handle a single transition event: update counters and emit Block if threshold is reached. -/// -/// `recovery_running`: when `true`, a mergemaster is currently in the pool for -/// the story and the failure is part of an in-flight recovery loop. We do NOT -/// increment the consecutive-failure counter in that case (bug 1025). -fn on_transition( - project_root: &Path, - fired: &crate::pipeline_state::TransitionFired, - counters: &mut HashMap, - recovery_running: bool, -) { - // Story 1086: gate on the typed `Status` projection — `Status::MergeFailure` - // is precisely the set of stages we count toward the block threshold. We - // still need the variant pattern below to read `kind`. - if fired.after.status() != Status::MergeFailure { - counters.remove(&fired.story_id); - return; - } - match &fired.after { - Stage::MergeFailure { kind, .. } => { - if recovery_running { - slog!( - "[merge-block-sub] Story '{}' MergeFailure while mergemaster is running; \ - not counting toward block threshold (recovery in progress).", - fired.story_id.0 - ); - return; - } - let entry = counters - .entry(fired.story_id.clone()) - .or_insert_with(|| (0, kind.clone())); - entry.0 += 1; - entry.1 = kind.clone(); - - let count = entry.0; - let threshold = load_threshold(project_root); - - if threshold == 0 { - return; - } - - if count >= threshold { - let kind_str = failure_kind_label(kind); - let reason = format!( - "Auto-blocked after {count} consecutive MergeFailure ({kind_str}) transitions." - ); - let story_id = fired.story_id.0.as_str(); - slog!( - "[merge-block-sub] Story '{story_id}' reached {count} consecutive \ - MergeFailure ({kind_str}); blocking." - ); - if let Err(e) = crate::pipeline_state::apply_transition( - story_id, - PipelineEvent::Block { reason }, - None, - ) { - slog_warn!("[merge-block-sub] Failed to block '{story_id}': {e}"); - } else { - counters.remove(&fired.story_id); - } - } - } - _ => { - counters.remove(&fired.story_id); - } - } -} - -/// Load the threshold from project config, falling back to the compiled default. -fn load_threshold(project_root: &Path) -> u32 { - crate::config::ProjectConfig::load(project_root) - .map(|c| c.merge_failure_block_threshold) - .unwrap_or(3) +/// Load the budget from project config off the async runtime (the read is +/// synchronous filesystem I/O — bug 1170 class), falling back to the compiled +/// default. +async fn load_threshold(project_root: &Path) -> u32 { + let root = project_root.to_path_buf(); + tokio::task::spawn_blocking(move || { + crate::config::ProjectConfig::load(&root) + .map(|c| c.merge_failure_block_threshold) + .unwrap_or(3) + }) + .await + .unwrap_or(3) } /// Short human-readable label for a [`MergeFailureKind`] variant. @@ -175,11 +463,7 @@ mod tests { use crate::pipeline_state::{BranchName, PipelineEvent, Stage, StoryId, TransitionFired}; use std::num::NonZeroU32; - fn setup_project(tmp: &tempfile::TempDir) { - let sk = tmp.path().join(".huskies"); - std::fs::create_dir_all(&sk).unwrap(); - std::fs::write(sk.join("project.toml"), "[[agent]]\nname = \"coder\"\n").unwrap(); - } + const THRESHOLD: u32 = 3; fn seed_at_merge(story_id: &str) { crate::crdt_state::init_for_test(); @@ -230,218 +514,232 @@ mod tests { } } - /// AC3 (threshold-not-reached): 2 consecutive failures below threshold of 3 must NOT block. - #[test] - fn below_threshold_does_not_block() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1018_below"; - seed_at_merge(story_id); - - // Transition to MergeFailure once to establish the stage. - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::GatesFailed("error".to_string()), - ) - .expect("initial MergeFailure transition"); - - let mut counters: HashMap = HashMap::new(); - let kind = MergeFailureKind::GatesFailed("error".to_string()); - - // Fire 2 MergeFailure events (default threshold is 3). - for _ in 0..2 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, false); + /// The MergeFailure → Merge bounce a retry causes. + fn make_retry_started_fired(story_id: &str) -> TransitionFired { + TransitionFired { + story_id: StoryId(story_id.to_string()), + before: Stage::MergeFailure { + kind: MergeFailureKind::GatesFailed("error".to_string()), + feature_branch: BranchName("feature/test".to_string()), + commits_ahead: NonZeroU32::new(1).unwrap(), + }, + after: Stage::Merge { + feature_branch: BranchName("feature/test".to_string()), + commits_ahead: NonZeroU32::new(1).unwrap(), + claim: None, + retries: 1, + server_start_time: None, + }, + event: PipelineEvent::MergeRetryStarted, + at: chrono::Utc::now(), } + } - // Story must still be in MergeFailure (not Blocked). - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - assert!( - matches!(item.stage, Stage::MergeFailure { .. }), - "story must still be in MergeFailure after 2 failures (threshold 3): {:?}", - item.stage + fn gates_failed() -> MergeFailureKind { + MergeFailureKind::GatesFailed("error".to_string()) + } + + /// Below the threshold, GatesFailed schedules a retry with the right + /// attempt numbering. + #[test] + fn gates_failed_below_threshold_schedules_retry() { + let mut counters = HashMap::new(); + let fired = make_merge_failure_fired("t_sched", gates_failed()); + + assert_eq!( + decide(&fired, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { + attempt: 1, + budget: THRESHOLD + } + ); + assert_eq!( + decide(&fired, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { + attempt: 2, + budget: THRESHOLD + } ); } - /// AC3 (threshold-reached): 3 consecutive failures at threshold of 3 must block. + /// 1185 review finding 1 (regression): the retry's own MergeFailure→Merge + /// bounce must NOT reset the counter — the third consecutive failure + /// blocks even though retries happened in between. #[test] - fn at_threshold_blocks_with_failure_kind_in_reason() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); + fn merge_retry_started_does_not_reset_counter() { + let mut counters = HashMap::new(); + let story = "t_no_reset"; + let fail = make_merge_failure_fired(story, gates_failed()); + let bounce = make_retry_started_fired(story); + + assert!(matches!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { attempt: 1, .. } + )); + assert_eq!( + decide(&bounce, &mut counters, false, THRESHOLD), + Decision::Nothing + ); + assert!(matches!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { attempt: 2, .. } + )); + assert_eq!( + decide(&bounce, &mut counters, false, THRESHOLD), + Decision::Nothing + ); + // Third consecutive failure: budget exhausted despite the bounces. + assert_eq!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::Block { count: 3 } + ); + } + + /// A real departure (FixupRequested → Coding) still resets the counter. + #[test] + fn real_departure_resets_counter() { + let mut counters = HashMap::new(); + let story = "t_reset"; + let fail = make_merge_failure_fired(story, gates_failed()); + + decide(&fail, &mut counters, false, THRESHOLD); + decide(&fail, &mut counters, false, THRESHOLD); + assert_eq!( + counters.get(&StoryId(story.to_string())).map(|e| e.0), + Some(2) + ); + + decide(&make_coding_fired(story), &mut counters, false, THRESHOLD); + assert!(!counters.contains_key(&StoryId(story.to_string()))); + + // Fresh failures start a fresh budget. + assert!(matches!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { attempt: 1, .. } + )); + } + + /// Non-GatesFailed kinds count toward the block budget but never schedule + /// a retry (ConflictDetected has its own mergemaster path; the rest need + /// humans). + #[test] + fn non_gates_failed_counts_but_does_not_retry() { + let mut counters = HashMap::new(); + let story = "t_conflict"; + let conflict = make_merge_failure_fired(story, MergeFailureKind::ConflictDetected(None)); + + assert_eq!( + decide(&conflict, &mut counters, false, THRESHOLD), + Decision::Nothing + ); + assert_eq!( + counters.get(&StoryId(story.to_string())).map(|e| e.0), + Some(1) + ); + + assert_eq!( + decide(&conflict, &mut counters, false, THRESHOLD), + Decision::Nothing + ); + assert_eq!( + decide(&conflict, &mut counters, false, THRESHOLD), + Decision::Block { count: 3 } + ); + } + + /// Mixed kinds share one budget: GatesFailed and ConflictDetected + /// interleavings block at the same total count (1185 review finding 5). + #[test] + fn mixed_kinds_share_one_budget() { + let mut counters = HashMap::new(); + let story = "t_mixed"; + let fail = make_merge_failure_fired(story, gates_failed()); + let conflict = make_merge_failure_fired(story, MergeFailureKind::ConflictDetected(None)); + + assert!(matches!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::ScheduleRetry { attempt: 1, .. } + )); + assert_eq!( + decide(&conflict, &mut counters, false, THRESHOLD), + Decision::Nothing + ); + assert_eq!( + decide(&fail, &mut counters, false, THRESHOLD), + Decision::Block { count: 3 } + ); + } + + /// Bug 1025: recovery in progress neither counts nor schedules. + #[test] + fn mergemaster_running_suppresses_counting_and_retry() { + let mut counters = HashMap::new(); + let story = "t_recovery"; + let fail = make_merge_failure_fired(story, gates_failed()); + + for _ in 0..3 { + assert_eq!( + decide(&fail, &mut counters, true, THRESHOLD), + Decision::Nothing + ); + } + assert!(!counters.contains_key(&StoryId(story.to_string()))); + } + + /// threshold == 0 disables both policies. + #[test] + fn threshold_zero_disables_block_and_retry() { + let mut counters = HashMap::new(); + let fail = make_merge_failure_fired("t_disabled", gates_failed()); + for _ in 0..5 { + assert_eq!(decide(&fail, &mut counters, false, 0), Decision::Nothing); + } + } + + /// Applying a Block decision transitions the story and names the kind. + #[test] + fn apply_block_blocks_with_failure_kind_in_reason() { let story_id = "1018_at_threshold"; seed_at_merge(story_id); - crate::agents::lifecycle::transition_to_merge_failure( story_id, MergeFailureKind::GatesFailed("fmt error".to_string()), ) .expect("initial MergeFailure transition"); - let mut counters: HashMap = HashMap::new(); - let kind = MergeFailureKind::GatesFailed("fmt error".to_string()); - - // Fire 3 MergeFailure events — the 3rd must trigger the block. - for _ in 0..3 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, false); - } + apply_block( + &StoryId(story_id.to_string()), + 3, + &MergeFailureKind::GatesFailed("fmt error".to_string()), + ); let item = crate::pipeline_state::read_typed(story_id) .expect("read") .expect("item"); - assert!( - matches!(item.stage, Stage::Blocked { .. }), - "story must be Blocked after 3 consecutive MergeFailures: {:?}", - item.stage - ); - - // The block reason must name the failure kind. - if let Stage::Blocked { reason } = &item.stage { - assert!( - reason.contains("GatesFailed"), - "block reason must name the failure kind: {reason}" - ); + match &item.stage { + Stage::Blocked { reason } => { + assert!( + reason.contains("GatesFailed"), + "block reason must name the failure kind: {reason}" + ); + } + other => panic!("story must be Blocked: {other:?}"), } } - /// AC3 (reset): counter clears after a non-MergeFailure transition. - /// - /// 2 failures → FixupRequested reset → 2 more failures: still below threshold, no block. + /// 1185 review finding 6 (regression): a newer schedule or a reset makes + /// earlier timers stale. #[test] - fn counter_resets_on_non_merge_failure_transition() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1018_reset"; - seed_at_merge(story_id); + fn stale_generations_are_not_current() { + let generations: Generations = Generations::default(); + let g1 = next_generation(&generations, "s"); + assert!(is_generation_current(&generations, "s", g1)); - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::ConflictDetected(None), - ) - .expect("initial MergeFailure transition"); + let g2 = next_generation(&generations, "s"); + assert!(!is_generation_current(&generations, "s", g1)); + assert!(is_generation_current(&generations, "s", g2)); - let mut counters: HashMap = HashMap::new(); - let kind = MergeFailureKind::ConflictDetected(None); - - // Fire 2 MergeFailure events. - for _ in 0..2 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, false); - } - assert_eq!( - counters.get(&StoryId(story_id.to_string())).map(|e| e.0), - Some(2), - "counter must be 2 after 2 failures" - ); - - // Simulate FixupRequested (non-MergeFailure transition). - let reset_fired = make_coding_fired(story_id); - on_transition(tmp.path(), &reset_fired, &mut counters, false); - assert!( - !counters.contains_key(&StoryId(story_id.to_string())), - "counter must be cleared after non-MergeFailure transition" - ); - - // Re-seed to MergeFailure so we can apply the block transition. - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::ConflictDetected(None), - ) - .expect("re-enter MergeFailure after reset"); - - // Fire 2 more MergeFailure events — still below threshold. - for _ in 0..2 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, false); - } - - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - assert!( - matches!(item.stage, Stage::MergeFailure { .. }), - "story must still be in MergeFailure after reset + 2 new failures: {:?}", - item.stage - ); - } - - /// Bug 1025: while a mergemaster is running, MergeFailure transitions are - /// recovery iterations, not consecutive give-ups. 3 failures with - /// `recovery_running=true` must NOT block. - #[test] - fn mergemaster_running_suppresses_block() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1025_recovery_running"; - seed_at_merge(story_id); - - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::ConflictDetected(None), - ) - .expect("initial MergeFailure transition"); - - let mut counters: HashMap = HashMap::new(); - let kind = MergeFailureKind::ConflictDetected(None); - - // Fire 3 MergeFailure events WHILE a mergemaster is running (gated). - for _ in 0..3 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, true); - } - - // Counter must NOT have incremented at all — recovery in progress. - assert!( - !counters.contains_key(&StoryId(story_id.to_string())), - "counter must not increment while mergemaster is running" - ); - - // And the story must still be in MergeFailure (not Blocked). - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - assert!( - matches!(item.stage, Stage::MergeFailure { .. }), - "story must NOT be blocked while mergemaster is running (recovery in progress): {:?}", - item.stage - ); - } - - /// Bug 1025 regression guard: the genuinely-stuck case (no mergemaster - /// running) still blocks at the threshold, so the original 1018 behaviour - /// is preserved. - #[test] - fn no_mergemaster_still_blocks_at_threshold() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1025_genuine_stuck"; - seed_at_merge(story_id); - - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::ConflictDetected(None), - ) - .expect("initial MergeFailure transition"); - - let mut counters: HashMap = HashMap::new(); - let kind = MergeFailureKind::ConflictDetected(None); - - // Fire 3 MergeFailure events with NO mergemaster (recovery_running=false). - for _ in 0..3 { - let fired = make_merge_failure_fired(story_id, kind.clone()); - on_transition(tmp.path(), &fired, &mut counters, false); - } - - // Story must be Blocked (genuine-stuck case unchanged). - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - assert!( - matches!(item.stage, Stage::Blocked { .. }), - "story must still block when no mergemaster is running: {:?}", - item.stage - ); + invalidate_generation(&generations, "s"); + assert!(!is_generation_current(&generations, "s", g2)); } } diff --git a/server/src/agents/pool/auto_assign/merge_failure_retry_subscriber.rs b/server/src/agents/pool/auto_assign/merge_failure_retry_subscriber.rs deleted file mode 100644 index a7163e87..00000000 --- a/server/src/agents/pool/auto_assign/merge_failure_retry_subscriber.rs +++ /dev/null @@ -1,570 +0,0 @@ -//! TransitionFired subscriber that auto-retries `GatesFailed` merge failures after a delay. -//! -//! Story 1185: ConflictDetected already gets an automatic recovery path via -//! [`super::merge_failure_subscriber`] (mergemaster auto-spawn). `GatesFailed` -//! previously just sat in `Stage::MergeFailure` until a human intervened or the -//! auto-block subscriber ([`super::merge_failure_block_subscriber`]) blocked the -//! story after `merge_failure_block_threshold` consecutive failures. -//! -//! This subscriber schedules a delayed re-trigger of the deterministic -//! server-side merge for `GatesFailed` failures, sharing the same -//! `merge_failure_block_threshold` budget the auto-block subscriber uses as its -//! retry ceiling. Once that many consecutive `GatesFailed` failures have -//! occurred, no further retry is scheduled and the story is left exactly where -//! it is — the untouched auto-block subscriber independently counts the same -//! consecutive-failure stream and transitions the story to `Stage::Blocked` at -//! the same threshold, so behaviour after budget exhaustion is unchanged from -//! today. - -use std::collections::HashMap; -use std::path::{Path, PathBuf}; -use std::sync::Arc; -use std::time::Duration; - -use crate::io::watcher::WatcherEvent; -use crate::pipeline_state::{MergeFailureKind, Stage, Status, StoryId, TransitionFired}; -use crate::slog; -use crate::slog_warn; - -use super::super::AgentPool; - -/// Delay before an auto-retry re-triggers the server-side merge for a -/// `GatesFailed` failure. Gives a moment for transient conditions (e.g. a -/// flaky test) to clear before retrying the exact same commit — retrying -/// instantly would just hammer the same deterministic failure. -const AUTO_RETRY_DELAY: Duration = Duration::from_secs(30); - -/// Reconcile: no-op, matching [`super::merge_failure_block_subscriber`]. -/// -/// The per-story attempt counter is in-memory bookkeeping that cannot be -/// reconstructed from CRDT state alone (only the current stage is stored, not -/// how many times a story has already been auto-retried). The live subscriber -/// reacting to each new transition is sufficient for eventual consistency. -pub(crate) fn reconcile_merge_failure_retry() {} - -/// Spawn a background task that auto-retries `GatesFailed` merge failures. -/// -/// Subscribes to the pipeline transition broadcast channel and tracks a -/// per-story auto-retry attempt counter, bounded by the same -/// `merge_failure_block_threshold` budget the auto-block subscriber uses. -pub(crate) fn spawn_merge_failure_retry_subscriber(pool: Arc, project_root: PathBuf) { - let mut rx = crate::pipeline_state::subscribe_transitions(); - tokio::spawn(async move { - let mut attempts: HashMap = HashMap::new(); - loop { - match rx.recv().await { - Ok(fired) => { - if let Some((attempt, budget)) = - decide_retry(&project_root, &fired, &mut attempts) - { - schedule_auto_retry( - Arc::clone(&pool), - project_root.clone(), - fired.story_id.0.clone(), - attempt, - budget, - ); - } - } - Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => { - slog_warn!( - "[merge-retry-sub] Subscriber lagged, skipped {n} event(s). \ - Some GatesFailed stories may need a manual retry." - ); - } - Err(tokio::sync::broadcast::error::RecvError::Closed) => break, - } - } - }); -} - -/// Pure decision function: given a fired transition and the running per-story -/// attempt counter, decide whether an auto-retry should be scheduled. -/// -/// Returns `Some((attempt, budget))` when a retry should fire — `attempt` is -/// this retry's 1-based sequence number and `budget` is the shared -/// `merge_failure_block_threshold`. Returns `None` when the story left -/// `MergeFailure`, the failure kind isn't `GatesFailed`, or the budget is -/// exhausted (`attempt >= budget`). -fn decide_retry( - project_root: &Path, - fired: &TransitionFired, - attempts: &mut HashMap, -) -> Option<(u32, u32)> { - if fired.after.status() != Status::MergeFailure { - attempts.remove(&fired.story_id); - return None; - } - let Stage::MergeFailure { kind, .. } = &fired.after else { - return None; - }; - if !matches!(kind, MergeFailureKind::GatesFailed(_)) { - attempts.remove(&fired.story_id); - return None; - } - - let entry = attempts.entry(fired.story_id.clone()).or_insert(0); - *entry += 1; - let attempt = *entry; - - let budget = load_threshold(project_root); - if budget == 0 || attempt >= budget { - slog!( - "[merge-retry-sub] Story '{}' exhausted GatesFailed auto-retry budget \ - ({attempt}/{budget}); parking in MergeFailure for human intervention.", - fired.story_id.0 - ); - return None; - } - - Some((attempt, budget)) -} - -/// Spawn the actual delayed retry as a background task. -/// -/// Kept separate from [`decide_retry`] so the counter/budget decision stays -/// synchronous and unit-testable without waiting on a real timer. -fn schedule_auto_retry( - pool: Arc, - project_root: PathBuf, - story_id: String, - attempt: u32, - budget: u32, -) { - slog!( - "[merge-retry-sub] Story '{story_id}' GatesFailed (attempt {attempt}/{budget}); \ - scheduling auto-retry in {AUTO_RETRY_DELAY:?}." - ); - tokio::spawn(async move { - tokio::time::sleep(AUTO_RETRY_DELAY).await; - perform_auto_retry(&pool, &project_root, &story_id, attempt, budget).await; - }); -} - -/// Re-check the story is still parked in `GatesFailed`, notify chat, and -/// re-trigger the deterministic server-side merge. -/// -/// The guard protects against the story having been fixed manually, blocked, -/// or otherwise moved on while the retry delay elapsed. -async fn perform_auto_retry( - pool: &Arc, - project_root: &Path, - story_id: &str, - attempt: u32, - budget: u32, -) { - let still_gates_failed = matches!( - crate::pipeline_state::read_typed(story_id), - Ok(Some(item)) if matches!( - item.stage, - Stage::MergeFailure { - kind: MergeFailureKind::GatesFailed(_), - .. - } - ) - ); - if !still_gates_failed { - slog!( - "[merge-retry-sub] Story '{story_id}' left GatesFailed before auto-retry \ - ({attempt}/{budget}) fired; skipping." - ); - return; - } - - slog!("[merge-retry-sub] Auto-retrying merge for '{story_id}' (attempt {attempt}/{budget})."); - let _ = pool.watcher_tx.send(WatcherEvent::MergeAutoRetry { - story_id: story_id.to_string(), - attempt, - budget, - }); - pool.trigger_server_side_merge(project_root, story_id); -} - -/// Load the auto-retry budget from project config, falling back to the -/// compiled default. Shares `merge_failure_block_threshold` with the -/// auto-block subscriber (story 1185, AC2). -fn load_threshold(project_root: &Path) -> u32 { - crate::config::ProjectConfig::load(project_root) - .map(|c| c.merge_failure_block_threshold) - .unwrap_or(3) -} - -// ── Tests ────────────────────────────────────────────────────────────────── - -#[cfg(test)] -mod tests { - use super::*; - use crate::pipeline_state::{BranchName, PipelineEvent, Stage, StoryId, TransitionFired}; - use std::num::NonZeroU32; - use std::process::Command; - - fn setup_project(tmp: &tempfile::TempDir) { - let sk = tmp.path().join(".huskies"); - std::fs::create_dir_all(&sk).unwrap(); - std::fs::write(sk.join("project.toml"), "[[agent]]\nname = \"coder\"\n").unwrap(); - } - - fn seed_at_merge(story_id: &str) { - crate::crdt_state::init_for_test(); - crate::db::ensure_content_store(); - crate::db::write_item_with_content( - story_id, - "4_merge", - "---\nname: Test\n---\n", - crate::db::ItemMeta::named("Test"), - ); - } - - fn make_merge_failure_fired(story_id: &str, kind: MergeFailureKind) -> TransitionFired { - TransitionFired { - story_id: StoryId(story_id.to_string()), - before: Stage::Merge { - feature_branch: BranchName("feature/test".to_string()), - commits_ahead: NonZeroU32::new(1).unwrap(), - claim: None, - retries: 0, - server_start_time: None, - }, - after: Stage::MergeFailure { - kind: kind.clone(), - feature_branch: BranchName("feature/test".to_string()), - commits_ahead: NonZeroU32::new(1).unwrap(), - }, - event: PipelineEvent::MergeFailed { kind }, - at: chrono::Utc::now(), - } - } - - fn make_coding_fired(story_id: &str) -> TransitionFired { - TransitionFired { - story_id: StoryId(story_id.to_string()), - before: Stage::MergeFailure { - kind: MergeFailureKind::GatesFailed("error".to_string()), - feature_branch: BranchName("feature/test".to_string()), - commits_ahead: NonZeroU32::new(1).unwrap(), - }, - after: Stage::Coding { - claim: None, - plan: Default::default(), - retries: 0, - }, - event: PipelineEvent::FixupRequested, - at: chrono::Utc::now(), - } - } - - // ── decide_retry: pure decision logic ─────────────────────────────────── - - /// AC1/AC6: a GatesFailed failure below the budget must be scheduled for retry. - #[test] - fn gates_failed_below_budget_schedules_retry() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1185_retry_basic"; - seed_at_merge(story_id); - - let mut attempts: HashMap = HashMap::new(); - let fired = make_merge_failure_fired( - story_id, - MergeFailureKind::GatesFailed("cargo test failed".to_string()), - ); - - let decision = decide_retry(tmp.path(), &fired, &mut attempts); - assert_eq!( - decision, - Some((1, 3)), - "first GatesFailed failure (default threshold 3) must schedule attempt 1/3" - ); - } - - /// AC2/AC6: once the attempt count reaches the shared threshold, no further - /// retry is scheduled — the story parks in MergeFailure exactly as today - /// (the untouched block subscriber handles the eventual Blocked transition). - #[test] - fn budget_exhaustion_stops_scheduling_retries() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1185_retry_exhausted"; - seed_at_merge(story_id); - - let mut attempts: HashMap = HashMap::new(); - let kind = MergeFailureKind::GatesFailed("cargo test failed".to_string()); - - // Default threshold is 3: attempts 1 and 2 schedule a retry, attempt 3 - // (== budget) does not. - let fired1 = make_merge_failure_fired(story_id, kind.clone()); - assert_eq!( - decide_retry(tmp.path(), &fired1, &mut attempts), - Some((1, 3)) - ); - - let fired2 = make_merge_failure_fired(story_id, kind.clone()); - assert_eq!( - decide_retry(tmp.path(), &fired2, &mut attempts), - Some((2, 3)) - ); - - let fired3 = make_merge_failure_fired(story_id, kind); - assert_eq!( - decide_retry(tmp.path(), &fired3, &mut attempts), - None, - "budget exhausted at attempt == threshold; must not schedule another retry" - ); - } - - /// AC3/AC6: ConflictDetected, EmptyDiff, NoCommits, and Other must never be - /// auto-retried by this subscriber (ConflictDetected has its own mergemaster - /// auto-spawn path; the rest require human intervention). - #[test] - fn non_gates_failed_kinds_are_not_retried() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - - for (story_id, kind) in [ - ( - "1185_retry_conflict", - MergeFailureKind::ConflictDetected(None), - ), - ("1185_retry_emptydiff", MergeFailureKind::EmptyDiff), - ("1185_retry_nocommits", MergeFailureKind::NoCommits), - ( - "1185_retry_other", - MergeFailureKind::Other("unknown".to_string()), - ), - ] { - seed_at_merge(story_id); - let mut attempts: HashMap = HashMap::new(); - let fired = make_merge_failure_fired(story_id, kind); - assert_eq!( - decide_retry(tmp.path(), &fired, &mut attempts), - None, - "non-GatesFailed kind must not be scheduled for auto-retry" - ); - } - } - - /// Counter resets when the story leaves MergeFailure (e.g. FixupRequested), - /// mirroring the auto-block subscriber's reset behaviour. - #[test] - fn counter_resets_on_non_merge_failure_transition() { - let tmp = tempfile::tempdir().unwrap(); - setup_project(&tmp); - let story_id = "1185_retry_reset"; - seed_at_merge(story_id); - - let mut attempts: HashMap = HashMap::new(); - let kind = MergeFailureKind::GatesFailed("error".to_string()); - - let fired1 = make_merge_failure_fired(story_id, kind.clone()); - decide_retry(tmp.path(), &fired1, &mut attempts); - assert_eq!(attempts.get(&StoryId(story_id.to_string())), Some(&1)); - - let reset_fired = make_coding_fired(story_id); - decide_retry(tmp.path(), &reset_fired, &mut attempts); - assert!( - !attempts.contains_key(&StoryId(story_id.to_string())), - "counter must be cleared after non-MergeFailure transition" - ); - - // A fresh failure after reset must start back at attempt 1. - let fired2 = make_merge_failure_fired(story_id, kind); - assert_eq!( - decide_retry(tmp.path(), &fired2, &mut attempts), - Some((1, 3)) - ); - } - - /// threshold == 0 disables auto-retry entirely, mirroring the auto-block - /// subscriber's "0 disables" convention. - #[test] - fn zero_threshold_disables_auto_retry() { - let tmp = tempfile::tempdir().unwrap(); - let sk = tmp.path().join(".huskies"); - std::fs::create_dir_all(&sk).unwrap(); - // Top-level keys must precede `[[agent]]` — TOML parses anything after - // the array-of-tables header as belonging to that table. - std::fs::write( - sk.join("project.toml"), - "merge_failure_block_threshold = 0\n[[agent]]\nname = \"coder\"\n", - ) - .unwrap(); - let story_id = "1185_retry_zero_threshold"; - seed_at_merge(story_id); - - let mut attempts: HashMap = HashMap::new(); - let fired = - make_merge_failure_fired(story_id, MergeFailureKind::GatesFailed("error".to_string())); - assert_eq!(decide_retry(tmp.path(), &fired, &mut attempts), None); - } - - // ── perform_auto_retry: end-to-end success path ───────────────────────── - - fn init_git_repo(repo: &std::path::Path) { - Command::new("git") - .args(["init"]) - .current_dir(repo) - .output() - .unwrap(); - Command::new("git") - .args(["config", "user.email", "test@test.com"]) - .current_dir(repo) - .output() - .unwrap(); - Command::new("git") - .args(["config", "user.name", "Test"]) - .current_dir(repo) - .output() - .unwrap(); - Command::new("git") - .args(["commit", "--allow-empty", "-m", "init"]) - .current_dir(repo) - .output() - .unwrap(); - } - - /// AC5/AC6: a mid-budget auto-retry that succeeds must drive the story to - /// `Stage::Done` via the same path a manual re-run uses - /// (`trigger_server_side_merge` → `start_merge_agent_work` → - /// `move_story_to_done`) — no 1178-style trap state. - #[tokio::test] - async fn perform_auto_retry_success_completes_the_story() { - use std::fs; - use tempfile::tempdir; - - crate::crdt_state::init_for_test(); - let tmp = tempdir().unwrap(); - let repo = tmp.path(); - init_git_repo(repo); - - let story_id = "1185_retry_success"; - let branch = format!("feature/story-{story_id}"); - Command::new("git") - .args(["checkout", "-b", &branch]) - .current_dir(repo) - .output() - .unwrap(); - fs::write(repo.join("feature.txt"), "feature content").unwrap(); - Command::new("git") - .args(["add", "."]) - .current_dir(repo) - .output() - .unwrap(); - Command::new("git") - .args(["commit", "-m", "add feature"]) - .current_dir(repo) - .output() - .unwrap(); - Command::new("git") - .args(["checkout", "master"]) - .current_dir(repo) - .output() - .unwrap(); - - crate::db::ensure_content_store(); - crate::db::write_item_with_content( - story_id, - "4_merge", - "---\nname: Retry Success Test\n---\n", - crate::db::ItemMeta::named("Retry Success Test"), - ); - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::GatesFailed("cargo test failed".to_string()), - ) - .expect("seed MergeFailure(GatesFailed)"); - - let pool = Arc::new(AgentPool::new_test(3199)); - perform_auto_retry(&pool, repo, story_id, 1, 3).await; - - // start_merge_agent_work transitions MergeFailure -> Merge synchronously - // (transition_merge_failure_to_retry) and then runs the actual squash - // merge as a background task; poll until the story reaches Done (or - // times out). Merge is an expected in-flight state, not a terminal one - // — only stop polling on Done or a re-observed MergeFailure/Blocked. - // Generous deadline: the real git subprocess pipeline (branch - // checkout, worktree add/remove, squash, cherry-pick, source-map - // regen) can be slow under a loaded, parallel `cargo test` run - // (matches the 5s+ headroom used by similar tests in - // `agents/pool/pipeline/merge/tests.rs` and `pipeline/advance/tests_regression.rs`). - let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); - loop { - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - if matches!(item.stage, Stage::Done { .. }) { - break; - } - if !matches!(item.stage, Stage::Merge { .. }) { - panic!( - "successful mid-budget auto-retry must complete the story via \ - the manual-rerun path; got {:?}", - item.stage - ); - } - if tokio::time::Instant::now() >= deadline { - panic!( - "auto-retry did not reach Done in time; still {:?}", - item.stage - ); - } - tokio::time::sleep(std::time::Duration::from_millis(50)).await; - } - } - - /// The guard in `perform_auto_retry` must skip re-triggering when the story - /// left GatesFailed while the retry delay was elapsing (e.g. blocked or - /// fixed manually in the meantime). - #[tokio::test] - async fn perform_auto_retry_skips_when_story_left_gates_failed() { - crate::crdt_state::init_for_test(); - let tmp = tempfile::tempdir().unwrap(); - let repo = tmp.path(); - - let story_id = "1185_retry_stale_guard"; - crate::db::ensure_content_store(); - crate::db::write_item_with_content( - story_id, - "4_merge", - "---\nname: Stale Guard Test\n---\n", - crate::db::ItemMeta::named("Stale Guard Test"), - ); - crate::agents::lifecycle::transition_to_merge_failure( - story_id, - MergeFailureKind::GatesFailed("error".to_string()), - ) - .expect("seed MergeFailure(GatesFailed)"); - - // Story moved on (e.g. blocked) before the delayed retry fired. - crate::pipeline_state::apply_transition( - story_id, - PipelineEvent::Block { - reason: "manual block".to_string(), - }, - None, - ) - .expect("block story"); - - let pool = Arc::new(AgentPool::new_test(3198)); - perform_auto_retry(&pool, repo, story_id, 1, 3).await; - - // No new merge job should have been started — the story is still - // Blocked. `transition_to_merge_failure` above already wrote a - // "failed" job entry as part of seeding; the guard skipping the - // retry means that entry must be untouched (in particular, never - // flipped to "running" by `trigger_server_side_merge`). - let item = crate::pipeline_state::read_typed(story_id) - .expect("read") - .expect("item"); - assert!( - matches!(item.stage, Stage::Blocked { .. }), - "guard must skip re-triggering once the story left GatesFailed: {:?}", - item.stage - ); - let job_status = crate::crdt_state::read_merge_job(story_id).map(|j| j.status); - assert_ne!( - job_status.as_deref(), - Some("running"), - "no new merge job should be started when the guard skips the retry" - ); - } -} diff --git a/server/src/agents/pool/auto_assign/mod.rs b/server/src/agents/pool/auto_assign/mod.rs index 416e9b6b..6085d3a2 100644 --- a/server/src/agents/pool/auto_assign/mod.rs +++ b/server/src/agents/pool/auto_assign/mod.rs @@ -4,10 +4,9 @@ mod auto_assign; mod backlog; mod merge; -/// TransitionFired subscriber that auto-blocks stories after N consecutive MergeFailure transitions. +/// TransitionFired subscriber owning the consecutive-MergeFailure budget: +/// auto-retries GatesFailed below the threshold, auto-blocks at it. pub(crate) mod merge_failure_block_subscriber; -/// TransitionFired subscriber that auto-retries GatesFailed merge failures after a delay. -pub(crate) mod merge_failure_retry_subscriber; /// TransitionFired subscriber that auto-spawns mergemaster on ConflictDetected merge failures. pub(crate) mod merge_failure_subscriber; mod pipeline; @@ -24,10 +23,6 @@ pub(crate) use merge_failure_block_subscriber::reconcile_merge_failure_block; /// Re-export for `startup::tick_loop`. pub(crate) use merge_failure_block_subscriber::spawn_merge_failure_block_subscriber; /// Re-export for `startup::tick_loop`. -pub(crate) use merge_failure_retry_subscriber::reconcile_merge_failure_retry; -/// Re-export for `startup::tick_loop`. -pub(crate) use merge_failure_retry_subscriber::spawn_merge_failure_retry_subscriber; -/// Re-export for `startup::tick_loop`. pub(crate) use merge_failure_subscriber::reconcile_merge_failure; /// Re-export for `startup::tick_loop`. pub(crate) use merge_failure_subscriber::spawn_merge_failure_subscriber; diff --git a/server/src/startup/tick_loop.rs b/server/src/startup/tick_loop.rs index 4394697d..693c418e 100644 --- a/server/src/startup/tick_loop.rs +++ b/server/src/startup/tick_loop.rs @@ -89,24 +89,17 @@ pub(crate) fn spawn_event_bridges( root.clone(), ); - // Consecutive-failure auto-block subscriber: blocks stories after N - // consecutive MergeFailure transitions (story 1018). Bug 1025: takes - // the agent pool so it can gate the counter on mergemaster presence — - // failures during active recovery iteration do not count toward block. + // Consecutive-failure budget subscriber: auto-retries GatesFailed + // failures below merge_failure_block_threshold (story 1185) and blocks + // stories at it (story 1018), with one shared counter. Bug 1025: + // takes the agent pool so it can gate both policies on mergemaster + // presence — failures during active recovery iteration neither count + // nor retry. crate::agents::pool::auto_assign::spawn_merge_failure_block_subscriber( Arc::clone(&agents), root.clone(), ); - // GatesFailed auto-retry subscriber: re-triggers the server-side merge - // after a delay for GatesFailed failures, bounded by the same - // merge_failure_block_threshold budget the auto-block subscriber above - // uses (story 1185). - crate::agents::pool::auto_assign::spawn_merge_failure_retry_subscriber( - Arc::clone(&agents), - root.clone(), - ); - // Content-store GC subscriber: purges all ContentKey::* entries for a // story when it reaches a terminal stage, preventing zombie entries from // accumulating in the process heap (story 996). @@ -548,12 +541,11 @@ pub(crate) async fn run_reconcile_pass( // Merge-failure: spawn mergemaster for ConflictDetected stories with no active agent. crate::agents::pool::auto_assign::reconcile_merge_failure(agents, root).await; - // Merge-block: no-op (in-memory counter cannot be reconstructed from CRDT). + // Merge-block: no-op for the periodic pass (in-memory counter cannot be + // reconstructed from CRDT); stranded-GatesFailed catch-up runs once at + // subscriber startup instead. crate::agents::pool::auto_assign::reconcile_merge_failure_block(); - // Merge-retry: no-op (in-memory attempt counter cannot be reconstructed from CRDT). - crate::agents::pool::auto_assign::reconcile_merge_failure_retry(); - // Audit-log: no-op (historical replay would produce misleading entries). crate::pipeline_state::reconcile_audit_log(); }