use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::time::Duration;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Cost {
pub took: Duration,
pub files: Option<usize>,
pub over_the_wire: Option<OverTheWire>,
}
fn combine_wire(a: OverTheWire, b: OverTheWire) -> OverTheWire {
OverTheWire {
requests: a.requests + b.requests,
bytes: a.bytes.zip(b.bytes).map(|(x, y)| x + y),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Total {
pub took: Duration,
pub over_the_wire: Option<OverTheWire>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct OverTheWire {
pub requests: usize,
pub bytes: Option<u64>,
}
#[derive(Debug, Default)]
struct Tally {
ran: AtomicBool,
nanos: AtomicU64,
files: AtomicUsize,
live_requests: AtomicUsize,
live_bytes: AtomicU64,
counted_bytes: AtomicBool,
counted: AtomicBool,
requests: AtomicUsize,
bytes: AtomicU64,
counted_files: AtomicBool,
}
impl Tally {
fn record(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
self.nanos.fetch_add(
took.as_nanos().min(u64::MAX as u128) as u64,
Ordering::Relaxed,
);
if let Some(files) = files {
self.files.fetch_add(files, Ordering::Relaxed);
self.counted_files.store(true, Ordering::Relaxed);
}
if over_the_wire {
self.requests.fetch_max(
self.live_requests.load(Ordering::Relaxed),
Ordering::Relaxed,
);
self.bytes
.fetch_max(self.live_bytes.load(Ordering::Relaxed), Ordering::Relaxed);
self.counted.store(true, Ordering::Relaxed);
}
self.ran.store(true, Ordering::Release);
}
fn replace(&self, took: Duration, files: Option<usize>) {
self.nanos.store(
took.as_nanos().min(u64::MAX as u128) as u64,
Ordering::Relaxed,
);
match files {
Some(files) => {
self.files.store(files, Ordering::Relaxed);
self.counted_files.store(true, Ordering::Relaxed);
}
None => self.counted_files.store(false, Ordering::Relaxed),
}
self.ran.store(true, Ordering::Release);
}
fn request(&self, bytes: u64) {
self.live_requests.fetch_add(1, Ordering::Relaxed);
self.live_bytes.fetch_add(bytes, Ordering::Relaxed);
self.counted_bytes.store(true, Ordering::Relaxed);
}
fn add_requests(&self, wire: OverTheWire) {
self.live_requests
.fetch_add(wire.requests, Ordering::Relaxed);
if let Some(bytes) = wire.bytes {
self.live_bytes.fetch_add(bytes, Ordering::Relaxed);
self.counted_bytes.store(true, Ordering::Relaxed);
}
}
fn cost(&self) -> Option<Cost> {
if !self.ran.load(Ordering::Acquire) {
return None;
}
Some(Cost {
took: Duration::from_nanos(self.nanos.load(Ordering::Relaxed)),
files: self
.counted_files
.load(Ordering::Relaxed)
.then(|| self.files.load(Ordering::Relaxed)),
over_the_wire: self.counted.load(Ordering::Relaxed).then(|| OverTheWire {
requests: self.requests.load(Ordering::Relaxed),
bytes: self
.counted_bytes
.load(Ordering::Relaxed)
.then(|| self.bytes.load(Ordering::Relaxed)),
}),
})
}
}
#[derive(Debug, Clone, Default)]
pub struct OpenReport {
pub progress: std::sync::Arc<crate::formats::schema_union::FooterProgress>,
pub meter: std::sync::Arc<Meter>,
pub remembered: Option<crate::cache::CacheManager>,
pub(crate) writes: crate::app::background::CacheWrites,
}
#[derive(Debug, Default)]
pub struct Meter {
listing: Tally,
footers: Tally,
last_page: Tally,
counted_rows: AtomicBool,
}
impl Meter {
pub fn listed(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
self.listing.record(took, files, over_the_wire);
}
pub fn read_footers(&self, took: Duration, files: Option<usize>, over_the_wire: bool) {
self.footers.record(took, files, over_the_wire);
}
pub fn counted_rows(
&self,
took: Duration,
files: Option<usize>,
wire: Option<OverTheWire>,
) -> bool {
if self.listing.cost().is_none() {
return false;
}
if self
.counted_rows
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return false;
}
if let Some(w) = wire {
self.footers.add_requests(w);
}
self.footers.record(took, files, wire.is_some());
true
}
pub fn read_page(&self, took: Duration, files: Option<usize>) {
self.last_page.replace(took, files);
}
pub fn last_page(&self) -> Option<Cost> {
self.last_page.cost()
}
pub fn footer_request(&self, bytes: u64) {
self.footers.request(bytes);
}
pub fn listing(&self) -> Option<Cost> {
self.listing.cost()
}
pub fn footers(&self) -> Option<Cost> {
self.footers.cost()
}
pub fn total(&self) -> Option<Total> {
let parts: Vec<Cost> = [self.listing(), self.footers()]
.into_iter()
.flatten()
.collect();
if parts.len() < 2 {
return None;
}
Some(Total {
took: parts.iter().map(|p| p.took).sum(),
over_the_wire: parts
.iter()
.filter_map(|p| p.over_the_wire)
.reduce(combine_wire),
})
}
}
#[derive(Debug, Default)]
pub struct LoopTimes {
pub frames: Durations,
pub handlers: Durations,
}
impl LoopTimes {
pub fn frame(&mut self, took: Duration) {
self.frames.record(took);
if self.frames.count.is_multiple_of(Durations::WINDOW as u64) {
log::debug!(
target: "datui",
"frames {}; handlers {}",
self.frames.summary(),
self.handlers.summary()
);
}
}
pub fn handler(&mut self, took: Duration) {
self.handlers.record(took);
}
}
#[derive(Debug, Default)]
pub struct Durations {
recent: std::collections::VecDeque<Duration>,
count: u64,
}
impl Durations {
pub const WINDOW: usize = 240;
pub fn record(&mut self, took: Duration) {
if self.recent.len() == Self::WINDOW {
self.recent.pop_front();
}
self.recent.push_back(took);
self.count += 1;
}
#[cfg(test)]
pub fn count(&self) -> u64 {
self.count
}
pub fn quantile(&self, q: f64) -> Option<Duration> {
let mut sorted: Vec<Duration> = self.recent.iter().copied().collect();
sorted.sort_unstable();
let last = sorted.len().checked_sub(1)?;
sorted.get(((last as f64) * q).round() as usize).copied()
}
pub fn summary(&self) -> String {
let ms = |d: Duration| format!("{:.1}ms", d.as_secs_f64() * 1000.0);
match (self.quantile(0.5), self.quantile(0.99), self.quantile(1.0)) {
(Some(p50), Some(p99), Some(max)) => {
format!("p50 {} p99 {} max {}", ms(p50), ms(p99), ms(max))
}
_ => "-".to_string(),
}
}
}
#[cfg(test)]
mod tests;