Skip to main content

tollgate_core/
usage.rs

1//! Usage events: the billing record.
2//!
3//! Leases *bound* spend; usage events *are* what gets billed. Reconciliation
4//! compares the two ledgers and steady-state drift is zero (INVARIANTS.md,
5//! ledger-roles note). Events are idempotent on `request_id`, so batched
6//! writers may retry whole batches freely (INVARIANTS.md GL-7).
7
8use std::sync::{
9    Arc,
10    atomic::{AtomicU64, Ordering},
11};
12
13use jiff::Timestamp;
14
15use crate::ids::{AccountId, FencingToken, KeyId, LeaseId, PolicyRevision, RequestId};
16use crate::units::CostUnits;
17
18/// Pre-reserved capacity for exactly one usage event.
19///
20/// Implementations must record without fallible I/O: obtaining a slot is the
21/// backpressure decision, while consuming it is the committed-charge path.
22pub trait UsageSlot: Send + 'static {
23    /// Record the charge into the capacity this slot reserved. Called from the
24    /// committed guard's drop, so it must not block, allocate fallibly or
25    /// fail.
26    fn record(self, event: UsageEvent);
27}
28
29/// A [`UsageSlot`] that discards the event and counts it.
30///
31/// The reference implementation for an integration that has not reached usage
32/// export yet, in the same spirit as `NoGate` for `CapacityGate`: it lets an
33/// embedder adopt admission one stage at a time instead of taking the client
34/// runtime and a background writer purely to satisfy a type.
35///
36/// It counts what it discarded, deliberately. A silently dropping slot makes
37/// "usage is not wired up" indistinguishable from "no usage happened", and
38/// those are very different operational stories — one of them means an
39/// integration is admitting billable work and losing the record of it.
40///
41/// Not for production billing. Discarded events are not recoverable; use the
42/// batching writer once usage must be durable.
43///
44/// Deliberately no `Default`. A zero-valued counter would make
45/// `Default::default()` and `new()` the same value, leaving a mutant no test
46/// can distinguish; an equivalent mutant is a design smell rather than a
47/// coverage gap. `new` returning `Arc<Self>` also keeps clippy's
48/// `new_without_default` inapplicable, so the two gates agree rather than
49/// pulling opposite ways.
50///
51/// ```
52/// # use tollgate_core::DiscardedUsage;
53/// let discarded = DiscardedUsage::new();
54/// // hand `discarded.slot()` to `admit`; assert on the count later
55/// assert_eq!(discarded.count(), 0);
56/// ```
57#[derive(Debug)]
58pub struct DiscardedUsage {
59    count: AtomicU64,
60}
61
62impl DiscardedUsage {
63    /// A shared counter, ready to hand out slots.
64    ///
65    /// Returns `Arc<Self>` because that is the only useful form: [`Self::slot`]
66    /// needs one, and a bare `DiscardedUsage` can do nothing. Making the
67    /// reachable shape the only constructible one removes a step an embedder
68    /// would otherwise have to know to take.
69    #[must_use]
70    pub fn new() -> Arc<Self> {
71        Arc::new(Self {
72            count: AtomicU64::new(0),
73        })
74    }
75
76    /// A slot that discards one event into this counter.
77    #[must_use]
78    pub fn slot(self: &Arc<Self>) -> DiscardedUsageSlot {
79        DiscardedUsageSlot(Arc::clone(self))
80    }
81
82    /// How many events have been discarded.
83    #[must_use]
84    pub fn count(&self) -> u64 {
85        self.count.load(Ordering::Relaxed)
86    }
87}
88
89/// One pre-reserved discard, obtained from [`DiscardedUsage::slot`].
90#[derive(Debug)]
91pub struct DiscardedUsageSlot(Arc<DiscardedUsage>);
92
93impl UsageSlot for DiscardedUsageSlot {
94    fn record(self, _event: UsageEvent) {
95        self.0.count.fetch_add(1, Ordering::Relaxed);
96    }
97}
98
99/// What funded the units in a [`UsageEvent`].
100///
101/// Two variants, because two things fund spend and they are validated by
102/// opposite rules. A leased charge arrives with a capability the sink checks
103/// before it touches the ledger; an overage charge has no capability to check,
104/// because no lease existed to issue one.
105///
106/// Making this a sum type rather than a pair of `Option`s is deliberate. A
107/// half-formed capability — a lease id with no token, or a token naming no
108/// lease — is the shape a sink would have to defend against, and here it
109/// cannot be written down. The storage schema mirrors the same constraint.
110#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
112pub enum UsageSource {
113    /// Spent from a lease. The sink requires the stored
114    /// `(lease_id, account_id, fencing_token)` triple to match before the
115    /// event may change either ledger; token age relative to another active
116    /// lease is irrelevant (INVARIANTS.md GL-4).
117    Leased {
118        /// The lease the units were debited from.
119        lease_id: LeaseId,
120        /// That lease's capability token.
121        fencing_token: FencingToken,
122    },
123    /// Admitted under [`EnforcementMode::Elastic`] with no lease behind it.
124    ///
125    /// **It carries no lease id on purpose, and must never be given one.**
126    /// Attributing unfunded spend to a real lease drives that lease's `used`
127    /// past its `granted`, which makes reclaim's `granted - used` credit
128    /// negative — a panic in the memory backend and a permanently failing
129    /// transaction in Postgres, so the account's expired leases would never be
130    /// reclaimed again. Overage stands outside lease accounting entirely and
131    /// is funded by its own ledger term.
132    ///
133    /// [`EnforcementMode::Elastic`]: crate::snapshot::EnforcementMode::Elastic
134    Overage,
135}
136
137impl UsageSource {
138    /// The lease this charge was spent from, or `None` for overage.
139    #[must_use]
140    pub const fn lease_id(self) -> Option<LeaseId> {
141        match self {
142            UsageSource::Leased { lease_id, .. } => Some(lease_id),
143            UsageSource::Overage => None,
144        }
145    }
146
147    /// The referenced lease's capability token, or `None` for overage.
148    #[must_use]
149    pub const fn fencing_token(self) -> Option<FencingToken> {
150        match self {
151            UsageSource::Leased { fencing_token, .. } => Some(fencing_token),
152            UsageSource::Overage => None,
153        }
154    }
155}
156
157/// One committed charge. Produced only from a committed
158/// [`crate::reservation::Reservation`]; there is deliberately no public
159/// constructor path for uncommitted work.
160///
161/// `#[non_exhaustive]` is what makes that sentence true outside this crate
162/// rather than merely stated. The type had said "produced only from a
163/// committed reservation" while remaining a plain struct literal any crate
164/// could fill in, which is a convention rather than a boundary; a caller could
165/// assemble an event for work that never committed and hand it to a sink.
166/// Construction now goes through [`UsageEvent::new`], and the sealing has a
167/// second benefit the workspace pays for once: a field added here no longer
168/// breaks every literal in every downstream test.
169#[derive(Debug, Clone, Copy, PartialEq, Eq)]
170#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
171#[non_exhaustive]
172pub struct UsageEvent {
173    /// Idempotency key: a sink must treat a replayed `request_id` as the same
174    /// charge, not a new one.
175    ///
176    /// The only field of this struct that is independent of how the units were
177    /// funded, which is what lets overage replay under the same rule as any
178    /// other event (INVARIANTS.md GL-7).
179    pub request_id: RequestId,
180    /// The account billed.
181    pub account_id: AccountId,
182    /// What funded the units, and the evidence the sink validates.
183    pub source: UsageSource,
184    /// The units charged: the quote reserved at admission.
185    pub units: CostUnits,
186    /// When the charge was committed, at execution start.
187    pub occurred_at: Timestamp,
188    /// The consuming application's policy identity, copied from the pinned
189    /// snapshot that priced this request (GL-94).
190    ///
191    /// Carried so a billing record can be traced to the exact product policy
192    /// that produced it. Tollgate never reads it, and an event from a
193    /// publisher that stated no revision carries
194    /// [`PolicyRevision::UNSTATED`].
195    ///
196    /// Defaults on the wire, so an ingest from a peer that predates the field
197    /// decodes rather than failing — the same rule the snapshot's optional
198    /// fields follow.
199    #[cfg_attr(feature = "serde", serde(default))]
200    pub policy_revision: PolicyRevision,
201    /// Credential named by the pinned, key-scoped snapshot. An absent value
202    /// preserves billing but supplies no credential activity. This metadata
203    /// is never authorization evidence; the sink checks account ownership.
204    #[cfg_attr(feature = "serde", serde(default))]
205    pub key_id: Option<KeyId>,
206}
207
208impl UsageEvent {
209    /// Build a committed charge.
210    ///
211    /// Called by [`Reservation::usage_event`], which is the only place that
212    /// can prove the charge committed. It is public because tests, benchmarks,
213    /// and store backends across the workspace need to construct events
214    /// without a live reservation; what `#[non_exhaustive]` buys is that every
215    /// such construction goes through one signature, so a new field reaches
216    /// them as a compile error at one call each rather than as a silent
217    /// default.
218    ///
219    /// [`Reservation::usage_event`]: crate::reservation::Reservation::usage_event
220    #[must_use]
221    pub const fn new(
222        request_id: RequestId,
223        account_id: AccountId,
224        source: UsageSource,
225        units: CostUnits,
226        occurred_at: Timestamp,
227        policy_revision: PolicyRevision,
228        key_id: Option<KeyId>,
229    ) -> Self {
230        Self {
231            request_id,
232            account_id,
233            source,
234            units,
235            occurred_at,
236            policy_revision,
237            key_id,
238        }
239    }
240}
241
242#[cfg(all(test, feature = "serde"))]
243mod revision_wire_tests {
244    use super::*;
245
246    fn event(revision: PolicyRevision) -> UsageEvent {
247        UsageEvent::new(
248            RequestId(9),
249            AccountId(1),
250            UsageSource::Overage,
251            CostUnits(70),
252            Timestamp::UNIX_EPOCH,
253            revision,
254            None,
255        )
256    }
257
258    /// A peer that predates GL-94 sends no revision, and its events must ingest
259    /// rather than fail. The absent key decodes to "unstated" — a value, not
260    /// an error — which is what makes the additive rollout safe in both
261    /// directions.
262    #[test]
263    fn an_event_without_a_revision_key_decodes_as_unstated() {
264        let mut value =
265            serde_json::to_value(event(PolicyRevision([0xc3; 32]))).expect("an event serializes");
266        assert_eq!(
267            value["policy_revision"],
268            serde_json::Value::String("c3".repeat(32)),
269            "a stated revision is on the wire in canonical form"
270        );
271        assert!(
272            value
273                .as_object_mut()
274                .expect("an event is a JSON object")
275                .remove("policy_revision")
276                .is_some()
277        );
278
279        let decoded: UsageEvent =
280            serde_json::from_value(value).expect("an older payload still decodes");
281        assert_eq!(decoded.policy_revision, PolicyRevision::UNSTATED);
282    }
283
284    /// The revision survives the wire byte for byte. It is compared for
285    /// equality by the consumer to select its own metadata, so a value that
286    /// round-tripped to something else would silently name a different policy.
287    #[test]
288    fn a_revision_survives_the_event_round_trip_exactly() {
289        let mut bytes = [0u8; 32];
290        for (index, byte) in bytes.iter_mut().enumerate() {
291            *byte = (index as u8).wrapping_mul(7).wrapping_add(1);
292        }
293        let original = event(PolicyRevision(bytes));
294        let encoded = serde_json::to_string(&original).expect("an event serializes");
295        let decoded: UsageEvent = serde_json::from_str(&encoded).expect("it decodes");
296        assert_eq!(decoded, original);
297        assert_eq!(decoded.policy_revision.as_bytes(), &bytes);
298    }
299}
300
301#[cfg(test)]
302mod discarded_usage_tests {
303    use super::*;
304
305    fn event(request: u128) -> UsageEvent {
306        UsageEvent::new(
307            RequestId(request),
308            AccountId(1),
309            UsageSource::Overage,
310            CostUnits(70),
311            Timestamp::UNIX_EPOCH,
312            PolicyRevision::UNSTATED,
313            None,
314        )
315    }
316
317    #[test]
318    fn discarding_is_counted_not_silent() {
319        let discarded = DiscardedUsage::new();
320        assert_eq!(discarded.count(), 0);
321
322        discarded.slot().record(event(1));
323        discarded.slot().record(event(2));
324
325        // The count is the whole point: without it, an integration that is
326        // admitting billable work and losing every record of it looks exactly
327        // like one that has served nothing.
328        assert_eq!(discarded.count(), 2);
329    }
330
331    #[test]
332    fn slots_are_independent_and_share_one_counter() {
333        let discarded = DiscardedUsage::new();
334        let first = discarded.slot();
335        let second = discarded.slot();
336
337        // A slot is pre-reserved capacity for exactly one event, so holding
338        // two and consuming them in either order must total two.
339        second.record(event(2));
340        first.record(event(1));
341
342        assert_eq!(discarded.count(), 2);
343    }
344
345    #[test]
346    fn counts_across_threads() {
347        let discarded = DiscardedUsage::new();
348        std::thread::scope(|scope| {
349            for index in 0..8u128 {
350                let discarded = Arc::clone(&discarded);
351                scope.spawn(move || discarded.slot().record(event(index)));
352            }
353        });
354        assert_eq!(discarded.count(), 8);
355    }
356}