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
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
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
//! Tests for manual resolve counter integration with parallel execution.

use super::support::{create_test_config, TestWorkspaceManager};
use crate::events::ExecutionEvent;
use crate::openspec::{Change, ProposalMetadata};
use crate::parallel::cleanup::WorkspaceCleanupGuard;
use crate::parallel::dynamic_queue::ReanalysisReason;
use crate::parallel::lifecycle_slots::SlotPhase;
use crate::parallel::queue_state::ReanalysisDispatchContext;
use crate::parallel::{ParallelExecutor, SchedulerLifetime, WorkspaceResult};
use crate::tui::queue::DynamicQueue;
use crate::vcs::VcsBackend;
use std::collections::{HashMap, HashSet};
use std::process::Command;
use std::sync::{atomic::AtomicUsize, Arc};
use tempfile::TempDir;
use tokio::task::JoinSet;

/// Wire the reducer the way production does before a dynamic queue push.
///
/// `OperatorCommandService::add_to_queue` applies `AddToQueue` first and only
/// then publishes the wake-up hint, and ingestion validates the hint against
/// that reducer intent — never against the catalog alone — so a dynamic-queue
/// test must record the same intent to exercise the accepted path.
fn shared_state_with_queue_intent(
    change_ids: &[&str],
) -> Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>> {
    use crate::orchestration::state::{OrchestratorState, ReducerCommand};

    let mut state = OrchestratorState::new(change_ids.iter().map(|id| id.to_string()).collect(), 1);
    for change_id in change_ids {
        state.apply_command(ReducerCommand::AddToQueue(change_id.to_string()));
    }
    Arc::new(tokio::sync::RwLock::new(state))
}

#[tokio::test]
async fn test_manual_resolve_counter_tracks_active_resolves() {
    // Create a temporary directory for the test repository
    let temp_dir = TempDir::new().unwrap();
    let repo_root = temp_dir.path().to_path_buf();

    // Create a basic config
    let config = create_test_config();

    // Create a manual resolve counter
    let manual_resolve_counter = Arc::new(AtomicUsize::new(0));

    // Create a ParallelExecutor with max_concurrent = 4
    let mut executor = ParallelExecutor::new(repo_root.clone(), config.clone(), None);

    // Set the manual resolve counter
    executor.set_manual_resolve_counter(manual_resolve_counter.clone());

    // Initially, counter should be 0
    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        0,
        "Manual resolve counter should start at 0"
    );

    // Simulate a manual resolve starting (TUI would increment this)
    manual_resolve_counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);

    // Verify counter is now 1
    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        1,
        "Manual resolve counter should be 1 after increment"
    );

    // The counter is observability only. Dispatch capacity comes from
    // lifecycle-slot membership, and a manual resolve runs inside the slot its
    // change already owns — see `parallel::tests::lifecycle_slot_ownership`.

    // Simulate resolve completing
    manual_resolve_counter.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);

    // Counter should be back to 0
    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        0,
        "Manual resolve counter should return to 0 after completion"
    );
}

#[tokio::test]
async fn test_multiple_manual_resolves_are_tracked_independently() {
    // Create a temporary directory for the test repository
    let temp_dir = TempDir::new().unwrap();
    let repo_root = temp_dir.path().to_path_buf();

    // Create a basic config
    let config = create_test_config();

    // Create a manual resolve counter
    let manual_resolve_counter = Arc::new(AtomicUsize::new(0));

    // Create a ParallelExecutor
    let mut executor = ParallelExecutor::new(repo_root.clone(), config.clone(), None);
    executor.set_manual_resolve_counter(manual_resolve_counter.clone());

    // Simulate 2 concurrent manual resolves
    manual_resolve_counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
    manual_resolve_counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);

    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        2,
        "Manual resolve counter should be 2 for concurrent resolves"
    );

    // Two resolves means two admitted changes, each already counted once by the
    // lifecycle slot it owns; the counter itself subtracts no capacity.

    // Simulate first resolve completing
    manual_resolve_counter.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        1,
        "Manual resolve counter should be 1 after one completes"
    );

    // Simulate second resolve completing
    manual_resolve_counter.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
    assert_eq!(
        manual_resolve_counter.load(std::sync::atomic::Ordering::SeqCst),
        0,
        "Manual resolve counter should be 0 after all complete"
    );
}

#[tokio::test]
async fn test_manual_resolve_completion_notifies_scheduler() {
    let queue = DynamicQueue::new();
    let notified = queue.notified();

    queue.notify_scheduler();

    tokio::time::timeout(std::time::Duration::from_secs(1), notified)
        .await
        .expect("scheduler notification should wake waiters");
}

#[test]
fn test_manual_resolve_counter_is_thread_safe() {
    // Create a counter
    let counter = Arc::new(AtomicUsize::new(0));

    // Spawn multiple threads to increment/decrement concurrently
    let handles: Vec<_> = (0..10)
        .map(|_| {
            let counter_clone = counter.clone();
            std::thread::spawn(move || {
                for _ in 0..100 {
                    counter_clone.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
                    counter_clone.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
                }
            })
        })
        .collect();

    // Wait for all threads to complete
    for handle in handles {
        handle.join().unwrap();
    }

    // Counter should be back to 0
    assert_eq!(
        counter.load(std::sync::atomic::Ordering::SeqCst),
        0,
        "Counter should be 0 after concurrent increment/decrement operations"
    );
}

