tollgate_client/usage_writer.rs
1//! The batched usage writer.
2//!
3//! A bounded queue separates the request path from billing I/O. The request
4//! path reserves a slot *before* admitting work
5//! ([`UsageRecorder::try_reserve`]); a full queue is
6//! `DenyReason::AccountingBackpressure` — shed with zero units charged,
7//! never a silent drop, never an unbounded block (INVARIANTS.md GL-8). The
8//! permit outlives execution, and the shutdown drain below waits for it, so
9//! a post-commit send is either ingested or explicitly counted — never
10//! silently dropped (INVARIANTS.md GL-13).
11//!
12//! The queue is a set of lanes, each a bounded mpsc channel, chosen by the
13//! request thread's sticky locality and sized to split `queue_capacity`
14//! exactly (GL-137). A request reserves in its own lane and tries the others
15//! before shedding, so the shed point is the whole queue. Lanes keep one
16//! thread's sender count, semaphore, tail and entry counter off every other
17//! thread's lines; one shared channel put all of them, for every account, on
18//! the same few. Events within a lane keep their order; events that overflow
19//! into another lane are delivered in that lane's order.
20//!
21//! The writer never parks on a lane, so a send wakes nothing. It drains every
22//! lane on its flush tick — when a partial batch was due anyway — and earlier
23//! when a lane reaches its ring point, which rings a doorbell once per fill
24//! rather than once per event.
25//!
26//! The writer task drains the lanes into batches and ingests them through
27//! the [`UsageSink`]. Ingest is idempotent on
28//! request id (INVARIANTS.md GL-7), so retrying a whole batch after a backend
29//! error is always safe. A failing backend is retried with backoff forever
30//! while the channel backs up and sheds upstream — memory stays bounded at
31//! one in-flight batch plus the channel.
32//!
33//! Every ingest is wall-clock bounded by
34//! [`UsageWriterConfig::ingest_timeout`] (INVARIANTS.md GL-18): a sink that
35//! hangs rather than erroring is indistinguishable from one that is merely
36//! slow, and neither may park the task. A timed-out call is a failed
37//! attempt — during shutdown it counts toward `lost`, never toward success.
38//!
39//! Shutdown is level-triggered: every loop consults the watch's *current*
40//! value, never only its edge notification, so a shutdown signalled during a
41//! retry backoff still reaches the bounded final flush (this was review
42//! finding GL-2 — the original edge-triggered design could consume the
43//! notification inside the retry loop and then wait forever for a second one
44//! that never came). Dropping the [`UsageWriter`] handle without calling
45//! [`shutdown`](UsageWriter::shutdown) aborts the task outright — enqueued
46//! events are lost in that path, which is why graceful code always calls
47//! `shutdown`.
48//!
49//! The final flush closes every lane (new reservations deny from that
50//! instant), then drains until every lane reports disconnected — empty *and*
51//! every outstanding permit resolved by sending or dropping, never merely
52//! momentarily empty — waking on each permit that resolves, bounded
53//! by [`UsageWriterConfig::shutdown_drain_deadline`]. Permits still
54//! unresolved at the deadline are reported in [`WriterStats::unresolved`];
55//! their charges are locally committed but unbilled, bounded thereafter by
56//! TTL reclaim (INVARIANTS.md GL-9).
57//!
58//! Every charge that enters the queue is counted until the writer gives it a
59//! billing outcome, in a counter held outside the task. A writer that dies
60//! instead of reporting therefore still says how many committed charges it
61//! was carrying: [`shutdown`](UsageWriter::shutdown) yields
62//! [`WriterShutdownError`], never a zeroed [`WriterStats`] that would read
63//! exactly like a clean run (INVARIANTS.md GL-8).
64//!
65//! Lifecycle order the embedder must follow: stop admitting, quiesce the
66//! request tasks holding permits or committed `Committed` guards, `shutdown()`
67//! this writer, and only then shut the lease manager down — events must land
68//! while their lease is live (INVARIANTS.md GL-12). Size the drain deadline
69//! within `expiry_safety_margin + reclaim_grace`, so a slow drain surfaces
70//! as `rejected` at the sink rather than silent loss.
71
72use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering};
73use std::sync::{Arc, OnceLock};
74
75use jiff::{SignedDuration, Timestamp};
76use tokio::sync::mpsc::error::{TryRecvError, TrySendError};
77use tokio::sync::{Notify, mpsc, watch};
78use tracing::Instrument as _;
79
80use tollgate_core::{DenyReason, LocalSharding, Locality, UsageEvent};
81use tollgate_store::Clock;
82use tollgate_store::{MAX_INGEST_BATCH, UsageSink};
83
84/// The usage queue's size and the writer's batching, retry and shutdown
85/// timing.
86///
87/// There are no defaults; every value is a deployment decision.
88/// [`validate`](Self::validate) runs before the task starts
89/// (INVARIANTS.md 16).
90#[derive(Debug, Clone, Copy)]
91pub struct UsageWriterConfig {
92 /// Channel capacity — the shed point. Size it to cover the sink's worst
93 /// tolerable outage at peak admission rate.
94 ///
95 /// Counted in usage events, one per in-flight or committed-but-unwritten
96 /// request. Every reserved permit holds a slot, so it also caps
97 /// concurrent admitted requests. Too small sheds
98 /// (`AccountingBackpressure`, zero charge) during brief sink hiccups or
99 /// bursts; too large lets a longer outage pile up unbilled events in
100 /// memory and in the shutdown drain. The writer splits it exactly into
101 /// lanes by host parallelism, with at least 64 slots a lane when there is
102 /// more than one. Must be positive.
103 pub queue_capacity: usize,
104 /// Largest batch handed to one `ingest` call.
105 ///
106 /// Counted in events. Too small multiplies sink calls under load; too
107 /// large makes each call, and each retry of a failing one, heavier. It
108 /// also sets how full a lane gets before it wakes the writer early. Must
109 /// be positive and at most [`MAX_INGEST_BATCH`], the ingest endpoint's
110 /// limit.
111 pub max_batch: usize,
112 /// A partial batch is flushed after at most this long.
113 ///
114 /// The latency between a charge being recorded and the sink seeing it at
115 /// low traffic. Too long delays billing and leaves more events in the
116 /// queue when a process dies; too short sends many small batches. Must be
117 /// positive.
118 pub flush_interval: std::time::Duration,
119 /// Backoff between retries of a failing ingest.
120 ///
121 /// A failing sink is retried forever at this pace while the queue fills
122 /// and sheds upstream. Too short hot-spins against a failing sink; too
123 /// long delays recovery after it returns. Must be positive.
124 pub retry_backoff: std::time::Duration,
125 /// Wall-clock bound on the shutdown drain: how long `shutdown` waits for
126 /// outstanding permits (slots reserved by in-flight requests or committed
127 /// committed guards) to resolve, *including* the ingest calls it makes along
128 /// the way. Must be positive. Size it within
129 /// `expiry_safety_margin + reclaim_grace` so an event landing at the end
130 /// of the drain is still billable against its lease.
131 ///
132 /// Too short reports permits as [`WriterStats::unresolved`] and batches
133 /// as [`WriterStats::lost`] that a longer drain would have billed; too
134 /// long risks events landing after their lease settled, which the sink
135 /// rejects. Under an [`InstanceRuntime`](crate::InstanceRuntime) the
136 /// runtime's `shutdown_deadline` must cover it.
137 pub shutdown_drain_deadline: std::time::Duration,
138 /// Wall-clock bound on one `UsageSink::ingest` call. A sink that hangs
139 /// rather than erroring would otherwise park the writer task forever, and
140 /// with it every later shutdown step. Must be positive.
141 ///
142 /// A timed-out call is a failed attempt and is retried; the sink may
143 /// still have recorded the batch, which the retry then reports as
144 /// `duplicate`. Set it above the sink's slowest legitimate ingest of a
145 /// full batch: too short turns a slow sink into endless retries; too long
146 /// holds the queue behind a hung call.
147 pub ingest_timeout: std::time::Duration,
148}
149
150/// Why a [`UsageWriterConfig`] was refused: the first rule it broke.
151#[derive(Debug, Clone, Copy, PartialEq, Eq)]
152pub struct UsageWriterConfigError(pub &'static str);
153
154impl std::fmt::Display for UsageWriterConfigError {
155 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
156 f.write_str(self.0)
157 }
158}
159
160impl std::error::Error for UsageWriterConfigError {}
161
162impl UsageWriterConfig {
163 /// Every field is a contract, so none of them is silently repaired
164 /// (INVARIANTS.md GL-16): a zero capacity or batch size has no sensible
165 /// coercion, a zero backoff hot-spins against a failing sink, and a zero
166 /// deadline or timeout reports failure without waiting at all.
167 pub fn validate(&self) -> Result<(), UsageWriterConfigError> {
168 if self.queue_capacity == 0 {
169 return Err(UsageWriterConfigError("queue_capacity must be positive"));
170 }
171 if self.max_batch > MAX_INGEST_BATCH {
172 return Err(UsageWriterConfigError(
173 "max_batch exceeds the ingest endpoint's documented limit",
174 ));
175 }
176 if self.max_batch == 0 {
177 return Err(UsageWriterConfigError("max_batch must be positive"));
178 }
179 if self.flush_interval.is_zero() {
180 return Err(UsageWriterConfigError("flush_interval must be positive"));
181 }
182 if self.retry_backoff.is_zero() {
183 return Err(UsageWriterConfigError("retry_backoff must be positive"));
184 }
185 if self.shutdown_drain_deadline.is_zero() {
186 return Err(UsageWriterConfigError(
187 "shutdown_drain_deadline must be positive",
188 ));
189 }
190 if self.ingest_timeout.is_zero() {
191 return Err(UsageWriterConfigError("ingest_timeout must be positive"));
192 }
193 Ok(())
194 }
195}
196
197/// A counter written from the request path, on a cache line of its own.
198///
199/// `#[repr(align(128))]` rounds the type's size up to its alignment, so no two
200/// of these — and nothing else in the struct — can share a line on the
201/// 128-byte Apple Silicon target or on 64-byte-line x86-64. The counters the
202/// writer task alone updates need no such treatment: one writer cannot
203/// contend with itself.
204#[repr(align(128))]
205#[derive(Debug)]
206struct Contended(AtomicU64);
207
208impl Contended {
209 const fn zero() -> Self {
210 Contended(AtomicU64::new(0))
211 }
212
213 #[inline]
214 fn bump(&self, by: u64) {
215 self.0.fetch_add(by, Ordering::Relaxed);
216 }
217
218 #[inline]
219 fn get(&self) -> u64 {
220 self.0.load(Ordering::Relaxed)
221 }
222}
223
224/// Every accounting number the writer produces, held outside the task.
225///
226/// The tally used to be a `WriterStats` local on the task's stack, materialised
227/// only when the task exited — so a process that crashed, was killed, or simply
228/// kept running reported nothing, and `lost` was observable only after the one
229/// kind of shutdown where loss is least likely (GL-38). Living out here it is
230/// readable at any time, and it survives the task's death exactly as the
231/// unaccounted count always has (INVARIANTS.md GL-8).
232///
233/// These are also the *only* copy: [`UsageWriter::shutdown`] returns a snapshot
234/// of these counters rather than a parallel tally, so the running totals and the
235/// final report cannot disagree.
236///
237/// `Relaxed` throughout, like the admission counters: nothing is published
238/// through them, and atomicity — no lost increments — is all they need.
239#[derive(Debug)]
240pub struct WriterCounters {
241 unattributed: AtomicU64,
242 attribution_unreported_batches: AtomicU64,
243 attribution_degraded: AtomicBool,
244 counter_overflow: AtomicBool,
245 // Written only by the writer task, between batches.
246 accepted: AtomicU64,
247 duplicate: AtomicU64,
248 rejected: AtomicU64,
249 lost: AtomicU64,
250 unresolved: AtomicU64,
251 /// When the sink last answered an ingest, in milliseconds since the epoch.
252 /// `i64::MIN` means "never": zero cannot be the sentinel, because the epoch
253 /// itself is a legitimate timestamp that tests use routinely.
254 last_ingest_ms: AtomicI64,
255 /// Charges the writer has given an outcome. Written only by the task.
256 settled: AtomicU64,
257 /// Each lane's count of charges that entered it (GL-137). Set once when the
258 /// writer spawns; a counter set nobody attached reads as holding nothing.
259 lanes: OnceLock<Arc<[LaneStats]>>,
260 // Written from the request path, by as many cores as serve requests.
261 shed: Contended,
262}
263
264impl WriterCounters {
265 /// A counter set that has recorded nothing and is attached to no lanes.
266 #[must_use]
267 pub const fn new() -> Self {
268 WriterCounters {
269 unattributed: AtomicU64::new(0),
270 attribution_unreported_batches: AtomicU64::new(0),
271 attribution_degraded: AtomicBool::new(false),
272 counter_overflow: AtomicBool::new(false),
273 accepted: AtomicU64::new(0),
274 duplicate: AtomicU64::new(0),
275 rejected: AtomicU64::new(0),
276 lost: AtomicU64::new(0),
277 unresolved: AtomicU64::new(0),
278 last_ingest_ms: AtomicI64::new(i64::MIN),
279 settled: AtomicU64::new(0),
280 lanes: OnceLock::new(),
281 shed: Contended::zero(),
282 }
283 }
284
285 /// The sink answered. Recorded even when every event in the batch was
286 /// refused: `rejected` is an answer, and what this timestamp distinguishes
287 /// is a reachable sink from an unreachable one.
288 fn record_ingest(&self, report: &tollgate_store::IngestReport, at: Timestamp) {
289 self.add_outcome(&self.accepted, report.accepted);
290 self.add_outcome(&self.duplicate, report.duplicate);
291 self.add_outcome(&self.rejected, report.rejected);
292 match report.unattributed {
293 Some(n) => self.add_outcome(&self.unattributed, n),
294 None => self.add_outcome(&self.attribution_unreported_batches, 1),
295 }
296 let degraded = report.unattributed != Some(0);
297 let previous = self.attribution_degraded.swap(degraded, Ordering::Relaxed);
298 if degraded && !previous {
299 tracing::warn!(coverage_complete = false, unattributed = ?report.unattributed,
300 "credential activity coverage is incomplete or unavailable");
301 } else if previous && !degraded {
302 tracing::info!(
303 coverage_complete = true,
304 "credential attribution reporting recovered for this batch"
305 );
306 }
307 self.last_ingest_ms
308 .store(at.as_millisecond(), Ordering::Relaxed);
309 }
310
311 fn record_lost(&self, events: u64) {
312 self.add_outcome(&self.lost, events);
313 }
314
315 // Off-path cumulative outcomes saturate visibly. The request-side queue
316 // accounting counters retain their existing mechanism and budget.
317 fn add_outcome(&self, counter: &AtomicU64, delta: u64) {
318 if counter
319 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |old| {
320 old.checked_add(delta)
321 })
322 .is_err()
323 {
324 counter.store(u64::MAX, Ordering::Relaxed);
325 if !self.counter_overflow.swap(true, Ordering::Relaxed) {
326 tracing::error!(
327 counter_overflow = true,
328 "usage outcome counter overflow; totals are saturated"
329 );
330 }
331 }
332 }
333
334 fn set_unresolved(&self, permits: u64) {
335 self.unresolved.store(permits, Ordering::Relaxed);
336 }
337
338 /// `events` charges have been given a billing outcome — delivered or
339 /// reported lost — so they no longer count as unaccounted.
340 fn settled(&self, events: u64) {
341 self.settled.fetch_add(events, Ordering::Relaxed);
342 }
343
344 /// Charges in the queue with no billing outcome yet: every lane's entries
345 /// less everything the writer has settled. Each entry is counted before its
346 /// event is sent, and settled only after it is received, so the difference
347 /// never undercounts what a dying writer holds.
348 fn unaccounted(&self) -> u64 {
349 let entered = self.lanes.get().map_or(0, |lanes| {
350 lanes.iter().fold(0u64, |total, lane| {
351 total.wrapping_add(lane.enqueued.load(Ordering::Relaxed))
352 })
353 });
354 entered.saturating_sub(self.settled.load(Ordering::Relaxed))
355 }
356
357 /// A backpressure refusal, counted where it happens rather than by the
358 /// caller: `try_reserve` is the only way to be refused, so counting here
359 /// cannot be forgotten by an embedder.
360 fn record_shed(&self) {
361 self.shed.bump(1);
362 }
363
364 /// The five numbers [`UsageWriter::shutdown`] reports.
365 #[must_use]
366 pub fn stats(&self) -> WriterStats {
367 WriterStats {
368 unattributed: self.unattributed.load(Ordering::Relaxed),
369 attribution_unreported_batches: self
370 .attribution_unreported_batches
371 .load(Ordering::Relaxed),
372 counter_overflow: self.counter_overflow.load(Ordering::Relaxed),
373 accepted: self.accepted.load(Ordering::Relaxed),
374 duplicate: self.duplicate.load(Ordering::Relaxed),
375 rejected: self.rejected.load(Ordering::Relaxed),
376 lost: self.lost.load(Ordering::Relaxed),
377 unresolved: self.unresolved.load(Ordering::Relaxed),
378 }
379 }
380
381 fn last_ingest_at(&self) -> Option<Timestamp> {
382 let millis = self.last_ingest_ms.load(Ordering::Relaxed);
383 (millis != i64::MIN)
384 .then(|| Timestamp::from_millisecond(millis).ok())
385 .flatten()
386 }
387
388 /// Everything readable while the writer runs, given the queue's shape —
389 /// which only a holder of the channel can supply.
390 fn health(&self, queue_depth: usize, queue_capacity: usize) -> WriterHealth {
391 WriterHealth {
392 stats: self.stats(),
393 unaccounted: self.unaccounted(),
394 shed: self.shed.get(),
395 queue_depth,
396 queue_capacity,
397 last_ingest_at: self.last_ingest_at(),
398 }
399 }
400}
401
402impl Default for WriterCounters {
403 fn default() -> Self {
404 Self::new()
405 }
406}
407
408/// A reading of the writer's accounting health, safe to serialise.
409///
410/// Note what `lost` does *not* tell you at runtime: the steady-state path
411/// retries an *unavailable* sink forever, so an outage declares nothing lost
412/// until the final flush gives up. A healthy-but-cut-off process reports
413/// `lost == 0` for its whole life. The one steady-state loss is a batch the
414/// sink *refuses* (a permanent error, such as an accounting overflow): it is
415/// counted in `lost` and dropped at once, because replaying it would earn the
416/// same refusal forever and block every event behind it. The leading indicators are [`WriterHealth::ingest_age`], a
417/// `queue_depth` approaching `queue_capacity`, and `stats.rejected` — which the
418/// sink has already refused, and which is bounded billing loss.
419#[derive(Debug, Clone, Copy, PartialEq, Eq)]
420pub struct WriterHealth {
421 /// The same accounting and coverage diagnostics a clean shutdown returns, read live.
422 pub stats: WriterStats,
423 /// Charges in the queue with no billing outcome yet.
424 pub unaccounted: u64,
425 /// Requests refused for want of queue capacity (INVARIANTS.md GL-8).
426 pub shed: u64,
427 /// Slots currently held by queued events and outstanding permits.
428 pub queue_depth: usize,
429 /// The shed point: `queue_depth` reaching this denies the next request.
430 pub queue_capacity: usize,
431 /// When the sink last answered, or `None` if it never has.
432 pub last_ingest_at: Option<Timestamp>,
433}
434
435impl WriterHealth {
436 /// How long since the sink last answered — the signal that separates "no
437 /// traffic" from "the sink has been unreachable for twenty minutes".
438 /// `None` while no ingest has ever succeeded, which is also the state of a
439 /// freshly started process.
440 #[must_use]
441 pub fn ingest_age(&self, now: Timestamp) -> Option<SignedDuration> {
442 self.last_ingest_at.map(|at| now.duration_since(at))
443 }
444}
445
446/// One lane's request-written state, on a cache line of its own (GL-137).
447#[repr(align(128))]
448#[derive(Debug)]
449struct LaneStats {
450 /// Charges that entered this lane. Read with every other lane's, less what
451 /// the writer settled, as the instance's unaccounted count.
452 enqueued: AtomicU64,
453 /// Set by the send that finds the lane at its ring point, cleared by the
454 /// writer before it drains the lane: one doorbell per fill, not per event.
455 rung: AtomicBool,
456}
457
458/// What the lanes and the writer share. Never cloned per request: a permit
459/// reaches it through its own lane, so no request touches this `Arc`'s count.
460#[derive(Debug)]
461struct Queue {
462 stats: Arc<[LaneStats]>,
463 /// Wakes the writer early: a lane reaching its ring point, a lane whose
464 /// last handle went away, and, while draining, every permit that resolves.
465 doorbell: Notify,
466 /// Set by the final flush before it closes the lanes. Read on every permit
467 /// release and never written again, so it costs a shared read, not a write.
468 draining: AtomicBool,
469 /// How full a lane gets before its send rings the doorbell.
470 ring_at: usize,
471 counters: Arc<WriterCounters>,
472}
473
474/// Rings the doorbell when a lane's last handle goes away. A field of its own,
475/// declared after the sender, so the ring comes *after* the sender drops and
476/// the woken writer sees the lane disconnected rather than merely empty.
477#[derive(Debug)]
478struct RingOnDrop(Arc<Queue>);
479
480impl Drop for RingOnDrop {
481 fn drop(&mut self) {
482 self.0.doorbell.notify_one();
483 }
484}
485
486/// One lane of the usage queue: a bounded channel a subset of request threads
487/// reserve in, on lines no other lane's threads write.
488#[repr(align(128))]
489#[derive(Debug)]
490struct Lane {
491 tx: mpsc::Sender<UsageEvent>,
492 index: usize,
493 queue: RingOnDrop,
494}
495
496/// Cheap-to-clone handle for request handlers.
497///
498/// The queue is partitioned into lanes by the request thread's sticky
499/// locality (GL-137). One shared channel put every request of every account on
500/// the same sender count, semaphore, tail and waker lines, and woke the writer
501/// once per event: eight threads measured 3.2 µs per reserve-and-record. A
502/// request reserves in its own lane and tries the others before shedding, so
503/// the shed point is still exactly `queue_capacity` — a partition, not a
504/// reservation.
505#[derive(Clone)]
506pub struct UsageRecorder {
507 lanes: Arc<[Arc<Lane>]>,
508 layout: LocalSharding,
509 queue: Arc<Queue>,
510}
511
512impl UsageRecorder {
513 /// Reserve accounting capacity for one request, *before* admission. A
514 /// full queue denies here — before any units are reserved or any work
515 /// runs.
516 pub fn try_reserve(&self) -> Result<UsagePermit, DenyReason> {
517 let count = self.lanes.len();
518 let mut index = Locality::current().index(self.layout);
519 for _ in 0..count {
520 let lane = &self.lanes[index];
521 match lane.tx.clone().try_reserve_owned() {
522 Ok(permit) => {
523 return Ok(UsagePermit {
524 permit: Some(permit),
525 lane: Arc::clone(lane),
526 });
527 }
528 // A full lane is not a full queue: another lane may have room.
529 Err(TrySendError::Full(_)) => {}
530 // Closed lanes close together, at shutdown.
531 Err(TrySendError::Closed(_)) => break,
532 }
533 index += 1;
534 if index == count {
535 index = 0;
536 }
537 }
538 self.queue.counters.record_shed();
539 Err(DenyReason::AccountingBackpressure)
540 }
541
542 /// Whether the writer task has exited and can no longer accept events.
543 #[must_use]
544 pub fn is_closed(&self) -> bool {
545 self.lanes[0].tx.is_closed()
546 }
547
548 pub(crate) async fn closed(&self) {
549 // Every lane closes together: the final flush closes them all, and the
550 // task holds every receiver.
551 self.lanes[0].tx.closed().await;
552 }
553
554 /// The writer's accounting health, readable at any time.
555 ///
556 /// Deliberately on the *recorder*: a service holds this handle in its
557 /// request state, while the [`UsageWriter`] is usually moved into whatever
558 /// owns shutdown — so exposing the numbers only there would put them out of
559 /// reach of the endpoint that needs to report them.
560 #[must_use]
561 pub fn health(&self) -> WriterHealth {
562 let (depth, capacity) = self.lanes.iter().fold((0, 0), |(depth, capacity), lane| {
563 (
564 depth + lane.tx.max_capacity() - lane.tx.capacity(),
565 capacity + lane.tx.max_capacity(),
566 )
567 });
568 self.queue.counters.health(depth, capacity)
569 }
570}
571
572/// One reserved accounting slot. Send the committed request's event with
573/// [`record`](UsagePermit::record); dropping the permit (deny, cancel,
574/// zero-charge path) releases the slot.
575pub struct UsagePermit {
576 /// `None` only after `record` has sent through it.
577 permit: Option<mpsc::OwnedPermit<UsageEvent>>,
578 lane: Arc<Lane>,
579}
580
581impl UsagePermit {
582 /// Send the committed request's event into the reserved slot. Never
583 /// blocks and cannot fail: the slot was reserved before admission. The
584 /// event counts as unaccounted until the writer gives it a billing
585 /// outcome.
586 pub fn record(mut self, event: UsageEvent) {
587 let permit = self
588 .permit
589 .take()
590 .expect("a permit is consumed only by record, which consumes the permit");
591 let queue = &self.lane.queue.0;
592 let stats = &queue.stats[self.lane.index];
593 // Counted from the moment it enters the queue until the writer gives
594 // it a billing outcome; a writer that dies in between is therefore
595 // able to say how many charges it was carrying. Counted before the
596 // send, so the writer can never settle an event that was not counted.
597 stats.enqueued.fetch_add(1, Ordering::Relaxed);
598 let tx = permit.send(event);
599 // No wake per event: the writer drains on its flush tick. A lane that
600 // reaches its ring point rings once, so a burst cannot outrun the tick.
601 if tx.max_capacity() - tx.capacity() >= queue.ring_at
602 && !stats.rung.load(Ordering::Relaxed)
603 && !stats.rung.swap(true, Ordering::AcqRel)
604 {
605 queue.doorbell.notify_one();
606 }
607 }
608}
609
610impl Drop for UsagePermit {
611 fn drop(&mut self) {
612 // Release the slot first, so a draining writer woken below finds it
613 // released rather than still outstanding.
614 drop(self.permit.take());
615 let queue = &self.lane.queue.0;
616 if queue.draining.load(Ordering::Acquire) {
617 queue.doorbell.notify_one();
618 }
619 }
620}
621
622impl tollgate_core::UsageSlot for UsagePermit {
623 fn record(self, event: UsageEvent) {
624 UsagePermit::record(self, event);
625 }
626}
627
628/// Terminal accounting of a writer's lifetime, returned by
629/// [`UsageWriter::shutdown`]. `lost` counts events a final flush could not
630/// deliver; `unresolved` counts permits still outstanding when the drain
631/// deadline expired — both reported, never silent.
632///
633/// Deliberately not [`Default`]: a zeroed report must never be conjurable
634/// from a failure (`unwrap_or_default` on a dead task's `JoinError` is the
635/// bug this type's history records — issue GL-41). Use [`WriterStats::ZERO`]
636/// when a starting value is genuinely meant.
637#[derive(Debug, Clone, Copy, PartialEq, Eq)]
638pub struct WriterStats {
639 /// Confirmed unattributed newly accepted events. Lost acknowledgements
640 /// followed by duplicate replies cannot reconstruct historical counts.
641 pub unattributed: u64,
642 /// Acknowledged batches whose sink did not report attribution support.
643 pub attribution_unreported_batches: u64,
644 /// At least one cumulative outcome exceeded u64; its value is saturated.
645 pub counter_overflow: bool,
646 /// Events the sink newly recorded: billed.
647 pub accepted: u64,
648 /// Events whose request id the sink had already recorded, from a retried
649 /// batch whose acknowledgement was lost. Ingest is idempotent
650 /// (INVARIANTS.md 7), so these are not double charges; a steady rate
651 /// points at acknowledgements lost to timeouts.
652 pub duplicate: u64,
653 /// Events the sink refused: unknown lease, lease-capability mismatch, no
654 /// remaining accounting capacity, or units outside its storage domain.
655 /// Bounded billing loss that has already happened, left for
656 /// reconciliation. Normal is zero; a trickle usually means events land
657 /// after their lease settled, so the drain deadline sits too close to
658 /// `expiry_safety_margin + reclaim_grace`.
659 pub rejected: u64,
660 /// Events that will never reach the ledger: batches the sink refused
661 /// with a non-retryable error, and events the final flush could not
662 /// deliver before its deadline. Any nonzero value is committed charges
663 /// that went unbilled. A healthy but unreachable sink leaves this at zero
664 /// until shutdown, so alert on [`WriterHealth::ingest_age`] instead.
665 pub lost: u64,
666 /// Permits (in-flight requests or committed guards) that neither sent
667 /// nor dropped before the drain deadline. Their charges are locally
668 /// committed but unbilled; TTL reclaim bounds them (INVARIANTS.md GL-9).
669 pub unresolved: u64,
670}
671
672impl WriterStats {
673 /// A writer that has accounted for nothing yet.
674 pub const ZERO: WriterStats = WriterStats {
675 unattributed: 0,
676 attribution_unreported_batches: 0,
677 counter_overflow: false,
678 accepted: 0,
679 duplicate: 0,
680 rejected: 0,
681 lost: 0,
682 unresolved: 0,
683 };
684}
685
686/// The writer task ended without reporting: it panicked, or it was aborted.
687/// `unaccounted` is a *lower bound* on committed charges left with no billing
688/// record — the events that had entered the queue but had not yet been given
689/// an outcome.
690#[derive(Debug, Clone, Copy, PartialEq, Eq)]
691pub struct WriterShutdownError {
692 /// Lower bound on committed charges that entered the queue and received
693 /// no billing outcome before the task ended.
694 pub unaccounted: u64,
695 /// True when the task panicked, false when it was cancelled or aborted.
696 pub panicked: bool,
697}
698
699impl std::fmt::Display for WriterShutdownError {
700 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
701 let cause = if self.panicked { "panicked" } else { "aborted" };
702 write!(
703 f,
704 "usage writer {cause} before reporting; at least {} committed charge(s) have no billing record",
705 self.unaccounted
706 )
707 }
708}
709
710impl std::error::Error for WriterShutdownError {}
711
712/// Handle to the writer task.
713pub struct UsageWriter {
714 shutdown: watch::Sender<bool>,
715 deadline: Arc<crate::ShutdownDeadline>,
716 handle: Option<tokio::task::JoinHandle<WriterStats>>,
717 counters: Arc<WriterCounters>,
718 /// Only to read the queue's depth for [`UsageWriter::health`] — weak, so
719 /// this handle never keeps a lane open by itself.
720 queue: Arc<[mpsc::WeakSender<UsageEvent>]>,
721 queue_capacity: usize,
722}
723
724/// The fewest slots a lane is given.
725///
726/// A lane smaller than a burst turns the sibling fallback into the routine
727/// path, and events that overflow into another lane are delivered in that
728/// lane's order rather than their thread's. Within one lane a thread's events
729/// keep their order; a queue too small for two such lanes keeps one, exactly
730/// the single channel it replaced.
731const MIN_LANE_CAPACITY: usize = 64;
732
733/// How many lanes a queue of `capacity` slots gets: the host's parallelism as
734/// a power of two, so the locality reduction is a mask, but never so many that
735/// a lane falls below [`MIN_LANE_CAPACITY`].
736fn lane_count(capacity: usize) -> usize {
737 let parallelism = LocalSharding::available_parallelism()
738 .get()
739 .next_power_of_two();
740 let affordable = capacity / MIN_LANE_CAPACITY;
741 if affordable < 2 {
742 return 1;
743 }
744 // A power of two no larger than what the capacity affords.
745 parallelism.min(1 << affordable.ilog2())
746}
747
748/// `total` slots split exactly across `count` lanes; the sizes sum to `total`.
749fn lane_capacity(total: usize, count: usize, index: usize) -> usize {
750 total / count + usize::from(index < total % count)
751}
752
753impl UsageWriter {
754 /// Validate `config` and start the writer task. Returns the cloneable
755 /// request-side [`UsageRecorder`] and this owning handle, which must be
756 /// retained and shut down; dropping it aborts the task and loses
757 /// whatever is queued. Must be called within a Tokio runtime.
758 ///
759 /// # Errors
760 ///
761 /// [`UsageWriterConfigError`] when `config` fails validation; nothing is
762 /// started.
763 pub fn spawn(
764 sink: Arc<dyn UsageSink>,
765 clock: Arc<dyn Clock>,
766 config: UsageWriterConfig,
767 ) -> Result<(UsageRecorder, UsageWriter), UsageWriterConfigError> {
768 config.validate()?;
769 Ok(Self::spawn_lanes(
770 sink,
771 clock,
772 config,
773 lane_count(config.queue_capacity),
774 ))
775 }
776
777 /// `spawn` with an explicit lane count, so the multi-lane paths can be
778 /// tested on any host. `config` is already validated and `count` is at
779 /// least one and at most `queue_capacity`.
780 fn spawn_lanes(
781 sink: Arc<dyn UsageSink>,
782 clock: Arc<dyn Clock>,
783 config: UsageWriterConfig,
784 count: usize,
785 ) -> (UsageRecorder, UsageWriter) {
786 let (senders, receivers): (Vec<_>, Vec<_>) = (0..count)
787 .map(|index| mpsc::channel(lane_capacity(config.queue_capacity, count, index)))
788 .unzip();
789 // Weak handles let the drain count still-outstanding permits at the
790 // deadline without holding any lane open itself.
791 let weak: Arc<[mpsc::WeakSender<UsageEvent>]> =
792 senders.iter().map(mpsc::Sender::downgrade).collect();
793 let (shutdown, shutdown_rx) = watch::channel(false);
794 let counters = Arc::new(WriterCounters::new());
795 let stats: Arc<[LaneStats]> = (0..count)
796 .map(|_| LaneStats {
797 enqueued: AtomicU64::new(0),
798 rung: AtomicBool::new(false),
799 })
800 .collect();
801 counters
802 .lanes
803 .set(Arc::clone(&stats))
804 .expect("a fresh counter set has no lanes yet");
805 let smallest_lane = config.queue_capacity / count;
806 let queue = Arc::new(Queue {
807 stats,
808 doorbell: Notify::new(),
809 draining: AtomicBool::new(false),
810 // Ring at a batch, or at half a small lane so a burst rings before
811 // the lane is full rather than only as it sheds.
812 ring_at: config.max_batch.min(smallest_lane / 2).max(1),
813 counters: Arc::clone(&counters),
814 });
815 let lanes: Arc<[Arc<Lane>]> = senders
816 .into_iter()
817 .enumerate()
818 .map(|(index, tx)| {
819 Arc::new(Lane {
820 tx,
821 index,
822 queue: RingOnDrop(Arc::clone(&queue)),
823 })
824 })
825 .collect();
826 let deadline = Arc::new(crate::ShutdownDeadline::default());
827 let handle = tokio::spawn(
828 run(
829 Writer {
830 sink,
831 clock,
832 config,
833 counters: Arc::clone(&counters),
834 deadline: Arc::clone(&deadline),
835 },
836 Lanes {
837 rx: receivers,
838 weak: Arc::clone(&weak),
839 queue: Arc::clone(&queue),
840 },
841 shutdown_rx,
842 )
843 // One writer serves every account, so the span carries the
844 // queue's shape; account and lease identify individual events.
845 .instrument(tracing::info_span!(
846 "usage_writer",
847 queue_capacity = config.queue_capacity,
848 lanes = count,
849 max_batch = config.max_batch
850 )),
851 );
852 (
853 UsageRecorder {
854 lanes,
855 layout: LocalSharding::new(
856 std::num::NonZeroUsize::new(count).expect("lane_count is at least one"),
857 ),
858 queue,
859 },
860 UsageWriter {
861 shutdown,
862 deadline,
863 handle: Some(handle),
864 counters,
865 queue: weak,
866 queue_capacity: config.queue_capacity,
867 },
868 )
869 }
870
871 /// The writer's accounting health, readable at any time — the same numbers
872 /// [`UsageRecorder::health`] reports, for an embedder that holds this half.
873 #[must_use]
874 pub fn health(&self) -> WriterHealth {
875 let depth = self
876 .queue
877 .iter()
878 .filter_map(mpsc::WeakSender::upgrade)
879 .map(|tx| tx.max_capacity() - tx.capacity())
880 .sum();
881 self.counters.health(depth, self.queue_capacity)
882 }
883
884 /// Flush everything already enqueued, then stop. Call this *before*
885 /// releasing leases — events must land while their lease is live.
886 ///
887 /// A writer task that died instead of reporting yields
888 /// [`WriterShutdownError`] carrying the charges it was still holding —
889 /// never a zeroed [`WriterStats`], which would be indistinguishable from
890 /// a clean shutdown (INVARIANTS.md GL-8).
891 pub async fn shutdown(mut self) -> Result<WriterStats, WriterShutdownError> {
892 crate::signal(&self.shutdown, true, "usage-writer shutdown");
893 let Some(handle) = self.handle.as_mut() else {
894 // Unreachable through the public API: `shutdown` consumes the
895 // handle, and `Drop` runs only afterwards.
896 return Err(self.died(false));
897 };
898 // Borrow the handle so cancelling this future cannot detach billing.
899 match handle.await {
900 Ok(stats) => Ok(stats),
901 Err(join) => Err(self.died(join.is_panic())),
902 }
903 }
904
905 /// Constrain cleanup before signalling it, so a runtime's total budget
906 /// governs the actual receives, ingests, and backoffs.
907 pub(crate) fn stop_at(&self, deadline: tokio::time::Instant) {
908 self.deadline.constrain(deadline);
909 crate::signal(&self.shutdown, true, "usage-writer shutdown");
910 }
911
912 fn died(&self, panicked: bool) -> WriterShutdownError {
913 WriterShutdownError {
914 unaccounted: self.counters.unaccounted(),
915 panicked,
916 }
917 }
918}
919
920impl Drop for UsageWriter {
921 fn drop(&mut self) {
922 // Dropped without shutdown(): abort rather than leak a detached
923 // task. Ungraceful by definition — enqueued events die with it.
924 if let Some(handle) = self.handle.take() {
925 handle.abort();
926 }
927 }
928}
929
930/// The writer's collaborators, fixed for the task's lifetime.
931struct Writer {
932 sink: Arc<dyn UsageSink>,
933 clock: Arc<dyn Clock>,
934 config: UsageWriterConfig,
935 counters: Arc<WriterCounters>,
936 deadline: Arc<crate::ShutdownDeadline>,
937}
938
939impl Writer {
940 /// Every event in `batch` now has a billing outcome, so it no longer
941 /// counts against what a dying writer would be holding.
942 fn account_for(&self, batch: &[UsageEvent]) {
943 self.counters.settled(batch.len() as u64);
944 }
945}
946
947/// The writer's half of the lanes.
948struct Lanes {
949 rx: Vec<mpsc::Receiver<UsageEvent>>,
950 weak: Arc<[mpsc::WeakSender<UsageEvent>]>,
951 queue: Arc<Queue>,
952}
953
954/// What one pass over the lanes found.
955enum Collected {
956 /// `batch` reached `max_batch`: deliver it before collecting more.
957 Full,
958 /// Every lane is empty for now. `disconnected` when no lane can ever yield
959 /// again: every handle is gone, or the lanes were closed and every permit
960 /// has resolved.
961 Empty { disconnected: bool },
962}
963
964impl Lanes {
965 /// Move whatever the lanes hold into `batch`, up to `max_batch`, without
966 /// waiting. A lane's doorbell flag is cleared *before* the lane is read, so
967 /// a send that lands after the read rings again rather than waiting a tick.
968 fn collect(&mut self, batch: &mut Vec<UsageEvent>, max_batch: usize) -> Collected {
969 let mut disconnected = true;
970 for (index, rx) in self.rx.iter_mut().enumerate() {
971 self.queue.stats[index].rung.store(false, Ordering::Release);
972 loop {
973 if batch.len() >= max_batch {
974 return Collected::Full;
975 }
976 match rx.try_recv() {
977 Ok(event) => batch.push(event),
978 Err(TryRecvError::Empty) => {
979 disconnected = false;
980 break;
981 }
982 Err(TryRecvError::Disconnected) => break,
983 }
984 }
985 }
986 Collected::Empty { disconnected }
987 }
988}
989
990async fn run(writer: Writer, mut lanes: Lanes, mut shutdown: watch::Receiver<bool>) -> WriterStats {
991 let config = writer.config;
992 let max_batch = config.max_batch;
993 let mut batch: Vec<UsageEvent> = Vec::with_capacity(max_batch);
994
995 loop {
996 // Level check at every loop boundary: a shutdown observed anywhere
997 // below (including inside the retry backoff) lands here.
998 if *shutdown.borrow() {
999 return final_flush(&writer, &mut lanes, &mut batch).await;
1000 }
1001
1002 // The writer never parks on a lane's receiver, so a send finds no
1003 // waker to wake (GL-137). It drains on this tick, which is when a
1004 // partial batch was due anyway, and earlier only when a lane rings.
1005 let deadline = tokio::time::sleep(config.flush_interval);
1006 tokio::pin!(deadline);
1007 let stop = loop {
1008 let due = tokio::select! {
1009 () = lanes.queue.doorbell.notified() => false,
1010 () = &mut deadline => true,
1011 changed = shutdown.changed() => {
1012 // Err = sender dropped without shutdown(); treat both as
1013 // stop so the task can never outlive its handle usefully.
1014 if changed.is_err() || *shutdown.borrow() {
1015 break true;
1016 }
1017 continue;
1018 }
1019 };
1020 let disconnected = loop {
1021 match lanes.collect(&mut batch, max_batch) {
1022 Collected::Full => {
1023 // Retry until delivered or shutdown interrupts; either
1024 // way the loop-top level check decides what happens next.
1025 flush_retrying(&writer, &mut batch, &mut shutdown).await;
1026 if *shutdown.borrow() || shutdown.has_changed().is_err() {
1027 break false;
1028 }
1029 }
1030 Collected::Empty { disconnected } => break disconnected,
1031 }
1032 };
1033 if *shutdown.borrow() || shutdown.has_changed().is_err() {
1034 break true;
1035 }
1036 if due && !batch.is_empty() {
1037 flush_retrying(&writer, &mut batch, &mut shutdown).await;
1038 }
1039 // Every handle is gone: nothing can arrive, so stop as the channel's
1040 // close used to (the recorder's lanes ring as they drop).
1041 if disconnected {
1042 break true;
1043 }
1044 if due {
1045 break false;
1046 }
1047 };
1048 if stop {
1049 return final_flush(&writer, &mut lanes, &mut batch).await;
1050 }
1051 }
1052}
1053
1054/// Ingest `batch`, retrying with backoff until it is delivered (batch
1055/// cleared) or shutdown is observed (batch left intact for the final flush).
1056async fn flush_retrying(
1057 writer: &Writer,
1058 batch: &mut Vec<UsageEvent>,
1059 shutdown: &mut watch::Receiver<bool>,
1060) {
1061 let Writer {
1062 sink,
1063 clock,
1064 config,
1065 counters,
1066 ..
1067 } = writer;
1068 // An outage is a *duration*, not an event: retrying forever is the
1069 // designed behavior, so the only way it becomes visible is by reporting
1070 // when it began, and how long it lasted once it ends. `WriterStats` sees
1071 // none of this — a recovered outage produces a perfectly clean report.
1072 let mut outage: Option<(tokio::time::Instant, u64)> = None;
1073 loop {
1074 // A sink that hangs is indistinguishable from one that is merely slow,
1075 // and neither may park this task: the timeout turns both into the
1076 // ordinary retry path.
1077 // The same instant the batch is ingested with is the one recorded as
1078 // the last time the sink answered, so the health reading and the
1079 // ledger agree about when this batch happened.
1080 let now = clock.now();
1081 let ingest =
1082 tokio::time::timeout(config.ingest_timeout, ingest_checked(&**sink, batch, now));
1083 let outcome = tokio::select! {
1084 outcome = ingest => outcome,
1085 _ = shutdown.changed() => return,
1086 };
1087 match outcome {
1088 Ok(Ok(report)) => {
1089 if let Some((began, attempts)) = outage {
1090 tracing::info!(
1091 attempts,
1092 outage_ms = began.elapsed().as_millis(),
1093 "usage sink recovered"
1094 );
1095 }
1096 counters.record_ingest(&report, now);
1097 writer.account_for(batch);
1098 batch.clear();
1099 return;
1100 }
1101 // A refusal is a fact about this batch, not about the sink's
1102 // availability: replaying it unchanged earns the same answer
1103 // forever, and every event queued behind it waits for a recovery
1104 // that cannot come. Counted lost and dropped, so the queue drains
1105 // and later events bill — a permanent configuration or contract
1106 // error must not become an unbounded billing outage (GL-61).
1107 //
1108 // `lost` is the honest word for it: these events entered the queue
1109 // and will never reach the ledger. INVARIANTS GL-8 asks that they be
1110 // counted rather than silently dropped, not that they be delivered
1111 // by a sink that refuses them.
1112 Ok(Err(refused)) if !refused.is_retryable() => {
1113 counters.record_lost(batch.len() as u64);
1114 tracing::error!(
1115 events = batch.len(),
1116 %refused,
1117 "usage sink refused this batch and will refuse it again; \
1118 counted lost so later events are not blocked behind it"
1119 );
1120 writer.account_for(batch);
1121 batch.clear();
1122 return;
1123 }
1124 outcome => {
1125 let timed_out = outcome.is_err();
1126 let attempts = match &mut outage {
1127 Some((_, attempts)) => {
1128 *attempts += 1;
1129 *attempts
1130 }
1131 none => {
1132 // First failure of this outage: say so once at warn,
1133 // then stay quiet at debug so a long outage does not
1134 // become a log flood.
1135 *none = Some((tokio::time::Instant::now(), 1));
1136 tracing::warn!(
1137 events = batch.len(),
1138 timed_out,
1139 "usage sink failing; batching up and retrying"
1140 );
1141 1
1142 }
1143 };
1144 if attempts > 1 {
1145 tracing::debug!(attempts, timed_out, "usage sink still failing");
1146 }
1147 tokio::select! {
1148 _ = tokio::time::sleep(config.retry_backoff) => {}
1149 _ = shutdown.changed() => {}
1150 }
1151 // Level check covers every wake-up path: backoff elapsed,
1152 // signal received, or sender dropped.
1153 if *shutdown.borrow() || shutdown.has_changed().is_err() {
1154 if let Some((began, attempts)) = outage {
1155 tracing::warn!(
1156 attempts,
1157 outage_ms = began.elapsed().as_millis(),
1158 events = batch.len(),
1159 "shutdown observed during a sink outage; \
1160 the final flush decides these events' fate"
1161 );
1162 }
1163 return;
1164 }
1165 }
1166 }
1167 }
1168}
1169
1170/// Drain the channel until every outstanding permit resolves or the
1171/// configured deadline expires, giving each batch a bounded number of
1172/// delivery attempts. Whatever cannot be delivered is *reported* lost;
1173/// permits still outstanding at the deadline are *reported* unresolved. A
1174/// deadline expiry can therefore never look like a clean flush, and the
1175/// drain can never block past its bound.
1176async fn final_flush(
1177 writer: &Writer,
1178 lanes: &mut Lanes,
1179 batch: &mut Vec<UsageEvent>,
1180) -> WriterStats {
1181 let config = &writer.config;
1182 let max_batch = config.max_batch;
1183 // Refuse new reservations from this instant. Permits already handed out
1184 // keep their slots and can still deliver into the drain below. A closed
1185 // lane reports disconnected only once it is empty and every one of its
1186 // permits has sent or dropped — never merely because it is momentarily
1187 // empty — which is the done signal the single channel's `recv` gave.
1188 // Every permit that resolves from here rings the doorbell, so the drain
1189 // waits on it rather than polling.
1190 lanes.queue.draining.store(true, Ordering::Release);
1191 for rx in &mut lanes.rx {
1192 rx.close();
1193 }
1194 let deadline = writer.deadline.within(config.shutdown_drain_deadline);
1195 let mut expired = false;
1196 loop {
1197 let drained = loop {
1198 match lanes.collect(batch, max_batch) {
1199 Collected::Full => flush_bounded(writer, batch, deadline).await,
1200 Collected::Empty { disconnected } => break disconnected,
1201 }
1202 };
1203 if drained || expired {
1204 if !batch.is_empty() {
1205 flush_bounded(writer, batch, deadline).await;
1206 }
1207 // A final sweep after the deadline has already collected anything
1208 // that was queued, so what remains outstanding is permits.
1209 if !drained {
1210 writer
1211 .counters
1212 .set_unresolved(outstanding_permits(&lanes.weak));
1213 }
1214 return writer.counters.stats();
1215 }
1216 // Each wake is a resolved permit, a dropped lane, or the deadline; the
1217 // finite permit set bounds the loop. One more sweep follows an expiry.
1218 if tokio::time::timeout_at(deadline, lanes.queue.doorbell.notified())
1219 .await
1220 .is_err()
1221 {
1222 expired = true;
1223 }
1224 }
1225}
1226
1227/// Ingest `batch` with a bounded number of attempts, none of which may run
1228/// past the drain `deadline`; an undeliverable batch is counted lost. The
1229/// batch is cleared either way.
1230async fn flush_bounded(
1231 writer: &Writer,
1232 batch: &mut Vec<UsageEvent>,
1233 deadline: tokio::time::Instant,
1234) {
1235 let Writer {
1236 sink,
1237 clock,
1238 config,
1239 counters,
1240 ..
1241 } = writer;
1242 const FINAL_FLUSH_ATTEMPTS: u32 = 3;
1243 let mut delivered = false;
1244 for attempt in 1..=FINAL_FLUSH_ATTEMPTS {
1245 // Two bounds, whichever is sooner: one call may not exceed the ingest
1246 // timeout, and the drain as a whole may not exceed its deadline. A
1247 // timed-out call is a failed attempt — never a silent success.
1248 let attempt_deadline = deadline.min(tokio::time::Instant::now() + config.ingest_timeout);
1249 let now = clock.now();
1250 match tokio::time::timeout_at(attempt_deadline, ingest_checked(&**sink, batch, now)).await {
1251 Ok(Ok(report)) => {
1252 counters.record_ingest(&report, now);
1253 delivered = true;
1254 break;
1255 }
1256 Ok(Err(_)) if attempt < FINAL_FLUSH_ATTEMPTS => {
1257 // The backoff sleeps into whatever budget is left, never past
1258 // it. Sleeping the full `retry_backoff` here overran the
1259 // deadline by up to `2 * retry_backoff`, because the comment
1260 // claiming otherwise was attached to the *timeout* arm while
1261 // this one — an ordinary store error, and the common case —
1262 // slept unconditionally (GL-63).
1263 //
1264 // That overrun is not merely a slow shutdown. The drain's
1265 // budget is sized inside `expiry_safety_margin + reclaim_grace`
1266 // (GL-12), so overrunning it releases leases past the window
1267 // that keeps a straggler billable: events that do land arrive
1268 // against a lease the allocator has re-granted, and are
1269 // refused. Bounding by attempt *count* alone is not a bound
1270 // (GL-18).
1271 tokio::time::sleep_until(
1272 deadline.min(tokio::time::Instant::now() + config.retry_backoff),
1273 )
1274 .await;
1275 if tokio::time::Instant::now() >= deadline {
1276 break;
1277 }
1278 }
1279 // Once the deadline has passed there is no budget left to back
1280 // off into, and none to make another attempt with.
1281 Err(_) => break,
1282 Ok(Err(_)) => {}
1283 }
1284 }
1285 if !delivered {
1286 counters.record_lost(batch.len() as u64);
1287 }
1288 // Delivered or lost, the batch has been reported either way.
1289 writer.account_for(batch);
1290 batch.clear();
1291}
1292
1293async fn ingest_checked(
1294 sink: &dyn UsageSink,
1295 events: &[UsageEvent],
1296 now: Timestamp,
1297) -> Result<tollgate_store::IngestReport, tollgate_store::IngestError> {
1298 let report = sink.ingest(events, now).await?;
1299 report
1300 .validate(events.len())
1301 .map_err(tollgate_store::IngestError::Unavailable)?;
1302 Ok(report)
1303}
1304
1305/// Slots still held at the drain deadline. The upgrade succeeds exactly
1306/// while some permit keeps the channel alive — which is when there is
1307/// something to report — and the momentary strong sender is dropped
1308/// immediately, so it cannot mask completion.
1309fn outstanding_permits(lanes: &[mpsc::WeakSender<UsageEvent>]) -> u64 {
1310 lanes
1311 .iter()
1312 .filter_map(mpsc::WeakSender::upgrade)
1313 .map(|tx| (tx.max_capacity() - tx.capacity()) as u64)
1314 .sum()
1315}
1316
1317#[cfg(test)]
1318mod layout_tests {
1319 #[test]
1320 fn attribution_and_existing_outcome_counters_saturate_visibly() {
1321 use super::*;
1322 use tracing_subscriber::layer::SubscriberExt;
1323 #[derive(Clone)]
1324 struct OverflowEvents(Arc<AtomicU64>);
1325 impl<S: tracing::Subscriber> tracing_subscriber::Layer<S> for OverflowEvents {
1326 fn on_event(
1327 &self,
1328 event: &tracing::Event<'_>,
1329 _: tracing_subscriber::layer::Context<'_, S>,
1330 ) {
1331 struct Fields(bool);
1332 impl tracing::field::Visit for Fields {
1333 fn record_debug(&mut self, _: &tracing::field::Field, _: &dyn std::fmt::Debug) {
1334 }
1335 fn record_bool(&mut self, field: &tracing::field::Field, value: bool) {
1336 if field.name() == "counter_overflow" {
1337 self.0 = value;
1338 }
1339 }
1340 }
1341 let mut fields = Fields(false);
1342 event.record(&mut fields);
1343 if fields.0 {
1344 self.0.fetch_add(1, Ordering::Relaxed);
1345 }
1346 }
1347 }
1348 let events = Arc::new(AtomicU64::new(0));
1349 // This unit-test binary has one subscriber. Overflow is reached only
1350 // here; global installation also makes tracing's callsite cache stable.
1351 tracing::subscriber::set_global_default(
1352 tracing_subscriber::registry().with(OverflowEvents(events.clone())),
1353 )
1354 .unwrap();
1355 let counters = WriterCounters::new();
1356 for counter in [
1357 &counters.accepted,
1358 &counters.duplicate,
1359 &counters.rejected,
1360 &counters.lost,
1361 &counters.unattributed,
1362 &counters.attribution_unreported_batches,
1363 ] {
1364 counter.store(u64::MAX - 1, Ordering::Relaxed);
1365 counters.add_outcome(counter, 1);
1366 assert_eq!(counter.load(Ordering::Relaxed), u64::MAX);
1367 counters.add_outcome(counter, 1);
1368 assert_eq!(counter.load(Ordering::Relaxed), u64::MAX);
1369 }
1370 assert!(counters.stats().counter_overflow);
1371 assert_eq!(events.load(Ordering::Relaxed), 1);
1372 let counters = WriterCounters::new();
1373 counters.record_ingest(
1374 &tollgate_store::IngestReport {
1375 accepted: 4,
1376 duplicate: 2,
1377 rejected: 1,
1378 unattributed: Some(3),
1379 },
1380 Timestamp::UNIX_EPOCH,
1381 );
1382 counters.record_ingest(
1383 &tollgate_store::IngestReport::default(),
1384 Timestamp::UNIX_EPOCH,
1385 );
1386 let stats = counters.stats();
1387 assert_eq!(
1388 (
1389 stats.accepted,
1390 stats.duplicate,
1391 stats.rejected,
1392 stats.unattributed,
1393 stats.attribution_unreported_batches
1394 ),
1395 (4, 2, 1, 3, 1)
1396 );
1397 assert!(!stats.counter_overflow);
1398 }
1399
1400 use super::{Contended, UsageWriter, UsageWriterConfig};
1401 use jiff::Timestamp;
1402 use std::sync::Arc;
1403
1404 #[tokio::test(start_paused = true)]
1405 async fn a_runtime_stop_closes_the_queue_and_bounds_the_actual_drain() {
1406 let (recorder, writer) = UsageWriter::spawn(
1407 tollgate_store::MemoryStore::new(tollgate_store::GrantPolicy::default()).unwrap(),
1408 Arc::new(crate::ManualClock::new(
1409 Timestamp::from_second(100).unwrap(),
1410 )),
1411 UsageWriterConfig {
1412 queue_capacity: 1,
1413 max_batch: 1,
1414 flush_interval: std::time::Duration::from_millis(1),
1415 retry_backoff: std::time::Duration::from_millis(1),
1416 shutdown_drain_deadline: std::time::Duration::from_millis(100),
1417 ingest_timeout: std::time::Duration::from_millis(1),
1418 },
1419 )
1420 .unwrap();
1421 let permit = recorder.try_reserve().unwrap();
1422 let began = tokio::time::Instant::now();
1423 writer.stop_at(began + std::time::Duration::from_millis(5));
1424 tokio::time::timeout(std::time::Duration::from_millis(1), recorder.closed())
1425 .await
1426 .unwrap();
1427 let stats = writer.shutdown().await.unwrap();
1428 assert_eq!(stats.unresolved, 1);
1429 assert_eq!(began.elapsed(), std::time::Duration::from_millis(5));
1430 drop(permit);
1431 }
1432
1433 #[test]
1434 fn request_path_counters_are_isolated_on_supported_cache_lines() {
1435 assert_eq!(align_of::<Contended>(), 128);
1436 assert_eq!(size_of::<Contended>(), 128);
1437 }
1438}
1439
1440#[cfg(test)]
1441mod lane_tests {
1442 use super::*;
1443 use std::sync::Mutex;
1444 use std::time::Duration;
1445 use tollgate_core::{AccountId, CostUnits, PolicyRevision, RequestId, UsageSource};
1446 use tollgate_store::{IngestError, IngestReport};
1447
1448 /// Accepts everything and remembers which request ids it was given.
1449 #[derive(Default)]
1450 struct Recording(Mutex<Vec<u128>>);
1451
1452 #[async_trait::async_trait]
1453 impl UsageSink for Recording {
1454 async fn ingest(
1455 &self,
1456 events: &[UsageEvent],
1457 _now: Timestamp,
1458 ) -> Result<IngestReport, IngestError> {
1459 self.0
1460 .lock()
1461 .unwrap()
1462 .extend(events.iter().map(|event| event.request_id.0));
1463 Ok(IngestReport {
1464 accepted: events.len() as u64,
1465 unattributed: Some(0),
1466 ..IngestReport::default()
1467 })
1468 }
1469 }
1470
1471 fn event(id: u128) -> UsageEvent {
1472 UsageEvent::new(
1473 RequestId(id),
1474 AccountId(1),
1475 UsageSource::Overage,
1476 CostUnits(1),
1477 Timestamp::from_second(100).unwrap(),
1478 PolicyRevision::UNSTATED,
1479 None,
1480 )
1481 }
1482
1483 fn config(queue_capacity: usize, max_batch: usize) -> UsageWriterConfig {
1484 UsageWriterConfig {
1485 queue_capacity,
1486 max_batch,
1487 flush_interval: Duration::from_secs(60),
1488 retry_backoff: Duration::from_millis(10),
1489 shutdown_drain_deadline: Duration::from_secs(5),
1490 ingest_timeout: Duration::from_secs(5),
1491 }
1492 }
1493
1494 fn spawn(
1495 sink: &Arc<Recording>,
1496 queue_capacity: usize,
1497 max_batch: usize,
1498 lanes: usize,
1499 ) -> (UsageRecorder, UsageWriter) {
1500 UsageWriter::spawn_lanes(
1501 Arc::clone(sink) as Arc<dyn UsageSink>,
1502 Arc::new(crate::ManualClock::new(
1503 Timestamp::from_second(100).unwrap(),
1504 )),
1505 config(queue_capacity, max_batch),
1506 lanes,
1507 )
1508 }
1509
1510 #[test]
1511 fn lanes_partition_the_capacity_exactly_and_never_go_below_their_floor() {
1512 for (total, count) in [(256, 4), (4_096, 16), (100, 3), (7, 7)] {
1513 let sizes: Vec<usize> = (0..count).map(|i| lane_capacity(total, count, i)).collect();
1514 assert_eq!(sizes.iter().sum::<usize>(), total, "{total}/{count}");
1515 assert!(sizes.iter().max().unwrap() - sizes.iter().min().unwrap() <= 1);
1516 }
1517 assert_eq!(lane_count(1), 1);
1518 assert_eq!(
1519 lane_count(MIN_LANE_CAPACITY * 2 - 1),
1520 1,
1521 "one lane below two floors"
1522 );
1523 for capacity in [128, 4_096, 65_536] {
1524 let count = lane_count(capacity);
1525 assert!(count.is_power_of_two());
1526 assert!(
1527 capacity / count >= MIN_LANE_CAPACITY,
1528 "{capacity} -> {count}"
1529 );
1530 }
1531 }
1532
1533 /// A full lane is not a full queue: one thread, whose own lane fills first,
1534 /// still reserves every slot of every lane, and the next request sheds at
1535 /// exactly `queue_capacity` (INVARIANTS.md GL-8).
1536 #[tokio::test(start_paused = true)]
1537 async fn the_shed_point_is_the_whole_queue_across_lanes() {
1538 let sink = Arc::new(Recording::default());
1539 let (recorder, writer) = spawn(&sink, 256, 64, 4);
1540 let permits: Vec<_> = (0..256).map(|_| recorder.try_reserve().unwrap()).collect();
1541 assert_eq!(recorder.health().queue_depth, 256);
1542 assert_eq!(recorder.health().queue_capacity, 256);
1543 assert_eq!(
1544 recorder.try_reserve().err(),
1545 Some(DenyReason::AccountingBackpressure)
1546 );
1547 assert_eq!(recorder.health().shed, 1);
1548 drop(permits);
1549 assert_eq!(recorder.health().queue_depth, 0);
1550 assert!(writer.shutdown().await.unwrap().unresolved == 0);
1551 }
1552
1553 /// Events in every lane are delivered by the drain, and permits still held
1554 /// in several lanes at the deadline are all reported unresolved.
1555 #[tokio::test(start_paused = true)]
1556 async fn the_drain_delivers_every_lane_and_reports_every_lanes_permits() {
1557 let sink = Arc::new(Recording::default());
1558 let (recorder, writer) = spawn(&sink, 256, 256, 4);
1559 // 200 records from one thread fill its lane and spill into the rest.
1560 for id in 0..200 {
1561 recorder.try_reserve().unwrap().record(event(id));
1562 }
1563 assert_eq!(recorder.health().unaccounted, 200);
1564 let held: Vec<_> = (0..40).map(|_| recorder.try_reserve().unwrap()).collect();
1565 let stats = writer.shutdown().await.unwrap();
1566 let mut delivered = sink.0.lock().unwrap().clone();
1567 delivered.sort_unstable();
1568 assert_eq!(
1569 delivered,
1570 (0..200).collect::<Vec<_>>(),
1571 "every lane drained"
1572 );
1573 assert_eq!(stats.accepted, 200);
1574 assert_eq!(
1575 stats.unresolved, 40,
1576 "permits held across lanes are all reported"
1577 );
1578 assert_eq!(recorder.health().unaccounted, 0);
1579 drop(held);
1580 }
1581
1582 /// A permit that resolves during the drain wakes it: the drain returns as
1583 /// soon as the last one does, not at its deadline.
1584 #[tokio::test(start_paused = true)]
1585 async fn a_resolving_permit_in_any_lane_completes_the_drain() {
1586 let sink = Arc::new(Recording::default());
1587 let (recorder, writer) = spawn(&sink, 256, 64, 4);
1588 let held: Vec<_> = (0..100).map(|_| recorder.try_reserve().unwrap()).collect();
1589 let began = tokio::time::Instant::now();
1590 let release = tokio::spawn(async move {
1591 tokio::time::sleep(Duration::from_millis(50)).await;
1592 for (id, permit) in (0..).zip(held) {
1593 permit.record(event(id));
1594 }
1595 });
1596 let stats = writer.shutdown().await.unwrap();
1597 release.await.unwrap();
1598 assert_eq!(stats.accepted, 100);
1599 assert_eq!(stats.unresolved, 0);
1600 assert_eq!(
1601 began.elapsed(),
1602 Duration::from_millis(50),
1603 "woken, not timed out"
1604 );
1605 }
1606
1607 /// A lane that reaches its ring point is delivered before the flush tick,
1608 /// so a burst cannot back up to the shed point waiting for it.
1609 #[tokio::test(start_paused = true)]
1610 async fn a_full_batch_is_delivered_before_the_tick() {
1611 let sink = Arc::new(Recording::default());
1612 let (recorder, writer) = spawn(&sink, 256, 16, 4);
1613 for id in 0..16 {
1614 recorder.try_reserve().unwrap().record(event(id));
1615 }
1616 for _ in 0..10 {
1617 tokio::task::yield_now().await;
1618 }
1619 assert_eq!(sink.0.lock().unwrap().len(), 16, "delivered without a tick");
1620 assert_eq!(writer.shutdown().await.unwrap().accepted, 16);
1621 }
1622
1623 /// Below the ring point nothing wakes the writer; the tick delivers.
1624 #[tokio::test(start_paused = true)]
1625 async fn a_partial_batch_waits_for_the_tick() {
1626 let sink = Arc::new(Recording::default());
1627 let (recorder, writer) = spawn(&sink, 256, 64, 4);
1628 for id in 0..3 {
1629 recorder.try_reserve().unwrap().record(event(id));
1630 }
1631 for _ in 0..10 {
1632 tokio::task::yield_now().await;
1633 }
1634 assert!(sink.0.lock().unwrap().is_empty(), "no wake per event");
1635 tokio::time::sleep(Duration::from_secs(61)).await;
1636 assert_eq!(sink.0.lock().unwrap().len(), 3, "the tick delivered it");
1637 assert_eq!(writer.shutdown().await.unwrap().accepted, 3);
1638 }
1639
1640 /// Dropping every recorder handle stops the writer without a tick: the
1641 /// lanes ring as their last handles go.
1642 #[tokio::test(start_paused = true)]
1643 async fn dropping_the_recorder_stops_the_writer_promptly() {
1644 let sink = Arc::new(Recording::default());
1645 let (recorder, mut writer) = spawn(&sink, 256, 64, 4);
1646 recorder.try_reserve().unwrap().record(event(1));
1647 drop(recorder);
1648 let handle = writer.handle.take().unwrap();
1649 let stats = tokio::time::timeout(Duration::from_millis(1), handle)
1650 .await
1651 .expect("stopped before any tick")
1652 .unwrap();
1653 assert_eq!(stats.accepted, 1);
1654 }
1655}