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
//! Process-local execution facts: typed phases, boundaries, and activity.
//!
//! This is an *observability* store, not a second lifecycle authority. It holds
//! nothing the workspace cannot re-derive, it is created fresh with the process,
//! and it is discarded at exit — under `openspec/CONSTITUTION.md` the next
//! workflow action after a restart is still recomputed from workspace and Git
//! evidence alone.
//!
//! Two rules keep it from becoming a second state machine:
//!
//! * **The reducer owns the current phase.** [`ExecutionFactsStore::observe`] is
//!   called from the authoritative typed-event dispatch boundary with the
//!   reducer state that event produced, and projects
//!   [`crate::orchestration::state::ActivityState`] into the closed wire
//!   vocabulary. Nothing here classifies a phase from a display string, a task
//!   count, a log line, or a commit subject.
//! * **Completion comes from typed completion events.** A phase is recorded as
//!   *completed* only when its own typed completion event is dispatched, so an
//!   Apply that failed is never reported as the last completed phase merely
//!   because the reducer left Applying.
//!
//! The facts are read by the `/api/v2` execution-status resource and by
//! stop-and-dequeue settlement, which is what lets both answer "what was
//! actually happening" with the same evidence.

use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::{Mutex, MutexGuard};

use chrono::{DateTime, Utc};

use crate::events::ExecutionEvent;
use crate::orchestration::state::{
    ActivityState, ChangeRuntimeState, OrchestratorState, QueueIntent, TerminalState, WaitState,
};

/// Closed per-change lifecycle phase vocabulary.
///
/// `Merge` is deliberately reachable only as a *completed* phase: the reducer
/// has no merging activity, so a per-change merge is observed through its typed
/// completion fact and never advertised as something currently running.
///
/// No `hook` value exists yet. Production emits typed hook start/completion
/// events, but they are not per-change lifecycle phases the reducer tracks, and
/// advertising a vocabulary value the server cannot classify truthfully would be
/// worse than reporting the phase the reducer really holds.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum ExecutionPhase {
    /// Admitted to a slot and preparing its managed workspace.
    Preparing,
    /// Running Apply.
    Apply,
    /// Running acceptance.
    Acceptance,
    /// Running dedicated rejection review.
    RejectionReview,
    /// Running archive.
    Archive,
    /// Running merge resolution.
    Resolve,
    /// A typed push episode is open.
    Push,
    /// A typed per-change merge completed (completed-phase only).
    Merge,
    /// No phase is active.
    #[default]
    None,
    /// Typed evidence exists but cannot be classified.
    Unknown,
}

impl ExecutionPhase {
    /// Stable wire token.
    pub fn as_str(self) -> &'static str {
        match self {
            Self::Preparing => "preparing",
            Self::Apply => "apply",
            Self::Acceptance => "acceptance",
            Self::RejectionReview => "rejection_review",
            Self::Archive => "archive",
            Self::Resolve => "resolve",
            Self::Push => "push",
            Self::Merge => "merge",
            Self::None => "none",
            Self::Unknown => "unknown",
        }
    }

    /// True when this value names a phase that can actually be running.
    pub fn is_active(self) -> bool {
        !matches!(self, Self::None | Self::Unknown)
    }
}

/// Closed per-change execution-state vocabulary.
///
/// Terminal outcomes take precedence, then the graceful-stop episode, then the
/// reducer's own activity and wait facts. A row the reducer tracks but has
/// neither queued, activated, held, nor finished is reported as `Unknown`
/// rather than mapped onto a state the reducer never claimed.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ChangeExecutionState {
    /// Requested to run and waiting for a slot.
    Queued,
    /// The reducer holds an active execution stage.
    Active,
    /// Held on a wait or blocker condition.
    Waiting,
    /// Active while a graceful process stop is in flight.
    Stopping,
    /// Stopped by operator request.
    Stopped,
    /// Reached a terminal failure or rejection.
    Failed,
    /// Reached a terminal success.
    Completed,
    /// Tracked, but no typed evidence classifies it.
    #[default]
    Unknown,
}