fn test_change(id: &str) -> Change {
    Change {
        id: id.to_string(),
        completed_tasks: 0,
        total_tasks: 1,
        last_modified: String::new(),
        dependencies: Vec::new(),
        metadata: ProposalMetadata::default(),
    }
}

fn create_active_change_fixture(repo_root: &std::path::Path, change_id: &str) {
    let change_dir = repo_root.join("openspec").join("changes").join(change_id);
    std::fs::create_dir_all(&change_dir).expect("create synthetic OpenSpec change directory");
    std::fs::write(
        change_dir.join("proposal.md"),
        format!("# Synthetic Change {change_id}\n\n## Why\n\nTest fixture.\n"),
    )
    .expect("write synthetic proposal");
    std::fs::write(
        change_dir.join("tasks.md"),
        "# Tasks\n\n- [ ] Synthetic fixture task\n",
    )
    .expect("write synthetic tasks");
}

fn init_minimal_git_repo(repo_root: &std::path::Path) {
    for args in [
        vec!["init", "-b", "main"],
        vec!["config", "user.email", "test@example.com"],
        vec!["config", "user.name", "Test User"],
    ] {
        let output = Command::new("git")
            .args(args)
            .current_dir(repo_root)
            .output()
            .expect("run git setup command");
        assert!(
            output.status.success(),
            "git setup command failed: {}",
            String::from_utf8_lossy(&output.stderr)
        );
    }
    std::fs::write(repo_root.join("README.md"), "base\n").expect("write base file");
    for args in [vec!["add", "-A"], vec!["commit", "-m", "Base"]] {
        let output = Command::new("git")
            .args(args)
            .current_dir(repo_root)
            .output()
            .expect("run git commit command");
        assert!(
            output.status.success(),
            "git commit command failed: {}",
            String::from_utf8_lossy(&output.stderr)
        );
    }
}

fn analysis_result<'a>(
    changes: &'a [Change],
    _in_flight: &'a [String],
    _iteration: u32,
) -> std::pin::Pin<
    Box<dyn std::future::Future<Output = crate::analyzer::AnalysisOutcome> + Send + 'a>,
> {
    let order = changes.iter().map(|change| change.id.clone()).collect();
    Box::pin(async move {
        crate::analyzer::AnalysisResult {
            order,
            dependencies: HashMap::new(),
            groups: None,
        }
        .into()
    })
}

/// Zero capacity gates the analyzer as well as dispatch.
///
/// An analysis order nothing can consume is wasted work, and recording that the
/// input was analysed is what would suppress the evaluation capacity recovery
/// depends on. Queue classification, reducer reconciliation, and diagnostics all
/// still run — they are just above this gate.
#[tokio::test]
async fn test_manual_resolve_zero_capacity_gates_analysis_and_apply_dispatch() {
    let temp_dir = TempDir::new().unwrap();
    let (tx, mut rx) = tokio::sync::mpsc::channel(16);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    // A manual resolve runs inside the lifecycle slot its change already owns,
    // so the occupancy that gates this pass is that admitted change.
    let manual_resolve_counter = Arc::new(AtomicUsize::new(1));
    executor.set_manual_resolve_counter(manual_resolve_counter);
    executor
        .lifecycle_slots
        .occupy_now("resolving-change", SlotPhase::Retained)
        .await;

    let mut queued = vec![test_change("queued-apply")];
    let mut in_flight = HashSet::new();
    let mut join_set: JoinSet<WorkspaceResult> = JoinSet::new();
    let mut cleanup_guard =
        WorkspaceCleanupGuard::new(VcsBackend::Git, temp_dir.path().to_path_buf());

    let (should_break, iteration) = executor
        .perform_reanalysis_and_dispatch(ReanalysisDispatchContext {
            queued: &mut queued,
            in_flight: &mut in_flight,
            max_parallelism: 1,
            iteration: 1,
            reanalysis_reason: ReanalysisReason::ResolveCompletion,
            analyzer: &analysis_result,
            join_set: &mut join_set,
            cleanup_guard: &mut cleanup_guard,
            work_snapshot: None,
        })
        .await
        .expect("re-analysis should not fail");

    assert!(!should_break);
    assert_eq!(
        iteration, 1,
        "suppressed dispatch must not advance iteration"
    );
    assert!(
        in_flight.is_empty(),
        "zero capacity must not start apply work"
    );
    assert_eq!(
        queued.len(),
        1,
        "queued change remains pending until capacity recovers"
    );
    assert!(
        join_set.is_empty(),
        "no workspace task should be spawned at zero capacity"
    );

    let mut saw_analysis_started = false;
    let mut saw_apply_started = false;
    while let Ok(event) = rx.try_recv() {
        match event {
            ExecutionEvent::AnalysisStarted { .. } => saw_analysis_started = true,
            ExecutionEvent::ApplyStarted { .. } => saw_apply_started = true,
            _ => {}
        }
    }

    assert!(
        !saw_analysis_started,
        "the expensive analyzer must not start while a manual resolve holds every slot"
    );
    assert!(
        !saw_apply_started,
        "ordinary apply must remain capacity-gated during active manual resolve"
    );
}

