stats.rsannotatedstats.rssource135 lines · 4.3 KB · raw
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}