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 | |
| 18 | use oxedyne_fe2o3_core::prelude::*; |
| 19 | |
| 20 | use 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 | |
| 39 | pub 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. |
| 42 | pub const MAX_PATHS_PER_VHOST: usize = 256; |
| 43 | pub 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. |
| 46 | pub const DEFAULT_HISTORY_CAPACITY: usize = 720; |
| 47 | pub 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)] |
| 57 | pub 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)] |
| 76 | pub 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)] |
| 83 | pub 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)] |
| 93 | pub 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)] |
| 107 | pub 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 | |
| 120 | impl 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 | |
| 241 | impl 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. |
| 252 | pub 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)] |
| 264 | mod 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 | } |