/// Repeated zero-capacity wakes stay inert rather than relaunching the analyzer.
#[tokio::test]
async fn repeated_capacity_zero_never_starts_analysis() {
    let temp_dir = TempDir::new().unwrap();
    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );

    let mut queued = vec![test_change("queued-apply")];
    let mut in_flight = HashSet::from(["active-apply".to_string()]);
    let mut join_set: JoinSet<WorkspaceResult> = JoinSet::new();
    let mut cleanup_guard =
        WorkspaceCleanupGuard::new(VcsBackend::Git, temp_dir.path().to_path_buf());

    for iteration in 1..=2 {
        let (should_break, returned_iteration) = executor
            .perform_reanalysis_and_dispatch(ReanalysisDispatchContext {
                queued: &mut queued,
                in_flight: &mut in_flight,
                max_parallelism: 1,
                iteration,
                reanalysis_reason: ReanalysisReason::ResolveCompletion,
                analyzer: &analysis_result,
                join_set: &mut join_set,
                cleanup_guard: &mut cleanup_guard,
                work_snapshot: None,
            })
            .await
            .expect("re-analysis should not fail");

        assert!(!should_break);
        assert_eq!(
            returned_iteration, iteration,
            "suppressed dispatch must not advance iteration"
        );
    }

    assert_eq!(
        queued.len(),
        1,
        "queued change remains pending while capacity is zero"
    );
    assert_eq!(
        in_flight.len(),
        1,
        "test must keep capacity at zero across repeated analysis iterations"
    );
    assert!(
        join_set.is_empty(),
        "no workspace task should be spawned at zero capacity"
    );

    let mut analysis_started_count = 0;
    let mut apply_started_count = 0;
    while let Ok(event) = rx.try_recv() {
        match event {
            ExecutionEvent::AnalysisStarted { .. } => analysis_started_count += 1,
            ExecutionEvent::ApplyStarted { .. } => apply_started_count += 1,
            _ => {}
        }
    }

    assert_eq!(
        analysis_started_count, 0,
        "a suppressed evaluation is not a distinct analysis attempt; saw {analysis_started_count}"
    );
    assert_eq!(
        apply_started_count, 0,
        "ordinary apply must remain capacity-gated"
    );
}

#[tokio::test]
async fn scheduler_loop_ingests_dynamic_queue_during_gated_manual_resolve() {
    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let seed_change_id = "synthetic-seed-gated";
    let synthetic_change_id = "synthetic-dynamic-gated-resolve";
    create_active_change_fixture(temp_dir.path(), seed_change_id);
    create_active_change_fixture(temp_dir.path(), synthetic_change_id);

    let (tx, mut rx) = tokio::sync::mpsc::channel(64);
    let dynamic_queue = Arc::new(DynamicQueue::new());

    let cancel_token = tokio_util::sync::CancellationToken::new();

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_cancel_token(cancel_token.clone());
    executor.set_dynamic_queue(dynamic_queue.clone());
    executor.set_scheduler_lifetime(SchedulerLifetime::Persistent);
    executor.set_manual_resolve_counter(Arc::new(AtomicUsize::new(1)));
    // The gate is admitted-change occupancy, held for the whole run: every
    // configured slot belongs to a change parked in manual resolution, so no
    // capacity can appear while the loop runs.
    for index in 0..executor.configured_max_concurrent() {
        executor
            .lifecycle_slots
            .occupy_now(&format!("gated-resolve-{index}"), SlotPhase::Retained)
            .await;
    }
    executor.set_shared_orchestrator_state(shared_state_with_queue_intent(&[
        seed_change_id,
        synthetic_change_id,
    ]));

    let scheduler_queue = dynamic_queue.clone();
    let scheduler = tokio::spawn(async move {
        executor
            .execute_with_order_based_reanalysis(vec![test_change(seed_change_id)], analysis_result)
            .await
    });
    scheduler_queue.push(synthetic_change_id.to_string()).await;

    let mut saw_dynamic_ingest = false;
    let mut saw_analysis_started = false;
    let mut saw_apply_started = false;
    let mut saw_capacity_diagnostic = false;
    let mut log_messages = Vec::new();

    tokio::time::timeout(std::time::Duration::from_millis(500), async {
        while !(saw_dynamic_ingest && saw_capacity_diagnostic) {
            match rx.recv().await {
                Some(ExecutionEvent::Log(entry))
                    if entry.message.contains(&format!(
                        "Dynamically added to parallel execution: {synthetic_change_id}"
                    )) =>
                {
                    saw_dynamic_ingest = true;
                }
                Some(ExecutionEvent::Log(entry))
                    if entry.message.contains("analysis_capacity_zero") =>
                {
                    saw_capacity_diagnostic = true;
                    log_messages.push(entry.message);
                }
                Some(ExecutionEvent::Log(entry)) => log_messages.push(entry.message),
                Some(ExecutionEvent::AnalysisStarted { .. }) => saw_analysis_started = true,
                Some(ExecutionEvent::ApplyStarted { .. }) => saw_apply_started = true,
                Some(_) => {}
                None => break,
            }
        }
    })
    .await
    .expect("scheduler loop should ingest and analyze bounded dynamic work");

    cancel_token.cancel();
    let _ = tokio::time::timeout(std::time::Duration::from_millis(500), scheduler)
        .await
        .expect("scheduler should stop after cancellation")
        .expect("scheduler task should not panic");

    assert!(
        saw_dynamic_ingest,
        "expected dynamic ingest log for {synthetic_change_id}; saw logs: {log_messages:?}"
    );
    assert!(
        !saw_analysis_started,
        "zero recalculated capacity must suppress the expensive analyzer too"
    );
    assert!(saw_capacity_diagnostic);
    assert!(
        !saw_apply_started,
        "zero recalculated capacity must suppress apply dispatch while gated resolve is active"
    );
}

