Skip to main content

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}