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
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
//! The one process-local operator application transaction.
//!
//! [`crate::orchestration::operator_command::OperatorCommandService`] and
//! [`crate::orchestration::run_control::RunControlService`] own *what* an
//! operator command does. This module owns *when* — the ordering that makes an
//! accepted command truthful to every frontend at once:
//!
//! ```text
//! acquire the application gate
//!   -> validate mode, target, and eligibility against Core state
//!   -> reserve every fallible runtime capability (nothing mutated yet)
//!   -> commit the staged decision state
//!   -> dispatch one authoritative outcome and capture its exact revision
//!   -> activate the infallible scheduler permit, or issue its wake
//! release the gate
//! ```
//!
//! Two properties follow from that shape and are the reason it exists.
//!
//! *Preparation precedes mutation*, so a runtime that refuses to launch leaves
//! no reducer transition, no mark, no queue intent, no explicit-retry edge, no
//! resolve reservation, no mode change, and no event behind. A failure is the
//! absence of an effect rather than a rollback anyone has to trust.
//!
//! *Activation follows dispatch*, so scheduler progress enabled by a command can
//! never emit an event before the command's own accepted effect. A frontend
//! observing the run therefore cannot see `ProcessingStarted` before it sees the
//! Start that caused it.
//!
//! Activation is nevertheless *inside* the gate, after the outcome dispatch and
//! revision capture, not after the release. Activation is infallible, so it
//! cannot fail the transaction; releasing first would only open a window in
//! which a force stop is admitted against a run whose permit has not spawned
//! yet, so its `cancel_run` finds no handle, no-ops, and the run spawns anyway
//! immediately afterwards. Holding the gate across the activation makes the
//! admitted-to-spawned interval indivisible to every other command.
//!
//! A command that must await confirmed runtime termination is the explicit
//! two-phase exception: it holds the gate only for admission and cancellation
//! issuance, waits outside it, then reacquires the gate and revalidates the
//! target's *current runtime state* — not the revision it was admitted at, which
//! unrelated commands may legitimately have advanced.
//!
//! Everything here is process-local and dropped at exit. Under
//! `openspec/CONSTITUTION.md` the next workflow action after a restart is
//! recomputed from workspace and Git evidence alone; a mode, a gate, or a
//! reservation held in memory is never durable workflow authority.

use std::sync::{Arc, Mutex};

use crate::events::{
    accepted_start_opens_idle_run_episode, all_completed_may_overwrite_mode,
    graceful_stop_is_idle_origin, is_admitted_work_start, persistent_idle_may_project_ready,
    EventDispatcher, ExecutionEvent, OperatorCommandEffect, OutcomeRevisions,
};
use crate::orchestration::apply_commit_evidence::ApplyCommitEvidencePort;
use crate::orchestration::mark_settlement::{
    MarkSettlementAction, MarkSettlementExclusion, MarkSettlementFailure, MarkSettlementPlan,
    MarkSettlementRuntime,
};
use crate::orchestration::operator_command::{
    OperatorMode, OperatorOutcome, PendingForceStop, PendingTermination,
};
use crate::orchestration::run_control::{
    ResolveReservation, RunControlError, RunControlOutcome, RunControlService, RunNoOpReason,
    SchedulerEffect,
};

// ============================================================================
// Core process lifecycle mode
// ============================================================================

/// The one ephemeral process lifecycle mode every frontend projects.
///
/// This is *not* reducer change lifecycle and *not* durable workflow state. It
/// answers exactly one question — which operator commands this process will
/// admit right now — and TUI `execution_mode` and Web `app_mode` are projections
/// of it rather than independent authorities.
///
/// It is transitioned from two sources and no others: an accepted command
/// outcome, and an authoritative lifecycle event delivered through the
/// process-lifetime dispatch boundary. The second source is what keeps the
/// process admissible: without it, natural completion would leave Core wedged in
/// `Running` and a later Start would be refused forever.
#[derive(Debug)]
pub struct CoreMode {
    state: Mutex<CoreModeState>,
}

/// Admission mode plus the one qualifier that changes what `Select` means.
#[derive(Debug, Clone, Copy)]
struct CoreModeState {
    mode: OperatorMode,
    /// True while `Select` describes a live persistent-scheduler idle episode
    /// rather than ordinary pre-run selection.
    ///
    /// Held next to the mode, under the same lock, because the two are read and
    /// written as one decision: a cancel-stop that saw a stale value of either
    /// would restore the wrong mode.
    persistent_idle: bool,
}

impl Default for CoreMode {
    fn default() -> Self {
        Self::new()
    }
}

impl CoreMode {
    /// A fresh process, in `Select`.
    ///
    /// Fresh on every construction on purpose: a restart begins with no
    /// admission history, and the next action is recomputed from the workspace.
    pub fn new() -> Self {
        Self {
            state: Mutex::new(CoreModeState {
                mode: OperatorMode::Select,
                persistent_idle: false,
            }),
        }
    }

