cflx 0.6.327

Conflux – a spec-driven parallel coding orchestrator that runs AI agents on git worktrees
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
//! Parallel execution coordinator for VCS workspace-based parallel change application.
//!
//! This module is the entry point for the parallel execution subsystem. It defines the
//! shared state container (`ParallelExecutor`) and re-exports the public API.
//!
//! Implementation is split into focused submodules:
//! - `builder`: construction and initialization
//! - `queue_state`: queue management and dispatch coordination
//! - `executor`: apply/acceptance/archive execution in workspaces
//! - `merge`: branch merge and conflict resolution
//! - `dispatch`: per-change dispatch logic
//! - `orchestration`: order-based re-analysis scheduler loop

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;

// Re-export unified event type as ParallelEvent for backward compatibility.
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,
};
// Only test code re-exports this helper through `crate::parallel`; production callers use
// `super::types` directly.
#[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;

/// Defines when the parallel scheduler should terminate.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SchedulerLifetime {
    /// Finite execution (CLI `run`): stop once no queued/in-flight work remains.
    Finite,
    /// Persistent execution (loop-based/TUI): keep waiting for queue notifications until stopped.
    Persistent,
}

/// Terminal report of one scheduler invocation that returned on its own.
///
/// A scheduler failure is the `Err` half of the run result, so this enum only
/// describes returns that were *not* failures. `CompletedWithErrors` exists so a
/// finite run that drained while change-local failures remain unresolved can
/// never be reported as plain success.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SchedulerRunReport {
    /// Every eligible change reached its terminal state without failure.
    Completed,
    /// Eligible work drained, but one or more changes ended in a change-local
    /// failure whose evidence is preserved for explicit retry.
    CompletedWithErrors,
    /// Operator cancellation stopped the run.
    Stopped,
    /// A finite run ended with queued work that is still blocked or stalled.
    ///
    /// Nothing drained: the remaining candidates are held by a failed
    /// dependency, an unresolved blocker, or another wait lane. This is neither
    /// a success nor an execution failure, and it must never be announced as
    /// completion.
    BlockedOrStalled,
}

impl SchedulerRunReport {
    /// Whether this run ended without completing every eligible change.
    ///
    /// Boundaries use this to withhold a success announcement: both
    /// change-local failures and blocked/stalled remainders leave work the
    /// operator still owns.
    pub fn is_incomplete(self) -> bool {
        matches!(self, Self::CompletedWithErrors | Self::BlockedOrStalled)
    }
}

/// Terminal action after a change has been archived.
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum PostArchiveAction {
    #[default]
    MergeToBase,
    PushToRemote {
        remote: String,
    },
}

/// Global lock for serializing all merge/resolve operations to base branch.
///
/// This ensures that only one merge operation can modify the base branch
/// at any given time, regardless of which `ParallelExecutor` instance initiates it.
static GLOBAL_MERGE_LOCK: OnceLock<Mutex<()>> = OnceLock::new();

/// Runtime-only set of post-archive merge tasks currently owning a change.
///
/// This is intentionally not durable workflow state. It only prevents the live
/// scheduler from redispatching the same archived dirty workspace while an
/// already-spawned post-archive merge task is still deriving the authoritative
/// outcome from repository/workspace git state.
static ACTIVE_POST_ARCHIVE_MERGES: OnceLock<StdMutex<HashSet<String>>> = OnceLock::new();

/// Get the global merge lock, initializing it if necessary.
fn global_merge_lock() -> &'static Mutex<()> {
    GLOBAL_MERGE_LOCK.get_or_init(|| Mutex::new(()))
}

/// Test-only serialization mutex for suites that manipulate the
/// [`global_merge_lock`] directly (acquiring it to exercise contention paths).
///
/// Such tests race against each other when the harness runs them concurrently:
/// one test holding the global lock makes another's `try_lock` fail. All tests
/// that touch the global merge lock must hold this mutex first so they run
/// serially regardless of which module they live in.
#[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()))
}