#[tokio::test]
async fn persistent_scheduler_dynamic_queue_push_after_initial_analysis_bypasses_debounce() {
    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let seed_change_id = "synthetic-seed-running";
    let dynamic_change_id = "synthetic-running-dynamic-queue";
    create_active_change_fixture(temp_dir.path(), seed_change_id);
    create_active_change_fixture(temp_dir.path(), dynamic_change_id);

    let (tx, mut rx) = tokio::sync::mpsc::channel(64);
    let dynamic_queue = Arc::new(DynamicQueue::new());
    let cancel_token = tokio_util::sync::CancellationToken::new();

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_cancel_token(cancel_token.clone());
    executor.set_dynamic_queue(dynamic_queue.clone());
    executor.set_scheduler_lifetime(SchedulerLifetime::Persistent);
    // Only the seed carries queue intent at start. The dynamic change acquires
    // it at push time, in the production order: reducer intent first, wake-up
    // hint second.
    let shared = shared_state_with_queue_intent(&[seed_change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let scheduler = tokio::spawn(async move {
        executor
            .execute_with_order_based_reanalysis(vec![test_change(seed_change_id)], analysis_result)
            .await
    });

    tokio::time::timeout(std::time::Duration::from_millis(500), async {
        loop {
            match rx.recv().await {
                Some(ExecutionEvent::AnalysisStarted { attempt_id, .. })
                    if attempt_id.contains(seed_change_id) =>
                {
                    break;
                }
                Some(_) => {}
                None => panic!("scheduler event stream closed before initial analysis"),
            }
        }
    })
    .await
    .expect("initial running scheduler analysis should start promptly");

    shared
        .write()
        .await
        .apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
            dynamic_change_id.to_string(),
        ));
    assert!(dynamic_queue.push(dynamic_change_id.to_string()).await);

    let mut saw_dynamic_ingest = false;
    let mut dynamic_analysis_attempt = None;
    tokio::time::timeout(std::time::Duration::from_millis(500), async {
        while dynamic_analysis_attempt.is_none() {
            match rx.recv().await {
                Some(ExecutionEvent::Log(entry))
                    if entry.message.contains(&format!(
                        "Dynamically added to parallel execution: {dynamic_change_id}"
                    )) =>
                {
                    saw_dynamic_ingest = true;
                }
                Some(ExecutionEvent::AnalysisStarted { attempt_id, .. })
                    if attempt_id.contains(dynamic_change_id) =>
                {
                    dynamic_analysis_attempt = Some(attempt_id);
                }
                Some(_) => {}
                None => panic!("scheduler event stream closed before dynamic queue analysis"),
            }
        }
    })
    .await
    .expect("dynamic queue push after initial analysis should trigger sub-second reanalysis");

    cancel_token.cancel();
    let _ = tokio::time::timeout(std::time::Duration::from_millis(500), scheduler)
        .await
        .expect("scheduler should stop after cancellation")
        .expect("scheduler task should not panic");

    assert!(saw_dynamic_ingest, "dynamic queue entry should be ingested");
    assert!(
        dynamic_analysis_attempt
            .as_deref()
            .is_some_and(|attempt_id| attempt_id.contains("trigger=queue")),
        "dynamic queue analysis must use explicit queue trigger, got {dynamic_analysis_attempt:?}"
    );
}

