#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum LogLevel {
Trace,
Debug,
Info,
Warn,
Error,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct EventTag {
pub attached: Option<crate::lifecycle::EventType>,
pub category: EventCategory,
}
impl EventTag {
pub const fn in_flight(category: EventCategory) -> Self {
Self {
attached: None,
category,
}
}
pub const fn at(moment: crate::lifecycle::EventType, category: EventCategory) -> Self {
Self {
attached: Some(moment),
category,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum EventCategory {
#[default]
General,
Outcome,
Evaluation,
Retry,
Throttle,
Resolution,
Report,
}
pub use crate::scene_tree::NodeKind as PreMapKind;
pub trait RunObserver: Send + Sync {
fn sysmon_update(&self, _sample: &crate::sysmon::SysmonSample) {}
fn session_dir_ready(&self, _dir: &std::path::Path) {}
fn phase_starting(
&self,
scene_node_id: crate::scene_tree::SceneNodeId,
name: &str,
labels: &str,
op_templates: usize,
total_cycles: u64,
concurrency: usize,
);
fn phase_completed(
&self,
scene_node_id: crate::scene_tree::SceneNodeId,
name: &str,
labels: &str,
duration_secs: f64,
);
fn phase_failed(
&self,
scene_node_id: crate::scene_tree::SceneNodeId,
name: &str,
labels: &str,
error: &str,
);
fn phase_progress(&self, update: &PhaseProgressUpdate);
fn op_starting(&self, _parent_phase: crate::scene_tree::SceneNodeId, _op_name: &str) {}
fn op_completed(
&self,
_parent_phase: crate::scene_tree::SceneNodeId,
_op_name: &str,
_duration_secs: f64,
) {
}
fn op_measure(
&self,
_parent_phase: crate::scene_tree::SceneNodeId,
_op_name: &str,
_text: &str,
) {
}
fn op_failed(
&self,
_parent_phase: crate::scene_tree::SceneNodeId,
_op_name: &str,
_error: &str,
) {
}
fn phase_render_attach(&self, _handle: PhaseRenderHandle) {}
fn run_finished(&self);
fn log(&self, level: LogLevel, message: &str);
fn log_tagged(&self, level: LogLevel, _tag: EventTag, message: &str) {
self.log(level, message);
}
fn suppresses_stderr(&self) -> bool {
false
}
fn live_suppress_flag(&self) -> Option<std::sync::Arc<std::sync::atomic::AtomicBool>> {
None
}
fn reporter(&self) -> Option<Box<dyn nmbrs_metrics::scheduler::Reporter>> {
None
}
fn reporters(
&self,
) -> Vec<(
std::time::Duration,
Box<dyn nmbrs_metrics::scheduler::Reporter>,
)> {
match self.reporter() {
Some(r) => vec![(std::time::Duration::from_secs(1), r)],
None => vec![],
}
}
fn cadences(&self) -> Option<nmbrs_metrics::cadence::Cadences> {
None
}
fn on_metrics_query(&self, _query: std::sync::Arc<nmbrs_metrics::metrics_query::MetricsQuery>) {
}
fn scenario_pre_mapped(&self, _tree: &crate::scene_tree::SceneTree) {}
}
#[derive(Clone, Debug)]
pub struct PhaseProgressUpdate {
pub exec_id: u64,
pub name: String,
pub labels: String,
pub cursor_name: String,
pub cursor_extent: u64,
pub daemon: bool,
pub rows_consumed: u64,
pub rows_total: u64,
pub fibers: usize,
pub ops_started: u64,
pub ops_finished: u64,
pub ops_ok: u64,
pub skips: u64,
pub errors: u64,
pub retries: u64,
pub ops_per_sec: f64,
pub adapter_counters: Vec<(String, u64, f64)>,
pub rows_per_batch: f64,
pub relevancy: Vec<crate::validation::RelevancyLive>,
}
#[derive(Clone)]
pub struct PhaseRenderHandle {
pub exec_id: u64,
pub name: String,
pub labels: String,
pub activity_name: String,
pub metrics: Arc<crate::activity::ActivityMetrics>,
pub bodies: Arc<Vec<crate::readouts::BakedBody>>,
pub memo: Arc<arc_swap::ArcSwap<String>>,
pub gutter: Arc<arc_swap::ArcSwapOption<crate::wrappers::gutter::GutterSpec>>,
pub status_metrics: Arc<[String]>,
pub concurrency: usize,
pub seq: Option<(usize, usize)>,
pub depth_indent: String,
}
impl std::fmt::Debug for PhaseRenderHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PhaseRenderHandle")
.field("exec_id", &self.exec_id)
.field("name", &self.name)
.field("labels", &self.labels)
.field("seq", &self.seq)
.field("bodies", &self.bodies.len())
.finish_non_exhaustive()
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
pub enum SkippedPhaseDisplay {
Elide,
#[default]
Mark,
Prune,
}
impl SkippedPhaseDisplay {
pub fn parse(s: &str) -> Option<Self> {
match s.trim().to_ascii_lowercase().as_str() {
"elide" => Some(Self::Elide),
"mark" => Some(Self::Mark),
"prune" => Some(Self::Prune),
_ => None,
}
}
}
static SKIPPED_PHASE_DISPLAY: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(1);
pub fn set_skipped_phase_display(mode: SkippedPhaseDisplay) {
let v = match mode {
SkippedPhaseDisplay::Elide => 0,
SkippedPhaseDisplay::Mark => 1,
SkippedPhaseDisplay::Prune => 2,
};
SKIPPED_PHASE_DISPLAY.store(v, std::sync::atomic::Ordering::Relaxed);
}
pub fn skipped_phase_display() -> SkippedPhaseDisplay {
match SKIPPED_PHASE_DISPLAY.load(std::sync::atomic::Ordering::Relaxed) {
0 => SkippedPhaseDisplay::Elide,
2 => SkippedPhaseDisplay::Prune,
_ => SkippedPhaseDisplay::Mark,
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
pub enum CompletedPhaseDisplay {
#[default]
Full,
Headers,
}
impl CompletedPhaseDisplay {
pub fn parse(s: &str) -> Option<Self> {
match s.trim().to_ascii_lowercase().as_str() {
"full" => Some(Self::Full),
"headers" => Some(Self::Headers),
_ => None,
}
}
}
static COMPLETED_PHASE_DISPLAY: std::sync::atomic::AtomicU8 = std::sync::atomic::AtomicU8::new(0);
pub fn set_completed_phase_display(mode: CompletedPhaseDisplay) {
let v = match mode {
CompletedPhaseDisplay::Full => 0,
CompletedPhaseDisplay::Headers => 1,
};
COMPLETED_PHASE_DISPLAY.store(v, std::sync::atomic::Ordering::Relaxed);
}
pub fn completed_phase_display() -> CompletedPhaseDisplay {
match COMPLETED_PHASE_DISPLAY.load(std::sync::atomic::Ordering::Relaxed) {
1 => CompletedPhaseDisplay::Headers,
_ => CompletedPhaseDisplay::Full,
}
}
static GLOBAL_OBSERVER: std::sync::OnceLock<Arc<dyn RunObserver>> = std::sync::OnceLock::new();
pub fn set_global_observer(observer: Arc<dyn RunObserver>) {
let _ = GLOBAL_OBSERVER.set(observer);
}
pub fn global_observer() -> Option<Arc<dyn RunObserver>> {
if let Some(obs) = crate::execution_context::current_observer() {
return Some(obs);
}
GLOBAL_OBSERVER.get().cloned()
}
pub fn set_log_file(path: &std::path::Path) -> std::io::Result<()> {
crate::log_sink::init(path)
}
pub fn use_color() -> bool {
use std::io::IsTerminal;
use std::sync::OnceLock;
static CACHE: OnceLock<bool> = OnceLock::new();
*CACHE.get_or_init(|| {
if std::env::var_os("NO_COLOR").is_some() {
return false;
}
std::io::stderr().is_terminal()
})
}
static EXPLAIN_HELD_UNTIL_NS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
static EXPLAIN_LAST_PRESS_NS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
const EXPLAIN_AUTO_REVERT_MS: u64 = 10_000;
const EXPLAIN_TOGGLE_DEBOUNCE_MS: u64 = 250;
fn now_nanos() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0)
}
pub fn toggle_explain() {
let now = now_nanos();
let last_press = EXPLAIN_LAST_PRESS_NS.load(std::sync::atomic::Ordering::Acquire);
if last_press != 0 && now.saturating_sub(last_press) < EXPLAIN_TOGGLE_DEBOUNCE_MS * 1_000_000 {
return;
}
EXPLAIN_LAST_PRESS_NS.store(now, std::sync::atomic::Ordering::Release);
let currently_on = {
let deadline = EXPLAIN_HELD_UNTIL_NS.load(std::sync::atomic::Ordering::Acquire);
deadline != 0 && now < deadline
};
if currently_on {
EXPLAIN_HELD_UNTIL_NS.store(0, std::sync::atomic::Ordering::Release);
} else {
let deadline = now.saturating_add(EXPLAIN_AUTO_REVERT_MS * 1_000_000);
EXPLAIN_HELD_UNTIL_NS.store(deadline, std::sync::atomic::Ordering::Release);
}
}
pub fn is_explain_held() -> bool {
let deadline = EXPLAIN_HELD_UNTIL_NS.load(std::sync::atomic::Ordering::Acquire);
deadline != 0 && now_nanos() < deadline
}
static READOUT_PAUSED: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
pub fn toggle_readout_pause() -> bool {
!READOUT_PAUSED.fetch_xor(true, std::sync::atomic::Ordering::AcqRel)
}
pub fn readouts_paused() -> bool {
READOUT_PAUSED.load(std::sync::atomic::Ordering::Acquire)
}
static RETAIN_LEVEL: std::sync::OnceLock<LogLevel> = std::sync::OnceLock::new();
pub fn set_retain_level(level: LogLevel) {
let _ = RETAIN_LEVEL.set(level);
}
pub fn retain_level() -> LogLevel {
*RETAIN_LEVEL.get().unwrap_or(&LogLevel::Debug)
}
static DISPLAY_LEVEL: std::sync::OnceLock<LogLevel> = std::sync::OnceLock::new();
pub fn set_display_level(level: LogLevel) {
let _ = DISPLAY_LEVEL.set(level);
}
pub fn display_level() -> LogLevel {
*DISPLAY_LEVEL.get().unwrap_or(&LogLevel::Info)
}
pub fn log(level: LogLevel, message: &str) {
log_tagged(level, EventTag::default(), message);
}
pub fn log_tagged(level: LogLevel, tag: EventTag, message: &str) {
if level >= retain_level()
&& let Some(sink) = crate::log_sink::global()
{
let tag = match level {
LogLevel::Trace => "TRC",
LogLevel::Debug => "DBG",
LogLevel::Info => "INF",
LogLevel::Warn => "WRN",
LogLevel::Error => "ERR",
};
let ts = crate::session::now_log_timestamp();
let line = format!(
"{ts} {tag} {}\n",
crate::readouts::snapshot::strip_ansi(message)
)
.into_bytes();
let _ = sink.try_send(line);
}
if let Some(obs) = global_observer() {
obs.log_tagged(level, tag, message);
} else {
crate::output_channel::log_to_surface(level, message);
}
}
pub fn op_output(line: &str) {
if let Some(channel) = crate::output_channel::installed() {
channel.op_output(line);
return;
}
op_output_raw(line);
}
pub(crate) fn op_output_raw(line: &str) {
if LogLevel::Info >= retain_level()
&& let Some(sink) = crate::log_sink::global()
{
let ts = crate::session::now_log_timestamp();
let bytes =
format!("{ts} INF {}\n", crate::readouts::snapshot::strip_ansi(line)).into_bytes();
let _ = sink.try_send(bytes);
}
use std::io::Write;
let mut out = std::io::stdout().lock();
let _ = writeln!(out, "{line}");
let _ = out.flush();
}
pub fn colorize_log_line(level: LogLevel, message: &str) -> String {
if !use_color() {
return message.to_string();
}
let (color, reset) = match level {
LogLevel::Trace => ("\x1b[2;90m", "\x1b[0m"),
LogLevel::Debug => ("\x1b[2m", "\x1b[0m"),
LogLevel::Info => ("", ""),
LogLevel::Warn => ("\x1b[33m", "\x1b[0m"),
LogLevel::Error => ("\x1b[1;31m", "\x1b[0m"),
};
if color.is_empty() {
message.to_string()
} else {
format!("{color}{message}{reset}")
}
}
#[macro_export]
macro_rules! diag {
($level:expr, $($arg:tt)*) => {
$crate::observer::log($level, &format!($($arg)*))
};
}
pub fn trace(labels: &nmbrs_metrics::labels::Labels, message: &str) {
crate::trace_router::log(labels, message);
}
pub fn trace_enabled() -> bool {
crate::trace_router::enabled()
}
#[macro_export]
macro_rules! trace_event {
($labels:expr, $($arg:tt)*) => {
if $crate::observer::trace_enabled() {
$crate::observer::trace($labels, &format!($($arg)*));
}
};
}
use std::sync::Arc;
pub struct StderrObserver {
pub min_level: LogLevel,
}
impl Default for StderrObserver {
fn default() -> Self {
Self {
min_level: LogLevel::Info,
}
}
}
impl StderrObserver {
pub fn with_min_level(min_level: LogLevel) -> Self {
Self { min_level }
}
}
impl RunObserver for StderrObserver {
fn phase_starting(
&self,
_scene_node_id: crate::scene_tree::SceneNodeId,
name: &str,
_labels: &str,
op_templates: usize,
total_cycles: u64,
concurrency: usize,
) {
let template_word = if op_templates == 1 {
"op template"
} else {
"op templates"
};
let cycle_word = if total_cycles == 1 { "cycle" } else { "cycles" };
crate::observer::log(
LogLevel::Info,
&format!(
"phase '{name}': {op_templates} {template_word}, {total_cycles} {cycle_word}, concurrency={concurrency}"
),
);
}
fn phase_completed(
&self,
_scene_node_id: crate::scene_tree::SceneNodeId,
_name: &str,
_labels: &str,
_duration_secs: f64,
) {
}
fn phase_failed(
&self,
_scene_node_id: crate::scene_tree::SceneNodeId,
_name: &str,
_labels: &str,
_error: &str,
) {
}
fn phase_progress(&self, _update: &PhaseProgressUpdate) {
}
fn run_finished(&self) {
crate::observer::log(LogLevel::Info, "all phases complete");
}
fn log(&self, level: LogLevel, message: &str) {
if level >= self.min_level {
if message.starts_with("session: graceful shutdown requested")
|| message.starts_with("session: force-exit")
{
eprintln!();
}
crate::output_channel::log_to_surface(level, message);
}
}
}
#[cfg(test)]
mod display_mode_tests {
use super::*;
#[test]
fn completed_phase_display_parses_and_rejects() {
assert_eq!(
CompletedPhaseDisplay::parse("full"),
Some(CompletedPhaseDisplay::Full)
);
assert_eq!(
CompletedPhaseDisplay::parse(" HEADERS "),
Some(CompletedPhaseDisplay::Headers)
);
assert_eq!(CompletedPhaseDisplay::parse("collapse"), None);
assert_eq!(
CompletedPhaseDisplay::default(),
CompletedPhaseDisplay::Full
);
}
}