impl ChangeExecutionState {
    /// Stable wire token.
    #[allow(dead_code)] // Read by execution-facts vocabulary coverage, not by the binary.
    pub fn as_str(self) -> &'static str {
        match self {
            Self::Queued => "queued",
            Self::Active => "active",
            Self::Waiting => "waiting",
            Self::Stopping => "stopping",
            Self::Stopped => "stopped",
            Self::Failed => "failed",
            Self::Completed => "completed",
            Self::Unknown => "unknown",
        }
    }
}

/// Closed process-level activity vocabulary.
///
/// These are the episodes that are real lifecycle work without belonging to any
/// single change. Each one is opened by its typed start event and closed by its
/// typed terminal event; a process terminal closes every open episode so a
/// crashed lane cannot leave the process reporting work forever.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ProcessActivity {
    /// Dependency analysis over the remaining changes.
    DependencyAnalysis,
    /// Sequential base-branch merge of a completed batch.
    BaseBranchMerge,
    /// Conflict resolution inside a base-branch merge.
    ConflictResolution,
    /// Worktree branch merge requested through the worktree surface.
    BranchMerge,
    /// Managed workspace cleanup.
    WorkspaceCleanup,
}

impl ProcessActivity {
    /// Stable wire token.
    pub fn as_str(self) -> &'static str {
        match self {
            Self::DependencyAnalysis => "dependency_analysis",
            Self::BaseBranchMerge => "base_branch_merge",
            Self::ConflictResolution => "conflict_resolution",
            Self::BranchMerge => "branch_merge",
            Self::WorkspaceCleanup => "workspace_cleanup",
        }
    }
}

/// One change's observed execution facts.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ChangeExecutionFacts {
    /// Process-local identity of the most recent admitted execution episode.
    ///
    /// `None` until any admission source moved this change into queued or
    /// active work in *this* incarnation. It survives the episode's terminal
    /// settlement so a late subscriber can still address the execution that
    /// just finished, and is replaced — never reused — by the next admission.
    pub execution_id: Option<String>,
    /// Closed execution state.
    pub execution_state: ChangeExecutionState,
    /// Phase the reducer currently holds.
    pub current_phase: ExecutionPhase,
    /// When the current phase became active; `None` when no phase is active or
    /// the boundary was not observed in this incarnation.
    pub phase_started_at: Option<DateTime<Utc>>,
    /// Last phase that published its own typed completion fact.
    pub last_completed_phase: Option<ExecutionPhase>,
    /// When that completion was observed.
    pub last_completed_at: Option<DateTime<Utc>>,
    /// Retained non-empty `ApplyCompleted.revision` OID for this incarnation.
    pub apply_commit_oid: Option<String>,
}

impl ChangeExecutionFacts {
    /// The value for a change this incarnation has no typed evidence about.
    ///
    /// Deliberately not the `Default`: "no phase is running" and "nothing was
    /// observed" are different claims, and only the second one is honest for a
    /// change the store has never seen.
    pub fn unknown() -> Self {
        Self {
            execution_id: None,
            execution_state: ChangeExecutionState::Unknown,
            current_phase: ExecutionPhase::Unknown,
            phase_started_at: None,
            last_completed_phase: None,
            last_completed_at: None,
            apply_commit_oid: None,
        }
    }
}

/// A coherent read of the whole store.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ExecutionFactsSnapshot {
    /// Per-change facts, keyed by change ID.
    pub changes: HashMap<String, ChangeExecutionFacts>,
    /// Process-level episodes that have started and not reached a terminal.
    pub activities: Vec<ProcessActivity>,
}

impl ExecutionFactsSnapshot {
    /// Facts for one change, explicitly unknown when nothing was observed.
    pub fn change(&self, change_id: &str) -> ChangeExecutionFacts {
        self.changes
            .get(change_id)
            .cloned()
            .unwrap_or_else(ChangeExecutionFacts::unknown)
    }

    /// True when a per-change phase or a process-level episode is running.
    ///
    /// Scheduler liveness is deliberately *not* an input: a parked persistent
    /// scheduler is alive without any admitted work, and conflating the two is
    /// exactly the ambiguity this resource exists to remove.
    pub fn has_active_work(&self) -> bool {
        !self.activities.is_empty()
            || self
                .changes
                .values()
                .any(|facts| facts.current_phase.is_active())
    }
}