#[tokio::test]
async fn dynamic_queue_ingestion_validates_candidates_against_executor_repo_root() {
    let temp_dir = TempDir::new().unwrap();
    let present_change_id = "synthetic-present-only-under-repo-root";
    let absent_change_id = "synthetic-absent-under-repo-root";
    create_active_change_fixture(temp_dir.path(), present_change_id);

    let (tx, mut rx) = tokio::sync::mpsc::channel(16);
    let dynamic_queue = Arc::new(DynamicQueue::new());
    dynamic_queue.push(present_change_id.to_string()).await;
    dynamic_queue.push(absent_change_id.to_string()).await;

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_dynamic_queue(dynamic_queue);
    // Both IDs carry accepted reducer queue intent, so the repo-root catalog
    // lookup is the only thing that can separate them here.
    executor.set_shared_orchestrator_state(shared_state_with_queue_intent(&[
        present_change_id,
        absent_change_id,
    ]));

    let mut queued = Vec::new();
    let in_flight = HashSet::new();
    let mut reanalysis_reason = ReanalysisReason::Initial;

    let queue_changed = executor
        .check_dynamic_queue_and_add_changes(&mut queued, &in_flight, &mut reanalysis_reason)
        .await;

    assert!(queue_changed, "present repo-root change should be ingested");
    assert_eq!(queued.len(), 1);
    assert_eq!(queued[0].id, present_change_id);
    assert_eq!(reanalysis_reason, ReanalysisReason::QueueNotification);

    let mut saw_present_ingest = false;
    let mut saw_absent_reconciliation = false;
    while let Ok(event) = rx.try_recv() {
        if let ExecutionEvent::Log(entry) = event {
            if entry.message.contains(&format!(
                "Dynamically added to parallel execution: {present_change_id}"
            )) {
                saw_present_ingest = true;
            }
            if entry.message.contains(&format!(
                "Queue reconciliation pending for '{absent_change_id}': candidate_not_found"
            )) {
                saw_absent_reconciliation = true;
            }
        }
    }

    assert!(
        saw_present_ingest,
        "ingestion log should name repo-root candidate"
    );
    assert!(
        saw_absent_reconciliation,
        "absent repo-root candidate should emit candidate_not_found reconciliation"
    );
    assert!(
        queued.iter().all(|change| change.id != absent_change_id),
        "absent candidate must not be queued"
    );
}

// ============================================================================
// Ghost queue prevention: accepted queue intent must converge with scheduler
// candidate discovery when the active OpenSpec catalog changes under a live
// owner.
//
// These are integration-scoped: they drive the real reducer, the real
// `DynamicQueue`, the real shared operator command boundary, and a real
// repository-visible `openspec/changes` tree on disk.
// ============================================================================

/// Reducer queue intent for one change, as an observer would read it.
async fn has_queue_intent(
    shared: &Arc<tokio::sync::RwLock<crate::orchestration::state::OrchestratorState>>,
    change_id: &str,
) -> bool {
    shared
        .read()
        .await
        .queued_change_ids()
        .iter()
        .any(|id| id == change_id)
}

/// Collect every scheduler log message emitted so far.
fn drain_log_messages(rx: &mut tokio::sync::mpsc::Receiver<ExecutionEvent>) -> Vec<String> {
    let mut messages = Vec::new();
    while let Ok(event) = rx.try_recv() {
        if let ExecutionEvent::Log(entry) = event {
            messages.push(entry.message);
        }
    }
    messages
}

/// The owner-before-proposal ordering, end to end on one executor.
///
/// The owner's first candidate lookup runs while `openspec/changes/<id>` does
/// not exist yet. The proposal then lands in the repository, and the *same*
/// executor — no restart, no new scheduler — must admit it from the refreshed
/// repository-visible view.
#[tokio::test]
async fn queued_intent_is_admitted_without_owner_restart_after_the_proposal_lands() {
    let temp_dir = TempDir::new().unwrap();
    let change_id = "synthetic-merged-after-owner-start";

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let dynamic_queue = Arc::new(DynamicQueue::new());
    dynamic_queue.push(change_id.to_string()).await;

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_dynamic_queue(dynamic_queue);
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let mut queued = Vec::new();
    let in_flight = HashSet::new();
    let mut reanalysis_reason = ReanalysisReason::Initial;

    // First lookup: the owner cannot see the proposal yet.
    let ingested = executor
        .check_dynamic_queue_and_add_changes(&mut queued, &in_flight, &mut reanalysis_reason)
        .await;
    assert!(!ingested, "an absent candidate cannot be ingested");
    assert!(queued.is_empty(), "no scheduler-local work exists yet");
    let messages = drain_log_messages(&mut rx);
    assert!(
        messages.iter().any(|message| message
            == &format!("Queue reconciliation pending for '{change_id}': candidate_not_found")),
        "the first miss must stay observable, got {messages:?}"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "hint ingestion alone must never revoke accepted queue intent"
    );

    // The proposal is merged into the base under the live owner.
    create_active_change_fixture(temp_dir.path(), change_id);

    let outcome = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;

    assert_eq!(
        outcome.queued_added, 1,
        "the same owner must admit the now-visible candidate"
    );
    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "a loadable candidate is never reconciled away"
    );
    assert!(
        queued.iter().any(|change| change.id == change_id),
        "the refreshed candidate must become scheduler-local queued work"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "admitted queue intent must survive the refresh"
    );
}

