cflx 0.6.322

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
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
//! Builder and initialization methods for [`super::ParallelExecutor`].
//!
//! This module provides the constructor and setter API for `ParallelExecutor`,
//! separating initialization concerns from execution logic.

use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;

use tokio::sync::{mpsc, Mutex};
use tokio_util::sync::CancellationToken;
use tracing::info;

use crate::ai_command_runner::{AiCommandRunner, RunCommandScope, SharedStaggerState};
use crate::config::OrchestratorConfig;
use crate::hooks::HookRunner;
use crate::vcs::{GitWorkspaceManager, VcsBackend, WorkspaceManager};

use super::{
    dedup::DiagnosticDeduplicationStore, FailedChangeTracker, ParallelEvent, ParallelExecutor,
    PostArchiveAction, SchedulerLifetime, DEFAULT_MAX_CONFLICT_RETRIES,
};

impl ParallelExecutor {
    /// Create a new parallel executor with automatic VCS detection
    pub fn new(
        repo_root: PathBuf,
        config: OrchestratorConfig,
        event_tx: Option<mpsc::Sender<ParallelEvent>>,
    ) -> Self {
        // Auto-detect VCS backend
        let vcs_backend = config.get_vcs_backend();
        Self::with_backend(repo_root, config, event_tx, vcs_backend)
    }

    /// Create a new parallel executor with a specific VCS backend and optional shared queue change
    /// timestamp
    pub fn with_backend_and_queue_state(
        repo_root: PathBuf,
        config: OrchestratorConfig,
        event_tx: Option<mpsc::Sender<ParallelEvent>>,
        vcs_backend: VcsBackend,
        shared_queue_change: Option<Arc<Mutex<Option<std::time::Instant>>>>,
    ) -> Self {
        Self::with_backend_and_queue_and_stagger(
            repo_root,
            config,
            event_tx,
            vcs_backend,
            shared_queue_change,
            None,
        )
    }

