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