Skip to main content

monoloop_loop/transaction/lifecycle/
supervisor.rs

1//! Supervisor command loop and terminal authority (v2 §10 / §13).
2//!
3//! Start commands use a dedicated bounded queue (capacity ≥ max_active).
4//! Cancel / shutdown use a separate control queue. WorkerExited uses a third
5//! worker queue so coordinators cannot starve start or control.
6
7use super::coordinator::{run_coordinator, CoordinatorParams, WorkerMessage};
8use super::event_publisher::run_event_publisher;
9use super::ledger::{LifecycleLedger, TransactionPhase};
10use super::task_spawner::{SpawnRequest, TransactionTaskSpawner};
11use super::task_supervisor::{TaskClass, TaskExit, TaskSupervisor};
12use super::terminal::{build_completion, end_event, TerminalDecision, TerminalProposal};
13use crate::transaction::bootstrap::{
14    ControlHoldGate, FinalizerHoldGate, JoinOnlySpillInject, StartHoldGate, StoppedGate,
15};
16use crate::transaction::channel_registry::LiveChannel;
17use crate::transaction::dispatcher::OrphanToolPermitSet;
18use crate::transaction::host_tools::HostToolRegistry;
19use crate::transaction::mcp::{McpGatewayHandle, PreparedMcpGateway};
20use crate::transaction::tool_capacity::SharedToolCapacity;
21use monoloop_contracts::{
22    ChannelId, CleanupStatus, CompletionPublishResult, ShutdownReport, ShutdownSnapshot,
23    TerminalEventDelivery, TransactionEndKind, TransactionId,
24};
25use std::collections::HashMap;
26use std::sync::atomic::{AtomicU32, AtomicU64, AtomicU8, Ordering};
27use std::sync::{Arc, Mutex};
28use std::time::Duration;
29use tokio::sync::{mpsc, oneshot, Notify};
30use tokio_util::sync::CancellationToken;
31
32/// Start-queue commands (admission only).
33#[derive(Clone, Debug, PartialEq, Eq)]
34pub enum StartCommand {
35    /// Start an admitted queued transaction.
36    Start(TransactionId),
37}
38
39/// Control-queue commands (cancel / shutdown).
40#[derive(Clone, Debug, PartialEq, Eq)]
41pub enum ControlCommand {
42    /// Request cooperative cancel.
43    Cancel(TransactionId),
44    /// Request forced terminate.
45    ForceTerminate(TransactionId),
46    /// Begin shutdown generation (idempotent).
47    BeginShutdown,
48    /// Driver thread asks the supervisor to stop after quiesce work.
49    StopSupervisor,
50}
51
52/// Backward-compatible aggregate name for exports.
53#[derive(Clone, Debug, PartialEq, Eq)]
54pub enum SupervisorCommand {
55    /// Start.
56    Start(TransactionId),
57    /// Cancel.
58    Cancel(TransactionId),
59    /// Force terminate.
60    ForceTerminate(TransactionId),
61    /// Begin shutdown.
62    BeginShutdown,
63    /// Stop supervisor.
64    StopSupervisor,
65}
66
67pub(crate) const STATE_STARTING: u8 = 0;
68pub(crate) const STATE_ACCEPTING: u8 = 1;
69pub(crate) const STATE_QUIESCING: u8 = 2;
70pub(crate) const STATE_STOPPED: u8 = 3;
71
72/// Shared runtime state visible to admission handle and owner.
73pub(crate) struct RuntimeShared {
74    pub state: AtomicU8,
75    pub ledger: Mutex<LifecycleLedger>,
76    pub start_tx: mpsc::Sender<StartCommand>,
77    pub control_tx: mpsc::Sender<ControlCommand>,
78    pub worker_tx: mpsc::Sender<WorkerMessage>,
79    pub wake: Notify,
80    pub channels: Arc<HashMap<ChannelId, LiveChannel>>,
81    pub default_deadline: Duration,
82    pub cleanup_deadline: Duration,
83    /// Independent budget for terminal `Ended` enqueue (spec §8 / D-047).
84    pub terminal_event_delivery_deadline: Duration,
85    pub task_spawner: TransactionTaskSpawner,
86    pub shutdown_generation: AtomicU64,
87    pub shutdown_report: Mutex<Option<ShutdownReport>>,
88    pub completions_published: AtomicU64,
89    pub completions_receiver_dropped: AtomicU64,
90    pub completions_invariant_failed: AtomicU64,
91    pub runtime_shutdown_terminals: AtomicU64,
92    /// Live TaskSupervisor task count (updated by supervisor loop).
93    pub owned_tasks: AtomicU32,
94    /// Live [`TaskClass::ConnectorOwner`] count (D-051 register-before-I/O).
95    pub live_connector_owners: AtomicU32,
96    /// When true, supervisor owns a loopback MCP gateway as RuntimeService (D-043 / §17).
97    pub enable_mcp_listener: bool,
98    /// Bound address once the MCP gateway task reports ready.
99    pub mcp_listen_addr: Mutex<Option<std::net::SocketAddr>>,
100    /// Cloneable MCP install/activate handle while the RuntimeService is live.
101    pub mcp_gateway: Mutex<Option<McpGatewayHandle>>,
102    /// Cancels the MCP axum serve future on quiesce.
103    pub mcp_cancel: Mutex<Option<CancellationToken>>,
104    /// When `Some`, defer drain-complete until the gate is released (§22.5).
105    pub block_stopped: Option<Arc<StoppedGate>>,
106    /// When `Some` and held, supervisor does not drain the start queue (D-040).
107    pub hold_start: Option<Arc<StartHoldGate>>,
108    /// When `Some` and held, supervisor does not drain the control queue (§23).
109    pub hold_control: Option<Arc<ControlHoldGate>>,
110    /// When `Some`, Finalizer waits after Seal before completion send (§22.2).
111    pub hold_finalizer_after_seal: Option<Arc<FinalizerHoldGate>>,
112    /// When `Some`, executor thread waits after drain before teardown (D-049).
113    pub hold_executor_teardown: Option<Arc<StoppedGate>>,
114    /// When `Some`, spawn a never-awaiting RuntimeService that signals before park.
115    pub inject_non_yielding_service: Option<Arc<std::sync::atomic::AtomicBool>>,
116    /// When `Some`, register TaskSupervisor-owned JoinOnly-style work at start.
117    pub inject_join_only_spill: Option<Arc<JoinOnlySpillInject>>,
118    /// Supervisor finished drain; OS thread may still be in executor teardown (D-049).
119    /// Public `Stopped` is published by the owner only after thread join.
120    pub drain_complete: std::sync::atomic::AtomicBool,
121    /// Host tool definitions available to admission / coordinators.
122    pub tools_registry: HostToolRegistry,
123    /// Process-wide concurrent tool execution budget.
124    pub shared_tool_capacity: Arc<SharedToolCapacity>,
125    /// Runtime-scoped unfinished tool joins/permits (Law 8 — not process-global).
126    pub tool_spill: Arc<OrphanToolPermitSet>,
127    /// Live ProcessIsolated OS children (§18.2 ShutdownSnapshot.owned_processes).
128    pub owned_processes: Arc<AtomicU32>,
129    /// ProcessIsolated children retained until OS exit is observed (D-048).
130    pub process_registry: Arc<crate::transaction::owned_process_registry::OwnedProcessRegistry>,
131    /// Runtime admission / capacity limits (D-035 input accounting uses these).
132    pub transaction_limits: monoloop_contracts::TransactionLimits,
133}
134
135impl RuntimeShared {
136    fn start_drain_enabled(&self) -> bool {
137        !self.hold_start.as_ref().is_some_and(|g| g.is_held())
138    }
139
140    fn control_drain_enabled(&self) -> bool {
141        !self.hold_control.as_ref().is_some_and(|g| g.is_held())
142    }
143
144    pub fn runtime_state(&self) -> super::super::state::RuntimeState {
145        match self.state.load(Ordering::SeqCst) {
146            STATE_ACCEPTING => super::super::state::RuntimeState::Accepting,
147            STATE_QUIESCING => super::super::state::RuntimeState::Quiescing,
148            STATE_STOPPED => super::super::state::RuntimeState::Stopped,
149            _ => super::super::state::RuntimeState::Starting,
150        }
151    }
152
153    pub fn snapshot(&self) -> ShutdownSnapshot {
154        let ledger_entries = self.ledger.lock().map(|l| l.len()).unwrap_or(0) as u32;
155        let mcp_routes = self
156            .mcp_gateway
157            .lock()
158            .ok()
159            .and_then(|g| g.as_ref().map(|h| h.routes().len() as u32))
160            .unwrap_or(0);
161        // Prefer registry live count (D-048); fall back to lease counter.
162        let owned_processes =
163            self.process_registry
164                .live_count()
165                .max(self.owned_processes.load(Ordering::SeqCst) as usize) as u32;
166        ShutdownSnapshot {
167            generation: self.shutdown_generation.load(Ordering::SeqCst),
168            ledger_entries,
169            owned_tasks: self.owned_tasks.load(Ordering::SeqCst),
170            owned_processes,
171            mcp_routes,
172            completions_published: self.completions_published.load(Ordering::SeqCst),
173        }
174    }
175
176    pub fn final_report(&self) -> ShutdownReport {
177        ShutdownReport {
178            completions_published: self.completions_published.load(Ordering::SeqCst),
179            completions_receiver_dropped: self.completions_receiver_dropped.load(Ordering::SeqCst),
180            completions_invariant_failed: self.completions_invariant_failed.load(Ordering::SeqCst),
181            runtime_shutdown_terminals: self.runtime_shutdown_terminals.load(Ordering::SeqCst),
182        }
183    }
184}
185
186/// Run the supervisor until stopped invariants hold.
187pub(crate) async fn run_supervisor(
188    shared: Arc<RuntimeShared>,
189    mut start_rx: mpsc::Receiver<StartCommand>,
190    mut control_rx: mpsc::Receiver<ControlCommand>,
191    mut worker_rx: mpsc::Receiver<WorkerMessage>,
192    mut spawn_rx: mpsc::Receiver<SpawnRequest>,
193    mcp_prepared: Option<PreparedMcpGateway>,
194) {
195    let mut tasks = TaskSupervisor::new();
196    let mut stopping = false;
197    let mut quiesce_started: Option<tokio::time::Instant> = None;
198
199    if let Some(prepared) = mcp_prepared {
200        // Handle/addr already published before start ready (§7.1).
201        let serve = super::mcp_listener::serve_runtime_mcp(Arc::clone(&shared), prepared);
202        tasks.spawn(TaskClass::RuntimeService, serve);
203    }
204    // §22.3 sacrificial: never-awaiting future pins one worker; abort cannot
205    // join it, so shutdown must stay Quiescing (never false Stopped).
206    if let Some(entered) = shared.inject_non_yielding_service.clone() {
207        tasks.spawn(TaskClass::RuntimeService, async move {
208            // Signal before park so the harness does not shut down while the
209            // task is still only registered (abort-before-poll would join cleanly
210            // and falsely allow Stopped).
211            entered.store(true, Ordering::SeqCst);
212            // Park the worker thread with no `.await` — Tokio abort cannot stop
213            // this task (spin would burn a core for the same proof).
214            loop {
215                std::thread::park();
216            }
217        });
218    }
219    // §22.4 / Law 23 / M5.4: JoinOnly-style work under TaskSupervisor (not spill,
220    // not ambient tokio::spawn). Park the worker thread so abort cannot join
221    // until release() unparks — same abort-resistance as §22.3 sacrificial.
222    if let Some(inject) = shared.inject_join_only_spill.clone() {
223        tasks.spawn(TaskClass::RuntimeService, async move {
224            inject.store_parked_thread(std::thread::current());
225            inject.mark_entered();
226            loop {
227                if inject.is_released() {
228                    break;
229                }
230                std::thread::park();
231            }
232        });
233    }
234    shared
235        .owned_tasks
236        .store(tasks.registered_count() as u32, Ordering::SeqCst);
237
238    loop {
239        // Authoritative WorkerExited proposals before reaping coordinators, so
240        // on_task_exit does not invent a terminal over a queued proposal.
241        while let Ok(msg) = worker_rx.try_recv() {
242            let WorkerMessage::WorkerExited {
243                transaction_id,
244                proposal,
245            } = msg;
246            accept_terminal(&shared, &mut tasks, transaction_id, proposal, false);
247        }
248
249        // Preferential control drain (D-039): process BeginShutdown/StopSupervisor
250        // and Cancels without waiting behind a biased select recv, so a terminate
251        // flood cannot delay observing Quiescing. Skipped while `hold_control`
252        // is held so §23 max_actor_commands plus-one can fill the bound.
253        if shared.control_drain_enabled() {
254            while let Ok(cmd) = control_rx.try_recv() {
255                match cmd {
256                    ControlCommand::Cancel(tx) => {
257                        accept_terminal(
258                            &shared,
259                            &mut tasks,
260                            tx,
261                            TerminalProposal::new(TransactionEndKind::Cancelled),
262                            false,
263                        );
264                    }
265                    ControlCommand::ForceTerminate(tx) => {
266                        accept_terminal(
267                            &shared,
268                            &mut tasks,
269                            tx,
270                            TerminalProposal::new(TransactionEndKind::Terminated),
271                            true,
272                        );
273                    }
274                    ControlCommand::BeginShutdown => {
275                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
276                    }
277                    ControlCommand::StopSupervisor => {
278                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
279                        tasks.abort_all();
280                    }
281                }
282            }
283        }
284
285        for (_id, class, exit) in tasks.try_reap_finished() {
286            note_connector_owner_exit(&shared, &class);
287            on_task_exit(&shared, &mut tasks, &class, exit);
288        }
289
290        // Preferential drain: never starve spawn registration behind join_next.
291        while let Ok(req) = spawn_rx.try_recv() {
292            note_connector_owner_spawn(&shared, &req.class);
293            let id = tasks.spawn(req.class, req.future);
294            let _ = req.reply.send(id);
295        }
296        shared
297            .owned_tasks
298            .store(tasks.registered_count() as u32, Ordering::SeqCst);
299
300        if !stopping && shared.state.load(Ordering::SeqCst) == STATE_QUIESCING {
301            begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
302        }
303        if stopping && quiesce_started.is_none() {
304            quiesce_started = Some(tokio::time::Instant::now());
305        }
306
307        if stopping {
308            // Release orphan tool permits (capacity); joins are TaskSupervisor-owned.
309            let _ = shared.tool_spill.shutdown_progress();
310            // D-048: kill+poll ProcessIsolated children; do not clear without reap.
311            let _ = shared.process_registry.shutdown_progress();
312
313            // Drive residual abort every lap so coordinator/tool work cannot
314            // strand Finalizer → tombstone → Stopped.
315            let tx_ids: Vec<TransactionId> = {
316                let ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
317                ledger.transaction_ids()
318            };
319            for tx in &tx_ids {
320                tasks.abort_transaction_residuals(tx);
321            }
322
323            // Failsafe: idle rows clear only after Finalizing/CleanupPending
324            // (one completion path started). Laws 8–9 + one-active key; admission
325            // already closed.
326            let idle: Vec<TransactionId> = tx_ids
327                .iter()
328                .copied()
329                .filter(|tx| tasks.tasks_for(tx).is_empty())
330                .collect();
331            for tx in idle {
332                try_remove_tombstone(&shared, &tx);
333            }
334
335            // Hard quiesce deadline (D-045 / §22.2): never abort Finalizer —
336            // EventPublisher may be aborted after grace (Seal already tried or
337            // stuck); clear row only when no tx tasks remain so Seal→completion
338            // cannot lose the one completion attempt.
339            let grace = shared.cleanup_deadline.max(Duration::from_secs(2));
340            if quiesce_started.is_some_and(|t| t.elapsed() >= grace) {
341                let leftover: Vec<TransactionId> = {
342                    let ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
343                    ledger.transaction_ids()
344                };
345                for tx in &leftover {
346                    tasks.abort_transaction_except_finalizer(tx);
347                    if tasks.tasks_for(tx).is_empty() {
348                        force_remove_tombstone(&shared, tx);
349                    }
350                }
351            }
352        }
353
354        let mut ready_to_stop = false;
355        if stopping && shared.ledger.lock().map(|l| l.is_empty()).unwrap_or(true) {
356            // Reject any queued spawns so workers cannot block on reply forever.
357            while let Ok(req) = spawn_rx.try_recv() {
358                drop(req);
359            }
360            while let Ok(msg) = worker_rx.try_recv() {
361                let WorkerMessage::WorkerExited {
362                    transaction_id,
363                    proposal,
364                } = msg;
365                accept_terminal(&shared, &mut tasks, transaction_id, proposal, false);
366            }
367            // Ledger may have been re-populated if a late WorkerExited arrived.
368            if shared.ledger.lock().map(|l| l.is_empty()).unwrap_or(true) {
369                // M5.4: orphan permits never block Stopped after quiesce release.
370                let _ = shared.tool_spill.shutdown_progress();
371                let _ = shared.process_registry.shutdown_progress();
372                // D-048: Stopped requires process registry empty (reaped), not
373                // merely owned_processes counter zero via lease Drop.
374                if tasks.is_empty() && shared.process_registry.is_empty() {
375                    ready_to_stop = true;
376                } else {
377                    // One bounded drain; on timeout fall through to select (50ms
378                    // tick) — never tight-loop abort_and_drain (§18.2 / Drop).
379                    if tasks.abort_and_drain().await {
380                        let _ = shared.tool_spill.shutdown_progress();
381                        let _ = shared.process_registry.shutdown_progress();
382                        ready_to_stop = shared.ledger.lock().map(|l| l.is_empty()).unwrap_or(true)
383                            && tasks.is_empty()
384                            && shared.process_registry.is_empty();
385                    } else {
386                        shared
387                            .owned_tasks
388                            .store(tasks.registered_count() as u32, Ordering::SeqCst);
389                    }
390                }
391            }
392        }
393        if ready_to_stop {
394            shared.owned_tasks.store(0, Ordering::SeqCst);
395            shared.live_connector_owners.store(0, Ordering::SeqCst);
396            // §22.5 test gate: hold Quiescing until release so TimedOut is deterministic.
397            if let Some(gate) = shared.block_stopped.as_ref() {
398                gate.wait_released().await;
399            }
400            // D-049: drain-complete ≠ public Stopped. Owner publishes Stopped
401            // only after the executor OS thread join is observed.
402            *shared
403                .shutdown_report
404                .lock()
405                .unwrap_or_else(|e| e.into_inner()) = Some(shared.final_report());
406            shared.drain_complete.store(true, Ordering::SeqCst);
407            shared.wake.notify_waiters();
408            return;
409        }
410
411        tokio::select! {
412            biased;
413            // Control / worker / start before join_next so Accepting is not
414            // starved when many tasks complete in a burst.
415            ctrl = control_rx.recv(), if shared.control_drain_enabled() => {
416                match ctrl {
417                    None => {
418                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
419                        tasks.abort_all();
420                    }
421                    Some(ControlCommand::Cancel(tx)) => {
422                        accept_terminal(&shared, &mut tasks, tx, TerminalProposal::new(TransactionEndKind::Cancelled), false);
423                    }
424                    Some(ControlCommand::ForceTerminate(tx)) => {
425                        accept_terminal(&shared, &mut tasks, tx, TerminalProposal::new(TransactionEndKind::Terminated), true);
426                    }
427                    Some(ControlCommand::BeginShutdown) => {
428                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
429                    }
430                    Some(ControlCommand::StopSupervisor) => {
431                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
432                        tasks.abort_all();
433                    }
434                }
435            }
436            msg = worker_rx.recv() => {
437                if let Some(WorkerMessage::WorkerExited { transaction_id, proposal }) = msg {
438                    accept_terminal(&shared, &mut tasks, transaction_id, proposal, false);
439                }
440            }
441            start = start_rx.recv(), if shared.start_drain_enabled() => {
442                if let Some(StartCommand::Start(tx)) = start {
443                    if shared.state.load(Ordering::SeqCst) == STATE_ACCEPTING {
444                        handle_start(&shared, &mut tasks, tx);
445                    } else if stopping {
446                        // Late Start after Quiescing: entry was already snapshotted
447                        // (or must still terminalize). Never drop a Queued admit.
448                        accept_terminal(
449                            &shared,
450                            &mut tasks,
451                            tx,
452                            TerminalProposal::new(TransactionEndKind::RuntimeShutdown),
453                            false,
454                        );
455                    }
456                }
457            }
458            // While start/control drain is held, poll so release is observed without a notify.
459            _ = tokio::time::sleep(Duration::from_millis(5)),
460                if !shared.start_drain_enabled() || !shared.control_drain_enabled() => {}
461            spawn = spawn_rx.recv() => {
462                match spawn {
463                    Some(req) => {
464                        note_connector_owner_spawn(&shared, &req.class);
465                        let id = tasks.spawn(req.class, req.future);
466                        let _ = req.reply.send(id);
467                    }
468                    None => {
469                        // Spawner dropped — treat as shutdown pressure.
470                        begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
471                    }
472                }
473            }
474            // While quiescing, wake periodically so grace / residual-abort
475            // checks run even if join_next stalls on a non-yielding task.
476            _ = tokio::time::sleep(Duration::from_millis(50)), if stopping => {}
477            _ = shared.wake.notified() => {
478                if shared.state.load(Ordering::SeqCst) == STATE_QUIESCING {
479                    begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
480                }
481            }
482            finished = tasks.join_next(), if !tasks.is_empty() => {
483                // Re-drain WorkerExited before inventing terminals. A coordinator
484                // can try_send + complete in the same wakeup as join_next; the
485                // loop-head try_recv alone does not cover that race under load.
486                while let Ok(msg) = worker_rx.try_recv() {
487                    let WorkerMessage::WorkerExited {
488                        transaction_id,
489                        proposal,
490                    } = msg;
491                    accept_terminal(&shared, &mut tasks, transaction_id, proposal, false);
492                }
493                if let Some((_id, class, exit)) = finished {
494                    note_connector_owner_exit(&shared, &class);
495                    on_task_exit(&shared, &mut tasks, &class, exit);
496                }
497            }
498        }
499    }
500}
501
502fn note_connector_owner_spawn(shared: &RuntimeShared, class: &TaskClass) {
503    if matches!(class, TaskClass::ConnectorOwner(_, _)) {
504        shared.live_connector_owners.fetch_add(1, Ordering::SeqCst);
505    }
506}
507
508fn note_connector_owner_exit(shared: &RuntimeShared, class: &TaskClass) {
509    if matches!(class, TaskClass::ConnectorOwner(_, _)) {
510        shared.live_connector_owners.fetch_sub(1, Ordering::SeqCst);
511    }
512}
513
514fn on_task_exit(
515    shared: &Arc<RuntimeShared>,
516    tasks: &mut TaskSupervisor,
517    class: &TaskClass,
518    exit: TaskExit,
519) {
520    let Some(tx) = class.transaction_id() else {
521        return;
522    };
523
524    // Recover lost WorkerExited: coordinator finished/panicked/aborted without
525    // an accepted terminal — install one so every admission gets a completion.
526    if matches!(class, TaskClass::TransactionCoordinator(_)) {
527        let (missing, cancelled, quiescing) = {
528            let ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
529            match ledger.get(&tx) {
530                Some(entry) => (
531                    entry.terminal.is_none(),
532                    entry.resources.cancel.is_cancelled(),
533                    shared.state.load(Ordering::SeqCst) == STATE_QUIESCING,
534                ),
535                None => (false, false, false),
536            }
537        };
538        if missing {
539            let stashed = {
540                let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
541                ledger
542                    .get_mut(&tx)
543                    .and_then(|e| e.pending_worker_proposal.take())
544            };
545            let proposal = if let Some(p) = stashed {
546                p
547            } else {
548                let kind = if matches!(exit, TaskExit::Panicked) {
549                    // §22.2: coordinator panic → one InvariantFailed completion.
550                    TransactionEndKind::InvariantFailed
551                } else if quiescing {
552                    TransactionEndKind::RuntimeShutdown
553                } else if cancelled || matches!(exit, TaskExit::Cancelled) {
554                    TransactionEndKind::Cancelled
555                } else {
556                    TransactionEndKind::InvariantFailed
557                };
558                TerminalProposal::new(kind)
559            };
560            accept_terminal(shared, tasks, tx, proposal, false);
561        }
562    }
563
564    if tasks.tasks_for(&tx).is_empty() {
565        try_remove_tombstone(shared, &tx);
566        return;
567    }
568    // Finalizer finished (published or aborted): abort residual tx-scoped work
569    // so the tombstone (and SessionKey) can clear once joins are observed.
570    if matches!(class, TaskClass::Finalizer(_)) {
571        tasks.abort_transaction(&tx);
572    }
573}
574
575fn handle_start(shared: &Arc<RuntimeShared>, tasks: &mut TaskSupervisor, tx: TransactionId) {
576    let (
577        cancel,
578        channel_id,
579        session_id,
580        input,
581        session_config,
582        effective_config,
583        event_tx,
584        absolute_deadline,
585        selected_tools,
586    ) = {
587        let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
588        let Some(entry) = ledger.get_mut(&tx) else {
589            return;
590        };
591        if entry.phase != TransactionPhase::Queued {
592            return;
593        }
594        let Some(delivery) = entry.delivery.take() else {
595            return;
596        };
597        entry.completion_tx = Some(delivery.completion_tx);
598        entry.phase = TransactionPhase::Running;
599        let session_id = entry.session_key.as_ref().map(|k| k.session_id.clone());
600        let selected_tools = entry.tools.clone();
601        // One absolute transaction deadline: invocation may shorten, not exceed
602        // RuntimeShared.default_deadline (== TransactionLimits.transaction_deadline).
603        // Cap duration so Instant::checked_add cannot panic on Duration::MAX-class values.
604        const MAX_TX_DEADLINE: std::time::Duration =
605            std::time::Duration::from_secs(365 * 24 * 3600);
606        let ceiling = shared.default_deadline.min(MAX_TX_DEADLINE);
607        let base = entry
608            .effective_config
609            .deadline
610            .unwrap_or(ceiling)
611            .min(ceiling)
612            .min(MAX_TX_DEADLINE);
613        let absolute_deadline = std::time::Instant::now()
614            .checked_add(base)
615            .expect("MAX_TX_DEADLINE is Instant-representable");
616        (
617            Arc::clone(&entry.resources.cancel),
618            entry.channel_id.clone(),
619            session_id,
620            entry.input.clone(),
621            entry.session_config.clone(),
622            entry.effective_config.clone(),
623            delivery.event_tx,
624            absolute_deadline,
625            selected_tools,
626        )
627    };
628
629    let (pub_admit, pub_rx) = super::event_publisher::OrdinaryCmdAdmit::channel(64);
630    // D-047: Seal never shares the ordinary command queue — capacity 1 priority path.
631    let (seal_tx, seal_rx) = mpsc::channel::<super::event_publisher::SealCommand>(1);
632    {
633        let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
634        if let Some(entry) = ledger.get_mut(&tx) {
635            entry.publisher_cmd_tx = Some(pub_admit.clone());
636            entry.publisher_seal_tx = Some(seal_tx);
637        }
638    }
639
640    let channel_id_pub = channel_id.clone();
641    let session_id_pub = session_id.clone();
642    let cancel_pub = Arc::clone(&cancel);
643    let deadline_pub = absolute_deadline;
644    let admit_pub = pub_admit.clone();
645    // Do not retain a Sender inside the publisher task — that prevented natural
646    // channel closure after Finalizer took the seal sender (D-047).
647    tasks.spawn(TaskClass::EventPublisher(tx), async move {
648        let _ = run_event_publisher(
649            tx,
650            channel_id_pub,
651            session_id_pub,
652            event_tx,
653            pub_rx,
654            admit_pub,
655            seal_rx,
656            cancel_pub,
657            deadline_pub,
658        )
659        .await;
660    });
661
662    let mcp_gateway = shared.mcp_gateway.lock().ok().and_then(|g| g.clone());
663    let params = CoordinatorParams {
664        transaction_id: tx,
665        cancel,
666        channel_id,
667        session_id,
668        input,
669        session_config,
670        effective_config,
671        channels: Arc::clone(&shared.channels),
672        publish_tx: pub_admit,
673        worker_tx: shared.worker_tx.clone(),
674        tasks: shared.task_spawner.clone(),
675        deadline: absolute_deadline,
676        cleanup_deadline: shared.cleanup_deadline,
677        selected_tools,
678        tools_registry: shared.tools_registry.clone(),
679        shared_tool_capacity: Arc::clone(&shared.shared_tool_capacity),
680        tool_spill: Arc::clone(&shared.tool_spill),
681        owned_processes: Arc::clone(&shared.owned_processes),
682        process_registry: Arc::clone(&shared.process_registry),
683        mcp_gateway,
684        shared: Arc::clone(shared),
685    };
686    tasks.spawn(TaskClass::TransactionCoordinator(tx), async move {
687        run_coordinator(params).await;
688    });
689}
690
691fn accept_terminal(
692    shared: &Arc<RuntimeShared>,
693    tasks: &mut TaskSupervisor,
694    tx: TransactionId,
695    proposal: TerminalProposal,
696    force_upgrade: bool,
697) {
698    let (first_decision, seal_tx, kind) = {
699        let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
700        let Some(entry) = ledger.get_mut(&tx) else {
701            return;
702        };
703        let first = entry.terminal.is_none();
704        entry.pending_worker_proposal = None;
705        if let Some(existing) = entry.terminal.as_ref() {
706            if force_upgrade
707                && existing.kind == TransactionEndKind::Cancelled
708                && proposal.kind == TransactionEndKind::Terminated
709            {
710                entry.terminal = Some(TerminalDecision::new(TransactionEndKind::Terminated));
711            }
712        } else {
713            entry.terminal = Some(TerminalDecision::new(proposal.kind));
714            entry.phase = TransactionPhase::Finalizing;
715        }
716        entry.resources.cancel.cancel();
717        let kind = entry
718            .terminal
719            .as_ref()
720            .map(|t| t.kind)
721            .unwrap_or(proposal.kind);
722        // Take Seal sender so Finalizer owns the only remaining publisher control.
723        (first, entry.publisher_seal_tx.take(), kind)
724    };
725
726    // Wake coordinator; do not abort publisher until Seal is sent.
727    if let Ok(mut ledger) = shared.ledger.lock() {
728        if let Some(entry) = ledger.get_mut(&tx) {
729            entry.resources.cancel.cancel();
730        }
731    }
732
733    if first_decision {
734        let shared2 = Arc::clone(shared);
735        // Tx-scoped so tombstone stays until finalizer (+ other tx tasks) exit.
736        // Kind is re-read at Seal time so Cancel→Terminated upgrade can win (§22.2).
737        let _ = kind;
738        tasks.spawn(TaskClass::Finalizer(tx), async move {
739            finalize_after_terminal(shared2, tx, seal_tx).await;
740        });
741    }
742}
743
744async fn finalize_after_terminal(
745    shared: Arc<RuntimeShared>,
746    tx: TransactionId,
747    seal_tx: Option<mpsc::Sender<super::event_publisher::SealCommand>>,
748) {
749    // Brief yield so a racing ForceTerminate can upgrade Cancelled → Terminated
750    // in the ledger before we snapshot the kind (§22.2).
751    tokio::task::yield_now().await;
752
753    // Take completion sender *before* Seal so hard-grace / force_remove cannot
754    // strand the one completion attempt after the terminal-event try (§22.2).
755    let (channel_id, session_id, kind, completion_tx) = {
756        let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
757        let Some(entry) = ledger.get_mut(&tx) else {
758            return;
759        };
760        let kind = entry
761            .terminal
762            .as_ref()
763            .map(|t| t.kind)
764            .unwrap_or(TransactionEndKind::InvariantFailed);
765        if entry.completion_tx.is_none() {
766            if let Some(delivery) = entry.delivery.take() {
767                entry.completion_tx = Some(delivery.completion_tx);
768                drop(delivery.event_tx);
769            }
770        }
771        entry.phase = TransactionPhase::CleanupPending;
772        // Close ordinary admission before Seal so parked/new sends cannot cross
773        // the fence after the publisher begins draining (D-047 linearization).
774        if let Some(admit) = entry.publisher_cmd_tx.take() {
775            admit.close();
776        }
777        (
778            entry.channel_id.clone(),
779            entry.session_key.as_ref().map(|k| k.session_id.clone()),
780            kind,
781            entry.completion_tx.take(),
782        )
783    };
784
785    // D-041: never-attempted is not Published (spec §6.4).
786    let mut terminal_delivery = TerminalEventDelivery::NotAttempted;
787    let mut last_seq = 0u64;
788    let mut kind = kind;
789    if let Some(seal_tx) = seal_tx {
790        let (reply_tx, reply_rx) = oneshot::channel();
791        let terminal = end_event(tx, channel_id.clone(), session_id.clone(), kind, 0);
792        // Dedicated seal channel (cap 1): not blocked by a full ordinary cmd queue.
793        // One authoritative terminal-delivery Instant for Finalizer + publisher
794        // (`terminal_event_delivery_deadline` exactly — never cleanup/tx deadline,
795        // never a silent floor).
796        let seal_budget = shared.terminal_event_delivery_deadline;
797        let seal_deadline = std::time::Instant::now() + seal_budget;
798        match seal_tx.try_send(super::event_publisher::SealCommand {
799            terminal,
800            reply: reply_tx,
801            deadline: seal_deadline,
802        }) {
803            Ok(()) => {
804                // Publisher uses the same Instant. Small reply slack only — does
805                // not extend the configured terminal delivery budget itself.
806                let wait_budget = seal_budget.saturating_add(Duration::from_millis(100));
807                match tokio::time::timeout(wait_budget, reply_rx).await {
808                    Ok(Ok(res)) => {
809                        terminal_delivery = res.delivery;
810                        last_seq = res.last_sequence;
811                    }
812                    _ => {
813                        // Publisher must have timed out or died; do not leave a
814                        // path where Ended can still publish after completion.
815                        terminal_delivery = TerminalEventDelivery::DeadlineExceeded;
816                    }
817                }
818            }
819            Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => {
820                terminal_delivery = TerminalEventDelivery::QueueClosed;
821            }
822            Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => {
823                // Prior Seal already in flight (should not happen with take()).
824                terminal_delivery = TerminalEventDelivery::DeadlineExceeded;
825            }
826        }
827    }
828
829    // D-047: sticky ordinary/establish publication failure must not report
830    // Completed / ContinuationRequired with a truncated event stream.
831    if matches!(
832        terminal_delivery,
833        TerminalEventDelivery::LimitExceeded
834            | TerminalEventDelivery::DeadlineExceeded
835            | TerminalEventDelivery::QueueClosed
836    ) && matches!(
837        kind,
838        TransactionEndKind::Completed | TransactionEndKind::ContinuationRequired
839    ) {
840        kind = TransactionEndKind::EventDeliveryFailed;
841        if let Ok(mut ledger) = shared.ledger.lock() {
842            if let Some(entry) = ledger.get_mut(&tx) {
843                entry.terminal = Some(TerminalDecision::new(kind));
844            }
845        }
846    }
847
848    // §22.2 test gate: hold after Seal so shutdown cannot drop completion.
849    if let Some(gate) = shared.hold_finalizer_after_seal.as_ref() {
850        gate.wait_released().await;
851    }
852
853    // Best-effort: refresh sequence on the ledger row if it still exists.
854    if let Ok(mut ledger) = shared.ledger.lock() {
855        if let Some(entry) = ledger.get_mut(&tx) {
856            entry.event_sequence = last_seq;
857        }
858    }
859
860    let end = end_event(tx, channel_id, session_id, kind, last_seq);
861    // Snapshot live ownership at completion publish (§18.2 honesty — no hardcodes).
862    let owned_tasks = shared.owned_tasks.load(Ordering::SeqCst);
863    let owned_processes = shared.owned_processes.load(Ordering::SeqCst);
864    let cooperative_tools = shared.tool_spill.pending_count() as u32;
865    let completion = build_completion(
866        end,
867        terminal_delivery,
868        CleanupStatus::Pending {
869            owned_tasks,
870            owned_processes,
871            cooperative_tools,
872        },
873    );
874    if let Some(sender) = completion_tx {
875        match sender.send(completion) {
876            CompletionPublishResult::Published => {
877                shared.completions_published.fetch_add(1, Ordering::SeqCst);
878            }
879            CompletionPublishResult::ReceiverDropped => {
880                shared.completions_published.fetch_add(1, Ordering::SeqCst);
881                shared
882                    .completions_receiver_dropped
883                    .fetch_add(1, Ordering::SeqCst);
884            }
885            CompletionPublishResult::InvariantFailed => {
886                shared
887                    .completions_invariant_failed
888                    .fetch_add(1, Ordering::SeqCst);
889            }
890        }
891    } else {
892        shared
893            .completions_invariant_failed
894            .fetch_add(1, Ordering::SeqCst);
895    }
896
897    // Tombstone removal is deferred to `on_task_exit` once all tx-scoped tasks
898    // (including this Finalizer) have exited — keeps SessionKey reserved while
899    // residual work is live.
900}
901
902fn begin_shutdown_inner(
903    shared: &Arc<RuntimeShared>,
904    tasks: &mut TaskSupervisor,
905    stopping: &mut bool,
906) {
907    // §18.2 / D-010: flip Quiescing and snapshot ids under the admit install lock.
908    let ids = {
909        let ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
910        let _ = shared.state.compare_exchange(
911            STATE_ACCEPTING,
912            STATE_QUIESCING,
913            Ordering::SeqCst,
914            Ordering::SeqCst,
915        );
916        if shared.state.load(Ordering::SeqCst) == STATE_STARTING {
917            shared.state.store(STATE_QUIESCING, Ordering::SeqCst);
918        }
919        ledger.transaction_ids()
920    };
921    *stopping = true;
922    // Revoke MCP routes + cancel axum serve so RuntimeService can join (§17).
923    super::mcp_listener::signal_mcp_shutdown(shared);
924    // Wake any other runtime-wide waiters (D-043).
925    shared.wake.notify_waiters();
926
927    for tx in ids {
928        let already = {
929            let ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
930            ledger.get(&tx).and_then(|e| e.terminal.as_ref()).is_some()
931        };
932        if !already {
933            shared
934                .runtime_shutdown_terminals
935                .fetch_add(1, Ordering::SeqCst);
936            accept_terminal(
937                shared,
938                tasks,
939                tx,
940                TerminalProposal::new(TransactionEndKind::RuntimeShutdown),
941                false,
942            );
943        } else {
944            // Ensure cancel wake + abort residual work; never abort Finalizer
945            // (must publish the one completion for this admission).
946            if let Ok(mut ledger) = shared.ledger.lock() {
947                if let Some(entry) = ledger.get_mut(&tx) {
948                    entry.resources.cancel.cancel();
949                }
950            }
951            tasks.abort_transaction_residuals(&tx);
952        }
953    }
954}
955
956fn try_remove_tombstone(shared: &Arc<RuntimeShared>, tx: &TransactionId) {
957    let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
958    let Some(entry) = ledger.get(tx) else {
959        return;
960    };
961    // Caller (`on_task_exit`) only invokes this when no tx-scoped tasks remain,
962    // so removing here releases SessionKey only after residual work is gone.
963    // CleanupPending: completion published. Finalizing with no tasks: Finalizer
964    // never ran or was lost — still clear so Stopped is reachable (fail-closed).
965    if matches!(
966        entry.phase,
967        TransactionPhase::CleanupPending | TransactionPhase::Finalizing
968    ) {
969        let _ = ledger.remove(tx);
970    }
971}
972
973fn force_remove_tombstone(shared: &Arc<RuntimeShared>, tx: &TransactionId) {
974    let mut ledger = shared.ledger.lock().unwrap_or_else(|e| e.into_inner());
975    let _ = ledger.remove(tx);
976}
977
978/// Wait until supervisor drain is complete (D-049) or the deadline elapses.
979///
980/// Does **not** mean public `Stopped` — the owner must still join the executor
981/// OS thread before publishing `STATE_STOPPED`.
982pub(crate) async fn wait_until_drain_complete(
983    shared: &Arc<RuntimeShared>,
984    deadline: Duration,
985) -> Result<(), monoloop_contracts::ShutdownWaitOutcome> {
986    let start = tokio::time::Instant::now();
987    loop {
988        if shared.state.load(Ordering::SeqCst) == STATE_STOPPED
989            || shared.drain_complete.load(Ordering::SeqCst)
990        {
991            return Ok(());
992        }
993        if start.elapsed() >= deadline {
994            return Err(monoloop_contracts::ShutdownWaitOutcome::TimedOut(
995                shared.snapshot(),
996            ));
997        }
998        // D-039: while Quiescing, re-send BeginShutdown + wake so a dropped
999        // control command or missed notify cannot strand shutdown.
1000        if shared.state.load(Ordering::SeqCst) == STATE_QUIESCING {
1001            let _ = shared.control_tx.try_send(ControlCommand::BeginShutdown);
1002            shared.wake.notify_waiters();
1003        }
1004        tokio::time::sleep(Duration::from_millis(5)).await;
1005    }
1006}