/// A miss inside one reconciliation pass is re-checked against a *fresh*
/// repository-visible view before anything is decided.
///
/// This is the in-pass half of the race: the pass's own catalog map missed, the
/// proposal landed, and the re-read admits it rather than settling accepted
/// intent from a stale observation.
#[tokio::test]
async fn a_missing_candidate_is_re_read_before_any_verdict_is_reached() {
    let temp_dir = TempDir::new().unwrap();
    let change_id = "synthetic-visible-only-on-refresh";

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    // The pass's first catalog read missed; the repository now has the proposal.
    create_active_change_fixture(temp_dir.path(), change_id);

    let snapshot = executor.capture_reducer_work_snapshot().await;
    let mut queued = Vec::new();
    let mut outcome = crate::parallel::queue_state::QueueReconciliationOutcome::default();
    executor
        .resolve_missing_queued_candidate(change_id, &mut queued, &mut outcome, &snapshot)
        .await;

    assert_eq!(
        outcome.queued_added, 1,
        "the fresh re-read must admit the candidate"
    );
    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "a candidate the refresh can load is not unavailable"
    );
    assert!(
        queued.iter().any(|change| change.id == change_id),
        "the refreshed candidate must become scheduler-local queued work"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "queue intent must be preserved when the refresh admits the candidate"
    );
    let messages = drain_log_messages(&mut rx);
    assert!(
        messages.iter().any(|message| message
            == &format!("Queue reconciliation admitted '{change_id}': candidate_refreshed")),
        "the refreshed-and-admitted result must be identifiable, got {messages:?}"
    );
}

/// A genuinely absent candidate leaves no queued row behind.
///
/// The queued projection has no scheduler-local work, no wake edge, and no
/// typed wait behind it, so reconciliation settles it through an explicit
/// reducer transition instead of reporting pending work forever. Diagnostics
/// stay bounded and nothing is dispatched.
///
/// The repository is a real Git repository on purpose: settlement requires
/// *conclusive* archived-dirty repair evidence, so a fixture where the base
/// branch cannot be resolved at all would prove deferral rather than the
/// settlement this test is about.
#[tokio::test]
async fn a_genuinely_absent_candidate_does_not_remain_a_ghost_queued_row() {
    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let change_id = "synthetic-never-created-anywhere";

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let mut queued = Vec::new();
    let in_flight = HashSet::new();

    let first = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;
    assert_eq!(
        first.unavailable_reconciled, 1,
        "unavailable queue intent must be reconciled explicitly"
    );
    assert_eq!(
        first.repair_evidence_deferred, 0,
        "conclusive repair evidence must not be reported as a deferral"
    );
    assert_eq!(first.queued_added, 0, "nothing loadable was added");
    assert!(
        queued.is_empty(),
        "an absent candidate must never become dispatchable work"
    );
    assert!(
        !has_queue_intent(&shared, change_id).await,
        "the queued projection must not survive as a ghost row"
    );

    // The reducer row is idle work again, not a dequeued or terminal outcome:
    // a later explicit Start or mark settlement can still admit the proposal.
    {
        let mut guard = shared.write().await;
        guard.apply_command(crate::orchestration::state::ReducerCommand::AddToQueue(
            change_id.to_string(),
        ));
    }
    assert!(
        has_queue_intent(&shared, change_id).await,
        "reconciliation must leave the change re-admittable"
    );

    // A second pass over the same unavailable intent repeats no warning.
    let _ = drain_log_messages(&mut rx);
    let second = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;
    assert_eq!(
        second.unavailable_reconciled, 1,
        "the re-added intent is settled again on its own evidence"
    );
    let messages = drain_log_messages(&mut rx);
    assert!(
        !messages
            .iter()
            .any(|message| message.contains("candidate_not_found")),
        "identical missing-candidate warnings must not repeat, got {messages:?}"
    );
}

/// The shared operator command boundary both frontends submit through.
///
/// `OperatorIntent::SetQueueIntent` from the WebUI/`/api/v2` route and from the
/// TUI route both reach `OperatorCommandService::add_to_queue`, so the boundary
/// is exercised directly here rather than inferred from either frontend's
/// helper. Neither route may require an owner restart after the base catalog
/// gains the proposal, and neither may lose the operator's execution mark to a
/// failed admission.
#[tokio::test]
async fn the_shared_queue_command_boundary_needs_no_owner_restart_after_a_catalog_update() {
    use crate::orchestration::operator_command::{
        ExecutionMarkStore, NoopQueueHooks, OperatorCommandService,
    };

    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let change_id = "synthetic-shared-boundary-late-proposal";

    let (tx, _rx) = tokio::sync::mpsc::channel(32);
    let dynamic_queue = DynamicQueue::new();
    let shared = Arc::new(tokio::sync::RwLock::new(
        crate::orchestration::state::OrchestratorState::new(vec![change_id.to_string()], 1),
    ));
    let marks = Arc::new(ExecutionMarkStore::new());
    let service = OperatorCommandService::new(
        shared.clone(),
        Arc::new(dynamic_queue.clone()),
        Arc::new(NoopQueueHooks),
        marks.clone(),
    );

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_dynamic_queue(Arc::new(dynamic_queue.clone()));
    executor.set_shared_orchestrator_state(shared.clone());

    // Operator selection and admission, through the shared boundary.
    marks.set(change_id, true);
    service
        .add_to_queue(change_id)
        .await
        .expect("shared queue command should be accepted");

    let mut queued = Vec::new();
    let in_flight = HashSet::new();
    let mut reanalysis_reason = ReanalysisReason::Initial;

    executor
        .check_dynamic_queue_and_add_changes(&mut queued, &in_flight, &mut reanalysis_reason)
        .await;
    let unavailable = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;

    assert_eq!(
        unavailable.unavailable_reconciled, 1,
        "an admission the catalog cannot satisfy must be settled, not left queued"
    );
    assert!(queued.is_empty(), "nothing dispatchable was produced");
    assert!(
        marks.is_marked(change_id),
        "a failed queue admission must not revoke the operator's execution mark"
    );

    // The proposal reaches the base catalog, and the same owner is asked again
    // through the same shared command.
    create_active_change_fixture(temp_dir.path(), change_id);
    service
        .add_to_queue(change_id)
        .await
        .expect("shared queue command should be accepted after the catalog update");

    let mut reanalysis_reason = ReanalysisReason::Initial;
    let ingested = executor
        .check_dynamic_queue_and_add_changes(&mut queued, &in_flight, &mut reanalysis_reason)
        .await;

    assert!(
        ingested,
        "the same owner must admit the change without a restart"
    );
    assert!(
        queued.iter().any(|change| change.id == change_id),
        "the shared boundary must produce scheduler-local queued work"
    );
    assert!(
        marks.is_marked(change_id),
        "the execution mark stays an independent axis throughout"
    );
}

