nmbrs_runtime/observer.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Run observer: callback trait for phase lifecycle events.
5//!
6//! The executor notifies observers when phases start, complete,
7//! or fail. The TUI implements this to update its display state.
8//! The default stderr observer prints phase progress lines.
9
10/// Log level for diagnostic messages.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
12pub enum LogLevel {
13 /// Highest-volume per-cycle/per-event tracing (recall pipeline,
14 /// adapter call traces). Below `Debug` so the existing file
15 /// sink default (Debug) does not pick up trace traffic. Trace
16 /// events instead route through the [`crate::trace_router`]
17 /// to optionally-filtered per-label files.
18 Trace,
19 /// Detailed diagnostics (parse notes, connection info)
20 Debug,
21 /// Normal operational messages (phase info, metrics paths)
22 Info,
23 /// Warnings (CQL driver warnings, recoverable errors)
24 Warn,
25 /// Errors (phase failures, binding errors)
26 Error,
27}
28
29/// Provenance tag carried alongside a log message so display
30/// sinks can route by *kind*, not by sniffing message text.
31///
32/// The only consumer today is the `tui=terminal` log sink: it
33/// owns a managed phase-history region (the idempotent catch-up
34/// projection of the scene tree), so the phase start/end readout
35/// lines that would otherwise scroll past in the log stream are
36/// tagged `attached: PhaseStart` and suppressed from that
37/// stream — they live in the region instead. Every other surface
38/// (`session.log`, the failure dump, the full TUI panel) treats
39/// all events alike; the tag is purely additive.
40///
41/// The tag is TWO ORTHOGONAL AXES (they were one conflated enum
42/// once, with `Diagnostic` — "anything else" — riding beside
43/// boundary-attached kinds, and sinks keying presentation rules
44/// off the conflation):
45///
46/// - [`EventTag::attached`] — WHERE on the execution lifecycle the
47/// event belongs, in the canonical lifecycle vocabulary
48/// ([`crate::lifecycle::EventType`]: session/scope/each/phase
49/// boundaries). `None` = in-flight — emitted from the running
50/// body, bound to no boundary.
51/// - [`EventTag::category`] — WHAT the event is about
52/// ([`EventCategory`]): the semantic domain, independent of when
53/// it fired.
54///
55/// Presentation rules derive from the axes instead of naming
56/// bundles: "hide phase-start renders from scrollback" keys on
57/// `attached == Some(PhaseStart)`; "the TUI log panel shows only
58/// in-flight diagnostics" keys on `attached.is_none()`; "detail
59/// rows fold with their block" keys on `attached == Some(PhaseEnd)
60/// && category != Outcome`.
61#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
62pub struct EventTag {
63 /// The execution-lifecycle boundary this event is attached to;
64 /// `None` = in-flight.
65 pub attached: Option<crate::lifecycle::EventType>,
66 /// Semantic category.
67 pub category: EventCategory,
68}
69
70impl EventTag {
71 /// An in-flight event of the given category (no boundary).
72 pub const fn in_flight(category: EventCategory) -> Self {
73 Self {
74 attached: None,
75 category,
76 }
77 }
78 /// An event attached to a lifecycle boundary.
79 pub const fn at(moment: crate::lifecycle::EventType, category: EventCategory) -> Self {
80 Self {
81 attached: Some(moment),
82 category,
83 }
84 }
85}
86
87/// The semantic domain of a structured event — orthogonal to the
88/// lifecycle moment it attaches to. Extend as producers appear;
89/// a category earns a variant when some consumer (sink filter,
90/// counter, panel) needs to dispatch on it without string-matching
91/// rendered prefixes.
92#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
93pub enum EventCategory {
94 /// Ordinary diagnostics with no more specific domain (the
95 /// overwhelming majority).
96 #[default]
97 General,
98 /// An outcome projection (the `✓`/`✗` phase_outcome block).
99 /// SRD-81: a display projection, not a diagnostic.
100 Outcome,
101 /// Evaluation detail attached to an outcome (relevancy
102 /// `mean=/p50=` rows, verify summaries). SRD-92: detail rows
103 /// belong to the block above them.
104 Evaluation,
105 /// Retry-plane visibility: first-sighting advisories and
106 /// sampled counter-exemplars (`exec_events`).
107 Retry,
108 /// Adaptive backpressure governance (SRD-83 Part 9
109 /// `throttle:`).
110 Throttle,
111 /// Reference resolution (SRD-85 nearest-first shadowing
112 /// warnings).
113 Resolution,
114 /// Report-block warnings (SRD-46).
115 Report,
116}
117
118/// Kind of pre-mapped scenario entry. Re-export of
119/// [`crate::scene_tree::NodeKind`] for callers that already
120/// imported it via the observer module.
121pub use crate::scene_tree::NodeKind as PreMapKind;
122
123/// Lifecycle events from the executor.
124pub trait RunObserver: Send + Sync {
125 /// A completed sysmon sample window (session-level host utilization).
126 /// Display surfaces override; everything else ignores it.
127 fn sysmon_update(&self, _sample: &crate::sysmon::SysmonSample) {}
128
129 /// The session directory now exists. Display surfaces that persist a
130 /// transcript open it here — the runner creates the directory well after
131 /// the observer is built, so this is the first moment the path is known.
132 fn session_dir_ready(&self, _dir: &std::path::Path) {}
133
134 /// A phase is about to start executing.
135 ///
136 /// `op_templates` is the count of op definitions in the phase
137 /// (typically 1 for query workloads). `total_cycles` is the
138 /// number of times the stanza will iterate. Both are reported
139 /// because they answer different questions: the first describes
140 /// the *shape* of the phase, the second describes the *amount*
141 /// of work it represents.
142 ///
143 /// `scene_node_id` is the phase's dispatch-time
144 /// [`crate::scene_tree::SceneNodeId`] (SRD-100 P1c) — the stable
145 /// row key for lifecycle routing. Observers flip *that* node
146 /// directly instead of re-matching by `(name, status)` in DFS
147 /// order, which races under concurrent same-name dispatch
148 /// (sweep cells, comprehension iterations, daemon+foreground).
149 fn phase_starting(
150 &self,
151 scene_node_id: crate::scene_tree::SceneNodeId,
152 name: &str,
153 labels: &str,
154 op_templates: usize,
155 total_cycles: u64,
156 concurrency: usize,
157 );
158
159 /// A phase completed successfully. `scene_node_id` keys the
160 /// node to flip (see [`Self::phase_starting`]).
161 fn phase_completed(
162 &self,
163 scene_node_id: crate::scene_tree::SceneNodeId,
164 name: &str,
165 labels: &str,
166 duration_secs: f64,
167 );
168
169 /// A phase failed. `scene_node_id` keys the node to flip (see
170 /// [`Self::phase_starting`]).
171 fn phase_failed(
172 &self,
173 scene_node_id: crate::scene_tree::SceneNodeId,
174 name: &str,
175 labels: &str,
176 error: &str,
177 );
178
179 /// Update live metrics for the active phase (called at progress tick rate).
180 fn phase_progress(&self, update: &PhaseProgressUpdate);
181
182 /// SRD-63 — an op-level execution leaf became visible (`readout: visible`,
183 /// enabled by the `readout` wrapper). Nests under `parent_phase` in the
184 /// hierarchic status view; the consumer assigns its sequence within the
185 /// phase from arrival order. Default no-op so observers that don't render
186 /// op-level leaves (and the zero-cost non-readout path) are unaffected.
187 fn op_starting(&self, _parent_phase: crate::scene_tree::SceneNodeId, _op_name: &str) {}
188
189 /// An op-level execution leaf completed. `duration_secs` is the op's own
190 /// cumulative execution time. Default no-op.
191 fn op_completed(
192 &self,
193 _parent_phase: crate::scene_tree::SceneNodeId,
194 _op_name: &str,
195 _duration_secs: f64,
196 ) {
197 }
198
199 /// The op's KEY MEASURABLE, rendered from its `measure:` template at
200 /// completion — "12 sstables", "1.4 GiB", "98.1%".
201 ///
202 /// A leaf row already carries how LONG a step took; this carries what it
203 /// actually produced, which is the number a reader is usually after. Sent
204 /// as its own hook (rather than widening `op_completed`) so every existing
205 /// observer keeps compiling and ignoring it costs nothing.
206 fn op_measure(
207 &self,
208 _parent_phase: crate::scene_tree::SceneNodeId,
209 _op_name: &str,
210 _text: &str,
211 ) {
212 }
213
214 /// An op-level execution leaf failed. Default no-op.
215 fn op_failed(
216 &self,
217 _parent_phase: crate::scene_tree::SceneNodeId,
218 _op_name: &str,
219 _error: &str,
220 ) {
221 }
222
223 /// SRD-100 P2 — attach the per-phase [`PhaseRenderHandle`] to the live
224 /// display fold, **once**, after the activity's metrics/binder exist
225 /// (the executor calls this on-task at progress-setup time). Observers
226 /// that own a run-state actor (TUI / log-only) route it to an
227 /// `AttachPhaseRender` mutation so the consumer can fold `active_phases`
228 /// and re-derive each phase's status line itself; the no-op default is
229 /// correct for surfaces with no live status fold (Stderr / Headless).
230 fn phase_render_attach(&self, _handle: PhaseRenderHandle) {}
231
232 /// The entire run is complete.
233 fn run_finished(&self);
234
235 /// Diagnostic log message. Routed to stderr in CLI mode,
236 /// to a ring buffer in TUI mode. All `eprintln!` in the
237 /// runtime should go through this instead.
238 fn log(&self, level: LogLevel, message: &str);
239
240 /// Log a message carrying an explicit [`EventTag`]. The
241 /// default ignores the tag and delegates to [`Self::log`] —
242 /// correct for every observer whose surface treats all events
243 /// alike. Observers that feed a tag-aware sink (the run-state
244 /// actor ring) override this to retain both axes.
245 fn log_tagged(&self, level: LogLevel, _tag: EventTag, message: &str) {
246 self.log(level, message);
247 }
248
249 // SRD-100 P2 — `set_status_line` is removed. The status line is no
250 // longer produced by submitting a pre-rendered string up to the actor;
251 // each display surface folds `active_phases` and re-derives the status
252 // itself (`nmbrs_tui::status_fold`). The producer-side render handle now
253 // travels via `phase_render_attach` instead.
254
255 /// Whether to suppress the inline stderr progress line
256 /// (because the TUI is handling display).
257 fn suppresses_stderr(&self) -> bool {
258 false
259 }
260
261 /// Optional shared flag mirroring [`Self::suppresses_stderr`]
262 /// that the runner threads into long-lived components
263 /// (e.g. the activity's inline status thread) so they can
264 /// react to dismissal mid-run rather than honoring a
265 /// snapshot taken at construction. When `None`, the
266 /// activity uses a fresh `AtomicBool(false)` (never
267 /// suppress). Implementations that go through a TUI
268 /// (and only those) typically expose their internal
269 /// "tui_active" flag here.
270 fn live_suppress_flag(&self) -> Option<std::sync::Arc<std::sync::atomic::AtomicBool>> {
271 None
272 }
273
274 /// Optional reporter to register on the metrics scheduler.
275 /// The runner calls this once during setup. Return None for
276 /// observers that don't need metrics frames (like StderrObserver).
277 ///
278 /// Kept for back-compat; for observers that want multiple reporters
279 /// at different cadences, override [`reporters`] instead — the
280 /// default impl forwards this single reporter as the base cadence.
281 fn reporter(&self) -> Option<Box<dyn nmbrs_metrics::scheduler::Reporter>> {
282 None
283 }
284
285 /// Multiple reporters with explicit cadences. The runner calls
286 /// this once during setup. Each `(interval, reporter)` entry is
287 /// registered with the scheduler at that interval. The default
288 /// implementation returns whatever [`reporter`] produced at the
289 /// base 1s cadence, so existing observers work unchanged.
290 fn reporters(
291 &self,
292 ) -> Vec<(
293 std::time::Duration,
294 Box<dyn nmbrs_metrics::scheduler::Reporter>,
295 )> {
296 match self.reporter() {
297 Some(r) => vec![(std::time::Duration::from_secs(1), r)],
298 None => vec![],
299 }
300 }
301
302 /// User-declared cadences for this observer's consumers (SRD-42).
303 ///
304 /// When present, the runner uses these to plan the cadence tree
305 /// passed to the scheduler's [`nmbrs_metrics::cadence_reporter::CadenceReporter`].
306 /// The reporter writes all windowed snapshots into a single store
307 /// that every consumer reads through [`nmbrs_metrics::metrics_query::MetricsQuery`].
308 ///
309 /// Observers that don't need windowed views (e.g. StderrObserver)
310 /// return `None` and the runner falls back to
311 /// `Cadences::defaults()`.
312 fn cadences(&self) -> Option<nmbrs_metrics::cadence::Cadences> {
313 None
314 }
315
316 /// Callback invoked once the runner has built the shared
317 /// [`nmbrs_metrics::metrics_query::MetricsQuery`]. Observers that
318 /// render metrics (TUI, CLI status) capture this handle to read
319 /// cadence windows, `now` values, and session-lifetime aggregates.
320 fn on_metrics_query(&self, _query: std::sync::Arc<nmbrs_metrics::metrics_query::MetricsQuery>) {
321 }
322
323 /// Pre-populated scenario tree.
324 ///
325 /// Called once before execution begins with the full
326 /// [`crate::scene_tree::SceneTree`] — synthetic root, every
327 /// concrete phase, and every scope header (`for_each`,
328 /// `for_combinations`, `do_while`, `do_until`) wired up by
329 /// parent / children pointers.
330 ///
331 /// The TUI uses this to show all phases as Pending from the
332 /// start; renderers that want hierarchical features (collapse,
333 /// scope-level aggregate status) walk the tree directly. The
334 /// callee may store the tree (e.g. behind an `RwLock`) and
335 /// mutate node statuses in place via the lifecycle callbacks.
336 fn scenario_pre_mapped(&self, _tree: &crate::scene_tree::SceneTree) {}
337}
338
339/// Live metrics snapshot for progress updates.
340#[derive(Clone, Debug)]
341pub struct PhaseProgressUpdate {
342 /// The execution this update belongs to (SRD-88 `exec_id`,
343 /// SRD-100 §4). With `name`+`labels` it forms the live routing
344 /// key, so concurrent executions of the same phase route to
345 /// distinct `ActivePhase` slots. `1` for the single-execution
346 /// case (SRD-88 A1).
347 pub exec_id: u64,
348 /// Phase name this update belongs to — matches the `name`
349 /// passed to [`RunObserver::phase_starting`]. Present so
350 /// observers that track multiple concurrent phases can route
351 /// the update to the correct per-phase slot.
352 pub name: String,
353 /// Phase dimensional labels (e.g. `profile=label_00, k=10`) —
354 /// together with `name` this uniquely identifies one phase
355 /// iteration.
356 pub labels: String,
357 pub cursor_name: String,
358 pub cursor_extent: u64,
359 /// SRD-82 Part 6 — whether this phase runs as a DAEMON (off the
360 /// foreground budget, open-extent cursor, stopped by its group's
361 /// foreground completing). Display consumers exclude daemon phases
362 /// from AGGREGATE progress — the shared footer gutter bar averages
363 /// the non-daemon phases — since a daemon's percent-of-budget is
364 /// not meaningful progress toward the run.
365 pub daemon: bool,
366 /// Cursor ordinals CONSUMED so far for a data-driven phase
367 /// (polydat `DataSourceFactory::global_consumed()`) — the
368 /// authoritative row-level progress. `0` for non-cursor phases
369 /// (plain `cycles:`), where the display keeps the op-denominated
370 /// `cycles:` chip.
371 pub rows_consumed: u64,
372 /// Cursor ordinal EXTENT for a data-driven phase
373 /// (`global_extent()`). `0` for non-cursor phases — distinct from
374 /// `cursor_extent`, which falls back to the phase's `cycles:`
375 /// bound for sourceless phases; this stays `0` so the display can
376 /// tell a declared cursor from the synthesized `range(0, cycles)`
377 /// every phase gets. `rows_total > 0` is the "show `rows:`" signal.
378 pub rows_total: u64,
379 pub fibers: usize,
380 pub ops_started: u64,
381 pub ops_finished: u64,
382 pub ops_ok: u64,
383 /// SKIPPED ops (`skips_total`) — excluded from the ok% denominator
384 /// (a skip is neither a success nor a failure).
385 pub skips: u64,
386 pub errors: u64,
387 pub retries: u64,
388 pub ops_per_sec: f64,
389 pub adapter_counters: Vec<(String, u64, f64)>,
390 pub rows_per_batch: f64,
391 /// Live relevancy aggregates — one entry per relevancy metric (e.g.
392 /// `recall@10`). Each has a moving-window mean over the last N
393 /// recall calculations and a whole-activity running mean.
394 pub relevancy: Vec<crate::validation::RelevancyLive>,
395}
396
397/// SRD-100 P2 — the per-phase **live render handle**, attached to the
398/// display fold's `ActivePhase` exactly once (after the `Activity` and its
399/// metrics exist — executor.rs creates the activity well after
400/// `phase_starting`, so this cannot ride that callback). It carries
401/// everything a display surface needs to re-derive *this phase's* status
402/// line **at the consumer** by folding the snapshot, replacing the retired
403/// per-phase inline-status producer threads (SRD-100 §6).
404///
405/// Why a handle of live shared state rather than a pre-rendered string:
406/// the consumer calls [`crate::readout_context::build_inline_refresh_context`]
407/// verbatim against `metrics` and fires `bodies`, so single-run output is
408/// **byte-identical** to the old producer path (SRD-100 §12 A1) by code
409/// reuse — and §11 mandates that "each phase owns its `Arc<ActivityMetrics>`".
410/// `BakedBody` is `Send + Sync` (its `Readout` handles are), so the format
411/// template rides the ArcSwap snapshot as pure data; only the `!Sync`
412/// *binder* is kept out of the snapshot (the consumer fires bodies with
413/// `&self`).
414#[derive(Clone)]
415pub struct PhaseRenderHandle {
416 /// Routing key — together with `name`+`labels`, selects the
417 /// `ActivePhase` slot to attach to (mirrors [`PhaseProgressUpdate`]).
418 pub exec_id: u64,
419 pub name: String,
420 pub labels: String,
421 /// The activity's display name (`activity.config.name`, which may carry
422 /// the leaf coord, e.g. `"phase (k=10)"`). Fed verbatim to
423 /// `build_inline_refresh_context` so the consumer's render is
424 /// byte-identical to the retired producer thread's.
425 pub activity_name: String,
426 /// Live atomic counters for this phase (SRD-100 §11). Read lock-free
427 /// at render time so the status stays fresh without a producer thread.
428 pub metrics: Arc<crate::activity::ActivityMetrics>,
429 /// Resolved `on_update`-slot render template (workload binding + CLI
430 /// override + built-in `phase_status` default). Immutable, fired with
431 /// `&self` by the consumer; shared, never cloned per render.
432 pub bodies: Arc<Vec<crate::readouts::BakedBody>>,
433 /// Live memo header (the memo-wrapper `before:`/`after:` state),
434 /// snapshotted into the context each render.
435 pub memo: Arc<arc_swap::ArcSwap<String>>,
436 /// Live gutter-cell spec (the gutter-wrapper state). `None` in
437 /// the slot ⇒ the display derives the cell automatically.
438 pub gutter: Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>,
439 /// `status_metrics:` selection for the per-phase status chips.
440 pub status_metrics: Arc<[String]>,
441 /// Fiber count at attach time (the inline context's `concurrency`).
442 pub concurrency: usize,
443 /// `(seq, total)` pre-map coordinate, resolved from the dispatch-time
444 /// `scene_node_id` (NOT a racy by-name DFS match — SRD-100 P1c).
445 pub seq: Option<(usize, usize)>,
446 /// Depth indent string (`" ".repeat(depth-1)`), resolved from the
447 /// scene node at attach time.
448 pub depth_indent: String,
449}
450
451impl std::fmt::Debug for PhaseRenderHandle {
452 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
453 // Deliberately shallow: `ActivityMetrics` / `ArcSwap` are not
454 // `Debug`, and a render handle in a command log only needs its
455 // identity, not its live counters.
456 f.debug_struct("PhaseRenderHandle")
457 .field("exec_id", &self.exec_id)
458 .field("name", &self.name)
459 .field("labels", &self.labels)
460 .field("seq", &self.seq)
461 .field("bodies", &self.bodies.len())
462 .finish_non_exhaustive()
463 }
464}
465
466/// SRD-? — how a FULLY-SKIPPED phase (every cycle `if:`-gated off, or
467/// pre-entry pruned) is represented across the plan, traversal, and
468/// readout surfaces. Selected once per run via the `skipped_phases`
469/// param; read by the runtime's completion fire, the executor's
470/// pre-entry gate, and the TUI's tree fold.
471#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
472pub enum SkippedPhaseDisplay {
473 /// Traversed as normal, but elided from the readout and the
474 /// completed-phase tree — a gated-off phase leaves no trace.
475 Elide,
476 /// Traversed as normal and kept in the plan/readout, explicitly
477 /// rendered as skipped (`⊘ [name] gated off`). The default:
478 /// visible but unmistakable.
479 #[default]
480 Mark,
481 /// Not even traversed: when every op declares `if:` and all
482 /// gates evaluate false at phase entry, the phase is skipped
483 /// before dispatch — elided from plan, traversal, and readout.
484 /// (Post-hoc fully-skipped phases — gates that opened false
485 /// mid-run — are elided as in `Elide`.)
486 Prune,
487}
488
489impl SkippedPhaseDisplay {
490 pub fn parse(s: &str) -> Option<Self> {
491 match s.trim().to_ascii_lowercase().as_str() {
492 "elide" => Some(Self::Elide),
493 "mark" => Some(Self::Mark),
494 "prune" => Some(Self::Prune),
495 _ => None,
496 }
497 }
498}
499
500static SKIPPED_PHASE_DISPLAY: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(1);
501
502/// Set the run's skipped-phase display mode (runner, at param parse).
503pub fn set_skipped_phase_display(mode: SkippedPhaseDisplay) {
504 let v = match mode {
505 SkippedPhaseDisplay::Elide => 0,
506 SkippedPhaseDisplay::Mark => 1,
507 SkippedPhaseDisplay::Prune => 2,
508 };
509 SKIPPED_PHASE_DISPLAY.store(v, std::sync::atomic::Ordering::Relaxed);
510}
511
512/// The run's skipped-phase display mode. Default [`SkippedPhaseDisplay::Mark`].
513pub fn skipped_phase_display() -> SkippedPhaseDisplay {
514 match SKIPPED_PHASE_DISPLAY.load(std::sync::atomic::Ordering::Relaxed) {
515 0 => SkippedPhaseDisplay::Elide,
516 2 => SkippedPhaseDisplay::Prune,
517 _ => SkippedPhaseDisplay::Mark,
518 }
519}
520
521/// SRD-92 R5 — how much of a completed node's block is retained in
522/// scrollback. Completion is never a full collapse: the header line
523/// always lands; this selects whether the detail rows and op leaves
524/// land with it. Selected once per run via the `completed_phases=`
525/// param; session.log keeps every line unconditionally.
526#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
527pub enum CompletedPhaseDisplay {
528 /// Header + detail rows (counters, key metrics, memo) + op
529 /// leaves — the full block, in contract shape. The default.
530 #[default]
531 Full,
532 /// Header row only.
533 Headers,
534}
535
536impl CompletedPhaseDisplay {
537 pub fn parse(s: &str) -> Option<Self> {
538 match s.trim().to_ascii_lowercase().as_str() {
539 "full" => Some(Self::Full),
540 "headers" => Some(Self::Headers),
541 _ => None,
542 }
543 }
544}
545
546static COMPLETED_PHASE_DISPLAY: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(0);
547
548/// Set the run's completed-phase display mode (runner, at param parse).
549pub fn set_completed_phase_display(mode: CompletedPhaseDisplay) {
550 let v = match mode {
551 CompletedPhaseDisplay::Full => 0,
552 CompletedPhaseDisplay::Headers => 1,
553 };
554 COMPLETED_PHASE_DISPLAY.store(v, std::sync::atomic::Ordering::Relaxed);
555}
556
557/// The run's completed-phase display mode. Default [`CompletedPhaseDisplay::Full`].
558pub fn completed_phase_display() -> CompletedPhaseDisplay {
559 match COMPLETED_PHASE_DISPLAY.load(std::sync::atomic::Ordering::Relaxed) {
560 1 => CompletedPhaseDisplay::Headers,
561 _ => CompletedPhaseDisplay::Full,
562 }
563}
564
565/// Global observer for code that can't thread the observer through.
566/// Set once at run start, remains for the process lifetime.
567static GLOBAL_OBSERVER: std::sync::OnceLock<Arc<dyn RunObserver>> = std::sync::OnceLock::new();
568
569/// Set the global observer. Called once by the runner at startup.
570pub fn set_global_observer(observer: Arc<dyn RunObserver>) {
571 let _ = GLOBAL_OBSERVER.set(observer);
572}
573
574/// The observer the current code should route through. Used by code that
575/// needs the observer without threading it through every call site (e.g. the
576/// activity's inline-status refresh thread publishing into
577/// [`RunObserver::set_status_line`]).
578///
579/// SRD-88: a fiber running inside an [`ExecutionContext`](crate::execution_context)
580/// resolves to ITS execution's observer (so concurrent executions route
581/// lifecycle/log independently); outside any execution scope — or when the
582/// scoped context set no observer — it falls back to the process-global
583/// `GLOBAL_OBSERVER` (the single-run / CLI / test default; axiom A1).
584pub fn global_observer() -> Option<Arc<dyn RunObserver>> {
585 if let Some(obs) = crate::execution_context::current_observer() {
586 return Some(obs);
587 }
588 GLOBAL_OBSERVER.get().cloned()
589}
590
591/// Direct the log sink to a file. Opens for append-writes,
592/// installs the [`crate::log_sink`] async writer thread.
593/// Producers thereafter `try_send` and never block — see SRD-02
594/// §"Display and Diagnostic Decoupling". Silently no-ops on a
595/// second call — the first session wins (one run per process).
596pub fn set_log_file(path: &std::path::Path) -> std::io::Result<()> {
597 crate::log_sink::init(path)
598}
599
600/// Log a diagnostic message through the global observer and
601/// append to the async log sink (if initialized). Safe to call
602/// from anywhere — falls back to stderr if no observer is set.
603///
604/// The file write is non-blocking: the line is enqueued onto a
605/// bounded channel consumed by the dedicated `log-sink` thread.
606/// On overflow the line is dropped and the sink's `dropped_count`
607/// is bumped — never blocks the caller, even on a stalled disk.
608/// Whether ANSI color escapes are appropriate for stderr.
609/// `true` only when stderr is a TTY and the operator hasn't
610/// disabled color via the conventional `NO_COLOR` env var
611/// (https://no-color.org). Pipelined / CI contexts return
612/// `false` so log archives stay readable. Cached on first
613/// call; the answer doesn't change over a process's lifetime.
614pub fn use_color() -> bool {
615 use std::io::IsTerminal;
616 use std::sync::OnceLock;
617 static CACHE: OnceLock<bool> = OnceLock::new();
618 *CACHE.get_or_init(|| {
619 if std::env::var_os("NO_COLOR").is_some() {
620 return false;
621 }
622 std::io::stderr().is_terminal()
623 })
624}
625
626/// SRD — explainer-overlay toggle state. Process-global so the
627/// TUI keystroke layer (which holds the watcher) and the
628/// readout binder (which holds the render thread) can rendezvous
629/// without threading a channel through every readout call.
630///
631/// Value semantics: wall-clock nanos at which the overlay
632/// auto-reverts. `0` means "off"; any value > `now_nanos()`
633/// means "render `ContentMode::Explanation` until the deadline."
634///
635/// Toggle model (not hold): a single `?` press flips the
636/// overlay on with a 10 s auto-revert deadline. A second press
637/// while on flips it back off immediately. Auto-revert ensures
638/// the operator can't leave the overlay stuck on after walking
639/// away.
640static EXPLAIN_HELD_UNTIL_NS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
641
642/// Wall-clock nanos of the most recent `?` press. Drives the
643/// auto-repeat debounce — terminals in raw mode send a stream
644/// of keystrokes while `?` is held, and without this the
645/// second auto-repeat would flip the overlay off again
646/// 30 ms after the operator's first press.
647static EXPLAIN_LAST_PRESS_NS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
648
649/// How long the overlay stays on after a `?` toggle. Picked
650/// to be long enough that the operator can read the explainer
651/// surface without timing out mid-read, short enough that a
652/// forgotten toggle reverts on its own. 10 s lands in the
653/// middle of "long enough to be useful" and "short enough that
654/// the operator notices the auto-revert."
655const EXPLAIN_AUTO_REVERT_MS: u64 = 10_000;
656
657/// Auto-repeat debounce window. 250 ms swallows the 30 Hz
658/// auto-repeat stream while still allowing a deliberate
659/// second tap to take effect.
660const EXPLAIN_TOGGLE_DEBOUNCE_MS: u64 = 250;
661
662fn now_nanos() -> u64 {
663 std::time::SystemTime::now()
664 .duration_since(std::time::UNIX_EPOCH)
665 .map(|d| d.as_nanos() as u64)
666 .unwrap_or(0)
667}
668
669/// Toggle the explainer overlay. First press → on (with a 10 s
670/// auto-revert deadline). Second press while on → off
671/// immediately. Auto-repeat-safe via a 250 ms debounce.
672pub fn toggle_explain() {
673 let now = now_nanos();
674 let last_press = EXPLAIN_LAST_PRESS_NS.load(std::sync::atomic::Ordering::Acquire);
675 if last_press != 0 && now.saturating_sub(last_press) < EXPLAIN_TOGGLE_DEBOUNCE_MS * 1_000_000 {
676 return;
677 }
678 EXPLAIN_LAST_PRESS_NS.store(now, std::sync::atomic::Ordering::Release);
679 let currently_on = {
680 let deadline = EXPLAIN_HELD_UNTIL_NS.load(std::sync::atomic::Ordering::Acquire);
681 deadline != 0 && now < deadline
682 };
683 if currently_on {
684 EXPLAIN_HELD_UNTIL_NS.store(0, std::sync::atomic::Ordering::Release);
685 } else {
686 let deadline = now.saturating_add(EXPLAIN_AUTO_REVERT_MS * 1_000_000);
687 EXPLAIN_HELD_UNTIL_NS.store(deadline, std::sync::atomic::Ordering::Release);
688 }
689}
690
691/// True iff the explainer overlay is currently on (toggled on
692/// within the last `EXPLAIN_AUTO_REVERT_MS` and not yet
693/// toggled off). Read by the readout binder on each `fire()`
694/// to decide whether to dispatch with
695/// `ContentMode::Explanation` instead of `Value`.
696pub fn is_explain_held() -> bool {
697 let deadline = EXPLAIN_HELD_UNTIL_NS.load(std::sync::atomic::Ordering::Acquire);
698 deadline != 0 && now_nanos() < deadline
699}
700
701/// Readout pause — the `p` key's copy-support freeze. While set,
702/// readout frontends stop writing to the terminal (the line-mode
703/// sink freezes its log drain and status redraw; buffered output
704/// flushes on resume) so the operator can select and copy text
705/// without the surface repainting underneath the selection. The
706/// full-TUI app keeps its own snapshot-based pause; this global
707/// serves the line-mode frontends, which have no App state.
708static READOUT_PAUSED: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
709
710/// Toggle the readout pause; returns the NEW state (`true` =
711/// now paused).
712pub fn toggle_readout_pause() -> bool {
713 !READOUT_PAUSED.fetch_xor(true, std::sync::atomic::Ordering::AcqRel)
714}
715
716/// True while the readout pause is on (see
717/// [`toggle_readout_pause`]).
718pub fn readouts_paused() -> bool {
719 READOUT_PAUSED.load(std::sync::atomic::Ordering::Acquire)
720}
721
722/// Minimum severity that reaches the file sink
723/// (`session.log`). Default `Debug` — the file gets every
724/// non-trivial entry. The display threshold (per-observer
725/// `min_level`) is the SEPARATE knob that controls what
726/// reaches stderr; it defaults to `Info`. The two are
727/// configured independently via:
728///
729/// - `loglevel-retain=` / `--log-retain-level` /
730/// `NMBRS_LOG_RETAIN_LEVEL` — what the file gets.
731/// - `loglevel=` / `--log-display-level` /
732/// `NMBRS_LOG_DISPLAY_LEVEL` — what stderr gets.
733///
734/// The runner installs the user-supplied retain level via
735/// [`set_retain_level`] at startup; readers consult
736/// [`retain_level`] before sending to the sink.
737static RETAIN_LEVEL: std::sync::OnceLock<LogLevel> = std::sync::OnceLock::new();
738
739/// Install the file-sink retention threshold. Called once
740/// by the runner; subsequent calls are silent no-ops
741/// (matching the "first wins" pattern of the rest of the
742/// global observer surface).
743pub fn set_retain_level(level: LogLevel) {
744 let _ = RETAIN_LEVEL.set(level);
745}
746
747/// Effective retention threshold. Defaults to `Debug` so
748/// pre-runner-init log calls (very early startup) still
749/// reach the file sink.
750pub fn retain_level() -> LogLevel {
751 *RETAIN_LEVEL.get().unwrap_or(&LogLevel::Debug)
752}
753
754/// Companion to [`RETAIN_LEVEL`]: the *display* threshold
755/// that gates console emission. Each observer carries its
756/// own `min_level` for the live log path; this global
757/// captures the effective value so secondary surfaces
758/// (the failure-dump path in `nmbrs-tui::observer`) can
759/// honour the same threshold without plumbing the observer
760/// reference everywhere.
761static DISPLAY_LEVEL: std::sync::OnceLock<LogLevel> = std::sync::OnceLock::new();
762
763/// Install the console display threshold. Called by the
764/// runner alongside [`set_retain_level`] at startup.
765pub fn set_display_level(level: LogLevel) {
766 let _ = DISPLAY_LEVEL.set(level);
767}
768
769/// Effective console display threshold. Defaults to
770/// `Info` — same default the live observers use.
771pub fn display_level() -> LogLevel {
772 *DISPLAY_LEVEL.get().unwrap_or(&LogLevel::Info)
773}
774
775pub fn log(level: LogLevel, message: &str) {
776 log_tagged(level, EventTag::default(), message);
777}
778
779/// [`log`] with an explicit [`EventTag`]. The session-log write
780/// (unconditional, all tags) and the fallback stderr path are
781/// identical to [`log`]; the tag only changes how a tag-aware
782/// display sink files the message — e.g. the readout engine
783/// attaches phase-start renders to `PhaseStart` so the terminal
784/// sink keeps them out of its scrollback (they show in its managed
785/// phase-history region instead).
786pub fn log_tagged(level: LogLevel, tag: EventTag, message: &str) {
787 if level >= retain_level()
788 && let Some(sink) = crate::log_sink::global()
789 {
790 let tag = match level {
791 LogLevel::Trace => "TRC",
792 LogLevel::Debug => "DBG",
793 LogLevel::Info => "INF",
794 LogLevel::Warn => "WRN",
795 LogLevel::Error => "ERR",
796 };
797 // Human-readable wall-clock timestamp from the session
798 // formatter — matches the session id's date/time style
799 // so log lines correlate visually with the session
800 // directory.
801 let ts = crate::session::now_log_timestamp();
802 // The durable session.log is plain text — strip any ANSI
803 // a colored readout render carried in `message` (notably
804 // the phase `✓` outcome, SRD-81 push 1b). The live ring
805 // keeps the colored version for the terminal scrollback,
806 // and the replay capture is untouched; only this file
807 // projection is stripped.
808 let line = format!(
809 "{ts} {tag} {}\n",
810 crate::readouts::snapshot::strip_ansi(message)
811 )
812 .into_bytes();
813 let _ = sink.try_send(line);
814 }
815 // SRD-88: route through the current execution's observer (task-local) or
816 // the process-global default; `global_observer()` resolves both.
817 if let Some(obs) = global_observer() {
818 obs.log_tagged(level, tag, message);
819 } else {
820 // No observer yet (bootstrap): project straight to the log bucket —
821 // the channel owns the fd (SRD-87 §5), falling back to stderr when no
822 // channel is installed either.
823 crate::output_channel::log_to_surface(level, message);
824 }
825}
826
827/// Emit a line of adapter **op output** (SRD-41 / "console belongs to the
828/// adapter"). Delegates to the installed [`crate::output_channel::OutputChannel`]
829/// (SRD-87): the **op-output bucket**'s impl decides where the bytes land —
830/// raw to the owned stdout (a console-owning adapter or a pipe, via
831/// [`op_output_raw`]) or routed through the live display (an interactive
832/// dashboard, avoiding the raw-mode staircase). Before any channel is installed
833/// (bootstrap / unit tests with no run), falls back to the raw path so early
834/// output is never lost.
835pub fn op_output(line: &str) {
836 if let Some(channel) = crate::output_channel::installed() {
837 channel.op_output(line);
838 return;
839 }
840 op_output_raw(line);
841}
842
843/// Write an op-output line RAW to the stdout the producer owns (so
844/// `nmbrs run | grep`, `> file`, AND a console-owning adapter's interactive
845/// screen all show it) and capture it to `session.log` at INFO (the same
846/// durable projection [`log_tagged`] writes). The SRD-87
847/// [`crate::output_channel::RawStdoutChannel`] op-output bucket and the
848/// no-channel bootstrap fallback both route through here.
849pub(crate) fn op_output_raw(line: &str) {
850 if LogLevel::Info >= retain_level()
851 && let Some(sink) = crate::log_sink::global()
852 {
853 let ts = crate::session::now_log_timestamp();
854 let bytes =
855 format!("{ts} INF {}\n", crate::readouts::snapshot::strip_ansi(line)).into_bytes();
856 let _ = sink.try_send(bytes);
857 }
858 use std::io::Write;
859 let mut out = std::io::stdout().lock();
860 let _ = writeln!(out, "{line}");
861 let _ = out.flush();
862}
863
864/// ANSI-colorize a log line by severity for console
865/// output. Always applied at the producer side — every
866/// console emission of a log entry runs through this so
867/// `DBG`/`INF`/`WRN`/`ERR` are visually distinct without
868/// the operator having to squint at message bodies.
869/// Falls through to the bare message when stderr isn't a
870/// TTY or `NO_COLOR` is set (per [`use_color`]); pipeline
871/// captures stay readable.
872pub fn colorize_log_line(level: LogLevel, message: &str) -> String {
873 if !use_color() {
874 return message.to_string();
875 }
876 let (color, reset) = match level {
877 // Faintest grey for trace — even more de-emphasized
878 // than debug; rarely reaches console anyway.
879 LogLevel::Trace => ("\x1b[2;90m", "\x1b[0m"),
880 // Dim grey for debug — present but de-emphasized.
881 LogLevel::Debug => ("\x1b[2m", "\x1b[0m"),
882 // Default-color for info — the baseline; no
883 // override so user-themed terminals show their
884 // preferred default.
885 LogLevel::Info => ("", ""),
886 // Yellow for warn.
887 LogLevel::Warn => ("\x1b[33m", "\x1b[0m"),
888 // Bold red for error.
889 LogLevel::Error => ("\x1b[1;31m", "\x1b[0m"),
890 };
891 if color.is_empty() {
892 message.to_string()
893 } else {
894 format!("{color}{message}{reset}")
895 }
896}
897
898/// Convenience macros for logging through the global observer.
899#[macro_export]
900macro_rules! diag {
901 ($level:expr, $($arg:tt)*) => {
902 $crate::observer::log($level, &format!($($arg)*))
903 };
904}
905
906/// Emit a [`LogLevel::Trace`] event carrying the component's
907/// labels through the [`crate::trace_router`]. Returns
908/// immediately when no `--trace=<spec>` was configured (the
909/// router is empty), so the hot-path cost of an unused trace
910/// site is one atomic load + branch.
911///
912/// Trace events DO NOT flow to `session.log` — they are routed
913/// to dedicated trace files configured by `--trace=` so the
914/// main session log stays usable as a Debug-and-up record.
915pub fn trace(labels: &nmbrs_metrics::labels::Labels, message: &str) {
916 crate::trace_router::log(labels, message);
917}
918
919/// True iff the trace router has at least one configured sink.
920/// Hot-path guard so callers can skip expensive message
921/// formatting when tracing is off.
922pub fn trace_enabled() -> bool {
923 crate::trace_router::enabled()
924}
925
926/// Format-and-emit a trace event. Same shape as [`diag!`] but
927/// takes a labels handle first and only fires when the trace
928/// router is active.
929#[macro_export]
930macro_rules! trace_event {
931 ($labels:expr, $($arg:tt)*) => {
932 if $crate::observer::trace_enabled() {
933 $crate::observer::trace($labels, &format!($($arg)*));
934 }
935 };
936}
937
938use std::sync::Arc;
939
940/// Default observer: prints to stderr.
941///
942/// `min_level` controls the minimum severity that reaches
943/// stderr — Info by default, matching the TUI log panel's
944/// default filter (so high-cadence Debug instrumentation
945/// doesn't drown the signal in either mode). Override via
946/// `loglevel=debug|info|warn|error` on the workload command
947/// line. The async log sink (session.log) still receives
948/// every level regardless of this filter.
949pub struct StderrObserver {
950 pub min_level: LogLevel,
951}
952
953impl Default for StderrObserver {
954 fn default() -> Self {
955 Self {
956 min_level: LogLevel::Info,
957 }
958 }
959}
960
961impl StderrObserver {
962 /// Build a stderr observer with the given min severity.
963 pub fn with_min_level(min_level: LogLevel) -> Self {
964 Self { min_level }
965 }
966}
967
968impl RunObserver for StderrObserver {
969 fn phase_starting(
970 &self,
971 _scene_node_id: crate::scene_tree::SceneNodeId,
972 name: &str,
973 _labels: &str,
974 op_templates: usize,
975 total_cycles: u64,
976 concurrency: usize,
977 ) {
978 // Route through the canonical event channel so the line
979 // lands in `session.log` AND on stderr (via the recursive
980 // call back into `StderrObserver::log` below).
981 let template_word = if op_templates == 1 {
982 "op template"
983 } else {
984 "op templates"
985 };
986 let cycle_word = if total_cycles == 1 { "cycle" } else { "cycles" };
987 crate::observer::log(
988 LogLevel::Info,
989 &format!(
990 "phase '{name}': {op_templates} {template_word}, {total_cycles} {cycle_word}, concurrency={concurrency}"
991 ),
992 );
993 }
994
995 fn phase_completed(
996 &self,
997 _scene_node_id: crate::scene_tree::SceneNodeId,
998 _name: &str,
999 _labels: &str,
1000 _duration_secs: f64,
1001 ) {
1002 // No-op — the executor's own diag emits a fully-formatted
1003 // "phase 'X' complete (Ns)" line via the log path. Doing
1004 // it here too produced a duplicate (and a less
1005 // informative one — no duration). The structured
1006 // callback stays for non-stderr consumers.
1007 }
1008
1009 fn phase_failed(
1010 &self,
1011 _scene_node_id: crate::scene_tree::SceneNodeId,
1012 _name: &str,
1013 _labels: &str,
1014 _error: &str,
1015 ) {
1016 // Same reasoning as phase_completed — the executor diags
1017 // already emit "phase 'X' stopped by error handler (Ns)"
1018 // (or other failure messages) right before calling this.
1019 // Re-emitting here was a duplicate.
1020 }
1021
1022 fn phase_progress(&self, _update: &PhaseProgressUpdate) {
1023 // The inline status line in activity.rs handles this
1024 }
1025
1026 fn run_finished(&self) {
1027 // Same routing as `phase_starting` — through `observer::log`
1028 // so session.log captures the run-end marker.
1029 crate::observer::log(LogLevel::Info, "all phases complete");
1030 }
1031
1032 fn log(&self, level: LogLevel, message: &str) {
1033 // Severity filter: only entries `>= min_level` reach
1034 // stderr. The session log file still gets every level
1035 // via the async log sink — this filter only affects
1036 // what the operator sees on screen.
1037 if level >= self.min_level {
1038 // Cosmetic: when the runtime announces a Ctrl-C-
1039 // initiated graceful shutdown, the terminal has just
1040 // echoed `^C` on the current line. A leading blank
1041 // line makes the announcement visually clear that
1042 // marker without leaving a stray newline in the
1043 // structured session.log (which never sees this
1044 // path). Same idea for force-exit on second Ctrl-C.
1045 if message.starts_with("session: graceful shutdown requested")
1046 || message.starts_with("session: force-exit")
1047 {
1048 eprintln!();
1049 }
1050 // SRD-87 §5: the live-surface write goes through the log bucket
1051 // (the channel owns the fd); the `min_level` gate above stays here.
1052 crate::output_channel::log_to_surface(level, message);
1053 }
1054 }
1055}
1056
1057#[cfg(test)]
1058mod display_mode_tests {
1059 use super::*;
1060
1061 #[test]
1062 fn completed_phase_display_parses_and_rejects() {
1063 assert_eq!(
1064 CompletedPhaseDisplay::parse("full"),
1065 Some(CompletedPhaseDisplay::Full)
1066 );
1067 assert_eq!(
1068 CompletedPhaseDisplay::parse(" HEADERS "),
1069 Some(CompletedPhaseDisplay::Headers)
1070 );
1071 assert_eq!(CompletedPhaseDisplay::parse("collapse"), None);
1072 // Default is Full — completion is never a full collapse
1073 // (SRD-92 R5).
1074 assert_eq!(
1075 CompletedPhaseDisplay::default(),
1076 CompletedPhaseDisplay::Full
1077 );
1078 }
1079}