    /// Create a new parallel executor with a specific VCS backend, optional shared queue change
    /// timestamp, and optional shared stagger state
    pub fn with_backend_and_queue_and_stagger(
        repo_root: PathBuf,
        config: OrchestratorConfig,
        event_tx: Option<mpsc::Sender<ParallelEvent>>,
        vcs_backend: VcsBackend,
        shared_queue_change: Option<Arc<Mutex<Option<std::time::Instant>>>>,
        shared_stagger_state: Option<SharedStaggerState>,
    ) -> Self {
        // Resolve workspace base directory
        let base_dir = if let Some(configured_dir) = config.get_workspace_base_dir() {
            // User configured a specific directory
            PathBuf::from(configured_dir)
        } else {
            // Use OS-specific default workspace directory
            crate::config::defaults::default_workspace_base_dir(Some(&repo_root))
        };
        info!("Using workspace base directory: {:?}", base_dir);

        let max_concurrent = config.get_max_concurrent_workspaces();
        let apply_command = config
            .get_apply_command()
            .expect("apply_command must be configured before creating ParallelExecutor")
            .to_string();
        let archive_command = config
            .get_archive_command()
            .expect("archive_command must be configured before creating ParallelExecutor")
            .to_string();

        // Resolve the VCS backend (handle Auto)
        let resolved_backend = Self::resolve_backend(vcs_backend, &repo_root);
        info!("Using VCS backend: {:?}", resolved_backend);

        let workspace_manager: Box<dyn WorkspaceManager> = match resolved_backend {
            VcsBackend::Git | VcsBackend::Auto => Box::new(GitWorkspaceManager::new(
                base_dir,
                repo_root.clone(),
                max_concurrent,
                config.clone(),
            )),
        };

        let last_queue_change_at =
            shared_queue_change.unwrap_or_else(|| Arc::new(Mutex::new(None)));

        // Use provided shared stagger state or create a new one
        let shared_stagger_state =
            shared_stagger_state.unwrap_or_else(|| Arc::new(Mutex::new(None)));

        // Create the invocation scope and the shared AI command runner bound to
        // it. Every run-owned command surface clones this runner, so no run path
        // can build one from the stagger timestamp alone. A run owner that
        // already created a scope replaces this one through
        // `set_run_command_scope` before dispatch.
        let run_command_scope = RunCommandScope::new();
        let ai_runner = AiCommandRunner::for_run(
            &config,
            shared_stagger_state.clone(),
            run_command_scope.clone(),
        );

        Self {
            workspace_manager,
            config,
            apply_command,
            archive_command,
            event_tx,
            max_conflict_retries: DEFAULT_MAX_CONFLICT_RETRIES,
            repo_root,
            no_resume: false,
            explicit_retry: false,
            failed_tracker: FailedChangeTracker::new(),
            change_dependencies: HashMap::new(),
            resolve_wait_changes: HashSet::new(),
            reject_wait_changes: HashSet::new(),
            merge_wait_changes: HashSet::new(),
            dependency_blocker_fingerprints: HashMap::new(),
            force_recreate_worktree: HashSet::new(),
            hooks: None,
            cancel_token: None,
            last_queue_change_at,
            last_available_slots: None,
            analyzer_capacity_suppressed: false,
            pending_retry_bypass: HashSet::new(),
            dynamic_queue: None,
            ai_runner,
            run_command_scope,
            shared_stagger_state,
            apply_history: Arc::new(Mutex::new(crate::history::ApplyHistory::new())),
            archive_history: Arc::new(Mutex::new(crate::history::ArchiveHistory::new())),
            acceptance_history: Arc::new(Mutex::new(crate::history::AcceptanceHistory::new())),
            acceptance_tail_injected: Arc::new(Mutex::new(std::collections::HashMap::new())),
            apply_budget: crate::execution::apply::ApplyBudget::new(),
            // Sized from the same configured limit the workspace manager
            // enforces, so admission and worktree capacity can never disagree.
            lifecycle_slots: super::lifecycle_slots::LifecycleSlots::new(max_concurrent),
            manual_resolve_count: None,
            auto_resolve_count: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
            pending_merge_count: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
            scheduler_lifetime: SchedulerLifetime::Finite,
            // Default-off: an executor no run owner bound a stop request to can
            // never observe one, so its loop behaves exactly as before.
            graceful_stop: None,
            persistent_idle_latched: Arc::new(std::sync::atomic::AtomicBool::new(false)),
            persistent_idle_baseline: Arc::new(std::sync::Mutex::new(None)),
            post_archive_action: PostArchiveAction::MergeToBase,
            shared_orchestrator_state: None,
            last_dispatched_resolve_wait_changes: HashSet::new(),
            last_dispatched_reject_wait_changes: HashSet::new(),
            resolve_wait_retry_triggered: false,
            last_resolve_wait_base_dirty: None,
            diagnostic_dedup: DiagnosticDeduplicationStore::new(),
            // A fresh executor holds no prior analysis signature, so the first eligible
            // evaluation always analyzes instead of trusting a previous process's state.
            last_completed_analysis_input: None,
            next_analysis_signature_probe_at: None,
            analysis_retry_throttle: None,
            analysis_input_probe: None,
            // Default-off: an executor built without an explicit `-u` invocation
            // installs no upstream coordinator and therefore no new behavior.
            upstream: None,
            // Default: explicit targets are resolved before construction.
            explicit_target_plan: None,
            // Invocation-scoped terminal-report bookkeeping starts clean.
            change_failures_this_run: HashSet::new(),
            run_fatal_abort: None,
            // Default: the scheduler loop owns its merge-result channel.
            #[cfg(test)]
            merge_result_channel_override: None,
            #[cfg(test)]
            run_command_cleanup_budget_override: None,
        }
    }

    /// Create a new parallel executor with a specific VCS backend
    pub fn with_backend(
        repo_root: PathBuf,
        config: OrchestratorConfig,
        event_tx: Option<mpsc::Sender<ParallelEvent>>,
        vcs_backend: VcsBackend,
    ) -> Self {
        Self::with_backend_and_queue_state(repo_root, config, event_tx, vcs_backend, None)
    }

    /// Set the hook runner for executing hooks during parallel execution.
    #[allow(dead_code)] // Public API for future integration with CLI/TUI
    pub fn set_hooks(&mut self, hooks: HookRunner) {
        self.hooks = Some(Arc::new(hooks));
    }

    /// Set whether to disable automatic workspace resume.
    ///
    /// When `no_resume` is true, existing workspaces are always deleted
    /// and new ones are created. When false (default), existing workspaces
    /// are reused to resume interrupted work.
    pub fn set_no_resume(&mut self, no_resume: bool) {
        self.no_resume = no_resume;
    }

    pub fn set_explicit_retry(&mut self, explicit_retry: bool) {
        self.explicit_retry = explicit_retry;
    }

    /// Install deferred explicit-target classification.
    ///
    /// The plan is evaluated once, immediately after the initial upstream
    /// checkpoint and before any change-worktree creation or reuse registration.
    pub fn set_explicit_target_plan(
        &mut self,
        plan: crate::orchestration::target_resolution::ExplicitTargetPlan,
    ) {
        self.explicit_target_plan = Some(plan);
    }