    fn lock(&self) -> std::sync::MutexGuard<'_, CoreModeState> {
        self.state
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
    }

    /// The current admission mode.
    pub fn get(&self) -> OperatorMode {
        self.lock().mode
    }

    /// Replace the admission mode, reporting whether it changed.
    ///
    /// The idle-episode fact qualifies `Select` and nothing else, so it travels
    /// with the mode it qualifies — the same rule the `/api/v2` snapshot applies
    /// when a caller publishes a different execution mode.
    ///
    /// Production moves the mode only through [`CoreMode::apply_event`]; this is
    /// how a test arranges the process it is about to command.
    #[cfg_attr(not(test), allow(dead_code))]
    pub fn set(&self, mode: OperatorMode) -> bool {
        let mut guard = self.lock();
        let changed = guard.mode != mode;
        guard.mode = mode;
        guard.persistent_idle &= mode == OperatorMode::Select;
        changed
    }

    /// Arrange the idle-episode half of the mode directly.
    ///
    /// Production only ever moves this through [`CoreMode::apply_event`]; a test
    /// that arranged a frontend in persistent-idle Ready without arranging Core
    /// would be exercising a split state the process cannot be in.
    #[cfg(test)]
    pub fn set_persistent_idle(&self, persistent_idle: bool) {
        self.lock().persistent_idle = persistent_idle;
    }

    /// Apply an authoritative lifecycle or command-outcome event.
    ///
    /// Returns the new mode when the event moved it. The guards are the same
    /// ones every frontend already applies, routed through one implementation so
    /// Core and the frontends cannot disagree about the same run:
    /// `AllCompleted` may not overwrite a retained terminal mode, and a
    /// command-outcome effect that carries no mode meaning (a mark delta, a
    /// queued resolve) moves nothing.
    ///
    /// The persistent-idle episode is tracked here for the same reason: `Select`
    /// over a live parked scheduler and `Select` before any run are different
    /// admission states, and a frontend that projects this mode must not have to
    /// re-derive which one it is holding.
    pub fn apply_event(&self, event: &ExecutionEvent) -> Option<OperatorMode> {
        let mut guard = self.lock();

        // One shared idle-episode rule, applied before the mode decision so an
        // arm below already reads the settled episode fact. It is the same rule
        // both frontends apply, over the same authoritative events, so the three
        // copies cannot drift apart.
        if is_admitted_work_start(event)
            || matches!(
                event,
                ExecutionEvent::Stopped | ExecutionEvent::Error { .. }
            )
        {
            guard.persistent_idle = false;
        }

        let next = match event {
            // Typed run activation. The scheduler really started work, so the
            // process is Running whether or not a command put it there.
            ExecutionEvent::ProcessingStarted(_) => OperatorMode::Running,
            // An accepted graceful stop. A stop admitted from Ready is an
            // idle-origin stop over a live parked scheduler, so the episode fact
            // is established here rather than assumed to already exist: Ready
            // reached through `AllCompleted` settlement owns no typed idle edge,
            // and without this a later cancel-stop would restore Running for an
            // episode nothing ever opened.
            ExecutionEvent::Stopping => {
                if graceful_stop_is_idle_origin(guard.mode.as_app_mode()) {
                    guard.persistent_idle = true;
                }
                OperatorMode::Stopping
            }
            ExecutionEvent::Stopped => OperatorMode::Stopped,
            ExecutionEvent::Error { .. } => OperatorMode::Error,
            // A change-scoped failure. The reducer records terminal Error on
            // that one change and mark reconciliation revokes its stale
            // execution intent; the *process* is not what failed, so admission
            // mode is left exactly where it was. Promoting it here is what made
            // one exhausted acceptance command disable ordinary mark controls
            // for every unrelated change still running.
            ExecutionEvent::ProcessingError { .. } => return None,
            // A persistent scheduler parked with nothing left to execute. Ready
            // here describes a *live* scheduler, so the shared guard applies:
            // only a running process becomes Ready, and pre-run Select, a
            // pending stop, and both terminal modes are left exactly as they
            // are.
            ExecutionEvent::PersistentSchedulerIdle => {
                if !persistent_idle_may_project_ready(guard.mode.as_app_mode()) {
                    return None;
                }
                guard.persistent_idle = true;
                OperatorMode::Select
            }
            ExecutionEvent::AllCompleted => {
                // Natural completion returns to Select and readmits a later
                // Start; a late completion after a stop or a fatal error does
                // not overwrite that terminal mode.
                if !all_completed_may_overwrite_mode(guard.mode.as_app_mode()) {
                    return None;
                }
                OperatorMode::Select
            }
            ExecutionEvent::OperatorCommandApplied { effect } => match effect {
                // A newly spawned scheduler proves execution has begun.
                OperatorCommandEffect::RunDispatched {
                    scheduler_started: true,
                    ..
                } => OperatorMode::Running,
                // A dispatch that woke a scheduler which was already alive.
                // Against a persistent-idle Ready episode that is the operator's
                // accepted Start: run control revalidated the live scheduler,
                // committed queue or explicit-retry intent, and notified it, so
                // the run episode is open even though nothing has been admitted
                // for execution yet. Every other wake — a Start into a live run,
                // a refusal, a no-op — moves nothing and waits for typed
                // work-start evidence as before.
                OperatorCommandEffect::RunDispatched {
                    change_ids,
                    scheduler_started,
                    ..
                } => {
                    if !accepted_start_opens_idle_run_episode(
                        guard.mode.as_app_mode(),
                        guard.persistent_idle,
                        *scheduler_started,
                        change_ids,
                    ) {
                        return None;
                    }
                    guard.persistent_idle = false;
                    OperatorMode::Running
                }
                // Withdrawing a stop restores where the stop came from. A stop
                // requested from persistent-idle Ready returns to Ready with its
                // episode intact; claiming Running would advertise execution no
                // typed event ever proved.
                OperatorCommandEffect::StopCancelled if guard.persistent_idle => {
                    OperatorMode::Select
                }
                OperatorCommandEffect::StopCancelled => OperatorMode::Running,
                OperatorCommandEffect::ForceStopAwaitingBoundary { .. } => OperatorMode::Stopping,
                OperatorCommandEffect::ResolveReserved { active: true, .. } => {
                    OperatorMode::Running
                }
                // A queued reservation starts nothing, and a mark or queue delta
                // is orthogonal to the run lifecycle.
                OperatorCommandEffect::ResolveReserved { active: false, .. }
                | OperatorCommandEffect::MarkDelta { .. }
                | OperatorCommandEffect::QueueDelta { .. } => return None,
            },
            // Work-start evidence other than `ProcessingStarted` resumes a
            // parked Ready and nothing else: a graceful stop that arrived first
            // is still owed, and a terminal mode stays terminal.
            _ if is_admitted_work_start(event) && guard.mode == OperatorMode::Select => {
                OperatorMode::Running
            }
            _ => return None,
        };
        let changed = guard.mode != next;
        guard.mode = next;
        changed.then_some(next)
    }
}