/// Inconclusive archived-dirty repair evidence is not proof of absence.
///
/// The settle path spends accepted queue intent, and mark settlement is
/// edge-triggered, so a revoked row does not come back without an explicit
/// operator Start. A base branch that cannot be resolved means the repair probe
/// never ran at all — the precondition "no archived-dirty repair candidate
/// applies" was never established — so the pass must defer exactly as an
/// unreadable catalog does.
#[tokio::test]
async fn undetermined_repair_evidence_defers_instead_of_settling_queue_intent() {
    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let change_id = "synthetic-repair-evidence-unavailable";

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    // Detached HEAD with no recorded original branch: base identity is
    // unreadable, so the archived-dirty repair probe has no base to compare to.
    executor.set_workspace_manager(Box::new(
        TestWorkspaceManager::new(Arc::new(AtomicUsize::new(0))).with_failing_original_branch(),
    ));
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let mut queued = Vec::new();
    let in_flight = HashSet::new();

    let outcome = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;

    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "queue intent must never be settled on evidence that was never gathered"
    );
    assert_eq!(
        outcome.repair_evidence_deferred, 1,
        "the deferral must be observable rather than look like an idle pass"
    );
    assert_eq!(outcome.queued_added, 0, "nothing loadable was added");
    assert_eq!(outcome.repair_added, 0, "no repair candidate was proven");
    assert!(
        queued.is_empty(),
        "an undetermined candidate must never become dispatchable work"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "the wake edge and the operator's queue intent must both survive"
    );

    let messages = drain_log_messages(&mut rx);
    assert!(
        messages.iter().any(|message| message
            == &format!(
                "Queue reconciliation deferred for '{change_id}': repair_evidence_unavailable"
            )),
        "the deferral must be identifiable in operator-facing logs, got {messages:?}"
    );
    assert!(
        !messages
            .iter()
            .any(|message| message.contains("candidate_unavailable")),
        "an undetermined probe must not publish a settled verdict, got {messages:?}"
    );

    // Repeating the pass repeats neither the settle nor the diagnostic.
    let second = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;
    assert_eq!(
        second.unavailable_reconciled, 0,
        "repeating an undetermined pass still settles nothing"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "queue intent survives every undetermined pass"
    );
    let repeated = drain_log_messages(&mut rx);
    assert!(
        !repeated
            .iter()
            .any(|message| message.contains("repair_evidence_unavailable")),
        "deferral diagnostics stay bounded, got {repeated:?}"
    );
}

/// The same rule for the other half of the repair probe.
///
/// A workspace lookup that errors asked the question and got no answer, so
/// "this change has no repairable workspace" is precisely what was not proven.
#[tokio::test]
async fn a_failed_workspace_lookup_defers_instead_of_settling_queue_intent() {
    let temp_dir = TempDir::new().unwrap();
    init_minimal_git_repo(temp_dir.path());
    let change_id = "synthetic-workspace-lookup-failed";

    let (tx, _rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    // Base identity resolves; only workspace discovery fails.
    executor.set_workspace_manager(Box::new(
        TestWorkspaceManager::new(Arc::new(AtomicUsize::new(0)))
            .with_failing_existing_workspace_lookup(),
    ));
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let mut queued = Vec::new();
    let in_flight = HashSet::new();

    let outcome = executor
        .reconcile_queued_candidates_from_shared_state(&mut queued, &in_flight)
        .await;

    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "a failed workspace lookup proves no absence"
    );
    assert_eq!(
        outcome.repair_evidence_deferred, 1,
        "the failed probe must be counted as a deferral"
    );
    assert!(queued.is_empty(), "nothing dispatchable was produced");
    assert!(
        has_queue_intent(&shared, change_id).await,
        "queue intent must be preserved for a pass that can read the evidence"
    );
}

