Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/src/srv/admin/traffic.rs

14.7 KiB, 140 runs

created by r1870400018:10334, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! In-memory traffic recorder.
2//!
3//! Holds a bounded ring buffer of recent HTTP requests and a small
4//! set of per-vhost / per-status counters. Populated from the request
5//! pipeline in `srv/https.rs`; read by [`handler`](super::handler)
6//! when the operator opens the traffic view.
7//!
8//! The buffer shape is chosen with a second consumer in mind: once
9//! `fe2o3_net::guard::AddressGuard` lands (after extraction from
10//! `fe2o3_shield`), it will feed from the same counters to drive
11//! rate-limiting and blacklist transitions. Counters are therefore
12//! updated on the hot path under a short write lock; snapshots for
13//! the dashboard copy out once under a read lock.
14//!
15//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
16//! Anthropic Claude
17
18use oxedyne_fe2o3_core::prelude::*;
19
20use std::{
21 collections::{
22 HashMap,
23 VecDeque,
24 },
25 sync::{
26 Arc,
27 RwLock,
28 atomic::{
29 AtomicU64,
30 Ordering,
31 },
32 },
33 time::{
34 SystemTime,
35 UNIX_EPOCH,
36 },
37};
38
39pub const DEFAULT_RING_CAPACITY: usize = 10_000; // roughly 1-2 MiB live
40// Past this many distinct paths, a vhost's further paths fold into the
41// `_other` bucket, bounding worst-case memory when a caller probes unique URLs.
42pub const MAX_PATHS_PER_VHOST: usize = 256;
43pub const OTHER_PATH_BUCKET: &str = "_other";
44// At the default sample interval the history spans one hour. Five seconds
45// trades chart smoothness against CPU cost on a quiet host.
46pub const DEFAULT_HISTORY_CAPACITY: usize = 720;
47pub const DEFAULT_SAMPLE_INTERVAL_SECS: u64 = 5;
48
49// ┌───────────────────────────────────────────────────────────────────────────┐
50// │ REQUEST RECORD │
51// └───────────────────────────────────────────────────────────────────────────┘
52
53/// Snapshot of a single request as it leaves the handler pipeline. All fields
54/// are owned, so the record can outlive the request without keeping borrows
55/// alive.
56#[derive(Clone, Debug)]
57pub struct RequestRecord {
58 pub when_ns: u64, // unix nanoseconds at completion
59 pub vhost: String, // lowercased hostname, as keyed in vhost_dbs
60 pub method: String,
61 pub path: String, // includes the query string
62 pub status: u16,
63 pub peer: String, // IP and port
64 pub bytes: Option<u64>, // response body length, when known
65 pub duration_us: u64, // accept to final write, microseconds
66}
67
68// ┌───────────────────────────────────────────────────────────────────────────┐
69// │ COUNTERS │
70// └───────────────────────────────────────────────────────────────────────────┘
71
72/// Per-vhost counters, aggregated since the recorder was created, i.e. since
73/// Steel started. Counts never decrement, so the dashboard can compute a rate
74/// of change by sampling and subtracting across two fetches.
75#[derive(Clone, Debug, Default)]
76pub struct VhostCounters {
77 pub total: u64,
78 pub by_status: HashMap<u16, u64>, // `{200 => 1234, 404 => 12, ...}`
79 pub by_path: HashMap<String, u64>, // capped, overflow into `_other`
80}
81
82#[derive(Clone, Debug, Default)]
83pub struct CountersSnapshot {
84 pub total: u64,
85 pub by_status: HashMap<u16, u64>,
86 pub by_vhost: HashMap<String, VhostCounters>,
87}
88
89/// One point in the bounded traffic history: the monotonic totals at a
90/// particular unix second, so the dashboard can difference adjacent samples and
91/// draw a requests-per-interval chart.
92#[derive(Clone, Debug)]
93pub struct TrafficSample {
94 pub when_secs: u64, // unix seconds
95 pub total: u64, // cumulative across every vhost
96 pub by_status: HashMap<u16, u64>, // cumulative at this instant
97}
98
99// ┌───────────────────────────────────────────────────────────────────────────┐
100// │ TRAFFIC RECORDER │
101// └───────────────────────────────────────────────────────────────────────────┘
102
103/// Thread-safe ring buffer plus counters, cheaply cloneable via `Arc`. A typical
104/// deployment constructs one and stores it both in the request pipeline (via
105/// `ServerContext`) and in the admin state, for the dashboard.
106#[derive(Debug)]
107pub struct TrafficRecorder {
108 capacity: usize, // older entries drop once at capacity
109 ring: RwLock<VecDeque<RequestRecord>>, // newest last
110 // A distinct lock from the ring, so dashboard reads of counters and of
111 // records do not contend on the same lock.
112 counters: RwLock<CountersSnapshot>,
113 total: AtomicU64, // also in `counters`; atomic for cheap reads
114 // Periodic counter samples, newest last, populated by a background sampling
115 // task and read by the dashboard when drawing time-series charts.
116 history: RwLock<VecDeque<TrafficSample>>,
117 history_capacity: usize,
118}
119
120impl TrafficRecorder {
121 /// A zero capacity is treated as [`DEFAULT_RING_CAPACITY`].
122 pub fn new(capacity: usize) -> Self {
123 let cap = if capacity == 0 {
124 DEFAULT_RING_CAPACITY
125 } else {
126 capacity
127 };
128 Self {
129 capacity: cap,
130 ring: RwLock::new(VecDeque::with_capacity(cap)),
131 counters: RwLock::new(CountersSnapshot::default()),
132 total: AtomicU64::new(0),
133 history: RwLock::new(
134 VecDeque::with_capacity(DEFAULT_HISTORY_CAPACITY),
135 ),
136 history_capacity: DEFAULT_HISTORY_CAPACITY,
137 }
138 }
139
140 pub fn new_shared(capacity: usize) -> Arc<Self> {
141 Arc::new(Self::new(capacity))
142 }
143
144 pub fn capacity(&self) -> usize {
145 self.capacity
146 }
147
148 pub fn history_capacity(&self) -> usize {
149 self.history_capacity
150 }
151
152 pub fn total(&self) -> u64 {
153 self.total.load(Ordering::Relaxed)
154 }
155
156 /// Takes two short write locks, one on the ring and one on the counters. Any
157 /// lock poisoning surfaces as an error; the hot-path call site logs it and
158 /// continues rather than aborting the request.
159 pub fn record(&self, rec: RequestRecord) -> Outcome<()> {
160 // Counters first so a poisoned ring does not leave us
161 // with a stale count.
162 {
163 let mut ctr = lock_write!(self.counters);
164 ctr.total = ctr.total.saturating_add(1);
165 *ctr.by_status.entry(rec.status).or_insert(0) += 1;
166 let vh = ctr.by_vhost
167 .entry(rec.vhost.clone())
168 .or_insert_with(VhostCounters::default);
169 vh.total = vh.total.saturating_add(1);
170 *vh.by_status.entry(rec.status).or_insert(0) += 1;
171 if vh.by_path.contains_key(&rec.path)
172 || vh.by_path.len() < MAX_PATHS_PER_VHOST
173 {
174 *vh.by_path.entry(rec.path.clone()).or_insert(0) += 1;
175 } else {
176 *vh.by_path.entry(OTHER_PATH_BUCKET.to_string())
177 .or_insert(0) += 1;
178 }
179 }
180 self.total.fetch_add(1, Ordering::Relaxed);
181 {
182 let mut ring = lock_write!(self.ring);
183 if ring.len() == self.capacity {
184 ring.pop_front();
185 }
186 ring.push_back(rec);
187 }
188 Ok(())
189 }
190
191 /// Up to `limit` most recent records, newest first. A `limit` of zero returns
192 /// everything currently in the ring.
193 pub fn recent(&self, limit: usize) -> Outcome<Vec<RequestRecord>> {
194 let ring = lock_read!(self.ring);
195 let take = if limit == 0 { ring.len() } else { limit.min(ring.len()) };
196 let mut out = Vec::with_capacity(take);
197 // Iterate newest-first by walking back from the end.
198 for rec in ring.iter().rev().take(take) {
199 out.push(rec.clone());
200 }
201 Ok(out)
202 }
203
204 pub fn counters_snapshot(&self) -> Outcome<CountersSnapshot> {
205 let ctr = lock_read!(self.counters);
206 Ok(ctr.clone())
207 }
208
209 /// Meant for a background task on a fixed interval. Trims the oldest entry
210 /// when the ring reaches `history_capacity`.
211 pub fn sample_now(&self) -> Outcome<()> {
212 let ctr = lock_read!(self.counters);
213 let sample = TrafficSample {
214 when_secs: SystemTime::now()
215 .duration_since(UNIX_EPOCH)
216 .map(|d| d.as_secs())
217 .unwrap_or(0),
218 total: ctr.total,
219 by_status: ctr.by_status.clone(),
220 };
221 drop(ctr);
222 let mut hist = lock_write!(self.history);
223 if hist.len() == self.history_capacity {
224 hist.pop_front();
225 }
226 hist.push_back(sample);
227 Ok(())
228 }
229
230 /// Chronological order, oldest first.
231 pub fn history_snapshot(&self) -> Outcome<Vec<TrafficSample>> {
232 let hist = lock_read!(self.history);
233 let mut out = Vec::with_capacity(hist.len());
234 for s in hist.iter() {
235 out.push(s.clone());
236 }
237 Ok(out)
238 }
239}
240
241impl Default for TrafficRecorder {
242 fn default() -> Self {
243 Self::new(DEFAULT_RING_CAPACITY)
244 }
245}
246
247// ┌───────────────────────────────────────────────────────────────────────────┐
248// │ HELPERS │
249// └───────────────────────────────────────────────────────────────────────────┘
250
251/// Current unix time in nanoseconds, clamped to zero on clock error.
252pub fn now_ns() -> u64 {
253 SystemTime::now()
254 .duration_since(UNIX_EPOCH)
255 .map(|d| d.as_nanos() as u64)
256 .unwrap_or(0)
257}
258
259// ┌───────────────────────────────────────────────────────────────────────────┐
260// │ TESTS │
261// └───────────────────────────────────────────────────────────────────────────┘
262
263#[cfg(test)]
264mod tests {
265 use super::*;
266
267 fn mkrec(vhost: &str, path: &str, status: u16) -> RequestRecord {
268 RequestRecord {
269 when_ns: now_ns(),
270 vhost: vhost.to_string(),
271 method: "GET".to_string(),
272 path: path.to_string(),
273 status,
274 peer: "127.0.0.1:0".to_string(),
275 bytes: Some(42),
276 duration_us: 123,
277 }
278 }
279
280 #[test]
281 fn record_and_recent() {
282 let r = TrafficRecorder::new(4);
283 r.record(mkrec("a", "/", 200)).expect("rec 1");
284 r.record(mkrec("a", "/x", 200)).expect("rec 2");
285 r.record(mkrec("a", "/y", 404)).expect("rec 3");
286 let recent = r.recent(0).expect("recent");
287 assert_eq!(recent.len(), 3);
288 // Newest first.
289 assert_eq!(recent[0].path, "/y");
290 assert_eq!(recent[2].path, "/");
291 }
292
293 #[test]
294 fn ring_evicts_oldest() {
295 let r = TrafficRecorder::new(2);
296 r.record(mkrec("a", "/1", 200)).expect("rec 1");
297 r.record(mkrec("a", "/2", 200)).expect("rec 2");
298 r.record(mkrec("a", "/3", 200)).expect("rec 3");
299 let recent = r.recent(0).expect("recent");
300 assert_eq!(recent.len(), 2);
301 // Oldest ("/1") must have been dropped.
302 let paths: Vec<&str> = recent.iter().map(|r| r.path.as_str()).collect();
303 assert!(!paths.contains(&"/1"));
304 assert!(paths.contains(&"/3"));
305 }
306
307 #[test]
308 fn counters_track_status_and_vhost() {
309 let r = TrafficRecorder::new(100);
310 r.record(mkrec("a", "/", 200)).expect("rec");
311 r.record(mkrec("a", "/", 200)).expect("rec");
312 r.record(mkrec("a", "/", 404)).expect("rec");
313 r.record(mkrec("b", "/", 500)).expect("rec");
314 let snap = r.counters_snapshot().expect("snap");
315 assert_eq!(snap.total, 4);
316 assert_eq!(snap.by_status.get(&200).copied(), Some(2));
317 assert_eq!(snap.by_status.get(&404).copied(), Some(1));
318 assert_eq!(snap.by_status.get(&500).copied(), Some(1));
319 let vh_a = snap.by_vhost.get("a").expect("vhost a");
320 assert_eq!(vh_a.total, 3);
321 let vh_b = snap.by_vhost.get("b").expect("vhost b");
322 assert_eq!(vh_b.total, 1);
323 assert_eq!(r.total(), 4);
324 }
325
326 #[test]
327 fn per_vhost_path_bucket_saturates() {
328 let r = TrafficRecorder::new(10_000);
329 // Fill the per-vhost path map to its cap.
330 for i in 0..(MAX_PATHS_PER_VHOST + 5) {
331 let p = fmt!("/p{}", i);
332 r.record(mkrec("a", &p, 200)).expect("rec");
333 }
334 let snap = r.counters_snapshot().expect("snap");
335 let vh_a = snap.by_vhost.get("a").expect("vhost a");
336 // Cap observed, plus the _other overflow bucket.
337 assert!(vh_a.by_path.len() <= MAX_PATHS_PER_VHOST + 1);
338 assert_eq!(
339 vh_a.by_path.get(OTHER_PATH_BUCKET).copied(),
340 Some(5),
341 );
342 }
343
344 #[test]
345 fn recent_honours_limit() {
346 let r = TrafficRecorder::new(100);
347 for i in 0..10 {
348 r.record(mkrec("a", &fmt!("/p{}", i), 200)).expect("rec");
349 }
350 let three = r.recent(3).expect("recent");
351 assert_eq!(three.len(), 3);
352 assert_eq!(three[0].path, "/p9");
353 assert_eq!(three[2].path, "/p7");
354 }
355}