/// The states that mean "this change is admitted work right now".
///
/// Admission — not activity — is what opens an execution episode: a change that
/// is queued behind a slot has already been admitted by whichever source asked
/// for it, and a subscriber that only learned about it once a phase started
/// would miss the window an enqueue caller is actually in.
pub fn is_admitted_execution_state(state: ChangeExecutionState) -> bool {
    matches!(
        state,
        ChangeExecutionState::Queued
            | ChangeExecutionState::Active
            | ChangeExecutionState::Waiting
            | ChangeExecutionState::Stopping
    )
}

/// How one admitted execution episode ended, in the owner's own typed terms.
///
/// Derived from the reducer's terminal state, never from a display string, an
/// error body, or a change disappearing from the snapshot. `Completed` here is a
/// *claim worth verifying*, not a completion: the repository oracle decides.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EpisodeTerminal {
    /// The reducer reached a terminal success for this change.
    Completed,
    /// The reducer reached a terminal failure or rejection.
    Failed,
    /// The episode was stopped or dequeued, including before any active work.
    Stopped,
}

impl EpisodeTerminal {
    /// Stable wire token.
    // Read by the episode-vocabulary assertions; the binary maps the enum onto
    // the API's own event type instead.
    #[cfg_attr(not(test), allow(dead_code))]
    pub fn as_str(self) -> &'static str {
        match self {
            Self::Completed => "completed",
            Self::Failed => "failed",
            Self::Stopped => "stopped",
        }
    }
}

/// What happened to one execution episode.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EpisodeTransitionKind {
    /// A non-admitted change became admitted work; the episode ID is new.
    Started,
    /// The open episode entered a typed blocked/waiting condition.
    BlockedEntered,
    /// The open episode left that condition, arming the next attention edge.
    BlockedLeft,
    /// The episode settled; no further transition can belong to this ID.
    Terminal(EpisodeTerminal),
}

/// One published episode transition.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EpisodeTransition {
    /// Change the episode belongs to.
    pub change_id: String,
    /// Process-local episode identity.
    pub execution_id: String,
    /// What happened.
    pub kind: EpisodeTransitionKind,
}

/// A process-local consumer of episode transitions.
///
/// Observability only. An observer cannot refuse, delay, or redirect a
/// transition, and nothing it does is read back as a workflow input — it is
/// called after the store's own lock is released precisely so it can never
/// influence the projection it is describing.
pub trait EpisodeObserver: std::fmt::Debug + Send + Sync {
    /// Absorb one transition. Must not block for long.
    fn observe_episode(&self, transition: &EpisodeTransition);
}

#[derive(Debug, Clone, Default)]
struct ChangeFactsState {
    /// Identity of the most recent admitted episode, retained after it settles.
    execution_id: Option<String>,
    /// True while that episode is still admitted, which is what makes the next
    /// admission a *new* episode rather than a continuation of this one.
    episode_open: bool,
    /// True while the open episode is inside a blocked attention edge.
    blocked_edge: bool,
    execution_state: ChangeExecutionState,
    current_phase: ExecutionPhase,
    phase_started_at: Option<DateTime<Utc>>,
    last_completed_phase: Option<ExecutionPhase>,
    last_completed_at: Option<DateTime<Utc>>,
    apply_commit_oid: Option<String>,
    push_open: bool,
}

/// Bounded set of dispatch identities this store has already absorbed.
///
/// Two boundaries feed the store — the authoritative dispatch owner and the web
/// projection sink — and in a process that has both, one dispatch reaches it
/// twice. Absorbing an identity once is what stops the second delivery from
/// restamping a completion boundary with a later instant.
#[derive(Debug, Default)]
struct AbsorbedDispatches {
    order: VecDeque<u64>,
    seen: HashSet<u64>,
}

impl AbsorbedDispatches {
    const CAPACITY: usize = 1024;

    fn admit(&mut self, id: u64) -> bool {
        if !self.seen.insert(id) {
            return false;
        }
        self.order.push_back(id);
        while self.order.len() > Self::CAPACITY {
            if let Some(evicted) = self.order.pop_front() {
                self.seen.remove(&evicted);
            }
        }
        true
    }
}

