use std::convert::Infallible;
use std::sync::Arc;
use axum::Extension;
use axum::extract::State;
use axum::http::StatusCode;
use axum::response::sse::{Event, KeepAlive, Sse};
use axum_extra::extract::Query as QueryExtra;
use futures::stream::{self, Stream, StreamExt};
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use crate::ingest::aggregate::{self, Scope};
use crate::ingest::decode::SchedulingDelayKind;
use crate::ingest::refine::{self, FoldErrors, FoldOutcome, RefineOpts, Resolved};
use crate::server::AppState;
use crate::server::credentials::MaybeCreds;
use crate::server::metrics::OperationMetrics;
use arrow::array::Array;
const DURATION_FLOOR_NS: i64 = 100_000;
const LONG_POLL_TOP: usize = 100;
const LONG_POLL_SOFT_CAP: usize = LONG_POLL_TOP * 4;
const SCHEDULING_DELAY_TOP: usize = 100;
const SCHEDULING_DELAY_SOFT_CAP: usize = SCHEDULING_DELAY_TOP * 4;
const HIGH_SCHEDULING_DELAY_NS: i64 = 1_000_000;
#[derive(Deserialize)]
pub struct TokioStatsParams {
pub bucket: Option<String>,
pub prefix: Option<String>,
pub aws_region: Option<String>,
pub service: Option<String>,
#[serde(default)]
pub host: Vec<String>,
pub start_ns: Option<i64>,
pub end_ns: Option<i64>,
pub max_files: Option<usize>,
}
#[derive(Serialize)]
pub struct TokioStatsResponse {
pub time_span_ns: i64,
pub total_polls: u64,
pub bucket: String,
pub by_spawn_loc: Vec<SpawnLocStats>,
pub top_long_polls: Vec<LongPoll>,
pub top_scheduling_delays: Vec<SchedulingDelay>,
pub scheduling_delay_coverage: SchedulingDelayCoverage,
pub worker_activity: Vec<WorkerStats>,
#[serde(skip_serializing_if = "Option::is_none")]
pub coverage: Option<aggregate::Coverage>,
}
#[derive(Serialize)]
pub struct SpawnLocStats {
pub spawn_loc: String,
pub total_polls: u64,
pub durations_ns: Vec<i64>,
pub classes: Vec<u8>,
pub exemplars: [Option<PollExemplar>; 4],
}
#[derive(Serialize, Clone)]
pub struct PollExemplar {
pub start_ns: i64,
pub end_ns: i64,
pub duration_ns: i64,
pub host: String,
pub source_key: String,
}
#[derive(Serialize, Clone)]
pub struct LongPoll {
pub duration_ns: i64,
pub worker_id: u32,
pub task_id: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub spawn_loc: Option<String>,
pub start_ns: i64,
pub end_ns: i64,
pub host: String,
pub source_key: String,
}
#[derive(Serialize, Clone)]
pub struct SchedulingDelay {
pub delay_ns: i64,
pub ready_at_ns: i64,
pub poll_start_ns: i64,
pub poll_end_ns: i64,
pub worker_id: u32,
pub task_id: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub spawn_loc: Option<String>,
pub kind: SchedulingDelayKind,
#[serde(skip_serializing_if = "Option::is_none")]
pub waker_task_id: Option<u64>,
pub host: String,
pub source_key: String,
}
#[derive(Serialize, Clone, Default)]
pub struct SchedulingDelayCoverage {
pub observed_polls: u64,
pub unmeasured_polls: u64,
pub over_1ms_polls: u64,
pub spawn_inferred_polls: u64,
pub wake_observed_polls: u64,
pub wake_during_poll_polls: u64,
pub uninstrumented_unmeasured_polls: u64,
pub instrumentation_unknown_unmeasured_polls: u64,
pub missing_readiness_unmeasured_polls: u64,
}
#[derive(Serialize, Clone)]
pub struct WorkerStats {
pub worker_id: u32,
pub host: String,
pub total_polls: u64,
pub busy_ns: i64,
pub span_ns: i64,
pub busy_pct: f64,
pub notable_polls: u64,
pub worst_poll_ns: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub worst_exemplar: Option<PollExemplar>,
}
pub async fn get_tokio_stats(
State(state): State<AppState>,
creds: MaybeCreds,
QueryExtra(params): QueryExtra<TokioStatsParams>,
) -> Result<
(
Extension<OperationMetrics>,
Sse<impl Stream<Item = Result<Event, Infallible>>>,
),
(StatusCode, String),
> {
let Some(agg) = state
.agg_context_for(
params.bucket.as_deref(),
params.prefix.as_deref(),
params.aws_region.as_deref(),
creds,
)
.await?
else {
return Err((
StatusCode::NOT_FOUND,
"tokio-stats requires aggregation (start with --agg or supply a bucket)".to_string(),
));
};
let scope = Scope {
start_ns: params.start_ns,
end_ns: params.end_ns,
service: params.service.clone(),
hosts: params.host.clone(),
};
tracing::debug!(
source_bucket = %agg.source_bucket,
source_prefixes = ?agg.source_prefixes,
service = scope.service.as_deref().unwrap_or("(all)"),
hosts = ?scope.hosts,
start_ns = ?scope.start_ns,
end_ns = ?scope.end_ns,
"tokio-stats: starting"
);
let opts = RefineOpts {
max_files: params.max_files,
};
let Some(resolved) = refine::resolve(&agg, &scope, opts).await else {
return Err((
StatusCode::NOT_FOUND,
"no source files match this scope".to_string(),
));
};
let op = OperationMetrics::tokio_stats(
resolved.files_matched as u32,
resolved.capped_files_folded_in(resolved.folded()) as u32,
None,
);
let stream = tokio_stats_stream(agg, resolved, scope, state.fold_limits.clone());
Ok((
Extension(op),
Sse::new(stream).keep_alive(KeepAlive::default()),
))
}
struct StreamCtx {
agg: Arc<crate::ingest::aggregate::AggContext>,
resolved: Resolved,
scope: Scope,
source_bucket: String,
}
enum Phase {
Start,
Folding {
acc: Box<TokioStatsAccum>,
folded: HashSet<String>,
errors: FoldErrors,
},
}
fn tokio_stats_stream(
agg: crate::ingest::aggregate::AggContext,
resolved: Resolved,
scope: Scope,
limits: aggregate::FoldLimits,
) -> impl Stream<Item = Result<Event, Infallible>> + use<> {
let agg = Arc::new(agg);
let ctx = Arc::new(StreamCtx {
agg: Arc::clone(&agg),
source_bucket: agg.source_bucket.clone(),
scope,
resolved,
});
let folds = Box::pin(refine::fold_stream(
agg,
limits,
ctx.resolved.unfolded_capped(),
));
stream::unfold(
(ctx, folds, Phase::Start),
|(ctx, mut folds, phase)| async move {
match phase {
Phase::Start => {
let polls_data = aggregate::read_polls_parts(
&*ctx.agg.output,
&ctx.agg.output_bucket,
&ctx.agg.output_prefix,
&ctx.resolved.capped,
ctx.resolved.folded(),
)
.await;
let mut acc = TokioStatsAccum::default();
for (raw_key, data) in &polls_data {
read_polls_part_lossy(data, &ctx.scope, raw_key, &mut acc);
}
let folded = ctx.resolved.folded().clone();
let errors = FoldErrors::default();
let event = snapshot_event(&ctx, &acc, &folded, &errors);
Some((
Ok(event),
(
ctx,
folds,
Phase::Folding {
acc: Box::new(acc),
folded,
errors,
},
),
))
}
Phase::Folding {
mut acc,
mut folded,
mut errors,
} => {
match folds.next().await? {
FoldOutcome::Folded(f) => {
if let Some(data) = aggregate::fetch_polls_part(
&*ctx.agg.output,
&ctx.agg.output_bucket,
&ctx.agg.output_prefix,
&f.full_key,
)
.await
{
read_polls_part_lossy(&data, &ctx.scope, &f.raw_key, &mut acc);
}
folded.insert(aggregate::part_leaf_of(&f.full_key));
}
FoldOutcome::Failed { raw_key, error } => {
errors.record(&raw_key, &error);
}
}
let event = snapshot_event(&ctx, &acc, &folded, &errors);
Some((
Ok(event),
(
ctx,
folds,
Phase::Folding {
acc,
folded,
errors,
},
),
))
}
}
},
)
}
fn read_polls_part_lossy(data: &[u8], scope: &Scope, source_key: &str, acc: &mut TokioStatsAccum) {
use dial9_core::rate_limited;
if let Err((_, e)) = read_polls_part(data, scope, source_key, acc) {
rate_limited!(std::time::Duration::from_secs(60), {
tracing::warn!(key = %source_key, error = %e, "tokio-stats: failed to read polls part");
});
}
}
fn snapshot_event(
ctx: &StreamCtx,
acc: &TokioStatsAccum,
folded: &HashSet<String>,
errors: &FoldErrors,
) -> Event {
let time_span_ns = match (acc.min_ts, acc.max_ts) {
(Some(min), Some(max)) => (max - min).max(1),
_ => 1,
};
let mut by_spawn_loc: Vec<SpawnLocStats> = acc
.by_loc
.iter()
.map(|(loc, la)| {
let mut durs = la.durations.clone();
durs.sort_unstable_by_key(|&(d, _)| std::cmp::Reverse(d));
let durations_ns = durs.iter().map(|(d, _)| *d).collect();
let classes = durs.iter().map(|(_, c)| *c).collect();
SpawnLocStats {
spawn_loc: loc.clone(),
total_polls: la.total,
durations_ns,
classes,
exemplars: la.worst_by_class.clone(),
}
})
.collect();
by_spawn_loc.sort_by_key(|l| std::cmp::Reverse(l.durations_ns.len()));
let files_folded = ctx.resolved.capped_files_folded_in(folded);
let resp = TokioStatsResponse {
time_span_ns,
total_polls: acc.total_polls,
bucket: ctx.source_bucket.clone(),
by_spawn_loc,
top_long_polls: acc.top_long_polls(),
top_scheduling_delays: acc.top_scheduling_delays(),
scheduling_delay_coverage: acc.scheduling_delay_coverage.clone(),
worker_activity: acc.worker_activity(),
coverage: Some(aggregate::Coverage {
files_matched: ctx.resolved.files_matched,
files_folded,
folded_set_id: None,
target_folded_set_id: None,
fold_work_cap: ctx.resolved.fold_work_cap(),
samples_folded: files_folded,
total_bytes: ctx.resolved.total_bytes,
hosts_matched: ctx.resolved.hosts_matched,
hosts_folded: ctx.resolved.capped_folded_hosts(folded),
fold_errors: errors.count,
fold_error_sample: errors.sample.clone(),
}),
};
Event::default().json_data(&resp).unwrap_or_else(|e| {
use dial9_core::rate_limited;
rate_limited!(std::time::Duration::from_secs(60), {
tracing::warn!(error = %e, "tokio-stats: event serialize failed");
});
Event::default().comment("serialize error")
})
}
#[derive(Default)]
struct TokioStatsAccum {
total_polls: u64,
min_ts: Option<i64>,
max_ts: Option<i64>,
by_loc: HashMap<String, LocAccum>,
long_polls: Vec<LongPoll>,
by_worker: HashMap<(String, u32), WorkerAccum>,
scheduling_delays: Vec<SchedulingDelay>,
scheduling_delay_coverage: SchedulingDelayCoverage,
}
impl TokioStatsAccum {
fn push_long_poll(&mut self, poll: LongPoll) {
self.long_polls.push(poll);
if self.long_polls.len() > LONG_POLL_SOFT_CAP {
self.long_polls
.sort_unstable_by_key(|p| std::cmp::Reverse(p.duration_ns));
self.long_polls.truncate(LONG_POLL_TOP);
}
}
fn top_long_polls(&self) -> Vec<LongPoll> {
let mut top = self.long_polls.clone();
top.sort_unstable_by_key(|p| std::cmp::Reverse(p.duration_ns));
top.truncate(LONG_POLL_TOP);
top
}
fn push_scheduling_delay(&mut self, delay: SchedulingDelay) {
self.scheduling_delays.push(delay);
if self.scheduling_delays.len() > SCHEDULING_DELAY_SOFT_CAP {
self.scheduling_delays
.sort_unstable_by_key(|d| std::cmp::Reverse(d.delay_ns));
self.scheduling_delays.truncate(SCHEDULING_DELAY_TOP);
}
}
fn top_scheduling_delays(&self) -> Vec<SchedulingDelay> {
let mut top = self.scheduling_delays.clone();
top.sort_unstable_by_key(|d| std::cmp::Reverse(d.delay_ns));
top.truncate(SCHEDULING_DELAY_TOP);
top
}
fn worker_activity(&self) -> Vec<WorkerStats> {
if self.by_worker.is_empty() {
return Vec::new();
}
let mut workers: Vec<WorkerStats> = self
.by_worker
.iter()
.map(|((host, wid), wa)| {
let busy_pct = if wa.observed_ns > 0 {
(wa.busy_ns as f64 / wa.observed_ns as f64) * 100.0
} else {
0.0
};
WorkerStats {
worker_id: *wid,
host: host.clone(),
total_polls: wa.total_polls,
busy_ns: wa.busy_ns,
span_ns: wa.observed_ns,
busy_pct,
notable_polls: wa.notable_polls,
worst_poll_ns: wa.worst_poll_ns,
worst_exemplar: wa.worst_exemplar.clone(),
}
})
.collect();
workers.sort_unstable_by(|a, b| {
b.busy_pct
.partial_cmp(&a.busy_pct)
.unwrap_or(std::cmp::Ordering::Equal)
});
workers
}
}
#[derive(Default)]
struct WorkerAccum {
total_polls: u64,
busy_ns: i64,
observed_ns: i64,
notable_polls: u64,
worst_poll_ns: i64,
worst_exemplar: Option<PollExemplar>,
}
struct LocAccum {
total: u64,
durations: Vec<(i64, u8)>, worst_by_class: [Option<PollExemplar>; 4],
}
const OFF_CPU_CONFIDENCE_NS: i64 = 10_000_000;
enum PollClass {
OnCpu,
OffCpu,
Mixed,
Unknown,
}
fn classify_poll(cpu_count: u32, sched_count: u32, duration_ns: i64) -> PollClass {
if cpu_count > 0 && sched_count > 0 {
return PollClass::Mixed;
}
if cpu_count > 0 {
return PollClass::OnCpu;
}
if duration_ns >= OFF_CPU_CONFIDENCE_NS {
PollClass::OffCpu
} else {
PollClass::Unknown
}
}
fn read_polls_part(
data: &[u8],
scope: &Scope,
source_key: &str,
acc: &mut TokioStatsAccum,
) -> Result<(), (StatusCode, String)> {
let reader = parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(
bytes::Bytes::from(data.to_vec()),
4096,
)
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
let mut seg_windows: HashMap<(String, u32), (i64, i64)> = HashMap::new();
for batch in reader {
let batch = batch.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
let duration_arr = batch
.column_by_name("duration_ns")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::Int64Array>());
let start_arr = batch
.column_by_name("start_ns")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::Int64Array>());
let end_arr = batch
.column_by_name("end_ns")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::Int64Array>());
let cpu_arr = batch
.column_by_name("cpu_sample_count")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt32Array>());
let sched_arr = batch
.column_by_name("sched_sample_count")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt32Array>());
let spawn_loc_arr = batch
.column_by_name("spawn_loc")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::StringArray>());
let host_arr = batch
.column_by_name("host")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::StringArray>());
let worker_arr = batch
.column_by_name("worker_id")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt32Array>());
let task_arr = batch
.column_by_name("task_id")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt64Array>());
let ready_at_arr = batch
.column_by_name("ready_at_ns")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::Int64Array>());
let scheduling_delay_arr = batch
.column_by_name("scheduling_delay_ns")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::Int64Array>());
let scheduling_kind_arr = batch
.column_by_name("scheduling_delay_kind")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt8Array>());
let waker_task_arr = batch
.column_by_name("waker_task_id")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::UInt64Array>());
let task_instrumented_arr = batch
.column_by_name("task_instrumented")
.and_then(|c| c.as_any().downcast_ref::<arrow::array::BooleanArray>());
let Some(duration_arr) = duration_arr else {
continue;
};
for i in 0..batch.num_rows() {
if let Some(sa) = start_arr
&& let Some(start) = scope.start_ns
&& sa.value(i) < start
{
continue;
}
if let Some(sa) = start_arr
&& let Some(end) = scope.end_ns
&& sa.value(i) >= end
{
continue;
}
let dur = duration_arr.value(i);
acc.total_polls += 1;
if let Some(workers) = worker_arr {
let host_val = host_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.unwrap_or("");
let key = (host_val.to_string(), workers.value(i));
let wa = acc.by_worker.entry(key.clone()).or_default();
wa.total_polls += 1;
wa.busy_ns += dur;
if let (Some(sa), Some(ea)) = (start_arr, end_arr) {
let (start, end) = (sa.value(i), ea.value(i));
let w = seg_windows.entry(key).or_insert((start, end));
w.0 = w.0.min(start);
w.1 = w.1.max(end);
}
}
let measured = scheduling_delay_arr
.zip(ready_at_arr)
.filter(|(delay, ready)| !delay.is_null(i) && !ready.is_null(i))
.map(|(delay, ready)| (delay.value(i), ready.value(i)));
if let Some((delay_ns, ready_at_ns)) = measured {
let cov = &mut acc.scheduling_delay_coverage;
cov.observed_polls += 1;
if delay_ns >= HIGH_SCHEDULING_DELAY_NS {
cov.over_1ms_polls += 1;
}
let kind = scheduling_kind_arr
.filter(|a| !a.is_null(i))
.and_then(|a| SchedulingDelayKind::from_u8(a.value(i)));
match kind {
Some(SchedulingDelayKind::Spawn) => cov.spawn_inferred_polls += 1,
Some(SchedulingDelayKind::Wake) => cov.wake_observed_polls += 1,
Some(SchedulingDelayKind::WakeDuringPoll) => cov.wake_during_poll_polls += 1,
None => {}
}
if let (Some(kind), Some(workers), Some(tasks)) = (kind, worker_arr, task_arr) {
let host = host_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.unwrap_or("");
let spawn_loc = spawn_loc_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.map(str::to_string);
let waker_task_id =
waker_task_arr.filter(|a| !a.is_null(i)).map(|a| a.value(i));
acc.push_scheduling_delay(SchedulingDelay {
delay_ns,
ready_at_ns,
poll_start_ns: start_arr.map_or(0, |a| a.value(i)),
poll_end_ns: end_arr.map_or(0, |a| a.value(i)),
worker_id: workers.value(i),
task_id: tasks.value(i),
spawn_loc,
kind,
waker_task_id,
host: host.to_string(),
source_key: source_key.to_string(),
});
}
} else {
let cov = &mut acc.scheduling_delay_coverage;
cov.unmeasured_polls += 1;
match task_instrumented_arr
.filter(|a| !a.is_null(i))
.map(|a| a.value(i))
{
Some(true) => cov.missing_readiness_unmeasured_polls += 1,
Some(false) => cov.uninstrumented_unmeasured_polls += 1,
None => cov.instrumentation_unknown_unmeasured_polls += 1,
}
}
if let Some(sa) = start_arr {
let ts = sa.value(i);
acc.min_ts = Some(acc.min_ts.map_or(ts, |m| m.min(ts)));
acc.max_ts = Some(acc.max_ts.map_or(ts, |m| m.max(ts)));
}
let loc = spawn_loc_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.unwrap_or("(unknown)");
let la = acc
.by_loc
.entry(loc.to_string())
.or_insert_with(|| LocAccum {
total: 0,
durations: Vec::new(),
worst_by_class: [None, None, None, None],
});
la.total += 1;
if dur < DURATION_FLOOR_NS {
continue;
}
let cpu_count = cpu_arr.map_or(0, |a| a.value(i));
let sched_count = sched_arr.map_or(0, |a| a.value(i));
let class = match classify_poll(cpu_count, sched_count, dur) {
PollClass::OffCpu => 0u8,
PollClass::OnCpu => 1,
PollClass::Mixed => 2,
PollClass::Unknown => 3,
};
la.durations.push((dur, class));
let host = host_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.unwrap_or("");
let slot = &mut la.worst_by_class[class as usize];
if slot.as_ref().is_none_or(|w| dur > w.duration_ns) {
*slot = Some(PollExemplar {
start_ns: start_arr.map_or(0, |a| a.value(i)),
end_ns: end_arr.map_or(0, |a| a.value(i)),
duration_ns: dur,
host: host.to_string(),
source_key: source_key.to_string(),
});
}
if let Some(workers) = worker_arr {
let wa = acc
.by_worker
.entry((host.to_string(), workers.value(i)))
.or_default();
wa.notable_polls += 1;
if dur > wa.worst_poll_ns {
wa.worst_poll_ns = dur;
wa.worst_exemplar = Some(PollExemplar {
start_ns: start_arr.map_or(0, |a| a.value(i)),
end_ns: end_arr.map_or(0, |a| a.value(i)),
duration_ns: dur,
host: host.to_string(),
source_key: source_key.to_string(),
});
}
}
if let (Some(workers), Some(tasks)) = (worker_arr, task_arr) {
let spawn_loc = spawn_loc_arr
.and_then(|a| if a.is_null(i) { None } else { Some(a.value(i)) })
.map(str::to_string);
acc.push_long_poll(LongPoll {
duration_ns: dur,
worker_id: workers.value(i),
task_id: tasks.value(i),
spawn_loc,
start_ns: start_arr.map_or(0, |a| a.value(i)),
end_ns: end_arr.map_or(0, |a| a.value(i)),
host: host.to_string(),
source_key: source_key.to_string(),
});
}
}
}
for (key, (start, end)) in seg_windows {
if let Some(wa) = acc.by_worker.get_mut(&key) {
wa.observed_ns += (end - start).max(0);
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ingest::decode::decode_samples;
use crate::ingest::parquet_writer;
#[test]
fn test_read_polls_from_demo_trace() {
let data = std::fs::read(concat!(
env!("CARGO_MANIFEST_DIR"),
"/ui/public/demo-trace.bin"
))
.unwrap();
let decompressed = {
use std::io::Read;
let mut dec = flate2::read::GzDecoder::new(data.as_slice());
let mut buf = Vec::new();
dec.read_to_end(&mut buf).unwrap();
buf
};
let (_, _, polls, _) = decode_samples(&decompressed, "demo-trace.bin").unwrap();
assert!(!polls.is_empty());
let mut buf = Vec::new();
parquet_writer::write_polls(&mut buf, &polls).unwrap();
let scope = Scope::default();
let mut acc = TokioStatsAccum::default();
read_polls_part(&buf, &scope, "test-key", &mut acc).unwrap();
assert_eq!(acc.total_polls, polls.len() as u64);
let notable: usize = acc.by_loc.values().map(|la| la.durations.len()).sum();
assert!(notable > 0, "expected polls above 100µs floor");
let with_exemplar = acc
.by_loc
.values()
.filter(|la| la.worst_by_class.iter().any(|e| e.is_some()))
.count();
assert!(with_exemplar > 0);
let top = acc.top_long_polls();
assert!(!top.is_empty(), "expected some notable long polls");
assert!(top.len() <= LONG_POLL_TOP);
assert!(
top.windows(2).all(|w| w[0].duration_ns >= w[1].duration_ns),
"top long polls must be ranked descending by duration"
);
assert!(
top.iter().all(|p| !p.source_key.is_empty()),
"each long poll must carry a source key for deep-linking"
);
assert!(
top[0].duration_ns >= DURATION_FLOOR_NS,
"long polls must clear the notable floor"
);
let workers = acc.worker_activity();
assert!(!workers.is_empty(), "expected at least one worker");
assert!(
workers.windows(2).all(|w| w[0].busy_pct >= w[1].busy_pct),
"worker activity must be ranked descending by busyness"
);
assert!(
workers.iter().all(|w| w.total_polls > 0),
"every worker must have at least one poll"
);
assert!(
workers.iter().all(|w| w.busy_ns > 0),
"every worker must have accumulated some busy time"
);
assert!(
workers.iter().all(|w| w.busy_pct > 0.0),
"every worker must have non-zero busyness"
);
assert!(
workers.iter().all(|w| w.busy_pct <= 100.0 + 1e-6),
"busyness must be bounded to 100% (got {:?})",
workers.iter().map(|w| w.busy_pct).fold(0.0_f64, f64::max)
);
assert!(
workers.iter().all(|w| w.span_ns >= w.busy_ns),
"observed span must be at least the busy time for every worker"
);
let hosts: HashSet<&str> = workers.iter().map(|w| w.host.as_str()).collect();
assert_eq!(
hosts.len(),
1,
"demo trace is single-host: all workers should share one host value"
);
assert!(
workers.iter().any(|w| w.worst_exemplar.is_some()),
"at least one worker should have a worst exemplar (from above-floor polls)"
);
assert!(
workers
.iter()
.filter_map(|w| w.worst_exemplar.as_ref())
.all(|ex| !ex.source_key.is_empty()),
"worker worst exemplars must carry a source key for deep-linking"
);
let total_worker_polls: u64 = workers.iter().map(|w| w.total_polls).sum();
assert_eq!(
total_worker_polls, acc.total_polls,
"sum of per-worker polls must equal total_polls"
);
let cov = &acc.scheduling_delay_coverage;
assert_eq!(
cov.observed_polls + cov.unmeasured_polls,
acc.total_polls,
"every poll must be counted as measured or unmeasured"
);
assert_eq!(
cov.spawn_inferred_polls + cov.wake_observed_polls + cov.wake_during_poll_polls,
cov.observed_polls,
"per-kind tallies must partition the observed polls"
);
assert_eq!(
cov.uninstrumented_unmeasured_polls
+ cov.instrumentation_unknown_unmeasured_polls
+ cov.missing_readiness_unmeasured_polls,
cov.unmeasured_polls,
"per-reason tallies must partition the unmeasured polls"
);
assert!(
cov.over_1ms_polls <= cov.observed_polls,
"over-1ms polls are a subset of the measured polls"
);
let delays = acc.top_scheduling_delays();
assert!(delays.len() <= SCHEDULING_DELAY_TOP);
assert!(delays.len() as u64 <= cov.observed_polls);
assert!(
delays.windows(2).all(|w| w[0].delay_ns >= w[1].delay_ns),
"top scheduling delays must be ranked descending by delay"
);
assert!(
delays.iter().all(|d| !d.source_key.is_empty()),
"each scheduling delay must carry a source key for deep-linking"
);
assert!(
delays.iter().all(|d| d.delay_ns >= 0),
"scheduling delay cannot be negative"
);
eprintln!(
"tokio-stats: {} total, {} above floor, {} locs with exemplars, {} long polls (worst {}ns), {} workers (busiest {:.1}%), {} measured delays / {} unmeasured",
acc.total_polls,
notable,
with_exemplar,
top.len(),
top[0].duration_ns,
workers.len(),
workers[0].busy_pct,
cov.observed_polls,
cov.unmeasured_polls,
);
}
#[test]
fn scheduling_delay_top_n_is_bounded_and_ranked() {
let mut acc = TokioStatsAccum::default();
for n in 0..(SCHEDULING_DELAY_SOFT_CAP + 50) {
acc.push_scheduling_delay(SchedulingDelay {
delay_ns: n as i64,
ready_at_ns: 0,
poll_start_ns: n as i64,
poll_end_ns: n as i64 + 1,
worker_id: 0,
task_id: n as u64,
spawn_loc: None,
kind: SchedulingDelayKind::Wake,
waker_task_id: Some(1),
host: "h".to_string(),
source_key: "k".to_string(),
});
}
let top = acc.top_scheduling_delays();
assert_eq!(top.len(), SCHEDULING_DELAY_TOP);
assert!(
top.windows(2).all(|w| w[0].delay_ns >= w[1].delay_ns),
"must be ranked descending by delay"
);
assert_eq!(
top[0].delay_ns,
(SCHEDULING_DELAY_SOFT_CAP + 50 - 1) as i64,
"the largest delay must survive compaction"
);
}
}