/// The typed scope of the two failure events, at the one place that decides it.
///
/// Unit-scoped: [`CoreMode`] alone, over constructed events. No reducer, no
/// frontend, no process, no repository.
#[cfg(test)]
mod core_mode_scope_tests {
    use super::*;

    fn processing_error() -> ExecutionEvent {
        ExecutionEvent::ProcessingError {
            id: "alpha".to_string(),
            error: "acceptance command attempts exhausted".to_string(),
        }
    }

    /// A change-scoped failure preserves whatever process mode it arrived in.
    ///
    /// Every mode is exercised, not just Running: the incident was Running, but
    /// a mapping that promoted the event from one mode and not another would be
    /// a second, quieter version of the same bug.
    #[test]
    fn processing_error_preserves_shared_mode() {
        for mode in [
            OperatorMode::Select,
            OperatorMode::Running,
            OperatorMode::Stopping,
            OperatorMode::Stopped,
            OperatorMode::Error,
        ] {
            let core = CoreMode::new();
            core.set(mode);

            assert_eq!(
                core.apply_event(&processing_error()),
                None,
                "a change-scoped failure moves no process mode (from {mode:?})"
            );
            assert_eq!(
                core.get(),
                mode,
                "the process mode must be exactly the one that existed before the event"
            );
        }
    }

    /// The idle-episode qualifier is a process fact too, so it survives as well.
    ///
    /// Observed through the command it actually qualifies: a stop withdrawn from
    /// persistent-idle Ready returns to Ready, and would claim Running instead if
    /// a change-local failure had closed the episode.
    #[test]
    fn processing_error_preserves_shared_mode_persistent_idle_episode() {
        let core = CoreMode::new();
        core.set(OperatorMode::Select);
        core.set_persistent_idle(true);

        assert_eq!(core.apply_event(&processing_error()), None);

        assert_eq!(
            core.apply_event(&ExecutionEvent::OperatorCommandApplied {
                effect: OperatorCommandEffect::StopCancelled,
            }),
            None,
            "the withdrawn stop returns to the Ready the episode still describes"
        );
        assert_eq!(core.get(), OperatorMode::Select);
    }

    /// The fatal control: suppressing every Error transition would also pass the
    /// test above, so the one event that really means the run died is asserted
    /// from the same starting modes.
    #[test]
    fn processing_error_preserves_shared_mode_fatal_control_still_transitions() {
        for mode in [
            OperatorMode::Select,
            OperatorMode::Running,
            OperatorMode::Stopping,
            OperatorMode::Stopped,
        ] {
            let core = CoreMode::new();
            core.set(mode);

            assert_eq!(
                core.apply_event(&ExecutionEvent::Error {
                    message: "the orchestrator could not start".to_string(),
                }),
                Some(OperatorMode::Error),
                "a global Error is process-fatal (from {mode:?})"
            );
            assert_eq!(core.get(), OperatorMode::Error);
        }
    }
}

// ============================================================================
// Intents and results
// ============================================================================

/// One operator intent, in the vocabulary both adapters submit.
///
/// The TUI and `/api/v2` translate their own request shapes into this and
/// nothing else, which is what makes "equivalent intent takes the same path" a
/// property of the type system rather than of two hand-kept implementations.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OperatorIntent {
    /// Start, resume, or (in `Error` mode) retry the marked set.
    Start,
    /// Request a graceful stop.
    Stop,
    /// Withdraw a pending graceful stop.
    CancelStop,
    /// Stop immediately.
    ForceStop,
    /// Retry one change.
    RetryChange {
        /// Target change.
        change_id: String,
    },
    /// Retry every change in the list that carries retryable evidence.
    RetryErrors {
        /// Candidate changes.
        change_ids: Vec<String>,
    },
    /// Reserve a merge resolution.
    ResolveMerge {
        /// Target change.
        change_id: String,
    },
    /// Set one change's execution mark.
    SetExecutionMark {
        /// Target change.
        change_id: String,
        /// Requested value.
        marked: bool,
    },
    /// Set one change's queue intent.
    SetQueueIntent {
        /// Target change.
        change_id: String,
        /// Requested membership.
        queued: bool,
    },
    /// Derive and apply one execution-mark value across every eligible change.
    SetAllExecutionMarks,
    /// Stop a change and dequeue it after confirmed termination.
    ///
    /// The two-phase intent: [`OperatorApplication::apply`] runs both phases and
    /// releases the gate between them.
    StopAndDequeue {
        /// Target change.
        change_id: String,
    },
    /// Kill one change's managed process group immediately and dequeue it.
    ///
    /// The other two-phase intent, and the only target-scoped control that
    /// bypasses the graceful SIGTERM escalation window. It names exactly one
    /// change and never touches the process-wide mode, the scheduler, or any
    /// other change's marks, queue intent, or processes.
    ForceStopChange {
        /// Target change.
        change_id: String,
    },
}

/// What an accepted command actually did, in the shared service vocabulary.
///
/// Reusing the two existing outcome enums rather than inventing a third is
/// deliberate: every adapter already projects them, and a third spelling would
/// be one more place the two frontends could drift.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ApplicationOutcome {
    /// A run-lifecycle command settled.
    Run(RunControlOutcome),
    /// A per-change operator command settled.
    Operator(OperatorOutcome),
}

/// One completed application transaction.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ApplicationResult {
    /// The typed outcome, or the typed refusal.
    pub outcome: Result<ApplicationOutcome, RunControlError>,
    /// The revision this command's own outcome dispatch produced.
    ///
    /// `None` when the command dispatched no outcome — an ordinary no-op or an
    /// ordinary failure — which is what a caller settles as "the unchanged
    /// admitted revision". A two-phase command that refuses *after* the wait
    /// returns `Some(unchanged revision observed at settlement)` instead,
    /// because by then the admitted revision is no longer the right answer.
    pub revision: Option<u64>,
}

impl ApplicationResult {
    fn failed(error: RunControlError) -> Self {
        Self {
            outcome: Err(error),
            revision: None,
        }
    }