    /// Install invocation-scoped upstream integration.
    ///
    /// Fetch, merge, verification, and push run as native Conflux operations
    /// through [`crate::upstream::git_ops::GitUpstreamOps`] and
    /// [`crate::upstream::verify::CommandVerifier`]; only bounded repair reaches
    /// the AI command harness through the existing `resolve_command` runner.
    pub fn set_upstream_integration(&mut self, runtime: crate::upstream::UpstreamRuntime) {
        use crate::upstream::coordinator::UpstreamCoordinator;
        use crate::upstream::git_ops::GitUpstreamOps;
        use crate::upstream::repair::ResolveCommandRepairAgent;
        use crate::upstream::verify::CommandVerifier;

        let coordinator = UpstreamCoordinator::new(
            runtime.config.clone(),
            runtime.branch.clone(),
            Arc::new(GitUpstreamOps::new(self.repo_root.clone())),
            Arc::new(CommandVerifier::new(
                runtime.config.verify_command.clone(),
                self.repo_root.clone(),
            )),
            Arc::new(ResolveCommandRepairAgent::new(
                self.config.clone(),
                self.repo_root.clone(),
                self.max_conflict_retries,
                self.ai_runner.clone(),
            )),
            Arc::new(
                crate::parallel::upstream_bridge::EventUpstreamObserver::new(self.event_tx.clone()),
            ),
        );
        self.upstream = Some(Arc::new(Mutex::new(coordinator)));
    }

    /// Whether upstream integration is installed for this executor.
    ///
    /// Used by default-off and propagation coverage; the executable itself never
    /// needs to ask, because every upstream entry point short-circuits already.
    #[allow(dead_code)]
    pub fn has_upstream_integration(&self) -> bool {
        self.upstream.is_some()
    }

    /// Override how dependency-analysis signature material is probed.
    ///
    /// Production leaves this unset and probes the real repository. Tests inject a double so
    /// unchanged-input suppression can be exercised deterministically, including probe counts
    /// and probe failures, without depending on VCS subprocesses.
    #[cfg(test)]
    pub(crate) fn set_analysis_input_probe(
        &mut self,
        probe: Arc<dyn crate::parallel::analysis_signature::AnalysisInputProbe>,
    ) {
        self.analysis_input_probe = Some(probe);
    }

    /// Replace the workspace manager with a double (tests only).
    ///
    /// This is what lets effective dependency-base resolution be verified without a real VCS:
    /// branch selection and ref revision lookup both go through the workspace manager, so a
    /// recording double can prove which one the scheduler asked for.
    #[cfg(test)]
    pub(crate) fn set_workspace_manager(&mut self, workspace_manager: Box<dyn WorkspaceManager>) {
        self.workspace_manager = workspace_manager;
    }

    /// Set shared reducer state for scheduler-owned resolve retry intent.
    pub fn set_shared_orchestrator_state(
        &mut self,
        shared_state: Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
    ) {
        self.shared_orchestrator_state = Some(shared_state);
    }

    /// Set the service-owned reducer only when the caller did not already provide one.
    pub(crate) fn ensure_shared_orchestrator_state(
        &mut self,
        shared_state: Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
    ) {
        if self.shared_orchestrator_state.is_none() {
            self.shared_orchestrator_state = Some(shared_state);
        }
    }

    /// Set the cancellation token for force stop cleanup.
    pub fn set_cancel_token(&mut self, cancel_token: CancellationToken) {
        self.cancel_token = Some(cancel_token);
    }

    /// The invocation scope owning every AI command this executor launches.
    #[allow(dead_code)] // Read by scope-ownership coverage; the loop uses the field directly.
    pub fn run_command_scope(&self) -> &RunCommandScope {
        &self.run_command_scope
    }

    /// The shared runner every run command surface clones.
    #[cfg(test)]
    pub(crate) fn ai_runner_for_test(&self) -> &AiCommandRunner {
        &self.ai_runner
    }

    /// Shorten the run-command cleanup barrier so scheduler-level coverage of
    /// an unproven cleanup stays inside the default-suite time budget.
    #[cfg(test)]
    pub(crate) fn set_cancellation_cleanup_budget_for_test(&mut self, budget: std::time::Duration) {
        self.run_command_cleanup_budget_override = Some(budget);
    }

