1use 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#[derive(Clone, Debug, PartialEq, Eq)]
34pub enum StartCommand {
35 Start(TransactionId),
37}
38
39#[derive(Clone, Debug, PartialEq, Eq)]
41pub enum ControlCommand {
42 Cancel(TransactionId),
44 ForceTerminate(TransactionId),
46 BeginShutdown,
48 StopSupervisor,
50}
51
52#[derive(Clone, Debug, PartialEq, Eq)]
54pub enum SupervisorCommand {
55 Start(TransactionId),
57 Cancel(TransactionId),
59 ForceTerminate(TransactionId),
61 BeginShutdown,
63 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
72pub(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 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 pub owned_tasks: AtomicU32,
94 pub live_connector_owners: AtomicU32,
96 pub enable_mcp_listener: bool,
98 pub mcp_listen_addr: Mutex<Option<std::net::SocketAddr>>,
100 pub mcp_gateway: Mutex<Option<McpGatewayHandle>>,
102 pub mcp_cancel: Mutex<Option<CancellationToken>>,
104 pub block_stopped: Option<Arc<StoppedGate>>,
106 pub hold_start: Option<Arc<StartHoldGate>>,
108 pub hold_control: Option<Arc<ControlHoldGate>>,
110 pub hold_finalizer_after_seal: Option<Arc<FinalizerHoldGate>>,
112 pub hold_executor_teardown: Option<Arc<StoppedGate>>,
114 pub inject_non_yielding_service: Option<Arc<std::sync::atomic::AtomicBool>>,
116 pub inject_join_only_spill: Option<Arc<JoinOnlySpillInject>>,
118 pub drain_complete: std::sync::atomic::AtomicBool,
121 pub tools_registry: HostToolRegistry,
123 pub shared_tool_capacity: Arc<SharedToolCapacity>,
125 pub tool_spill: Arc<OrphanToolPermitSet>,
127 pub owned_processes: Arc<AtomicU32>,
129 pub process_registry: Arc<crate::transaction::owned_process_registry::OwnedProcessRegistry>,
131 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 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
186pub(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 let serve = super::mcp_listener::serve_runtime_mcp(Arc::clone(&shared), prepared);
202 tasks.spawn(TaskClass::RuntimeService, serve);
203 }
204 if let Some(entered) = shared.inject_non_yielding_service.clone() {
207 tasks.spawn(TaskClass::RuntimeService, async move {
208 entered.store(true, Ordering::SeqCst);
212 loop {
215 std::thread::park();
216 }
217 });
218 }
219 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 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 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 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 let _ = shared.tool_spill.shutdown_progress();
310 let _ = shared.process_registry.shutdown_progress();
312
313 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 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 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 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 if shared.ledger.lock().map(|l| l.is_empty()).unwrap_or(true) {
369 let _ = shared.tool_spill.shutdown_progress();
371 let _ = shared.process_registry.shutdown_progress();
372 if tasks.is_empty() && shared.process_registry.is_empty() {
375 ready_to_stop = true;
376 } else {
377 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 if let Some(gate) = shared.block_stopped.as_ref() {
398 gate.wait_released().await;
399 }
400 *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 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 accept_terminal(
449 &shared,
450 &mut tasks,
451 tx,
452 TerminalProposal::new(TransactionEndKind::RuntimeShutdown),
453 false,
454 );
455 }
456 }
457 }
458 _ = 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 begin_shutdown_inner(&shared, &mut tasks, &mut stopping);
471 }
472 }
473 }
474 _ = 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 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 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 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 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 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 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 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 (first, entry.publisher_seal_tx.take(), kind)
724 };
725
726 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 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 tokio::task::yield_now().await;
752
753 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 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 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 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 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 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 terminal_delivery = TerminalEventDelivery::DeadlineExceeded;
825 }
826 }
827 }
828
829 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 if let Some(gate) = shared.hold_finalizer_after_seal.as_ref() {
850 gate.wait_released().await;
851 }
852
853 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 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 }
901
902fn begin_shutdown_inner(
903 shared: &Arc<RuntimeShared>,
904 tasks: &mut TaskSupervisor,
905 stopping: &mut bool,
906) {
907 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 super::mcp_listener::signal_mcp_shutdown(shared);
924 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 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 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
978pub(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 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}