    fn run(outcome: RunControlOutcome, revision: Option<u64>) -> Self {
        Self {
            outcome: Ok(ApplicationOutcome::Run(outcome)),
            revision,
        }
    }

    fn operator(outcome: OperatorOutcome, revision: Option<u64>) -> Self {
        Self {
            outcome: Ok(ApplicationOutcome::Operator(outcome)),
            revision,
        }
    }
}

/// Exclusive ownership of the application gate, held across one transaction.
pub type ApplicationGuard = tokio::sync::OwnedMutexGuard<()>;

// ============================================================================
// The coordinator
// ============================================================================

/// The single process-local operator application coordinator.
///
/// It lives outside the `web-monitoring` feature on purpose: a TUI-only build
/// runs the same transaction with the same ordering and simply has no revision
/// to report. Binding Web adds a sink and a revision source; it never creates
/// the authority.
pub struct OperatorApplication {
    gate: Arc<tokio::sync::Mutex<()>>,
    mode: Arc<CoreMode>,
    run_control: Arc<RunControlService>,
    dispatch: Arc<EventDispatcher>,
    /// Where an outcome dispatch's exact revision is read back from.
    ///
    /// `None` for a process with no `/api/v2` projection, which is also a
    /// process with no command record to bind a revision to.
    revisions: Option<Arc<dyn OutcomeRevisions>>,
    /// Where explanatory Apply-commit evidence is read from at stop settlement.
    ///
    /// `None` for a process with no repository observation port, which reports
    /// unknown rather than guessing.
    apply_commit_evidence: Option<Arc<dyn ApplyCommitEvidencePort>>,
}

impl OperatorApplication {
    /// Build the coordinator over the shared services and the process-lifetime
    /// dispatch owner.
    pub fn new(
        mode: Arc<CoreMode>,
        run_control: Arc<RunControlService>,
        dispatch: Arc<EventDispatcher>,
    ) -> Self {
        Self {
            gate: Arc::new(tokio::sync::Mutex::new(())),
            mode,
            run_control,
            dispatch,
            revisions: None,
            apply_commit_evidence: None,
        }
    }

    /// Bind the projection an outcome dispatch's revision is read back from.
    pub fn with_revisions(mut self, revisions: Option<Arc<dyn OutcomeRevisions>>) -> Self {
        self.revisions = revisions;
        self
    }

    /// Bind the repository port stop settlement reads Apply-commit evidence from.
    pub fn with_apply_commit_evidence(
        mut self,
        port: Option<Arc<dyn ApplyCommitEvidencePort>>,
    ) -> Self {
        self.apply_commit_evidence = port;
        self
    }

    /// The process-local serialization gate.
    ///
    /// A frontend that must serialize its *own* admission against command
    /// execution — `/api/v2`, whose optimistic revision check would otherwise let
    /// two new commands consume one revision — holds this across its whole
    /// submission and calls [`Self::apply_held`].
    pub fn gate(&self) -> Arc<tokio::sync::Mutex<()>> {
        self.gate.clone()
    }

    /// The shared run-lifecycle service.
    pub fn run_control(&self) -> Arc<RunControlService> {
        self.run_control.clone()
    }

    /// Apply one intent, acquiring the gate.
    pub async fn apply(&self, intent: OperatorIntent) -> ApplicationResult {
        if let OperatorIntent::StopAndDequeue { change_id } = &intent {
            return self.apply_stop_and_dequeue(change_id, None).await;
        }
        if let OperatorIntent::ForceStopChange { change_id } = &intent {
            return self.apply_force_stop_change(change_id, None).await;
        }
        let guard = self.gate.clone().lock_owned().await;
        self.apply_ordinary(intent, &guard).await
    }

    /// Apply one intent under a gate the caller already holds.
    ///
    /// A two-phase intent takes ownership of the guard: it must release it
    /// before waiting for confirmed termination and reacquire it to commit.
    pub async fn apply_held(
        &self,
        intent: OperatorIntent,
        guard: ApplicationGuard,
    ) -> ApplicationResult {
        if let OperatorIntent::StopAndDequeue { change_id } = &intent {
            let change_id = change_id.clone();
            return self.apply_stop_and_dequeue(&change_id, Some(guard)).await;
        }
        if let OperatorIntent::ForceStopChange { change_id } = &intent {
            let change_id = change_id.clone();
            return self.apply_force_stop_change(&change_id, Some(guard)).await;
        }
        self.apply_ordinary(intent, &guard).await
    }

    /// Begin a two-phase stop-and-dequeue under a held gate.
    ///
    /// Phase one only: validation and cancellation issuance. The caller releases
    /// the gate and later drives [`Self::settle_stop_and_dequeue`], which is what
    /// keeps force stop, unrelated commands, event fan-out, and rendering live
    /// while confirmation is pending.
    pub async fn begin_stop_and_dequeue(
        &self,
        change_id: &str,
    ) -> Result<PendingTermination, RunControlError> {
        self.run_control
            .operator()
            .begin_stop_and_dequeue(change_id)
            .await
            .map_err(RunControlError::Operator)
    }

