Skip to main content

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}