pub(crate) mod acceptance_state;
pub(crate) mod analysis_signature;
mod archive_state;
mod builder;
mod cleanup;
mod conflict;
pub(crate) mod dedup;
mod dependency;
mod dispatch;
mod dynamic_queue;
mod events;
mod executor;
mod lifecycle_slots;
mod manual_continuation;
mod merge;
mod orchestration;
mod output_bridge;
pub(super) mod queue_state;
pub(crate) mod resolve_state;
mod target_plan;
mod types;
pub(crate) mod upstream_bridge;
mod upstream_lane;
mod work_snapshot;
mod workspace;
pub use crate::events::ExecutionEvent as ParallelEvent;
#[cfg(all(test, feature = "heavy-tests"))]
#[allow(unused_imports)]
pub use merge::{base_dirty_reason, resolve_deferred_merge};
pub use types::{
AlreadyReportedFailureKind, FailedChangeTracker, MergeResult, MergeResultDisposition,
MergeResultOrigin, MergeTaskOutcome, ResolveFailureClassification, WorkspaceResult,
};
#[allow(unused_imports)]
pub use types::resolve_failure_detail;
#[cfg(all(test, feature = "heavy-tests"))]
#[allow(unused_imports)]
pub use merge::MergeAttempt;
use crate::ai_command_runner::{AiCommandRunner, RunCommandScope, SharedStaggerState};
use crate::config::OrchestratorConfig;
use crate::hooks::HookRunner;
use crate::parallel::analysis_signature::{
AnalysisInputProbe, BoundedAnalysisRetry, CompletedAnalysisInput,
};
use crate::parallel::dedup::{DiagnosticDeduplicationKey, DiagnosticDeduplicationStore};
use crate::vcs::WorkspaceManager;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::{Arc, Mutex as StdMutex, OnceLock};
use tokio::sync::{mpsc, Mutex, RwLock};
use tokio_util::sync::CancellationToken;
use crate::orchestration::state::OrchestratorState;
type DependencyBlockerFingerprint = Vec<(String, String)>;
const DEFAULT_MAX_CONFLICT_RETRIES: u32 = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SchedulerLifetime {
Finite,
Persistent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SchedulerRunReport {
Completed,
CompletedWithErrors,
Stopped,
BlockedOrStalled,
}
impl SchedulerRunReport {
pub fn is_incomplete(self) -> bool {
matches!(self, Self::CompletedWithErrors | Self::BlockedOrStalled)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum PostArchiveAction {
#[default]
MergeToBase,
PushToRemote {
remote: String,
},
}
static GLOBAL_MERGE_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
static ACTIVE_POST_ARCHIVE_MERGES: OnceLock<StdMutex<HashSet<String>>> = OnceLock::new();
fn global_merge_lock() -> &'static Mutex<()> {
GLOBAL_MERGE_LOCK.get_or_init(|| Mutex::new(()))
}
#[cfg(test)]
pub(crate) fn merge_lock_test_mutex() -> &'static Mutex<()> {
static TEST_MUTEX: OnceLock<Mutex<()>> = OnceLock::new();
TEST_MUTEX.get_or_init(|| Mutex::new(()))
}
fn active_post_archive_merges() -> &'static StdMutex<HashSet<String>> {
ACTIVE_POST_ARCHIVE_MERGES.get_or_init(|| StdMutex::new(HashSet::new()))
}
pub struct ParallelExecutor {
workspace_manager: Box<dyn WorkspaceManager>,
config: OrchestratorConfig,
apply_command: String,
archive_command: String,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
max_conflict_retries: u32,
repo_root: PathBuf,
no_resume: bool,
explicit_retry: bool,
failed_tracker: FailedChangeTracker,
change_dependencies: HashMap<String, Vec<String>>,
resolve_wait_changes: HashSet<String>,
reject_wait_changes: HashSet<String>,
merge_wait_changes: HashSet<String>,
dependency_blocker_fingerprints: HashMap<String, DependencyBlockerFingerprint>,
force_recreate_worktree: HashSet<String>,
hooks: Option<Arc<HookRunner>>,
cancel_token: Option<CancellationToken>,
last_queue_change_at: Arc<Mutex<Option<std::time::Instant>>>,
last_available_slots: Option<usize>,
analyzer_capacity_suppressed: bool,
pending_retry_bypass: HashSet<String>,
dynamic_queue: Option<Arc<crate::tui::queue::DynamicQueue>>,
ai_runner: AiCommandRunner,
run_command_scope: RunCommandScope,
#[allow(dead_code)]
shared_stagger_state: SharedStaggerState,
apply_history: Arc<Mutex<crate::history::ApplyHistory>>,
archive_history: Arc<Mutex<crate::history::ArchiveHistory>>,
acceptance_history: Arc<Mutex<crate::history::AcceptanceHistory>>,
acceptance_tail_injected: Arc<Mutex<std::collections::HashMap<String, bool>>>,
apply_budget: crate::execution::apply::ApplyBudget,
lifecycle_slots: lifecycle_slots::LifecycleSlots,
manual_resolve_count: Option<Arc<std::sync::atomic::AtomicUsize>>,
auto_resolve_count: Arc<std::sync::atomic::AtomicUsize>,
pending_merge_count: Arc<std::sync::atomic::AtomicUsize>,
scheduler_lifetime: SchedulerLifetime,
graceful_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
persistent_idle_latched: Arc<std::sync::atomic::AtomicBool>,
persistent_idle_baseline: Arc<std::sync::Mutex<Option<HashSet<String>>>>,
post_archive_action: PostArchiveAction,
shared_orchestrator_state: Option<Arc<RwLock<OrchestratorState>>>,
last_dispatched_resolve_wait_changes: HashSet<String>,
last_dispatched_reject_wait_changes: HashSet<String>,
resolve_wait_retry_triggered: bool,
last_resolve_wait_base_dirty: Option<bool>,
diagnostic_dedup: DiagnosticDeduplicationStore<DiagnosticDeduplicationKey>,
last_completed_analysis_input: Option<CompletedAnalysisInput>,
next_analysis_signature_probe_at: Option<tokio::time::Instant>,
analysis_retry_throttle: Option<BoundedAnalysisRetry>,
upstream: Option<Arc<Mutex<crate::upstream::UpstreamCoordinator>>>,
analysis_input_probe: Option<Arc<dyn AnalysisInputProbe>>,
#[cfg(test)]
merge_result_channel_override: Option<(mpsc::Sender<MergeResult>, mpsc::Receiver<MergeResult>)>,
#[cfg(test)]
run_command_cleanup_budget_override: Option<std::time::Duration>,
explicit_target_plan: Option<crate::orchestration::target_resolution::ExplicitTargetPlan>,
change_failures_this_run: HashSet<String>,
run_fatal_abort: Option<String>,
}
#[cfg(test)]
mod tests;