/// Parallel executor for running changes in VCS workspaces (git worktrees today).
///
/// All execution logic lives in submodules as `impl ParallelExecutor` blocks.
pub struct ParallelExecutor {
    /// Workspace manager (VCS-agnostic)
    workspace_manager: Box<dyn WorkspaceManager>,
    /// Configuration (used for AgentRunner and resolve operations)
    config: OrchestratorConfig,
    /// Apply command template
    apply_command: String,
    /// Archive command template
    archive_command: String,
    /// Event sender
    event_tx: Option<mpsc::Sender<ParallelEvent>>,
    /// Maximum retries for conflict resolution
    max_conflict_retries: u32,
    /// Repository root path for archive operations
    repo_root: PathBuf,
    /// Disable automatic workspace resume (always create new workspaces)
    no_resume: bool,
    /// Release the automatic repair budget for an explicitly retried change.
    explicit_retry: bool,
    /// Tracker for failed changes to enable skipping dependent changes
    failed_tracker: FailedChangeTracker,
    /// Change-level dependencies (change_id -> dependency ids)
    change_dependencies: HashMap<String, Vec<String>>,
    /// Changes waiting for auto-resumable resolve retry (ResolveWait)
    resolve_wait_changes: HashSet<String>,
    /// Changes waiting for rejection review to run once the base-mutating lane is free (RejectWait)
    reject_wait_changes: HashSet<String>,
    /// Changes waiting for manual user intervention before merge can continue (MergeWait)
    merge_wait_changes: HashSet<String>,
    /// Last emitted dependency blocker fingerprint per change.
    ///
    /// This runtime-only observability state drives blocked/resolved diagnostics and
    /// worktree recreation after dependency resolution. It MUST NOT decide dispatch
    /// eligibility; `select_changes_for_dispatch` still derives executable work from
    /// the current analysis result and repository/workspace state on every pass.
    dependency_blocker_fingerprints: HashMap<String, DependencyBlockerFingerprint>,
    /// Changes that need forced worktree recreation (dependency just resolved)
    force_recreate_worktree: HashSet<String>,
    /// Hook runner for executing hooks (optional)
    hooks: Option<Arc<HookRunner>>,
    /// Cancellation token for force stop cleanup
    cancel_token: Option<CancellationToken>,
    /// Last queue change timestamp for debouncing re-analysis
    last_queue_change_at: Arc<Mutex<Option<std::time::Instant>>>,
    /// Last observed number of available execution slots.
    ///
    /// Used to bypass queue-edit debounce when capacity recovers from zero to positive,
    /// so queued changes dispatch immediately after a running task or manual resolve frees a slot.
    last_available_slots: Option<usize>,
    /// Whether the most recent reanalysis pass skipped the analyzer for lack of
    /// capacity alone.
    ///
    /// Per-pass and process-local. It exists so the loop step that owns
    /// one-shot trigger lifetime can tell "this edge was evaluated" from "this
    /// edge never got an evaluation", without every direct caller of the pass
    /// having to thread a richer return value through.
    analyzer_capacity_suppressed: bool,
    /// Retried targets whose explicit-retry edge has left the queue but whose
    /// analysis-bypass authority no eligible evaluation has spent yet.
    ///
    /// Draining the queue at Step 0 is not the same event as consuming the edge.
    /// A pass can end before the dependency-analysis evaluation it authorized —
    /// cancellation, an incomplete reducer view, an early break, a candidate list
    /// that is momentarily empty — and the authority must outlive that pass
    /// rather than be discarded with it. It is spent when the analyzer really
    /// runs, and only then.
    ///
    /// Ephemeral process-local control state: it is never persisted and a restart
    /// recomputes routing from the workspace alone.
    pending_retry_bypass: HashSet<String>,
    /// Dynamic queue for runtime change additions (TUI mode)
    dynamic_queue: Option<Arc<crate::tui::queue::DynamicQueue>>,
    /// Shared AI command runner for stagger coordination
    ai_runner: AiCommandRunner,
    /// Invocation-scoped ownership of every AI command this run launches.
    ///
    /// Ephemeral and process-local: it is created for one run, discarded when
    /// that run ends, and never persisted or consulted for restart routing.
    run_command_scope: RunCommandScope,
    /// Shared stagger state for resolve operations
    #[allow(dead_code)]
    shared_stagger_state: SharedStaggerState,
    /// History of apply attempts per change for context injection
    apply_history: Arc<Mutex<crate::history::ApplyHistory>>,
    /// History of archive attempts per change for context injection
    archive_history: Arc<Mutex<crate::history::ArchiveHistory>>,
    /// History of acceptance attempts per change for context injection
    acceptance_history: Arc<Mutex<crate::history::AcceptanceHistory>>,
    /// Tracks which changes have had acceptance tail injected (to prevent re-injection)
    acceptance_tail_injected: Arc<Mutex<std::collections::HashMap<String, bool>>>,
    /// The sole per-change `max_iterations` Apply-dispatch budget owner for this
    /// parallel run.
    ///
    /// Shared by every dispatched workspace task, so a change that re-enters
    /// Apply after an Acceptance FAIL keeps counting from where it left off.
    /// Active-run memory only: a restart starts from zero and re-derives routing
    /// from workspace and Git evidence.
    apply_budget: crate::execution::apply::ApplyBudget,
    /// The authoritative lifecycle-slot membership for this scheduler.
    ///
    /// One admitted change owns one slot continuously, from just before
    /// workspace preparation until repository-visible settlement — merged,
    /// terminal error, rejected, explicit dequeue, or the configured
    /// push/publication terminal settlement. Apply/acceptance/archive
    /// completion, spawning a background merge, detecting a conflict, starting
    /// a resolve, and entering `merge wait` all *transfer* the slot rather than
    /// release it, which is what keeps a queued change from taking capacity that
    /// is still owned.
    ///
    /// Ephemeral process-local state: nothing is persisted, and a restart
    /// reconstructs occupancy from workspace, Git, and reducer evidence
    /// (`openspec/CONSTITUTION.md`, law 1).
    lifecycle_slots: lifecycle_slots::LifecycleSlots,
    /// Counter for active manual resolve operations (TUI mode)
    ///
    /// Observability and base-lane serialization only. It must never subtract
    /// dispatch capacity: a manual resolve runs inside the lifecycle slot its
    /// change already owns, so subtracting it as well would double-count one
    /// admitted change.
    manual_resolve_count: Option<Arc<std::sync::atomic::AtomicUsize>>,
    /// Counter for active automatic resolve operations
    ///
    /// Observability and base-lane serialization only, for the same reason as
    /// `manual_resolve_count`: an automatic resolve runs inside the background
    /// merge of a change that already owns its lifecycle slot.
    auto_resolve_count: Arc<std::sync::atomic::AtomicUsize>,
    /// Counter for background merge tasks that have been spawned but not yet handled by scheduler.
    pending_merge_count: Arc<std::sync::atomic::AtomicUsize>,
    /// Scheduler lifetime policy (finite for CLI run, persistent for loop-based frontends).
    scheduler_lifetime: SchedulerLifetime,
    /// The run owner's pending graceful-stop request, when one is bound.
    ///
    /// Shared with the run supervisor that shared run control sets it through, so
    /// the scheduler reads the *same* request the operator's accepted stop
    /// recorded rather than a copy taken at launch. Process-local and
    /// invocation-scoped: a launch clears it, and a restart begins with no
    /// pending stop.
    graceful_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
    /// Edge-trigger latch for the persistent-idle transition.
    ///
    /// Set when the scheduler emits its one idle event for an episode and
    /// cleared only when admitted work actually begins, so repeated loop
    /// evaluations and wake notifications that start nothing cannot re-emit.
    /// Shared with the dispatch paths that own those two rearm points, and
    /// invocation-scoped like every other runtime latch here: a restart begins
    /// with no open episode.
    persistent_idle_latched: Arc<std::sync::atomic::AtomicBool>,
    /// Reducer-visible queue intent this idle episode already parked on.
    ///
    /// `None` while no episode is open. Recorded from the same coherent view the
    /// park was decided from, so a later pass can answer the only question the
    /// rearm needs — "is there intent this episode has not already evaluated?" —
    /// as a *level* observation rather than from an individual reducer command's
    /// outcome. That is what closes the prepare/commit race where a concurrent
    /// queue addition makes an accepted Start's own `AddToQueue` a no-op while
    /// coherent queue intent plainly exists.
    ///
    /// Comparing against the episode baseline rather than against emptiness is
    /// what keeps a blocked-only park quiet: its rows were already there, so a
    /// generic wake observes nothing new and emits no second idle edge.
    persistent_idle_baseline: Arc<std::sync::Mutex<Option<HashSet<String>>>>,
    /// Post-archive terminal action.
    post_archive_action: PostArchiveAction,
    /// Optional reducer shared state used for scheduler-owned resolve/merge retry intent.
    shared_orchestrator_state: Option<Arc<RwLock<OrchestratorState>>>,
    /// Last resolve-wait snapshot that was dispatched via scheduler retry.
    last_dispatched_resolve_wait_changes: HashSet<String>,
    /// Last reject-wait snapshot that was dispatched via scheduler retry.
    last_dispatched_reject_wait_changes: HashSet<String>,
    /// One-shot flag to allow retry dispatch on explicit wake/completion triggers.
    resolve_wait_retry_triggered: bool,
    /// Last scheduler-observed base dirtiness while reducer-owned base-lane waiters existed.
    ///
    /// This is runtime-only dedupe state. It is not durable workflow state; retry routing is
    /// recalculated from reducer state plus current base git/workspace state on each scheduler run.
    last_resolve_wait_base_dirty: Option<bool>,
    /// Runtime-only observability dedupe for operator-visible diagnostics.
    ///
    /// This state is intentionally in-memory and MUST NOT participate in scheduling decisions.
    diagnostic_dedup: DiagnosticDeduplicationStore<DiagnosticDeduplicationKey>,
    /// Last dependency-analysis input that completed with a usable result.
    ///
    /// Ordinary timer-driven analysis is skipped while the current input still matches this
    /// record, which is what stops an unchanged scheduler state from relaunching costly
    /// analysis agents forever. It is deliberately process-local: it is never persisted, so a
    /// restart always performs an initial analysis, and it never authorizes dispatch — the
    /// previous `AnalysisResult` is not retained.
    last_completed_analysis_input: Option<CompletedAnalysisInput>,
    /// Earliest instant at which a suppressed input may probe repository-visible signature
    /// material again.
    ///
    /// The scheduler wakes every 500 ms; without this bound each wake would re-read proposal
    /// files and spawn a VCS revision subprocess just to confirm nothing changed.
    next_analysis_signature_probe_at: Option<tokio::time::Instant>,
    /// Bounded retry deadline for an ordinary timer evaluation that established no completed
    /// input.
    ///
    /// This is deliberately separate from `last_completed_analysis_input`: a failed signature
    /// probe or an unusable analyzer result proves nothing was analyzed, so it must never be
    /// recorded as a completion. It only rate-limits the retry, so fail-open stays open without
    /// becoming a 500 ms probe/analyzer loop.
    analysis_retry_throttle: Option<BoundedAnalysisRetry>,
    /// Invocation-scoped upstream integration coordinator.
    ///
    /// `None` is the hard compatibility boundary: no checkpoint runs, no
    /// upstream fetch/merge/verification/push happens, and no upstream lifecycle
    /// evidence is emitted. It is installed only from an explicit `cflx run`
    /// `-u`/`--integrate-upstream` invocation, never from persistent config.
    upstream: Option<Arc<Mutex<crate::upstream::UpstreamCoordinator>>>,
    /// Override for dependency-analysis signature probing.
    ///
    /// `None` uses the real repository probe (VCS revision plus
    /// `openspec/changes/<id>/proposal.md` content). Tests inject a deterministic double so
    /// suppression coverage does not depend on real VCS subprocesses, and so probe counts and
    /// probe failures can be observed directly.
    analysis_input_probe: Option<Arc<dyn AnalysisInputProbe>>,
    /// Override for the scheduler loop's background merge/base-lane result channel.
    ///
    /// `None` is the ordinary path: the loop creates its own channel per run. Tests
    /// inject a double so the cancellation cleanup barrier can be driven
    /// deterministically — a pending merge can be left outstanding across cancellation
    /// and completed on demand — without spawning real detached merge tasks.
    #[cfg(test)]
    merge_result_channel_override: Option<(mpsc::Sender<MergeResult>, mpsc::Receiver<MergeResult>)>,
    /// Override for the bounded run-command cleanup barrier.
    ///
    /// `None` uses the fixed 30-second production budget. Tests shorten it so
    /// an unproven-cleanup path can be driven inside the default-suite time
    /// budget; there is no user-facing configuration surface for it.
    #[cfg(test)]
    run_command_cleanup_budget_override: Option<std::time::Duration>,
    /// Explicit-target classification deferred to the post-checkpoint boundary.
    ///
    /// `None` is the ordinary path: explicit targets were already classified
    /// against the captured local base before the executor was constructed. A
    /// real `-u` run installs a plan so classification reads the cumulative base
    /// produced by the mandatory initial upstream checkpoint instead.
    explicit_target_plan: Option<crate::orchestration::target_resolution::ExplicitTargetPlan>,
    /// Changes that ended this invocation in a change-local base-lane failure.
    ///
    /// Invocation-scoped in-memory bookkeeping only: it decides the terminal
    /// report of *this* scheduler run and nothing else. It is never persisted,
    /// so a restart recomputes the next action from workspace and Git evidence
    /// alone, as `openspec/CONSTITUTION.md` requires.
    change_failures_this_run: HashSet<String>,
    /// Detail of the run-fatal outcome that requested a scheduler abort.
    ///
    /// `Some` means the single global Error was already emitted by the
    /// queue/orchestration owner and the loop owes a bounded drain followed by
    /// scheduler failure. Also invocation-scoped and never persisted.
    run_fatal_abort: Option<String>,
}

#[cfg(test)]
mod tests;