use std::cell::{Cell, RefCell};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Instant;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamSubPhase {
ColdFault,
Decompress,
Merge,
Encode,
EncodeFraming,
GrpcWrite,
}
#[derive(Debug, Default)]
pub struct StreamSubPhaseTimings {
cold_fault_nanos: AtomicU64,
decompress_nanos: AtomicU64,
merge_nanos: AtomicU64,
encode_nanos: AtomicU64,
encode_framing_nanos: AtomicU64,
grpc_write_nanos: AtomicU64,
}
impl StreamSubPhaseTimings {
fn counter(&self, phase: StreamSubPhase) -> &AtomicU64 {
match phase {
StreamSubPhase::ColdFault => &self.cold_fault_nanos,
StreamSubPhase::Decompress => &self.decompress_nanos,
StreamSubPhase::Merge => &self.merge_nanos,
StreamSubPhase::Encode => &self.encode_nanos,
StreamSubPhase::EncodeFraming => &self.encode_framing_nanos,
StreamSubPhase::GrpcWrite => &self.grpc_write_nanos,
}
}
pub fn add_nanos(&self, phase: StreamSubPhase, nanos: u64) {
let _ = self
.counter(phase)
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |v| {
Some(v.saturating_add(nanos))
});
}
pub fn nanos(&self, phase: StreamSubPhase) -> u64 {
self.counter(phase).load(Ordering::Relaxed)
}
}
thread_local! {
static CURRENT_SINK: RefCell<Option<Arc<StreamSubPhaseTimings>>> =
const { RefCell::new(None) };
static SINK_ACTIVE: Cell<bool> = const { Cell::new(false) };
static PULL_WAIT_NANOS: Cell<u64> = const { Cell::new(0) };
}
#[must_use = "the sink is uninstalled when the guard is dropped"]
pub struct StreamSubPhaseGuard {
prev: Option<Arc<StreamSubPhaseTimings>>,
prev_active: bool,
}
impl Drop for StreamSubPhaseGuard {
fn drop(&mut self) {
let prev = self.prev.take();
CURRENT_SINK.with(|c| *c.borrow_mut() = prev);
SINK_ACTIVE.with(|c| c.set(self.prev_active));
}
}
pub fn install(sink: Option<Arc<StreamSubPhaseTimings>>) -> StreamSubPhaseGuard {
let active = sink.is_some();
let prev = CURRENT_SINK.with(|c| std::mem::replace(&mut *c.borrow_mut(), sink));
let prev_active = SINK_ACTIVE.with(|c| c.replace(active));
StreamSubPhaseGuard { prev, prev_active }
}
pub fn sink_active() -> bool {
SINK_ACTIVE.with(|c| c.get())
}
pub fn add_pull_wait_nanos(nanos: u64) {
PULL_WAIT_NANOS.with(|c| c.set(c.get().saturating_add(nanos)));
}
pub fn pull_wait_nanos() -> u64 {
PULL_WAIT_NANOS.with(|c| c.get())
}
pub fn current() -> Option<Arc<StreamSubPhaseTimings>> {
CURRENT_SINK.with(|c| c.borrow().clone())
}
pub fn elapsed_nanos(start: Instant) -> u64 {
start.elapsed().as_nanos().min(u64::MAX as u128) as u64
}
pub fn timed<T>(phase: StreamSubPhase, f: impl FnOnce() -> T) -> T {
let sink = current();
match sink {
None => f(),
Some(sink) => {
let start = Instant::now();
let out = f();
sink.add_nanos(phase, elapsed_nanos(start));
out
}
}
}
pub fn record_nanos(phase: StreamSubPhase, nanos: u64) {
CURRENT_SINK.with(|c| {
if let Some(sink) = c.borrow().as_ref() {
sink.add_nanos(phase, nanos);
}
});
}
pub struct SubPhaseTimer {
phase: StreamSubPhase,
start: Instant,
sink: Arc<StreamSubPhaseTimings>,
}
impl Drop for SubPhaseTimer {
fn drop(&mut self) {
self.sink.add_nanos(self.phase, elapsed_nanos(self.start));
}
}
pub fn scoped(phase: StreamSubPhase) -> Option<SubPhaseTimer> {
scoped_captured(¤t(), phase)
}
pub fn scoped_captured(
sink: &Option<Arc<StreamSubPhaseTimings>>,
phase: StreamSubPhase,
) -> Option<SubPhaseTimer> {
sink.clone().map(|sink| SubPhaseTimer {
phase,
start: Instant::now(),
sink,
})
}
pub fn time_recv<T>(f: impl FnOnce() -> T) -> T {
if sink_active() || super::read_phase::sink_active() {
let start = Instant::now();
let out = f();
add_pull_wait_nanos(elapsed_nanos(start));
out
} else {
f()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicBool;
#[test]
fn timed_is_noop_when_no_sink_installed() {
assert!(current().is_none());
let ran = AtomicBool::new(false);
let out = timed(StreamSubPhase::ColdFault, || {
ran.store(true, Ordering::Relaxed);
42
});
assert_eq!(out, 42);
assert!(ran.load(Ordering::Relaxed));
assert!(current().is_none(), "no sink leaked onto the thread");
}
#[test]
fn timed_accumulates_into_the_installed_sink() {
let sink = Arc::new(StreamSubPhaseTimings::default());
let _g = install(Some(sink.clone()));
timed(StreamSubPhase::Decompress, || {
std::thread::sleep(std::time::Duration::from_millis(2));
});
assert!(
sink.nanos(StreamSubPhase::Decompress) > 0,
"a wrapped op accumulates into its bucket"
);
assert_eq!(
sink.nanos(StreamSubPhase::ColdFault),
0,
"an unentered bucket stays zero"
);
}
#[test]
fn guard_restores_previous_sink_on_drop() {
assert!(current().is_none());
{
let sink = Arc::new(StreamSubPhaseTimings::default());
let _g = install(Some(sink));
assert!(current().is_some());
}
assert!(
current().is_none(),
"dropping the guard uninstalls the sink (no leak across RPCs)"
);
}
#[test]
fn sink_propagates_across_a_spawned_thread() {
let sink = Arc::new(StreamSubPhaseTimings::default());
let _g = install(Some(sink.clone()));
let captured = current();
std::thread::spawn(move || {
let _child = install(captured);
record_nanos(StreamSubPhase::ColdFault, 1234);
})
.join()
.expect("child thread joins");
assert_eq!(
sink.nanos(StreamSubPhase::ColdFault),
1234,
"the child thread accumulated into the propagated per-request sink"
);
}
#[test]
fn record_nanos_is_noop_without_a_sink() {
assert!(current().is_none());
record_nanos(StreamSubPhase::Merge, 999); assert!(current().is_none());
}
#[test]
fn sink_active_mirrors_install_and_restores_on_drop() {
assert!(!sink_active(), "no sink installed at start");
{
let sink = Arc::new(StreamSubPhaseTimings::default());
let _g = install(Some(sink));
assert!(sink_active(), "sink_active tracks an installed sink");
{
let _none = install(None);
assert!(!sink_active(), "nested None install deactivates");
}
assert!(sink_active(), "outer sink restored after nested drop");
}
assert!(!sink_active(), "sink_active cleared once the guard drops");
}
#[test]
fn pull_wait_accumulates_on_this_thread() {
std::thread::spawn(|| {
assert_eq!(pull_wait_nanos(), 0);
add_pull_wait_nanos(100);
add_pull_wait_nanos(50);
assert_eq!(
pull_wait_nanos(),
150,
"recv-wait accumulates monotonically"
);
})
.join()
.expect("thread joins");
}
#[test]
fn add_nanos_saturates_rather_than_wrapping() {
let sink = StreamSubPhaseTimings::default();
sink.add_nanos(StreamSubPhase::Merge, u64::MAX);
sink.add_nanos(StreamSubPhase::Merge, 10);
assert_eq!(
sink.nanos(StreamSubPhase::Merge),
u64::MAX,
"a second add saturates at u64::MAX instead of wrapping to a tiny value"
);
}
}