#[derive(Debug, Default)]
struct Inner {
    changes: HashMap<String, ChangeFactsState>,
    activities: HashSet<ProcessActivity>,
    stop_requested: bool,
    absorbed: AbsorbedDispatches,
    /// Episode transitions produced by the current absorption, drained and
    /// published *after* the store lock is released so an observer can never
    /// re-enter the store from inside its own mutex.
    pending: Vec<EpisodeTransition>,
}

/// The shared process-local execution-facts store.
#[derive(Debug, Default)]
pub struct ExecutionFactsStore {
    inner: Mutex<Inner>,
    /// Late-bound episode consumer. Unbound is the ordinary case: a build or a
    /// frontend with no completion-sink dispatcher still tracks episodes, it
    /// simply has nobody to tell.
    observer: std::sync::RwLock<Option<std::sync::Arc<dyn EpisodeObserver>>>,
}

impl ExecutionFactsStore {
    /// A store for a fresh process incarnation.
    pub fn new() -> Self {
        Self::default()
    }

    /// Bind the process-local episode consumer.
    ///
    /// Idempotent replacement rather than a list: exactly one dispatcher owns
    /// completion sinks in a process, and a second one would double-deliver.
    pub fn bind_episode_observer(&self, observer: std::sync::Arc<dyn EpisodeObserver>) {
        *self
            .observer
            .write()
            .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(observer);
    }

    /// Publish drained transitions without holding the store lock.
    fn publish(&self, transitions: Vec<EpisodeTransition>) {
        if transitions.is_empty() {
            return;
        }
        let observer = self
            .observer
            .read()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
            .clone();
        let Some(observer) = observer else {
            return;
        };
        for transition in &transitions {
            observer.observe_episode(transition);
        }
    }

    /// The most recent episode identity for one change, if this incarnation
    /// admitted it at all.
    // A convenience read over the same field `snapshot()` and `change()`
    // publish; the API resources reach it through those.
    #[cfg_attr(not(test), allow(dead_code))]
    pub fn execution_id(&self, change_id: &str) -> Option<String> {
        self.lock()
            .changes
            .get(change_id)
            .and_then(|facts| facts.execution_id.clone())
    }

    /// The change one episode identity belongs to, if this incarnation owns it.
    ///
    /// Linear over tracked changes on purpose: the map is bounded by the number
    /// of changes a process tracks, and a second index would be one more thing
    /// to keep consistent with the reducer.
    // The sink registry keeps its own binding, so this exists for assertions
    // that the store's own view agrees with it.
    #[cfg_attr(not(test), allow(dead_code))]
    pub fn change_of_execution(&self, execution_id: &str) -> Option<String> {
        self.lock()
            .changes
            .iter()
            .find(|(_, facts)| facts.execution_id.as_deref() == Some(execution_id))
            .map(|(id, _)| id.clone())
    }