/// An unreadable active-change catalog keeps the dynamic queue hint.
///
/// A popped hint may be the only wake edge a queued change has. A catalog read
/// that fails produced no lookup at all, so ingestion puts the hint back
/// instead of spending it, and the reducer's queue intent is untouched.
#[tokio::test]
async fn an_unreadable_catalog_retains_the_dynamic_queue_hint_and_queue_intent() {
    let temp_dir = TempDir::new().unwrap();
    let change_id = "synthetic-catalog-unreadable";

    // `openspec/changes` exists but is not a directory, so the catalog read
    // fails rather than reporting an empty active-change set.
    std::fs::create_dir_all(temp_dir.path().join("openspec")).expect("create openspec directory");
    std::fs::write(
        temp_dir.path().join("openspec").join("changes"),
        "not a dir\n",
    )
    .expect("write catalog blocker");

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let dynamic_queue = Arc::new(DynamicQueue::new());
    dynamic_queue.push(change_id.to_string()).await;

    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    executor.set_dynamic_queue(dynamic_queue.clone());
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let mut queued = Vec::new();
    let in_flight = HashSet::new();
    let mut reanalysis_reason = ReanalysisReason::Initial;

    let ingested = executor
        .check_dynamic_queue_and_add_changes(&mut queued, &in_flight, &mut reanalysis_reason)
        .await;

    assert!(!ingested, "an unreadable catalog ingests nothing");
    assert!(queued.is_empty(), "no scheduler-local work was produced");
    assert_eq!(
        dynamic_queue.len().await,
        1,
        "the wake edge must be requeued, not spent on a lookup that never happened"
    );
    assert_eq!(
        dynamic_queue.pop().await.as_deref(),
        Some(change_id),
        "the retained hint must be the same change, at the front"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "an unreadable catalog must not touch reducer queue intent"
    );

    let messages = drain_log_messages(&mut rx);
    assert!(
        messages.iter().any(|message| message.starts_with(&format!(
            "Queue reconciliation pending for '{change_id}': candidate_load_failed"
        ))),
        "the unreadable read must stay observable, got {messages:?}"
    );
    assert!(
        !messages
            .iter()
            .any(|message| message.contains("candidate_not_found")),
        "unreadable is not absent, got {messages:?}"
    );
}

/// The reconciliation half of the same rule.
///
/// `resolve_missing_queued_candidate` re-reads the catalog itself. When that
/// re-read fails, absence is unproven, so queue intent is left exactly as it
/// was and the diagnostic stays bounded.
#[tokio::test]
async fn an_unreadable_refresh_leaves_queue_intent_exactly_as_it_was() {
    let temp_dir = TempDir::new().unwrap();
    let change_id = "synthetic-refresh-unreadable";

    std::fs::create_dir_all(temp_dir.path().join("openspec")).expect("create openspec directory");
    std::fs::write(
        temp_dir.path().join("openspec").join("changes"),
        "not a dir\n",
    )
    .expect("write catalog blocker");

    let (tx, mut rx) = tokio::sync::mpsc::channel(32);
    let mut executor = ParallelExecutor::new(
        temp_dir.path().to_path_buf(),
        create_test_config(),
        Some(tx),
    );
    let shared = shared_state_with_queue_intent(&[change_id]);
    executor.set_shared_orchestrator_state(shared.clone());

    let snapshot = executor.capture_reducer_work_snapshot().await;
    let mut queued = Vec::new();
    let mut outcome = crate::parallel::queue_state::QueueReconciliationOutcome::default();

    executor
        .resolve_missing_queued_candidate(change_id, &mut queued, &mut outcome, &snapshot)
        .await;

    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "an unreadable refresh must never settle accepted queue intent"
    );
    assert_eq!(outcome.queued_added, 0, "nothing loadable was found");
    assert!(queued.is_empty(), "nothing dispatchable was produced");
    assert!(
        has_queue_intent(&shared, change_id).await,
        "queue intent must survive for a pass that can read the repository"
    );

    let messages = drain_log_messages(&mut rx);
    assert!(
        messages.iter().any(|message| message.starts_with(&format!(
            "Queue reconciliation pending for '{change_id}': candidate_load_failed"
        ))),
        "the unreadable refresh must stay observable, got {messages:?}"
    );

    // A second unreadable refresh repeats neither the settle nor the warning.
    executor
        .resolve_missing_queued_candidate(change_id, &mut queued, &mut outcome, &snapshot)
        .await;
    assert_eq!(
        outcome.unavailable_reconciled, 0,
        "repetition still proves no absence"
    );
    assert!(
        has_queue_intent(&shared, change_id).await,
        "queue intent survives every unreadable refresh"
    );
    let repeated = drain_log_messages(&mut rx);
    assert!(
        !repeated
            .iter()
            .any(|message| message.contains("candidate_load_failed")),
        "catalog read failure diagnostics stay bounded, got {repeated:?}"
    );
}