    /// Settle a two-phase stop-and-dequeue after its cancellation was issued.
    ///
    /// Reacquires the gate, revalidates the target's current runtime state, and
    /// commits — or settles with the explicit unchanged revision observed under
    /// the reacquired boundary, without publishing a dequeue event.
    ///
    /// Between confirmation and reacquisition the terminated worktree is
    /// quiescent, which is the one window where explanatory Git evidence can be
    /// read safely *and* cheaply: doing it here rather than under the boundary
    /// keeps a Git subprocess from monopolizing operator admission, and doing it
    /// after termination rather than before keeps it from racing a worker that
    /// is still committing.
    pub async fn settle_stop_and_dequeue(&self, pending: PendingTermination) -> ApplicationResult {
        let change_id = pending.change_id().to_string();
        if let Err(error) = pending.confirm_termination().await {
            // Cancellation was already issued and is not rollbackable, but no
            // dequeue decision state may be committed. The settlement revision
            // is read explicitly rather than assumed to be the admitted one:
            // unrelated commands may have advanced it during the wait. It is
            // read under the *reacquired* gate for the same reason the
            // post-wait commit path holds it — a revision sampled outside the
            // boundary is a sample of mutable global state, which is exactly
            // what a settled record may not contain.
            let _guard = self.gate.clone().lock_owned().await;
            return ApplicationResult {
                outcome: Err(RunControlError::Operator(error)),
                revision: Some(self.current_revision()),
            };
        }

        // Outside the boundary, over a quiescent worktree.
        let apply_commit = crate::orchestration::apply_commit_evidence::observe_apply_commit(
            self.run_control.operator().execution_facts().as_deref(),
            self.apply_commit_evidence.as_deref(),
            &change_id,
        )
        .await;

        let _guard = self.gate.clone().lock_owned().await;
        match self
            .run_control
            .operator()
            .commit_stop_and_dequeue(&change_id, apply_commit)
            .await
        {
            Ok(OperatorOutcome::Dequeued {
                change_id,
                settlement,
            }) => {
                // `ChangeDequeued` already means exactly this; a second spelling
                // would give the same fact two authorities.
                let revision = self
                    .dispatch_outcome(ExecutionEvent::ChangeDequeued {
                        change_id: change_id.clone(),
                    })
                    .await;
                ApplicationResult::operator(
                    OperatorOutcome::Dequeued {
                        change_id,
                        settlement,
                    },
                    Some(revision),
                )
            }
            Ok(outcome) => ApplicationResult::operator(outcome, Some(self.current_revision())),
            Err(error) => ApplicationResult {
                outcome: Err(RunControlError::Operator(error)),
                revision: Some(self.current_revision()),
            },
        }
    }

    /// Begin a two-phase targeted force-stop under a held gate.
    ///
    /// Phase one only: eligibility validation against the fresh authoritative
    /// state, then the immediate SIGKILL and its reaping proof. The caller
    /// releases the gate and later drives [`Self::settle_force_stop_change`], so
    /// unrelated commands, event fan-out, and rendering stay live while the
    /// killed task unwinds.
    pub async fn begin_force_stop_change(
        &self,
        change_id: &str,
    ) -> Result<PendingForceStop, RunControlError> {
        self.run_control
            .operator()
            .begin_force_stop_change(change_id)
            .await
            .map_err(RunControlError::Operator)
    }

    /// Settle a two-phase targeted force-stop after its target was killed.
    ///
    /// Structurally the same shape as [`Self::settle_stop_and_dequeue`], and for
    /// the same reasons: the settlement revision is read under the reacquired
    /// gate, and explanatory Git evidence is read outside the boundary over the
    /// now-quiescent worktree.
    ///
    /// The event it dispatches is the ordinary `ChangeDequeued`. A targeted
    /// force-stop ends one execution episode exactly the way a stop does, so a
    /// subscription's terminal `stopped` classification and `wait`'s
    /// `change_requires_action` release come from the existing edge rather than
    /// from a second spelling of the same fact. The edge cannot undo the
    /// terminal `stopped` row the commit just settled: the reducer treats an
    /// already-stopped change as settled, exactly as it treats a merged, pushed,
    /// or rejected one.
    pub async fn settle_force_stop_change(&self, pending: PendingForceStop) -> ApplicationResult {
        let change_id = pending.change_id().to_string();
        if let Err(error) = pending.confirm_termination().await {
            // The kill was already issued and proven; only the task handshake
            // failed. No dequeue decision state may be committed on that.
            let _guard = self.gate.clone().lock_owned().await;
            return ApplicationResult {
                outcome: Err(RunControlError::Operator(error)),
                revision: Some(self.current_revision()),
            };
        }

        let apply_commit = crate::orchestration::apply_commit_evidence::observe_apply_commit(
            self.run_control.operator().execution_facts().as_deref(),
            self.apply_commit_evidence.as_deref(),
            &change_id,
        )
        .await;

        let _guard = self.gate.clone().lock_owned().await;
        match self
            .run_control
            .operator()
            .commit_force_stop_change(&pending, apply_commit)
            .await
        {
            Ok(outcome @ OperatorOutcome::ForceStopped { .. }) => {
                let revision = self
                    .dispatch_outcome(ExecutionEvent::ChangeDequeued {
                        change_id: change_id.clone(),
                    })
                    .await;
                ApplicationResult::operator(outcome, Some(revision))
            }
            Ok(outcome) => ApplicationResult::operator(outcome, Some(self.current_revision())),
            Err(error) => ApplicationResult {
                outcome: Err(RunControlError::Operator(error)),
                revision: Some(self.current_revision()),
            },
        }
    }

    async fn apply_force_stop_change(
        &self,
        change_id: &str,
        held: Option<ApplicationGuard>,
    ) -> ApplicationResult {
        let guard = match held {
            Some(guard) => guard,
            None => self.gate.clone().lock_owned().await,
        };
        let pending = match self.begin_force_stop_change(change_id).await {
            Ok(pending) => pending,
            Err(error) => return ApplicationResult::failed(error),
        };
        drop(guard);
        self.settle_force_stop_change(pending).await
    }

    async fn apply_stop_and_dequeue(
        &self,
        change_id: &str,
        held: Option<ApplicationGuard>,
    ) -> ApplicationResult {
        let guard = match held {
            Some(guard) => guard,
            None => self.gate.clone().lock_owned().await,
        };
        let pending = match self.begin_stop_and_dequeue(change_id).await {
            Ok(pending) => pending,
            Err(error) => return ApplicationResult::failed(error),
        };
        // The gate is released *before* the wait, never during it.
        drop(guard);
        self.settle_stop_and_dequeue(pending).await
    }

    // ------------------------------------------------------------------
    // Ordinary transaction
    // ------------------------------------------------------------------

