Skip to main content

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}