    fn lock(&self) -> MutexGuard<'_, Inner> {
        self.inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
    }

    /// Absorb one typed dispatch and the reducer state it produced.
    ///
    /// `state` is `Some` only for a state-owning dispatch, which is the only
    /// kind that can have moved a phase. A log or presentation event still
    /// reaches here because a few typed episodes (analysis, cleanup, branch
    /// merge, conflict resolution) are presentation-owned yet are real work.
    ///
    /// `dispatch_id` makes absorption exactly-once: a repeated delivery of the
    /// same dispatch returns without restamping any boundary. Reports whether
    /// this delivery was the first.
    pub fn observe(
        &self,
        dispatch_id: u64,
        event: &ExecutionEvent,
        state: Option<&OrchestratorState>,
        now: DateTime<Utc>,
    ) -> bool {
        let transitions = {
            let mut inner = self.lock();
            if !inner.absorbed.admit(dispatch_id) {
                return false;
            }
            Self::observe_completion(&mut inner, event, now);
            Self::observe_push(&mut inner, event);
            Self::observe_process(&mut inner, event);
            if let Some(state) = state {
                Self::refresh_from_reducer(&mut inner, state, now);
            }
            std::mem::take(&mut inner.pending)
        };
        self.publish(transitions);
        true
    }

    /// Typed completion facts. Only these set the last completed phase.
    fn observe_completion(inner: &mut Inner, event: &ExecutionEvent, now: DateTime<Utc>) {
        use ExecutionEvent as E;
        let (change_id, phase) = match event {
            E::WorkspacePreparationEnded { change_id } => (change_id, ExecutionPhase::Preparing),
            E::ApplyCompleted {
                change_id,
                revision,
            } => {
                // The OID is retained *before* the completion is recorded so a
                // settlement racing this dispatch cannot see the phase without
                // the evidence that explains it. An empty revision is no
                // evidence at all and is never stored as one.
                if !revision.trim().is_empty() {
                    inner
                        .changes
                        .entry(change_id.clone())
                        .or_default()
                        .apply_commit_oid = Some(revision.trim().to_string());
                }
                (change_id, ExecutionPhase::Apply)
            }
            E::AcceptanceCompleted { change_id } => (change_id, ExecutionPhase::Acceptance),
            E::RejectionReviewCompleted { change_id, .. } => {
                (change_id, ExecutionPhase::RejectionReview)
            }
            E::ChangeArchived(change_id) => (change_id, ExecutionPhase::Archive),
            E::ResolveCompleted { change_id, .. } => (change_id, ExecutionPhase::Resolve),
            E::MergeCompleted { change_id, .. } => (change_id, ExecutionPhase::Merge),
            E::PushCompleted { change_id, .. } => (change_id, ExecutionPhase::Push),
            _ => return,
        };
        let facts = inner.changes.entry(change_id.clone()).or_default();
        facts.last_completed_phase = Some(phase);
        facts.last_completed_at = Some(now);
    }

    /// The typed push episode, which the reducer does not track as an activity.
    fn observe_push(inner: &mut Inner, event: &ExecutionEvent) {
        use ExecutionEvent as E;
        let (change_id, open) = match event {
            E::PushStarted { change_id, .. } => (change_id, true),
            E::PushCompleted { change_id, .. } | E::PushFailed { change_id, .. } => {
                (change_id, false)
            }
            _ => return,
        };
        inner
            .changes
            .entry(change_id.clone())
            .or_default()
            .push_open = open;
    }

    /// Process-level episodes and the graceful-stop qualifier.
    fn observe_process(inner: &mut Inner, event: &ExecutionEvent) {
        use ExecutionEvent as E;
        match event {
            E::AnalysisStarted { .. } => {
                inner.activities.insert(ProcessActivity::DependencyAnalysis);
            }
            E::AnalysisCompleted { .. } => {
                inner
                    .activities
                    .remove(&ProcessActivity::DependencyAnalysis);
            }
            E::MergeStarted { .. } => {
                inner.activities.insert(ProcessActivity::BaseBranchMerge);
            }
            // The base-lane episode has no terminal event of its own: it ends
            // when the batch it merges reaches a per-change outcome. Closing it
            // on any of those is what stops a finished lane from reporting work
            // forever.
            E::MergeCompleted { .. } | E::MergeDeferred { .. } | E::ResolveFailed { .. } => {
                inner.activities.remove(&ProcessActivity::BaseBranchMerge);
            }
            E::ConflictResolutionStarted => {
                inner.activities.insert(ProcessActivity::ConflictResolution);
            }
            E::ConflictResolutionCompleted | E::ConflictResolutionFailed { .. } => {
                inner
                    .activities
                    .remove(&ProcessActivity::ConflictResolution);
            }
            E::BranchMergeStarted { .. } => {
                inner.activities.insert(ProcessActivity::BranchMerge);
            }
            E::BranchMergeCompleted { .. } | E::BranchMergeFailed { .. } => {
                inner.activities.remove(&ProcessActivity::BranchMerge);
            }
            E::CleanupStarted { .. } => {
                inner.activities.insert(ProcessActivity::WorkspaceCleanup);
            }
            E::CleanupCompleted { .. } => {
                inner.activities.remove(&ProcessActivity::WorkspaceCleanup);
            }
            E::Stopping => inner.stop_requested = true,
            // A process terminal ends every episode. Anything still open at that
            // point is a lane that died with the process, not running work.
            E::Stopped | E::Error { .. } | E::AllCompleted => {
                inner.activities.clear();
                inner.stop_requested = false;
            }
            E::ProcessingStarted(_) => inner.stop_requested = false,
            _ => {}
        }
    }

    /// Re-project every tracked change from the reducer state.
    ///
    /// Every change is refreshed, not only the event's target: one typed
    /// transition can move another row (a released merge wait, a revoked
    /// dequeue), and a store that only followed the addressed change would
    /// publish a phase the reducer had already left.
    fn refresh_from_reducer(inner: &mut Inner, state: &OrchestratorState, now: DateTime<Utc>) {
        let stop_requested = inner.stop_requested;
        // Split the borrow so one pass can both update a change's facts and
        // append the episode transition that update produced.
        let Inner {
            changes, pending, ..
        } = inner;
        for change_id in state.tracked_change_ids() {
            let Some(runtime) = state.change_runtime(&change_id) else {
                continue;
            };
            let change_key = change_id.clone();
            let facts = changes.entry(change_id).or_default();
            let phase = project_phase(runtime, facts.push_open);
            if phase != facts.current_phase {
                facts.current_phase = phase;
                facts.phase_started_at = phase.is_active().then_some(now);
            }
            let next = project_execution_state(runtime, phase, stop_requested);
            facts.execution_state = next;
            Self::advance_episode(pending, &change_key, facts, next);
        }
    }

    /// Move one change's execution episode in step with its projected state.
    ///
    /// The rule is deliberately narrow: admission opens an episode, leaving
    /// admission settles it, and nothing else creates identity. A change that
    /// simply stops being tracked never settles here — disappearance proves
    /// nothing, and inventing a terminal for it is exactly the lie the
    /// completion contract exists to prevent.
    fn advance_episode(
        pending: &mut Vec<EpisodeTransition>,
        change_id: &str,
        facts: &mut ChangeFactsState,
        next: ChangeExecutionState,
    ) {
        let admitted = is_admitted_execution_state(next);

        if admitted && !facts.episode_open {
            let execution_id = crate::ids::new_hex_id();
            facts.execution_id = Some(execution_id.clone());
            facts.episode_open = true;
            facts.blocked_edge = false;
            pending.push(EpisodeTransition {
                change_id: change_id.to_string(),
                execution_id,
                kind: EpisodeTransitionKind::Started,
            });
        }

        let Some(execution_id) = facts.execution_id.clone() else {
            return;
        };

        if !facts.episode_open {
            return;
        }

        if admitted {
            // Attention is edge-triggered: an unchanged waiting state publishes
            // nothing, while leaving and re-entering it arms a new edge.
            let blocked = matches!(next, ChangeExecutionState::Waiting);
            if blocked != facts.blocked_edge {
                facts.blocked_edge = blocked;
                pending.push(EpisodeTransition {
                    change_id: change_id.to_string(),
                    execution_id,
                    kind: if blocked {
                        EpisodeTransitionKind::BlockedEntered
                    } else {
                        EpisodeTransitionKind::BlockedLeft
                    },
                });
            }
            return;
        }

        let terminal = match next {
            ChangeExecutionState::Completed => EpisodeTerminal::Completed,
            ChangeExecutionState::Failed => EpisodeTerminal::Failed,
            // `Stopped` is settled stop/dequeue removal. `Unknown` reaches here
            // only by leaving admission without a typed terminal — a revoked
            // queue intent — which is the same fact from the caller's side.
            ChangeExecutionState::Stopped | ChangeExecutionState::Unknown => {
                EpisodeTerminal::Stopped
            }
            // Unreachable: every remaining variant is an admitted state.
            _ => return,
        };
        facts.episode_open = false;
        facts.blocked_edge = false;
        pending.push(EpisodeTransition {
            change_id: change_id.to_string(),
            execution_id,
            kind: EpisodeTransitionKind::Terminal(terminal),
        });
    }

    /// A coherent read of every fact this incarnation has observed.
    pub fn snapshot(&self) -> ExecutionFactsSnapshot {
        let inner = self.lock();
        let mut activities: Vec<ProcessActivity> = inner.activities.iter().copied().collect();
        activities.sort();
        ExecutionFactsSnapshot {
            changes: inner
                .changes
                .iter()
                .map(|(id, facts)| {
                    (
                        id.clone(),
                        ChangeExecutionFacts {
                            execution_id: facts.execution_id.clone(),
                            execution_state: facts.execution_state,
                            current_phase: facts.current_phase,
                            phase_started_at: facts.phase_started_at,
                            last_completed_phase: facts.last_completed_phase,
                            last_completed_at: facts.last_completed_at,
                            apply_commit_oid: facts.apply_commit_oid.clone(),
                        },
                    )
                })
                .collect(),
            activities,
        }
    }

    /// One change's facts without cloning the whole store.
    pub fn change(&self, change_id: &str) -> ChangeExecutionFacts {
        let inner = self.lock();
        inner
            .changes
            .get(change_id)
            .map(|facts| ChangeExecutionFacts {
                execution_id: facts.execution_id.clone(),
                execution_state: facts.execution_state,
                current_phase: facts.current_phase,
                phase_started_at: facts.phase_started_at,
                last_completed_phase: facts.last_completed_phase,
                last_completed_at: facts.last_completed_at,
                apply_commit_oid: facts.apply_commit_oid.clone(),
            })
            .unwrap_or_else(ChangeExecutionFacts::unknown)
    }

    /// The retained Apply-completion OID for a change, if this incarnation saw one.
    ///
    /// Empty after a restart by construction, which is the whole reason Apply
    /// commit presence is nullable: a process that never observed the completion
    /// has no typed evidence and must not guess from the repository alone.
    pub fn apply_commit_oid(&self, change_id: &str) -> Option<String> {
        self.lock()
            .changes
            .get(change_id)
            .and_then(|facts| facts.apply_commit_oid.clone())
    }
}