    async fn apply_ordinary(
        &self,
        intent: OperatorIntent,
        _guard: &ApplicationGuard,
    ) -> ApplicationResult {
        let mode = self.mode.get();
        match intent {
            OperatorIntent::Start => {
                self.apply_prepared(self.run_control.prepare_start(mode).await)
                    .await
            }
            OperatorIntent::RetryChange { change_id } => {
                self.apply_prepared(self.run_control.prepare_retry_change(&change_id).await)
                    .await
            }
            OperatorIntent::RetryErrors { change_ids } => {
                self.apply_prepared(self.run_control.prepare_retry_errors(&change_ids).await)
                    .await
            }
            OperatorIntent::ResolveMerge { change_id } => self.apply_resolve(&change_id).await,
            OperatorIntent::Stop => self.apply_stop(mode).await,
            OperatorIntent::CancelStop => self.apply_cancel_stop(mode).await,
            OperatorIntent::ForceStop => self.apply_force_stop(mode).await,
            OperatorIntent::SetExecutionMark { change_id, marked } => {
                self.apply_mark(&change_id, marked).await
            }
            OperatorIntent::SetQueueIntent { change_id, queued } => {
                self.apply_queue_intent(&change_id, queued).await
            }
            OperatorIntent::SetAllExecutionMarks => self.apply_bulk_marks().await,
            // Routed before the gate was taken.
            OperatorIntent::StopAndDequeue { .. } | OperatorIntent::ForceStopChange { .. } => {
                unreachable!("a two-phase intent is routed before the ordinary transaction")
            }
        }
    }

    /// Commit a prepared run-lifecycle command, publish its outcome, then
    /// activate its scheduler dispatch.
    async fn apply_prepared(
        &self,
        prepared: Result<crate::orchestration::run_control::PreparedRunCommand, RunControlError>,
    ) -> ApplicationResult {
        let prepared = match prepared {
            Ok(prepared) => prepared,
            // Preparation never mutated anything, so this refusal has nothing to
            // undo and nothing to publish.
            Err(error) => return ApplicationResult::failed(error),
        };

        let committed = match self.run_control.commit(prepared).await {
            Ok(committed) => committed,
            Err(error) => return ApplicationResult::failed(error),
        };

        let outcome = committed.outcome.clone();
        let Some(event) = run_outcome_event(&outcome) else {
            // A no-op publishes nothing and advances no revision.
            committed
                .activate(self.run_control.scheduler().as_ref())
                .await;
            return ApplicationResult::run(outcome, None);
        };

        let revision = self.dispatch_outcome(event).await;
        // Only now may scheduler activity enabled by this command exist.
        committed
            .activate(self.run_control.scheduler().as_ref())
            .await;
        ApplicationResult::run(outcome, Some(revision))
    }

    async fn apply_resolve(&self, change_id: &str) -> ApplicationResult {
        match self.run_control.prepare_resolve(change_id).await {
            Ok(Some(prepared)) => self.apply_prepared(Ok(prepared)).await,
            Ok(None) => ApplicationResult::run(
                RunControlOutcome::NoOp {
                    reason: RunNoOpReason::ResolveAlreadyReserved {
                        change_id: change_id.to_string(),
                    },
                },
                None,
            ),
            Err(error) => ApplicationResult::failed(error),
        }
    }

    async fn apply_stop(&self, mode: OperatorMode) -> ApplicationResult {
        match self.run_control.stop(mode).await {
            // `Stopping` already means exactly this, so no new vocabulary is
            // introduced for it.
            Ok(outcome @ RunControlOutcome::StopRequested) => {
                let revision = self.dispatch_outcome(ExecutionEvent::Stopping).await;
                ApplicationResult::run(outcome, Some(revision))
            }
            Ok(other) => ApplicationResult::run(other, None),
            Err(error) => ApplicationResult::failed(error),
        }
    }

    async fn apply_cancel_stop(&self, mode: OperatorMode) -> ApplicationResult {
        match self.run_control.cancel_stop(mode).await {
            Ok(outcome @ RunControlOutcome::StopCancelled) => {
                let revision = self
                    .dispatch_outcome(ExecutionEvent::OperatorCommandApplied {
                        effect: OperatorCommandEffect::StopCancelled,
                    })
                    .await;
                ApplicationResult::run(outcome, Some(revision))
            }
            Ok(other) => ApplicationResult::run(other, None),
            Err(error) => ApplicationResult::failed(error),
        }
    }

    async fn apply_force_stop(&self, mode: OperatorMode) -> ApplicationResult {
        let outcome = match self.run_control.force_stop(mode).await {
            Ok(outcome) => outcome,
            Err(error) => return ApplicationResult::failed(error),
        };
        let RunControlOutcome::ForceStopped {
            classification,
            awaiting_safe_boundary,
        } = &outcome
        else {
            return ApplicationResult::run(outcome, None);
        };

        let event = if *awaiting_safe_boundary {
            // Terminal state must not be published before the required cleanup:
            // the scheduler owns the eventual `Stopped`.
            ExecutionEvent::OperatorCommandApplied {
                effect: OperatorCommandEffect::ForceStopAwaitingBoundary {
                    force_stop: classification.process_report.is_force_stop(),
                },
            }
        } else {
            // No live scheduler will emit this run's terminal stop, so this
            // command is its dispatch owner and `Stopped` is exact.
            ExecutionEvent::Stopped
        };
        let revision = self.dispatch_outcome(event).await;
        ApplicationResult::run(outcome, Some(revision))
    }

    async fn apply_mark(&self, change_id: &str, marked: bool) -> ApplicationResult {
        match self
            .run_control
            .operator()
            .set_execution_mark(change_id, marked)
            .await
        {
            Ok(outcome) => self.publish_operator_outcome(outcome).await,
            Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
        }
    }

    async fn apply_queue_intent(&self, change_id: &str, queued: bool) -> ApplicationResult {
        let service = self.run_control.operator();
        let result = if queued {
            service.add_to_queue(change_id).await
        } else {
            service.remove_from_queue(change_id).await
        };
        match result {
            Ok(queue) => {
                self.publish_operator_outcome(OperatorOutcome::Queue(queue))
                    .await
            }
            Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
        }
    }

