1//! This backend's counters behind `jev_stats()`, and the DEBUG1 lines for 2//! each retry, redial, connection and answer (review D4). 3//! 4//! Per backend, not cluster-wide: the counters are thread-locals on the 5//! backend thread, so a session sees its own work and nobody else's. 6 7use std::cell::Cell; 8use std::fmt; 9 10use jev_client::{Event, Observer}; 11 12/// A snapshot of this backend's counters since it started. 13#[derive(Clone, Copy, Debug, Default, PartialEq)] 14pub struct Stats { 15 /// HTTP attempts that reached the server, retries included. 16 pub requests: u64, 17 /// Attempts that failed and were sent again under the retry policy. 18 pub retries: u64, 19 /// Requests that never left and were resent on a fresh connection. 20 pub redials: u64, 21 /// Connections opened to the endpoint. 22 pub connections: u64, 23 /// Judgments served from the durable cache, and those it lacked. 24 pub cache_hits: u64, 25 pub cache_misses: u64, 26 /// Judgments that waited on an equal one already asked in the same 27 /// scan, instead of sending their own. 28 pub dedupe_hits: u64, 29 pub input_tokens: u64, 30 pub output_tokens: u64, 31 /// Dollars, from input tokens at the price in force when answered. 32 pub cost: f64, 33 /// HTTP/2 streams open right now: sent and not yet answered, failed 34 /// or abandoned. A gauge, not a counter. 35 pub in_flight: u64, 36} 37 38thread_local! { 39 static STATS: Cell<Stats> = Cell::new(Stats::default()); 40} 41 42pub fn snapshot() -> Stats { 43 STATS.with(Cell::get) 44} 45 46fn update(f: impl FnOnce(&mut Stats)) { 47 STATS.with(|s| { 48 let mut v = s.get(); 49 f(&mut v); 50 s.set(v); 51 }); 52} 53 54/// Counts one judgment the cache answered, or did not. 55pub fn cached(hit: bool) { 56 update(|s| if hit { s.cache_hits += 1 } else { s.cache_misses += 1 }); 57} 58 59/// Counts one judgment answered by an equal one's request. 60pub fn deduped() { 61 update(|s| s.dedupe_hits += 1); 62} 63 64/// One open HTTP/2 stream, counted in `in_flight` for as long as it 65/// lives. The count leaves on `Drop`, so a stream abandoned by a timeout, 66/// a cancel or an error leaves it as surely as an answered one. 67pub(crate) struct InFlight(()); 68 69impl InFlight { 70 pub(crate) fn open() -> Self { 71 update(|s| s.in_flight += 1); 72 InFlight(()) 73 } 74} 75 76impl Drop for InFlight { 77 fn drop(&mut self) { 78 // Never panics (cancellation unwinds through here): `try_with` 79 // tolerates a thread-local already torn down at backend exit. 80 let _ = STATS.try_with(|s| { 81 let mut v = s.get(); 82 v.in_flight = v.in_flight.saturating_sub(1); 83 s.set(v); 84 }); 85 } 86} 87 88pub(crate) fn connected(uri: &dyn fmt::Display) { 89 update(|s| s.connections += 1); 90 pgrx::debug1!("jev: connected to {uri}"); 91} 92 93/// Counts and logs one client's events. `price_per_mtok` is the input 94/// price read when the scan began. 95pub struct Telemetry { 96 pub price_per_mtok: f64, 97} 98 99impl Observer for Telemetry { 100 fn event(&self, event: Event<'_>) { 101 match event { 102 Event::Sending { .. } => update(|s| s.requests += 1), 103 Event::Retrying { attempt, delay, cause, request_id } => { 104 update(|s| s.retries += 1); 105 pgrx::debug1!("jev: retrying after attempt {attempt} in {delay:?}: {cause} (request id {})", id(request_id)); 106 } 107 Event::Redialing { cause } => { 108 // The attempt counted in `Sending` never left. 109 update(|s| { 110 s.redials += 1; 111 s.requests -= 1; 112 }); 113 pgrx::debug1!("jev: resending on a new connection, never sent: {cause}"); 114 } 115 Event::Answered { request_id, attempts, usage } => { 116 let price = self.price_per_mtok; 117 update(|s| { 118 s.input_tokens += usage.input_tokens; 119 s.output_tokens += usage.output_tokens; 120 s.cost += usage.input_tokens as f64 * price / 1e6; 121 }); 122 pgrx::debug1!( 123 "jev: answered, request id {}, {attempts} attempt(s), {} input tokens", 124 id(request_id), 125 usage.input_tokens 126 ); 127 } 128 _ => {} 129 } 130 } 131} 132 133fn id(request_id: Option<&str>) -> &str { 134 request_id.unwrap_or("none") 135}