/// Project the reducer's activity onto the closed phase vocabulary.
///
/// The reducer is the sole authority for everything except `Push`. Publication
/// has no activity of its own there — it reuses `Resolving`, because both hold
/// the base-mutating lane — so the typed push episode is what tells the two
/// apart. It overrides only the two activities publication can legitimately be
/// running under; any other activity is a newer reducer transition and wins.
pub fn project_phase(runtime: &ChangeRuntimeState, push_open: bool) -> ExecutionPhase {
    if push_open
        && matches!(
            runtime.activity,
            ActivityState::Idle | ActivityState::Resolving
        )
    {
        return ExecutionPhase::Push;
    }
    match runtime.activity {
        ActivityState::Preparing => ExecutionPhase::Preparing,
        ActivityState::Applying => ExecutionPhase::Apply,
        ActivityState::Accepting => ExecutionPhase::Acceptance,
        ActivityState::Rejecting => ExecutionPhase::RejectionReview,
        ActivityState::Archiving => ExecutionPhase::Archive,
        ActivityState::Resolving => ExecutionPhase::Resolve,
        ActivityState::Idle => ExecutionPhase::None,
    }
}

/// Project the reducer's runtime facts onto the closed execution-state vocabulary.
pub fn project_execution_state(
    runtime: &ChangeRuntimeState,
    phase: ExecutionPhase,
    stop_requested: bool,
) -> ChangeExecutionState {
    match &runtime.terminal {
        TerminalState::Merged | TerminalState::Pushed => return ChangeExecutionState::Completed,
        TerminalState::Rejected(_) | TerminalState::Error(_) => {
            return ChangeExecutionState::Failed
        }
        TerminalState::Stopped => return ChangeExecutionState::Stopped,
        TerminalState::None => {}
    }
    if runtime.dequeued {
        return ChangeExecutionState::Stopped;
    }
    if phase.is_active() {
        return if stop_requested {
            ChangeExecutionState::Stopping
        } else {
            ChangeExecutionState::Active
        };
    }
    // A dependency wait is a dispatch exclusion applied to admitted queue
    // intent: no execution episode has started, and the retained intent is what
    // will dispatch the change once the dependency resolves. It therefore stays
    // `queued` here, and the `blocked` display plus the structured dependency
    // blocker are what explain *why* the slot is empty. Every other wait state
    // holds an episode that already began, so it stays `waiting`.
    if matches!(runtime.wait_state, WaitState::DependencyBlocked)
        && matches!(runtime.queue_intent, QueueIntent::Queued)
    {
        return ChangeExecutionState::Queued;
    }
    if !matches!(runtime.wait_state, WaitState::None) {
        return ChangeExecutionState::Waiting;
    }
    if matches!(runtime.queue_intent, QueueIntent::Queued) {
        return ChangeExecutionState::Queued;
    }
    ChangeExecutionState::Unknown
}

#[cfg(test)]
mod tests;