    async fn apply_bulk_marks(&self) -> ApplicationResult {
        match self.run_control.operator().set_all_execution_marks().await {
            Ok(outcome) => self.publish_operator_outcome(outcome).await,
            Err(error) => ApplicationResult::failed(RunControlError::Operator(error)),
        }
    }

    /// Publish the delta an accepted per-change command produced.
    async fn publish_operator_outcome(&self, outcome: OperatorOutcome) -> ApplicationResult {
        let Some(event) = operator_outcome_event(&outcome) else {
            return ApplicationResult::operator(outcome, None);
        };
        let revision = self.dispatch_outcome(event).await;
        ApplicationResult::operator(outcome, Some(revision))
    }

    // ------------------------------------------------------------------
    // Dispatch
    // ------------------------------------------------------------------

    /// Dispatch one accepted outcome and report the exact revision it produced.
    ///
    /// The revision is looked up by *dispatch identity*, not sampled from global
    /// state afterwards: sampling would let unrelated later progress become the
    /// command's recorded `result_revision`.
    async fn dispatch_outcome(&self, event: ExecutionEvent) -> u64 {
        let dispatch_id = self.dispatch.dispatch(event).await;
        match &self.revisions {
            Some(revisions) => revisions
                .revision_for_dispatch(dispatch_id)
                .unwrap_or_else(|| revisions.current_revision()),
            None => 0,
        }
    }

    fn current_revision(&self) -> u64 {
        self.revisions
            .as_ref()
            .map_or(0, |revisions| revisions.current_revision())
    }
}

/// Bind this coordinator as the process's mark-settlement runtime.
///
/// Weakly, and deliberately: the mark store is reachable from the coordinator
/// through run control, so a strong handle would close a cycle that never drops.
///
/// Call it once, after the coordinator is in its `Arc`. Until it is called every
/// mark write is mark-only, which is the correct behaviour for a process that
/// has no application transaction to admit work through.
pub fn bind_mark_settlement(application: &Arc<OperatorApplication>) {
    let owned: std::sync::Weak<OperatorApplication> = Arc::downgrade(application);
    let runtime: std::sync::Weak<dyn MarkSettlementRuntime> = owned;
    application
        .run_control
        .operator()
        .marks()
        .settlement()
        .bind_runtime(runtime);
}

/// The mark-stability policy's view of this process.
///
/// The coordinator is the right implementation because settlement must produce
/// exactly what a frontend queue command produces: the same reducer transition,
/// the same `DynamicQueue` mutation, the same queue hook, the same authoritative
/// outcome, and the same scheduler wake. Routing it through
/// [`OperatorIntent::SetQueueIntent`] is what makes that identity structural
/// rather than a second implementation that has to be kept in step.
#[async_trait::async_trait]
impl MarkSettlementRuntime for OperatorApplication {
    fn admits_dynamic_queue(&self) -> bool {
        // Scheduler-task liveness, never presentation mode. A persistent
        // scheduler parked in Select is still perfectly able to admit work, and
        // a finite run that already exited is not.
        self.run_control.scheduler().is_running()
    }

    async fn settle_marks(&self, targets: Vec<String>) -> MarkSettlementPlan {
        // Classification happens outside the application gate; only the derived
        // plan is applied under it. Settlement therefore never holds the gate
        // across a reducer read, and never waits on a base-mutating lane.
        let mut plan = self
            .run_control
            .operator()
            .plan_mark_settlement(&targets)
            .await;

        let mutations: Vec<(String, MarkSettlementAction)> = plan
            .additions
            .iter()
            .map(|change_id| (change_id.clone(), MarkSettlementAction::Add))
            .chain(
                plan.removals
                    .iter()
                    .map(|change_id| (change_id.clone(), MarkSettlementAction::Remove)),
            )
            .collect();

        let mut applied_membership_change = false;
        // Reasons the write boundary produced for targets classification had
        // planned a mutation for. They belong in the returned plan: a target the
        // guard refused gained no queue effect, so leaving it listed as an
        // applied addition or removal would be a false claim — and a refusal the
        // plan never carried is a refusal the caller cannot act on.
        let mut guard_skipped: Vec<(String, MarkSettlementExclusion)> = Vec::new();
        for (change_id, action) in mutations {
            let application = {
                // One target, one gate acquisition. Holding the gate across the
                // whole batch would let a settled mark starve every operator
                // command behind it, and the per-target guard is what makes each
                // mutation correct on its own anyway.
                let _guard = self.gate.clone().lock_owned().await;
                self.run_control
                    .operator()
                    .apply_settlement_queue_intent(&change_id, action)
                    .await
            };
            if application.applied() {
                applied_membership_change = true;
            }
            if let Some(reason) = application.skipped {
                guard_skipped.push((change_id.clone(), reason));
            }
            // Publish the same queue delta an explicit frontend command
            // publishes, so the TUI and `/api/v2` project `queued` and
            // `not queued` from authoritative reducer intent rather than from
            // the mark set.
            self.publish_operator_outcome(OperatorOutcome::Queue(application.outcome))
                .await;
        }

        for (change_id, reason) in guard_skipped {
            plan.additions.retain(|planned| planned != &change_id);
            plan.removals.retain(|planned| planned != &change_id);
            plan.excluded.push((change_id, reason));
        }

        // Exactly one wake for the whole settled batch, and only when membership
        // really moved. A no-op batch changes no analysis input, so waking the
        // scheduler for it would manufacture an analysis attempt with no operator
        // action behind it.
        if applied_membership_change {
            self.run_control
                .operator()
                .notify_scheduler_after_settlement()
                .await;
        }

        // Observability only, and deliberately not operator-facing: a named row
        // that settlement skipped is otherwise indistinguishable from a deadline
        // that never expired, which is the one question this policy is hard to
        // answer from the outside.
        if !plan.excluded.is_empty() {
            let skipped = plan
                .excluded
                .iter()
                .map(|(change_id, reason)| format!("{change_id}={}", reason.as_str()))
                .collect::<Vec<_>>()
                .join(", ");
            tracing::debug!("Mark settlement changed no queue intent for: {skipped}");
        }
        plan
    }

