nmbrs_runtime/concurrent.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-88 — the **headless observer + single-run headless helper**.
5//!
6//! [`HeadlessObserver`] folds an execution's lifecycle + log into an
7//! [`ExecutionOutcome`] with no live display surface (SRD-88 §4), and
8//! [`run_workload_headless`] runs ONE workload inside a scoped
9//! [`ExecutionContext`](crate::execution_context) for that capture.
10//!
11//! The real CONCURRENT harness — N executions sharing ONE session — is
12//! [`crate::runner::run_executions`]: one `SessionHost`, forked under a
13//! `ScheduleSpec`/semaphore, each execution deriving its own `exec_id` under
14//! the shared session and flushing its metrics to the shared store. The bespoke
15//! task-based `run_executions_concurrent` that used to live here was a DUPLICATE
16//! concurrency path (SRD-02 One Concurrency Path) and has been retired.
17
18use std::sync::{Arc, Mutex};
19
20use crate::execution_context::{self, ExecutionContext};
21use crate::observer::{LogLevel, PhaseProgressUpdate, RunObserver};
22
23/// One phase's terminal record captured by a [`HeadlessObserver`].
24#[derive(Clone, Debug, PartialEq)]
25pub enum PhaseRecord {
26 Completed {
27 name: String,
28 labels: String,
29 duration_secs: f64,
30 },
31 Failed {
32 name: String,
33 labels: String,
34 error: String,
35 },
36}
37
38/// What one execution produced, returned by [`run_executions_concurrent`].
39#[derive(Clone, Debug)]
40pub struct ExecutionOutcome {
41 /// The execution's process-unique id (SRD-77 / §A2).
42 pub exec_id: u64,
43 /// Terminal phase records, in completion order.
44 pub phases: Vec<PhaseRecord>,
45 /// Diagnostic log lines this execution emitted (`level >= Info` kept;
46 /// captured headless rather than displayed).
47 pub logs: Vec<String>,
48}
49
50/// A headless [`RunObserver`]: captures lifecycle + log into an outcome, draws
51/// nothing. Each concurrent execution gets its own, so events route here via
52/// the task-local context (`observer::global_observer()` resolves to it) and
53/// never to a shared global.
54pub struct HeadlessObserver {
55 phases: Mutex<Vec<PhaseRecord>>,
56 logs: Mutex<Vec<String>>,
57}
58
59impl HeadlessObserver {
60 pub fn new() -> Self {
61 Self {
62 phases: Mutex::new(Vec::new()),
63 logs: Mutex::new(Vec::new()),
64 }
65 }
66
67 fn take(&self) -> (Vec<PhaseRecord>, Vec<String>) {
68 let p = std::mem::take(&mut *self.phases.lock().unwrap_or_else(|e| e.into_inner()));
69 let l = std::mem::take(&mut *self.logs.lock().unwrap_or_else(|e| e.into_inner()));
70 (p, l)
71 }
72}
73
74impl Default for HeadlessObserver {
75 fn default() -> Self {
76 Self::new()
77 }
78}
79
80impl RunObserver for HeadlessObserver {
81 fn phase_starting(
82 &self,
83 _scene_node_id: crate::scene_tree::SceneNodeId,
84 _name: &str,
85 _labels: &str,
86 _ops: usize,
87 _cycles: u64,
88 _conc: usize,
89 ) {
90 }
91
92 fn phase_completed(
93 &self,
94 _scene_node_id: crate::scene_tree::SceneNodeId,
95 name: &str,
96 labels: &str,
97 duration_secs: f64,
98 ) {
99 self.phases
100 .lock()
101 .unwrap_or_else(|e| e.into_inner())
102 .push(PhaseRecord::Completed {
103 name: name.to_string(),
104 labels: labels.to_string(),
105 duration_secs,
106 });
107 }
108
109 fn phase_failed(
110 &self,
111 _scene_node_id: crate::scene_tree::SceneNodeId,
112 name: &str,
113 labels: &str,
114 error: &str,
115 ) {
116 self.phases
117 .lock()
118 .unwrap_or_else(|e| e.into_inner())
119 .push(PhaseRecord::Failed {
120 name: name.to_string(),
121 labels: labels.to_string(),
122 error: error.to_string(),
123 });
124 }
125
126 fn phase_progress(&self, _update: &PhaseProgressUpdate) {}
127
128 fn run_finished(&self) {}
129
130 fn log(&self, _level: LogLevel, message: &str) {
131 self.logs
132 .lock()
133 .unwrap_or_else(|e| e.into_inner())
134 .push(message.to_string());
135 }
136}
137
138/// Run ONE workload headless inside its own [`ExecutionContext`]: route the
139/// run's lifecycle + log through a [`HeadlessObserver`] (no display surface),
140/// then return the captured [`ExecutionOutcome`] alongside the run result.
141///
142/// This is the building block that proves the de-globalized run path works
143/// **inside a scoped context** — the run's observer/scene-tree/stop resolve to
144/// the execution's own (task-local), and the propagated per-cycle fibers
145/// (SRD-88 fiber propagation) carry the context end-to-end.
146///
147/// Single-execution headless helper: lets `run_with_observer` create its own
148/// session. To run N executions CONCURRENTLY against ONE shared session, use
149/// [`crate::runner::run_executions`] (the session-tier harness: one
150/// `SessionHost`, forked under a `ScheduleSpec`/semaphore, each execution
151/// deriving its own `exec_id` under the shared session). The bespoke
152/// task-based `run_executions_concurrent` that used to live here was a
153/// DUPLICATE concurrency path (SRD-02 One Concurrency Path) and has been
154/// retired in favour of `run_executions`.
155pub async fn run_workload_headless(args: &[String]) -> (ExecutionOutcome, Result<(), String>) {
156 let obs = Arc::new(HeadlessObserver::new());
157 let ctx = ExecutionContext::with_observer(obs.clone() as Arc<dyn RunObserver>);
158 let exec_id = ctx.exec_id;
159 let result = execution_context::scope(
160 ctx,
161 crate::runner::run_with_observer(args, obs.clone() as Arc<dyn RunObserver>),
162 )
163 .await;
164 let (phases, logs) = obs.take();
165 (
166 ExecutionOutcome {
167 exec_id,
168 phases,
169 logs,
170 },
171 result,
172 )
173}