    /// Adopt the run owner's scope in place of the one built here.
    ///
    /// Rebinding the shared runner is what keeps analyze, Apply, Archive,
    /// Acceptance, cleanup review, rejection review, conflict resolve, and
    /// upstream repair on a single scope identity: they all clone this runner.
    pub fn set_run_command_scope(&mut self, scope: RunCommandScope) {
        self.ai_runner.set_run_command_scope(scope.clone());
        self.run_command_scope = scope;
    }

    pub fn set_post_archive_action(&mut self, action: PostArchiveAction) {
        self.post_archive_action = action;
    }

    /// The terminal action this executor takes after a successful archive.
    ///
    /// Read-only observation so coverage can prove a frontend's configured
    /// action actually reached the executor that frontend built, instead of
    /// inferring it from a downstream side effect.
    #[cfg(test)]
    pub(crate) fn post_archive_action_for_test(&self) -> &PostArchiveAction {
        &self.post_archive_action
    }

    /// Set the dynamic queue for runtime change additions (TUI mode).
    pub fn set_dynamic_queue(&mut self, dynamic_queue: Arc<crate::tui::queue::DynamicQueue>) {
        self.dynamic_queue = Some(dynamic_queue);
    }

    /// Configure scheduler lifetime policy.
    pub fn set_scheduler_lifetime(&mut self, lifetime: SchedulerLifetime) {
        self.scheduler_lifetime = lifetime;
    }

    /// Keep scheduler alive until explicit stop (loop-based frontends).
    pub fn set_persistent_lifetime(&mut self) {
        self.set_scheduler_lifetime(SchedulerLifetime::Persistent);
    }

    /// The lifetime policy this executor was configured with.
    ///
    /// A finite CLI run and a persistent loop-based frontend differ only by
    /// this field, so coverage reads it directly rather than waiting on a
    /// scheduler that would never return.
    #[cfg(test)]
    pub(crate) fn scheduler_lifetime_for_test(&self) -> SchedulerLifetime {
        self.scheduler_lifetime
    }

    /// Bind the run owner's pending graceful-stop request.
    ///
    /// The same flag shared run control writes through the scheduler port, so a
    /// stop the operator already had accepted is visible to the loop that has to
    /// honour it. Nothing else about the request is duplicated here: the flag is
    /// read, never written.
    pub fn set_graceful_stop_flag(&mut self, graceful_stop: Arc<std::sync::atomic::AtomicBool>) {
        self.graceful_stop = Some(graceful_stop);
    }

    /// The run owner's pending graceful-stop request bound to this executor.
    ///
    /// Coverage compares the returned handle by pointer identity, which is the
    /// only way to prove the executor observes the *same* flag the run owner
    /// writes through rather than an equal-valued copy.
    #[cfg(test)]
    pub(crate) fn graceful_stop_flag_for_test(
        &self,
    ) -> Option<&Arc<std::sync::atomic::AtomicBool>> {
        self.graceful_stop.as_ref()
    }

    /// Set the manual resolve counter for tracking active manual resolve operations (TUI mode).
    /// This allows manual resolves to consume parallel execution slots.
    pub fn set_manual_resolve_counter(&mut self, counter: Arc<std::sync::atomic::AtomicUsize>) {
        self.manual_resolve_count = Some(counter);
    }

    /// Get a clone of the automatic resolve counter for testing or external tracking.
    #[cfg(test)]
    pub fn get_auto_resolve_counter(&self) -> Arc<std::sync::atomic::AtomicUsize> {
        self.auto_resolve_count.clone()
    }

    /// The configured concurrency limit this executor's scheduler will enforce.
    ///
    /// The same value the run loop passes as `max_parallelism`, so a test can
    /// occupy exactly as many lifecycle slots as the run has.
    #[cfg(test)]
    pub fn configured_max_concurrent(&self) -> usize {
        self.workspace_manager.max_concurrent()
    }

    /// Get the VCS backend type
    #[allow(dead_code)] // Public API for external callers
    pub fn backend_type(&self) -> VcsBackend {
        self.workspace_manager.backend_type()
    }

    /// Check if VCS is available for parallel execution
    #[allow(dead_code)] // Public API, used via ParallelRunService
    pub async fn check_vcs_available(&self) -> crate::error::Result<bool> {
        self.workspace_manager
            .check_available()
            .await
            .map_err(Into::into)
    }

    /// Resolve VCS backend (convert Auto to concrete backend)
    pub(super) fn resolve_backend(backend: VcsBackend, _repo_root: &Path) -> VcsBackend {
        match backend {
            VcsBackend::Auto => VcsBackend::Git,
            other => other,
        }
    }
}