use std::sync::Arc;
use std::sync::atomic::Ordering;
use crate::lifecycle::EventType;
use crate::readouts::{LifecycleState, ReadoutContext};
pub struct ActivityReadoutContext {
pub phase_name: String,
pub phase_seq: Option<(usize, usize)>,
pub phase_labels: String,
pub cycles_completed: u64,
pub cycles_total: u64,
pub ops_ok: u64,
pub skips: u64,
pub errors: u64,
pub retries: u64,
pub concurrency: usize,
pub elapsed_secs: f64,
pub consumed: u64,
pub status_metric_chips: String,
pub depth_indent: String,
pub use_color: bool,
pub memo: String,
pub outcome: crate::phase_outcome::Outcome,
pub outcome_errors: Vec<crate::phase_outcome::PhaseErrorDetail>,
pub outcome_resume_cursor: Option<crate::phase_outcome::ResumeCursor>,
pub open_ended: bool,
}
impl ReadoutContext for ActivityReadoutContext {
fn subject_name(&self) -> &str {
&self.phase_name
}
fn subject_seq(&self) -> Option<(usize, usize)> {
self.phase_seq
}
fn subject_labels(&self) -> &str {
&self.phase_labels
}
fn open_ended(&self) -> bool {
self.open_ended
}
fn cycles_completed(&self) -> u64 {
self.cycles_completed
}
fn cycles_total(&self) -> u64 {
self.cycles_total
}
fn ops_ok(&self) -> u64 {
self.ops_ok
}
fn skips(&self) -> u64 {
self.skips
}
fn errors(&self) -> u64 {
self.errors
}
fn retries(&self) -> u64 {
self.retries
}
fn concurrency(&self) -> usize {
self.concurrency
}
fn elapsed_secs(&self) -> f64 {
self.elapsed_secs
}
fn consumed(&self) -> u64 {
self.consumed
}
fn status_metric_chips(&self) -> String {
self.status_metric_chips.clone()
}
fn depth_indent(&self) -> &str {
&self.depth_indent
}
fn use_color(&self) -> bool {
self.use_color
}
fn event(&self) -> EventType {
EventType::PhaseEnd
}
fn subject_state(&self) -> LifecycleState {
match self.outcome.validity {
crate::phase_outcome::Validity::Succeeded => LifecycleState::Completed,
crate::phase_outcome::Validity::Failed => LifecycleState::Failed(
self.outcome_errors
.first()
.map(|e| e.message.clone())
.unwrap_or_else(|| "phase failed".into()),
),
}
}
fn phase_memo(&self) -> &str {
&self.memo
}
fn outcome(&self) -> crate::phase_outcome::Outcome {
self.outcome.clone()
}
fn outcome_errors(&self) -> &[crate::phase_outcome::PhaseErrorDetail] {
&self.outcome_errors
}
fn outcome_resume_cursor(&self) -> Option<&crate::phase_outcome::ResumeCursor> {
self.outcome_resume_cursor.as_ref()
}
}
pub struct LifecycleContext {
pub event: crate::lifecycle::EventType,
pub subject_name: String,
pub subject_labels: String,
pub depth_indent: String,
pub use_color: bool,
pub stick_reattached: String,
}
impl ReadoutContext for LifecycleContext {
fn subject_name(&self) -> &str {
&self.subject_name
}
fn subject_seq(&self) -> Option<(usize, usize)> {
None
}
fn subject_labels(&self) -> &str {
&self.subject_labels
}
fn cycles_completed(&self) -> u64 {
0
}
fn cycles_total(&self) -> u64 {
0
}
fn ops_ok(&self) -> u64 {
0
}
fn errors(&self) -> u64 {
0
}
fn retries(&self) -> u64 {
0
}
fn concurrency(&self) -> usize {
0
}
fn elapsed_secs(&self) -> f64 {
0.0
}
fn consumed(&self) -> u64 {
0
}
fn status_metric_chips(&self) -> String {
String::new()
}
fn depth_indent(&self) -> &str {
&self.depth_indent
}
fn use_color(&self) -> bool {
self.use_color
}
fn event(&self) -> crate::lifecycle::EventType {
self.event
}
fn stick_reattached_session(&self) -> &str {
&self.stick_reattached
}
fn subject_state(&self) -> LifecycleState {
LifecycleState::Running
}
}
pub struct InlineRefreshContext {
pub phase_name: String,
pub activity_name: String,
pub phase_seq: Option<(usize, usize)>,
pub phase_labels: String,
pub cycles_completed: u64,
pub cycles_total: u64,
pub ops_started: u64,
pub ops_finished: u64,
pub ops_ok: u64,
pub skips: u64,
pub errors: u64,
pub retries: u64,
pub attempt_ok: u64,
pub attempt_failed: u64,
pub concurrency: usize,
pub elapsed_secs: f64,
pub consumed: u64,
pub rows_consumed: u64,
pub rows_total: u64,
pub status_metric_chips: String,
pub adapter_counters_text: String,
pub batch_info_text: String,
pub depth_indent: String,
pub refresh_tick: u64,
pub use_color: bool,
pub memo: String,
pub progress_override: Option<f64>,
pub progress_override_elapsed: Option<f64>,
pub open_ended: bool,
pub lat_p50_nanos: u64,
pub lat_p99_nanos: u64,
}
impl ReadoutContext for InlineRefreshContext {
fn subject_name(&self) -> &str {
&self.phase_name
}
fn activity_name(&self) -> &str {
&self.activity_name
}
fn subject_seq(&self) -> Option<(usize, usize)> {
self.phase_seq
}
fn subject_labels(&self) -> &str {
&self.phase_labels
}
fn cycles_completed(&self) -> u64 {
self.cycles_completed
}
fn cycles_total(&self) -> u64 {
self.cycles_total
}
fn ops_started(&self) -> u64 {
self.ops_started
}
fn ops_finished(&self) -> u64 {
self.ops_finished
}
fn ops_ok(&self) -> u64 {
self.ops_ok
}
fn skips(&self) -> u64 {
self.skips
}
fn errors(&self) -> u64 {
self.errors
}
fn retries(&self) -> u64 {
self.retries
}
fn attempt_ok(&self) -> u64 {
self.attempt_ok
}
fn attempt_failed(&self) -> u64 {
self.attempt_failed
}
fn concurrency(&self) -> usize {
self.concurrency
}
fn elapsed_secs(&self) -> f64 {
self.elapsed_secs
}
fn consumed(&self) -> u64 {
self.consumed
}
fn rows_consumed(&self) -> u64 {
self.rows_consumed
}
fn rows_total(&self) -> u64 {
self.rows_total
}
fn status_metric_chips(&self) -> String {
self.status_metric_chips.clone()
}
fn adapter_counters_text(&self) -> String {
self.adapter_counters_text.clone()
}
fn batch_info_text(&self) -> String {
self.batch_info_text.clone()
}
fn depth_indent(&self) -> &str {
&self.depth_indent
}
fn use_color(&self) -> bool {
self.use_color
}
fn event(&self) -> EventType {
EventType::Update
}
fn refresh_tick(&self) -> u64 {
self.refresh_tick
}
fn phase_memo(&self) -> &str {
&self.memo
}
fn progress_override(&self) -> Option<f64> {
self.progress_override
}
fn open_ended(&self) -> bool {
self.open_ended
}
fn latency_p50_nanos(&self) -> u64 {
self.lat_p50_nanos
}
fn latency_p99_nanos(&self) -> u64 {
self.lat_p99_nanos
}
fn eta_secs(&self) -> Option<f64> {
if self.open_ended {
return None;
}
if let (Some(f), Some(e)) = (self.progress_override, self.progress_override_elapsed)
&& f > 0.0
&& f < 1.0
&& e > 0.0
{
return Some(e * (1.0 - f) / f);
}
if self.rows_total > 0 && self.elapsed_secs > 0.0 {
if self.rows_consumed == 0 {
return None;
}
let rate = self.rows_consumed as f64 / self.elapsed_secs;
let remaining = self.rows_total.saturating_sub(self.rows_consumed) as f64;
return Some(remaining / rate);
}
if self.cycles_total == 0 || self.elapsed_secs <= 0.0 {
return None;
}
let rate = self.ops_finished as f64 / self.elapsed_secs;
if rate <= 0.0 {
return None;
}
let remaining = self.cycles_total.saturating_sub(self.ops_finished) as f64;
Some(remaining / rate)
}
}
pub fn fire_lifecycle(
event: crate::lifecycle::EventType,
bindings: &nmbrs_workload::model::ReadoutsBindings,
default: Option<crate::readouts::BakedBody>,
ctx: &dyn crate::readouts::ReadoutContext,
snapshot_writer: Option<&crate::readouts::snapshot::SnapshotWriter>,
) {
use crate::readouts::ReadoutBinder;
let seed = default.unwrap_or_default();
let mut binder = match crate::readouts::build_event_binder(bindings, event, seed) {
Ok(b) => b,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"readouts: failed to bind {slot} — {e}",
slot = event.slot_name()
);
return;
}
};
let mut sink = crate::readouts::StringSink::with_capacity(128);
binder.fire(event, ctx, &mut sink);
let rendered = sink.take();
if rendered.trim().is_empty() {
return; }
let tag = crate::observer::EventTag::at(event, crate::observer::EventCategory::General);
crate::observer::log_tagged(crate::observer::LogLevel::Info, tag, &rendered);
let subject_id = ctx.subject_id();
crate::readouts::snapshot::capture(
snapshot_writer,
event.slot_name(),
ctx.subject_exec_id(),
event.subject_kind().as_str(),
&subject_id,
"binder",
crate::readouts::snapshot::lod_str(crate::readouts::Lod::Labeled),
&rendered,
);
}
pub fn resolve_phase_coord_by_name(activity_name: &str) -> (Option<(usize, usize)>, String) {
let bare_name = activity_name
.split_once(" (")
.map(|(n, _)| n)
.unwrap_or(activity_name);
crate::scene_tree::current()
.and_then(|t| {
let node = t
.dfs_phases()
.find(|n| {
n.name == bare_name
&& matches!(n.status, crate::scene_tree::PhaseStatus::Running)
})?
.clone();
let seq = node.seq?;
let depth = node.depth.saturating_sub(1);
Some((Some((seq, t.total_phases())), " ".repeat(depth)))
})
.unwrap_or((None, String::new()))
}
pub fn resolve_phase_coord_by_id(
scene_node_id: crate::scene_tree::SceneNodeId,
) -> (Option<(usize, usize)>, String) {
crate::scene_tree::current()
.and_then(|t| {
let node = t.nodes.get(scene_node_id)?;
let seq = node.seq?;
let depth = node.depth.saturating_sub(1);
Some((Some((seq, t.total_phases())), " ".repeat(depth)))
})
.unwrap_or((None, String::new()))
}
pub fn is_internal_counter(name: &str) -> bool {
name.starts_with('_')
}
pub(crate) fn rows_per_batch(
rows_inserted: Option<u64>,
batch_writes: Option<u64>,
stanzas: u64,
) -> Option<f64> {
let rows = rows_inserted?;
match batch_writes {
Some(batches) if batches > 0 => (rows > batches).then(|| rows as f64 / batches as f64),
_ => (stanzas > 0 && rows > stanzas).then(|| rows as f64 / stanzas as f64),
}
}
#[allow(clippy::too_many_arguments)]
pub fn build_inline_refresh_context(
progress_metrics: &Arc<crate::activity::ActivityMetrics>,
activity_name: &str,
concurrency: usize,
total_extent: u64,
rows_consumed: u64,
rows_total: u64,
elapsed_secs: f64,
refresh_tick: u64,
status_metrics: &[String],
memo: &arc_swap::ArcSwap<String>,
phase_seq: Option<(usize, usize)>,
depth_indent: String,
open_ended: bool,
) -> InlineRefreshContext {
let started = progress_metrics.ops_started.load(Ordering::Relaxed);
let finished = progress_metrics.ops_finished.load(Ordering::Relaxed);
let ops_completed = progress_metrics.cycles_completed();
let successes = progress_metrics.result_success.count();
let errors = progress_metrics.errors_total.get();
let failed_ops = ops_completed
.saturating_sub(successes)
.saturating_sub(progress_metrics.skips_total.get());
let consumed = finished;
let attempt_ok = progress_metrics.attempt_success.count();
let attempt_failed = progress_metrics.attempt_failure.count();
let retries = attempt_failed.saturating_sub(failed_ops);
let mut adapter_counters_text = String::new();
let counters = progress_metrics.collect_status_counters();
for (name, total) in &counters {
if is_internal_counter(name) {
continue;
}
let item_rate = if elapsed_secs > 0.0 {
*total as f64 / elapsed_secs
} else {
0.0
};
let rate_str = if item_rate >= 1_000_000.0 {
format!("{:.1}M", item_rate / 1_000_000.0)
} else if item_rate >= 1_000.0 {
format!("{:.1}K", item_rate / 1_000.0)
} else {
format!("{:.0}", item_rate)
};
adapter_counters_text.push_str(&format!(" {name}:{rate_str}/s"));
}
let stanzas = progress_metrics.stanzas_total.get();
let find_counter = |want: &str| counters.iter().find(|(n, _)| n == want).map(|(_, t)| *t);
let batch_info_text = rows_per_batch(
find_counter("rows_inserted"),
find_counter("_batch_writes"),
stanzas,
)
.map(|avg| format!(" rows/batch:{avg:.1}"))
.unwrap_or_default();
let status_metric_chips = progress_metrics
.collect_status_values(status_metrics)
.concat();
let bare_name = activity_name
.split_once(" (")
.map(|(n, _)| n)
.unwrap_or(activity_name);
let memo_snapshot: String = memo.load().as_str().to_string();
let (lat_p50, lat_p99) = {
let snap = progress_metrics.service_time.peek_snapshot();
let h = &snap.histogram;
if h.is_empty() {
(0, 0)
} else {
(h.value_at_quantile(0.50), h.value_at_quantile(0.99))
}
};
InlineRefreshContext {
phase_name: bare_name.to_string(),
activity_name: activity_name.to_string(),
phase_seq,
phase_labels: String::new(),
cycles_completed: ops_completed,
cycles_total: total_extent,
ops_started: started,
ops_finished: finished,
ops_ok: successes,
skips: progress_metrics.skips_total.get(),
errors,
retries,
attempt_ok,
attempt_failed,
concurrency,
elapsed_secs,
consumed,
rows_consumed,
rows_total,
status_metric_chips,
adapter_counters_text,
batch_info_text,
depth_indent,
refresh_tick,
use_color: crate::observer::use_color(),
memo: memo_snapshot,
progress_override: progress_metrics.progress_override(),
progress_override_elapsed: progress_metrics.progress_override_elapsed_secs(),
open_ended,
lat_p50_nanos: lat_p50,
lat_p99_nanos: lat_p99,
}
}
#[cfg(test)]
mod tests {
#[test]
fn eta_prefers_measured_basis_when_override_present() {
let ctx = super::InlineRefreshContext {
phase_name: String::new(),
activity_name: String::new(),
phase_seq: None,
phase_labels: String::new(),
cycles_completed: 3,
cycles_total: 4,
ops_started: 4,
ops_finished: 3,
ops_ok: 3,
skips: 0,
errors: 0,
retries: 0,
attempt_ok: 3,
attempt_failed: 0,
concurrency: 1,
elapsed_secs: 600.0,
consumed: 3,
rows_consumed: 0,
rows_total: 0,
status_metric_chips: String::new(),
adapter_counters_text: String::new(),
batch_info_text: String::new(),
depth_indent: String::new(),
refresh_tick: 0,
use_color: false,
memo: String::new(),
progress_override: Some(0.25),
progress_override_elapsed: Some(60.0),
open_ended: false,
lat_p50_nanos: 0,
lat_p99_nanos: 0,
};
use crate::readouts::ReadoutContext;
let eta = ctx.eta_secs().expect("measured ETA");
assert!((eta - 180.0).abs() < 1e-9, "eta={eta}");
let ctx2 = super::InlineRefreshContext {
progress_override_elapsed: None,
..ctx
};
let eta2 = ctx2.eta_secs().expect("cycle ETA");
assert!((eta2 - 200.0).abs() < 1e-9, "eta2={eta2}");
}
#[test]
fn eta_uses_row_basis_for_cursor_phases() {
let ctx = super::InlineRefreshContext {
phase_name: String::new(),
activity_name: String::new(),
phase_seq: None,
phase_labels: String::new(),
cycles_completed: 10_000,
cycles_total: 10_000_000,
ops_started: 10_000,
ops_finished: 10_000,
ops_ok: 10_000,
skips: 0,
errors: 0,
retries: 0,
attempt_ok: 10_000,
attempt_failed: 0,
concurrency: 1,
elapsed_secs: 200.0,
consumed: 10_000,
rows_consumed: 2_000_000,
rows_total: 10_000_000,
status_metric_chips: String::new(),
adapter_counters_text: String::new(),
batch_info_text: String::new(),
depth_indent: String::new(),
refresh_tick: 0,
use_color: false,
memo: String::new(),
progress_override: None,
progress_override_elapsed: None,
open_ended: false,
lat_p50_nanos: 0,
lat_p99_nanos: 0,
};
use crate::readouts::ReadoutContext;
let eta = ctx.eta_secs().expect("row-basis ETA");
assert!((eta - 800.0).abs() < 1e-6, "eta={eta}");
let ctx2 = super::InlineRefreshContext {
rows_consumed: 0,
..ctx
};
assert!(
ctx2.eta_secs().is_none(),
"zero-row cursor phase must have no ETA"
);
}
use super::{is_internal_counter, rows_per_batch};
#[test]
fn internal_counter_is_underscore_prefixed() {
assert!(is_internal_counter("_batch_writes"));
assert!(!is_internal_counter("rows_inserted"));
assert!(!is_internal_counter("queries"));
}
#[test]
fn prefers_batch_writes_over_stanzas() {
let avg = rows_per_batch(Some(1000), Some(5), 40);
assert_eq!(avg, Some(200.0));
}
#[test]
fn falls_back_to_stanzas_without_batch_writes() {
let avg = rows_per_batch(Some(1000), None, 40);
assert_eq!(avg, Some(25.0));
}
#[test]
fn zero_batch_writes_uses_fallback() {
let avg = rows_per_batch(Some(1000), Some(0), 40);
assert_eq!(avg, Some(25.0));
}
#[test]
fn no_batch_observed_is_none() {
assert_eq!(rows_per_batch(Some(5), Some(5), 40), None);
assert_eq!(rows_per_batch(Some(40), None, 40), None);
assert_eq!(rows_per_batch(None, Some(5), 40), None);
}
}