    async fn report_abandoned_settlement(&self, pending: Vec<String>) {
        let targets = if pending.is_empty() {
            "no marked change".to_string()
        } else {
            pending.join(", ")
        };
        self.dispatch
            .dispatch(ExecutionEvent::Log(crate::events::LogEntry::info(format!(
                "Mark settlement abandoned because the scheduler ended: {targets}"
            ))))
            .await;
    }

    async fn report_settlement_failure(
        &self,
        failure: MarkSettlementFailure,
        targets: Vec<String>,
    ) {
        // Operator-facing, unlike the per-row exclusion diagnostic: a row that
        // settlement reasoned about keeps its mark for a reason the operator can
        // see on the row itself, but a batch the owner could not carry to a
        // decision at all is invisible everywhere else.
        let named = if targets.is_empty() {
            "no marked change".to_string()
        } else {
            targets.join(", ")
        };
        self.dispatch
            .dispatch(ExecutionEvent::Log(crate::events::LogEntry::warn(format!(
                "Mark settlement did not admit (reason={}): {named}",
                failure.as_str()
            ))))
            .await;
    }
}

/// The authoritative event one accepted run-lifecycle outcome publishes.
///
/// `None` for a no-op, which publishes nothing and advances no revision.
fn run_outcome_event(outcome: &RunControlOutcome) -> Option<ExecutionEvent> {
    match outcome {
        RunControlOutcome::RunDispatched {
            change_ids,
            explicit_retry,
            scheduler,
            // Exclusions are operator-facing reporting on the command result,
            // not run state: the authoritative event carries only what was
            // actually dispatched.
            excluded: _,
        } => {
            // This effect *means* scheduler-wake evidence for the targets it
            // names: a persistent-idle frontend reads a woken dispatch as the
            // accepted Start that opens its run episode. Run control only ever
            // builds this outcome from a reserved dispatch, so the invariant
            // holds by construction — asserting it is what stops a future
            // no-dispatch path from silently teaching frontends to claim
            // Running for a scheduler nobody woke.
            debug_assert!(
                scheduler.dispatched(),
                "an accepted run dispatch must carry scheduler evidence"
            );
            Some(ExecutionEvent::OperatorCommandApplied {
                effect: OperatorCommandEffect::RunDispatched {
                    change_ids: change_ids.clone(),
                    explicit_retry: *explicit_retry,
                    scheduler_started: matches!(scheduler, SchedulerEffect::Started),
                },
            })
        }
        RunControlOutcome::ResolveReserved {
            change_id,
            reservation,
            ..
        } => Some(ExecutionEvent::OperatorCommandApplied {
            effect: OperatorCommandEffect::ResolveReserved {
                change_id: change_id.clone(),
                active: matches!(reservation, ResolveReservation::Active),
            },
        }),
        // Published by their own dedicated paths, which choose between the exact
        // existing event and the awaiting-boundary variant.
        RunControlOutcome::StopRequested
        | RunControlOutcome::StopCancelled
        | RunControlOutcome::ForceStopped { .. } => None,
        RunControlOutcome::NoOp { .. } => None,
    }
}

/// The authoritative event one accepted per-change outcome publishes.
///
/// Crate-visible so a frontend-projection test can publish exactly what this
/// coordinator publishes for the same outcome. Hand-constructing the event there
/// would be a second answer to "what does an accepted mark broadcast", and the
/// projection under test would stop being exercised by the real one.
pub(crate) fn operator_outcome_event(outcome: &OperatorOutcome) -> Option<ExecutionEvent> {
    match outcome {
        OperatorOutcome::MarkSet { change_id, marked } => {
            Some(ExecutionEvent::OperatorCommandApplied {
                effect: OperatorCommandEffect::MarkDelta {
                    change_ids: vec![change_id.clone()],
                    marked: *marked,
                },
            })
        }
        OperatorOutcome::BulkMarks {
            marked, changed, ..
        } => (!changed.is_empty()).then(|| ExecutionEvent::OperatorCommandApplied {
            effect: OperatorCommandEffect::MarkDelta {
                change_ids: changed.clone(),
                marked: *marked,
            },
        }),
        OperatorOutcome::Queue(queue) => (queue.reducer_changed || queue.dynamic_queue_mutated)
            .then(|| ExecutionEvent::OperatorCommandApplied {
                effect: OperatorCommandEffect::QueueDelta {
                    change_id: queue.change_id.clone(),
                    queued: matches!(
                        queue.mutation,
                        crate::orchestration::operator_command::QueueMutation::Added
                    ),
                },
            }),
        // A plan is the reducer transition, not the dispatch that makes it real:
        // it carries no scheduler evidence at all. `RunDispatched` now *means*
        // scheduler-wake evidence for its committed targets — that is what lets
        // a persistent-idle frontend read one as the accepted Start that opens
        // the run episode — so a publisher holding no such evidence must not
        // reuse the shape. Every production retry reaches the run-control path,
        // which does hold the scheduler effect and does publish the outcome.
        OperatorOutcome::Retry(_) => None,
        // `ChangeDequeued` already means exactly this and is published by the
        // two-phase settlement paths. A targeted force-stop ends one execution
        // episode the same way a stop does, so it reuses that edge rather than
        // giving the same fact a second authority.
        OperatorOutcome::Dequeued { .. } | OperatorOutcome::ForceStopped { .. } => None,
        OperatorOutcome::NoOp { .. } => None,
    }
}

#[cfg(test)]
mod tests;

/// Phase-aware stop settlement, over deterministic termination ordering.
#[cfg(test)]
mod stop_settlement_tests;

/// The dispatch boundary with a real `/api/v2` projection attached.
#[cfg(all(test, feature = "web-monitoring"))]
mod mode_matrix_tests;