Skip to main content

tollgate_store_postgres/
lib.rs

1//! PostgreSQL backend.
2//!
3//! Reproduces `MemoryStore`'s settlement rules exactly — the shared
4//! correctness suite (`tests/postgres_suite.rs`, mirroring
5//! `tollgate-store/tests/store_suite.rs` by name) is the proof. Concurrency
6//! control is row-level: `SELECT ... FOR UPDATE` on the account funds
7//! acquire; on the lease it serializes release/ingest/reclaim against each
8//! other. Fencing tokens come from the account row's `next_fence` counter,
9//! so allocation is strictly monotonic per account across any number of
10//! servers sharing the database. Each token remains a capability for its own
11//! lease record; the sequence is not an account-wide validity epoch.
12//!
13//! Representation choices (PoC-pragmatic, documented):
14//! - u128 ids as 16-byte `BYTEA` (big-endian);
15//! - units as `BIGINT` with checked u64↔i64 conversion in both directions
16//!   (a balance beyond i64::MAX is refused, not wrapped; a negative stored
17//!   value is refused, not clamped) and schema-level CHECK constraints
18//!   keeping every unit column non-negative;
19//! - lease and credential expiry as floor `BIGINT` microseconds plus a
20//!   `SMALLINT` nanosecond remainder; informational timestamps as `BIGINT`
21//!   microseconds since the Unix epoch;
22//! - snapshots as storage-local `JSONB`: ids in the legacy u64 range remain
23//!   numeric for rollback, larger ids use canonical text, and the public
24//!   HTTP/Serde contract always uses text.
25//!
26//! Snapshot pushes broadcast in-process only; cross-process push
27//! (LISTEN/NOTIFY or the server's future SSE) is a documented seam in
28//! `docs/DESIGN.md`.
29
30#![deny(missing_docs)]
31
32pub(crate) mod instant;
33
34#[cfg(feature = "test-support")]
35pub mod test_support;
36
37use instant::StoredInstant;
38use std::num::{NonZeroI64, NonZeroUsize};
39use std::sync::Arc;
40
41use async_trait::async_trait;
42use jiff::{SignedDuration, Timestamp};
43use serde::{Deserialize, Serialize};
44use sqlx::postgres::PgPoolOptions;
45use sqlx::{PgPool, Postgres, Row, Transaction};
46use tokio::sync::broadcast;
47
48use tollgate_core::{
49    AccountId, AccountSnapshot, AccountStatus, BudgetSchedule, BudgetView, CapacityClass,
50    CostTable, CostUnits, EnforcementMode, FencingToken, Generation, KeyId, LeaseGrant, LeaseId,
51    Period, PermissionBits, PolicyRevision, Principal, PublishableSnapshot, ResolvedLimits,
52    Rollover, UsageEvent,
53};
54use tollgate_store::{
55    AccountConfig, AccountView, AdminReceipt, AdminState, AdminStore, AllocateError, Allocation,
56    BudgetError, Conservation, CreateAccountError, GrantPolicy, IngestError, IngestReport,
57    KeyDirectory, KeyError, KeyRecord, KeySnapshotError, KeySummary, LeaseAllocator,
58    PUSH_CHANNEL_CAPACITY, PublishSnapshotError, ReclaimBatch, ReclaimedLease, Revocation,
59    RolledAccount, RolloverBatch, SetStatusError, SnapshotPush, SnapshotResolution, SnapshotSource,
60    StatusChange, StoreError, StoreHealth, UsageSink, pushes_exceed_capacity,
61    validate_key_page_limit,
62};
63
64const STATE_ACTIVE: i16 = 0;
65const STATE_RELEASED: i16 = 1;
66const STATE_EXPIRED: i16 = 2;
67
68/// One storage-local identifier. Values the previous codec could represent
69/// stay numeric for rollback; the rest of the promised u128 domain uses
70/// canonical text because serde_json's default number type rejects it.
71#[derive(Debug, Clone, Copy)]
72struct StoredId(u128);
73
74impl Serialize for StoredId {
75    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
76    where
77        S: serde::Serializer,
78    {
79        match u64::try_from(self.0) {
80            Ok(value) => serializer.serialize_u64(value),
81            Err(_) => serializer.collect_str(&format_args!("{:032x}", self.0)),
82        }
83    }
84}
85
86impl<'de> Deserialize<'de> for StoredId {
87    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
88    where
89        D: serde::Deserializer<'de>,
90    {
91        struct StoredIdVisitor;
92
93        impl serde::de::Visitor<'_> for StoredIdVisitor {
94            type Value = StoredId;
95
96            fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
97                formatter.write_str("a legacy u64 number or a canonical 128-bit identifier string")
98            }
99
100            fn visit_u64<E>(self, value: u64) -> Result<Self::Value, E> {
101                Ok(StoredId(u128::from(value)))
102            }
103
104            fn visit_str<E>(self, value: &str) -> Result<Self::Value, E>
105            where
106                E: serde::de::Error,
107            {
108                value
109                    .parse::<AccountId>()
110                    .map(|id| StoredId(id.0))
111                    .map_err(E::custom)
112            }
113        }
114
115        deserializer.deserialize_any(StoredIdVisitor)
116    }
117}
118
119/// A digest column that is not exactly 32 bytes is corruption, never
120/// something to pad or truncate into shape: the comparison it feeds decides
121/// authentication, and a silently reshaped digest would either refuse a
122/// legitimate credential forever or, worse, shorten what is compared.
123fn digest_from(bytes: &[u8], key_id: KeyId) -> Result<[u8; 32], StoreError> {
124    <[u8; 32]>::try_from(bytes).map_err(|_| {
125        StoreError(format!(
126            "credential {key_id} has a {}-byte digest in storage, expected 32",
127            bytes.len()
128        ))
129    })
130}
131
132#[async_trait]
133impl KeyDirectory for PostgresStore {
134    async fn credential_activity(
135        &self,
136        keys: &[KeyId],
137    ) -> Result<Vec<tollgate_store::CredentialActivity>, StoreError> {
138        use tollgate_store::{CredentialActivity, CredentialActivityState};
139        let mut result = Vec::with_capacity(keys.len());
140        // Reuse the bulk-work budget, not a total operator-list cap. Each
141        // input has one row, including duplicates and missing credentials.
142        for chunk in keys.chunks(tollgate_store::MAX_INGEST_BATCH) {
143            let ids: Vec<_> = chunk.iter().map(|id| id_bytes(id.0)).collect();
144            let rows = sqlx::query(
145                // Each lateral lookup is bounded by the credential PK. The
146                // explicit one-row bound prevents a bulk join from choosing
147                // a catalogue-wide hash scan for a moderately sized directory.
148                "SELECT evidence.key_id IS NOT NULL, evidence.last_committed_at_us
149                 FROM UNNEST($1::bytea[]) WITH ORDINALITY AS requested(key_id, ordinal)
150                 LEFT JOIN LATERAL (
151                    SELECT k.key_id, a.last_committed_at_us
152                    FROM tollgate_credential_keys k
153                    LEFT JOIN tollgate_credential_activity a ON a.key_id = k.key_id
154                    WHERE k.key_id = requested.key_id LIMIT 1
155                 ) AS evidence ON true
156                 ORDER BY requested.ordinal",
157            )
158            .bind(&ids)
159            .fetch_all(&self.pool)
160            .await
161            .map_err(storage)?;
162            if rows.len() != chunk.len() {
163                return Err(StoreError("incomplete credential activity read".into()));
164            }
165            for (&key_id, row) in chunk.iter().zip(rows) {
166                let state = if !row.get::<bool, _>(0) {
167                    CredentialActivityState::Unknown
168                } else if let Some(at) = row.get::<Option<i64>, _>(1) {
169                    CredentialActivityState::Committed {
170                        last_committed_at: micros_ts(at, "credential activity")?,
171                    }
172                } else {
173                    CredentialActivityState::Unobserved
174                };
175                result.push(CredentialActivity { key_id, state });
176            }
177        }
178        Ok(result)
179    }
180
181    async fn insert_key(&self, record: KeyRecord) -> Result<(), KeyError> {
182        self.insert_credential(record, None)
183            .await
184            .map(|receipt| receipt.outcome)
185    }
186
187    async fn revoke_key(&self, key_id: KeyId, now: Timestamp) -> Result<Revocation, KeyError> {
188        self.revoke_key_audited(key_id, now)
189            .await
190            .map(|receipt| receipt.outcome)
191    }
192
193    async fn revoke_key_audited(
194        &self,
195        key_id: KeyId,
196        now: Timestamp,
197    ) -> Result<AdminReceipt<Revocation>, KeyError> {
198        let mut tx = self.pool.begin().await.map_err(storage)?;
199        let result = async {
200            // The row lock captures the actual predecessor, including a prior
201            // retirement committed while this call waited. No account lock is
202            // acquired after the credential lock.
203            let row = sqlx::query(
204                "SELECT account_id, revoked_at_us IS NOT NULL
205                 FROM tollgate_credential_keys WHERE key_id = $1 FOR UPDATE",
206            )
207            .bind(id_bytes(key_id.0))
208            .fetch_optional(&mut *tx)
209            .await
210            .map_err(storage)?
211            .ok_or(KeyError::UnknownKey)?;
212            let account_bytes: Vec<u8> = row.get(0);
213            let account_bytes: [u8; 16] = account_bytes
214                .try_into()
215                .map_err(|_| StoreError("credential account identifier is not 16 bytes".into()))?;
216            let account_id = AccountId(u128::from_be_bytes(account_bytes));
217            let revoked: bool = row.get(1);
218            let before = AdminState::Credential {
219                account_id,
220                key_id,
221                revoked,
222            };
223            let outcome = if revoked {
224                Revocation::AlreadyRetired
225            } else {
226                sqlx::query(
227                    "UPDATE tollgate_credential_keys SET revoked_at_us = $2 WHERE key_id = $1",
228                )
229                .bind(id_bytes(key_id.0))
230                .bind(ts_micros(now))
231                .execute(&mut *tx)
232                .await
233                .map_err(storage)?;
234                Revocation::Retired
235            };
236            Ok(AdminReceipt::new(
237                outcome,
238                before,
239                AdminState::Credential {
240                    account_id,
241                    key_id,
242                    revoked: true,
243                },
244            ))
245        }
246        .await;
247        finish_transaction(tx, result).await
248    }
249
250    async fn publish_key_snapshot(
251        &self,
252        account: AccountId,
253        key: KeyId,
254        snapshot: PublishableSnapshot,
255    ) -> Result<AdminReceipt<()>, KeySnapshotError> {
256        let generation = i64::try_from(snapshot.generation.0).map_err(|_| {
257            StoreError("snapshot generation exceeds PostgreSQL BIGINT range".into())
258        })?;
259        let mut tx = self.pool.begin().await.map_err(storage)?;
260        let result = async {
261            let (principal, revoked) = lock_account_key(&mut tx, account, key).await?;
262            if revoked {
263                return Err(KeySnapshotError::Retired { key_id: key });
264            }
265            if snapshot.key_id != Some(key) {
266                return Err(PublishSnapshotError::CredentialMismatch { key_id: key }.into());
267            }
268            let published = publish_in_tx(&mut tx, principal, generation, snapshot).await?;
269            Ok((principal, published))
270        }
271        .await;
272        let (principal, (written, published, before, after)) =
273            finish_transaction(tx, result).await?;
274        if written {
275            self.push_to_subscribers(SnapshotPush {
276                principal,
277                resolution: SnapshotResolution::Present(published),
278            });
279        }
280        Ok(AdminReceipt::new((), before, after))
281    }
282
283    async fn remove_key_snapshot(
284        &self,
285        account: AccountId,
286        key: KeyId,
287    ) -> Result<AdminReceipt<()>, KeySnapshotError> {
288        let mut tx = self.pool.begin().await.map_err(storage)?;
289        let result = async {
290            let (principal, _revoked) = lock_account_key(&mut tx, account, key).await?;
291            Ok::<_, KeySnapshotError>((principal, remove_in_tx(&mut tx, principal).await?))
292        }
293        .await;
294        let (principal, receipt) = finish_transaction(tx, result).await?;
295        self.announce_removal(principal, &receipt);
296        Ok(receipt)
297    }
298
299    async fn active_keys(&self, now: Timestamp) -> Result<Vec<KeyRecord>, StoreError> {
300        let cutoff = StoredInstant::from(now);
301        // Expiry is applied here, beside revocation, so this backend answers
302        // "active" exactly as `MemoryStore` does and a projection built from
303        // either sees the same live set. Ordering is explicit for the same
304        // reason: two instances must not build tables that differ by row
305        // order alone.
306        let rows = sqlx::query(
307            "SELECT key_id, account_id, principal, digest, not_after_floor_us, not_after_submicro_ns, not_after_is_lower_bound
308             FROM tollgate_credential_keys
309             WHERE revoked_at_us IS NULL
310               AND (not_after_floor_us IS NULL OR not_after_submicro_ns IS NULL
311                    OR (not_after_floor_us, not_after_submicro_ns) > ($1, $2))
312             ORDER BY key_id",
313        )
314        .bind(cutoff.micros)
315        .bind(cutoff.submicro_nanos)
316        .fetch_all(&self.pool)
317        .await
318        .map_err(storage)?;
319
320        rows.into_iter().map(credential_from_row).collect()
321    }
322
323    async fn account_keys(
324        &self,
325        account: AccountId,
326        after: Option<KeyId>,
327        limit: NonZeroUsize,
328    ) -> Result<Vec<KeySummary>, KeyError> {
329        validate_key_page_limit(limit)?;
330        // The account anchors the result in one statement snapshot: no rows
331        // means no account, while a NULL key marks an existing, empty page.
332        // LATERAL keeps the bounded credential scan on the account/key index;
333        // separate cursor shapes preserve an indexable range in prepared plans.
334        let sql = if after.is_some() {
335            "SELECT page.* FROM tollgate_accounts AS account
336             LEFT JOIN LATERAL (
337                 SELECT key_id, not_after_floor_us, not_after_submicro_ns,
338                        not_after_is_lower_bound, revoked_at_us
339                 FROM tollgate_credential_keys
340                 WHERE account_id = account.account_id AND key_id > $3
341                 ORDER BY key_id LIMIT $2
342             ) AS page ON TRUE
343             WHERE account.account_id = $1 ORDER BY page.key_id"
344        } else {
345            "SELECT page.* FROM tollgate_accounts AS account
346             LEFT JOIN LATERAL (
347                 SELECT key_id, not_after_floor_us, not_after_submicro_ns,
348                        not_after_is_lower_bound, revoked_at_us
349                 FROM tollgate_credential_keys
350                 WHERE account_id = account.account_id
351                 ORDER BY key_id LIMIT $2
352             ) AS page ON TRUE
353             WHERE account.account_id = $1 ORDER BY page.key_id"
354        };
355        let mut query = sqlx::query(sql)
356            .bind(id_bytes(account.0))
357            .bind(i64::try_from(limit.get()).unwrap_or(i64::MAX));
358        if let Some(cursor) = after {
359            query = query.bind(id_bytes(cursor.0));
360        }
361        let rows = query.fetch_all(&self.pool).await.map_err(storage)?;
362        if rows.is_empty() {
363            return Err(KeyError::UnknownAccount);
364        }
365        rows.into_iter()
366            .filter(|row| row.get::<Option<&[u8]>, _>(0).is_some())
367            .map(|row| summary_from_row(row).map_err(KeyError::Storage))
368            .collect()
369    }
370
371    async fn insert_key_within(
372        &self,
373        record: KeyRecord,
374        max_active: NonZeroUsize,
375        now: Timestamp,
376    ) -> Result<(), KeyError> {
377        self.insert_key_within_audited(record, max_active, now)
378            .await
379            .map(|receipt| receipt.outcome)
380    }
381
382    async fn insert_key_within_audited(
383        &self,
384        record: KeyRecord,
385        max_active: NonZeroUsize,
386        now: Timestamp,
387    ) -> Result<AdminReceipt<()>, KeyError> {
388        self.insert_credential(record, Some((max_active, now)))
389            .await
390    }
391}
392
393impl PostgresStore {
394    /// Both issuance APIs enter the same account-first transaction. The lock
395    /// covers the optional bound check, unique-index insertion, foreign-key
396    /// check and revision trigger through commit. Acquiring it after inserting
397    /// would invert the bounded issuer's order and permit a deadlock.
398    async fn insert_credential(
399        &self,
400        record: KeyRecord,
401        bound: Option<(NonZeroUsize, Timestamp)>,
402    ) -> Result<AdminReceipt<()>, KeyError> {
403        let mut tx = self
404            .pool
405            .begin()
406            .await
407            .map_err(|e| KeyError::Storage(storage(e)))?;
408        // Take the account row first. Counting and inserting without it is the
409        // race this method exists to prevent: under READ COMMITTED neither
410        // transaction sees the other's uncommitted credential, so both count
411        // `max_active - 1`, both insert, and the account ends up over the
412        // bound with no error raised anywhere. This is the same row
413        // `set_account_status` and `acquire` serialise on, so an issuance in
414        // flight also orders against a suspension.
415        let account = sqlx::query(ACCOUNT_LOCK_SQL)
416            .bind(id_bytes(record.account_id.0))
417            .fetch_optional(&mut *tx)
418            .await
419            .map_err(|e| KeyError::Storage(storage(e)))?;
420        if account.is_none() {
421            return Err(KeyError::UnknownAccount);
422        }
423
424        if let Some((max_active, now)) = bound {
425            // Identity before the bound, and the order is load-bearing. A caller
426            // that lost the response resends the same `key_id`; by then its own
427            // successful write may have filled the bound, and answering
428            // `ActiveKeyLimit` would tell it to retire a credential when in fact
429            // its first call worked. `AlreadyExists` is both true and what makes
430            // the retry safe (GL-121).
431            let existing: Option<i32> = sqlx::query_scalar(
432                "SELECT 1 FROM tollgate_credential_keys WHERE key_id = $1 OR principal = $2",
433            )
434            .bind(id_bytes(record.key_id.0))
435            .bind(id_bytes(record.principal.0))
436            .fetch_optional(&mut *tx)
437            .await
438            .map_err(|e| KeyError::Storage(storage(e)))?;
439            if existing.is_some() {
440                return Err(KeyError::AlreadyExists);
441            }
442
443            let cutoff = StoredInstant::from(now);
444            let live: i64 = sqlx::query_scalar(LIVE_KEY_COUNT_SQL)
445                .bind(id_bytes(record.account_id.0))
446                .bind(cutoff.micros)
447                .bind(cutoff.submicro_nanos)
448                .fetch_one(&mut *tx)
449                .await
450                .map_err(|e| KeyError::Storage(storage(e)))?;
451            if u128::from(live.max(0).unsigned_abs())
452                >= u128::try_from(max_active.get()).unwrap_or(u128::MAX)
453            {
454                return Err(KeyError::ActiveKeyLimit { limit: max_active });
455            }
456        }
457
458        let expiry = record.not_after.map(StoredInstant::from);
459        let result = sqlx::query(
460            "INSERT INTO tollgate_credential_keys
461             (key_id, account_id, principal, digest, not_after_floor_us,
462              not_after_submicro_ns, not_after_is_lower_bound, revoked_at_us)
463             VALUES ($1, $2, $3, $4, $5, $6, FALSE, NULL)
464             ON CONFLICT (key_id) DO NOTHING",
465        )
466        .bind(id_bytes(record.key_id.0))
467        .bind(id_bytes(record.account_id.0))
468        .bind(id_bytes(record.principal.0))
469        .bind(record.digest.to_vec())
470        .bind(expiry.map(|expiry| expiry.micros))
471        .bind(expiry.map(|expiry| expiry.submicro_nanos))
472        .execute(&mut *tx)
473        .await;
474        let outcome = match result {
475            Ok(done) if done.rows_affected() == 0 => Err(KeyError::AlreadyExists),
476            Ok(_) => Ok(()),
477            Err(sqlx::Error::Database(e)) if e.is_unique_violation() => {
478                Err(KeyError::AlreadyExists)
479            }
480            // No arm for the foreign key, deliberately. The account row was
481            // taken `FOR UPDATE` above and its absence already answered
482            // `UnknownAccount`; nothing in this crate deletes an account, so a
483            // violation here cannot mean "no such account". It would mean the
484            // row vanished under a held lock, and reporting that as a caller
485            // error would absorb storage corruption as a routine refusal --
486            // the caller would retire a credential over a broken database.
487            // `Storage` is the truthful answer, and the mutation gate is what
488            // noticed the old arm could not be reached to be tested.
489            Err(e) => Err(KeyError::Storage(storage(e))),
490        };
491        outcome?;
492        tx.commit()
493            .await
494            .map_err(|e| KeyError::Storage(storage(e)))?;
495        Ok(AdminReceipt::new(
496            (),
497            AdminState::Absent,
498            AdminState::Credential {
499                account_id: record.account_id,
500                key_id: record.key_id,
501                revoked: false,
502            },
503        ))
504    }
505}
506
507fn summary_from_row(row: sqlx::postgres::PgRow) -> Result<KeySummary, StoreError> {
508    let bytes: Vec<u8> = row.get(0);
509    let fixed: [u8; 16] = bytes
510        .as_slice()
511        .try_into()
512        .map_err(|_| StoreError("credential identifier is not 16 bytes".into()))?;
513    let not_after = match (row.get::<Option<i64>, _>(1), row.get::<Option<i16>, _>(2)) {
514        (None, None) if !row.get::<bool, _>(3) => None,
515        (Some(micros), Some(submicro_nanos)) => Some(
516            StoredInstant {
517                micros,
518                submicro_nanos,
519            }
520            .timestamp()?,
521        ),
522        _ => return Err(StoreError("incomplete stored credential expiry".into())),
523    };
524    Ok(KeySummary {
525        key_id: KeyId(u128::from_be_bytes(fixed)),
526        not_after,
527        revoked_at: row
528            .get::<Option<i64>, _>(4)
529            .map(|micros| {
530                StoredInstant {
531                    micros,
532                    submicro_nanos: 0,
533                }
534                .timestamp()
535            })
536            .transpose()?,
537    })
538}
539
540fn credential_from_row(row: sqlx::postgres::PgRow) -> Result<KeyRecord, StoreError> {
541    let id = |index| -> Result<u128, StoreError> {
542        let bytes: Vec<u8> = row.get(index);
543        let fixed: [u8; 16] = bytes
544            .as_slice()
545            .try_into()
546            .map_err(|_| StoreError("credential identifier is not 16 bytes".into()))?;
547        Ok(u128::from_be_bytes(fixed))
548    };
549    let key_id = KeyId(id(0)?);
550    let not_after = match (row.get::<Option<i64>, _>(4), row.get::<Option<i16>, _>(5)) {
551        (None, None) if !row.get::<bool, _>(6) => None,
552        (Some(micros), Some(submicro_nanos)) => Some(
553            StoredInstant {
554                micros,
555                submicro_nanos,
556            }
557            .timestamp()?,
558        ),
559        _ => return Err(StoreError("incomplete stored credential expiry".into())),
560    };
561    Ok(KeyRecord {
562        key_id,
563        account_id: AccountId(id(1)?),
564        principal: Principal(id(2)?),
565        digest: digest_from(row.get::<Vec<u8>, _>(3).as_slice(), key_id)?,
566        not_after,
567    })
568}
569
570#[async_trait]
571impl tollgate_store::KeySource for PostgresStore {
572    async fn active_keys_page(
573        &self,
574        now: Timestamp,
575        after: Option<KeyId>,
576        limit: NonZeroUsize,
577    ) -> Result<tollgate_store::KeyPage, StoreError> {
578        validate_key_page_limit(limit)?;
579        let cutoff = StoredInstant::from(now);
580        // One short snapshot per page, never a transaction held across HTTP
581        // requests. Revision and records therefore cannot describe two commits.
582        let mut tx = self.pool.begin().await.map_err(storage)?;
583        sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
584            .execute(&mut *tx)
585            .await
586            .map_err(storage)?;
587        let revision: i64 =
588            sqlx::query_scalar("SELECT revision FROM tollgate_credential_revision WHERE singleton")
589                .fetch_one(&mut *tx)
590                .await
591                .map_err(storage)?;
592        let revision = u64::try_from(revision)
593            .map_err(|_| StoreError("credential revision is negative".into()))?;
594        // Separate SQL shapes preserve an indexable range in prepared plans.
595        let sql = if after.is_some() {
596            "SELECT key_id, account_id, principal, digest, not_after_floor_us, not_after_submicro_ns, not_after_is_lower_bound
597             FROM tollgate_credential_keys WHERE revoked_at_us IS NULL
598             AND (not_after_floor_us IS NULL OR not_after_submicro_ns IS NULL
599                  OR (not_after_floor_us, not_after_submicro_ns) > ($1, $2)) AND key_id > $3
600             ORDER BY key_id LIMIT $4"
601        } else {
602            "SELECT key_id, account_id, principal, digest, not_after_floor_us, not_after_submicro_ns, not_after_is_lower_bound
603             FROM tollgate_credential_keys WHERE revoked_at_us IS NULL
604             AND (not_after_floor_us IS NULL OR not_after_submicro_ns IS NULL
605                  OR (not_after_floor_us, not_after_submicro_ns) > ($1, $2)) AND key_id >= $3
606             ORDER BY key_id LIMIT $4"
607        };
608        let rows = sqlx::query(sql)
609            .bind(cutoff.micros)
610            .bind(cutoff.submicro_nanos)
611            .bind(id_bytes(after.unwrap_or(KeyId(0)).0))
612            .bind((limit.get() + 1) as i64)
613            .fetch_all(&mut *tx)
614            .await
615            .map_err(storage)?;
616        tx.commit().await.map_err(storage)?;
617        // Decode lookahead too: corrupt data cannot masquerade as a normal
618        // continuation and silently defer a failure to another request.
619        let mut records = rows
620            .into_iter()
621            .map(credential_from_row)
622            .map(|row| row.map(tollgate_store::CredentialRecord::from))
623            .collect::<Result<Vec<_>, _>>()?;
624        let next_after = if records.len() > limit.get() {
625            records.pop();
626            records.last().map(|record| record.key_id)
627        } else {
628            None
629        };
630        tollgate_store::KeyPage::try_new(revision, now, after, limit, records, next_after)
631    }
632}
633
634#[cfg(test)]
635mod stored_id_tests {
636    use super::StoredId;
637
638    #[test]
639    fn malformed_storage_id_explains_both_accepted_representations() {
640        let error = serde_json::from_value::<StoredId>(serde_json::Value::Bool(true)).unwrap_err();
641        assert!(
642            error
643                .to_string()
644                .contains("a legacy u64 number or a canonical 128-bit identifier string"),
645            "unexpected diagnostic: {error}"
646        );
647    }
648}
649
650/// PostgreSQL's snapshot JSON predates the portable HTTP identifier contract.
651/// Keep this boundary explicit so existing rows remain readable and ordinary
652/// legacy-range writes remain rollback-safe while high-bit ids finally work.
653///
654/// **Two facts in this table follow opposite rules, deliberately.** The
655/// generation is *not* here — it lives only in the `generation` column, because
656/// that column is what the `ON CONFLICT ... WHERE` monotonicity guard compares
657/// and what a tombstone reports after its snapshot is gone (GL-54). `account_id`
658/// *is* here and must stay: its column is `GENERATED ALWAYS AS` a function of
659/// this JSON (migration 0006), so deleting it here would silently NULL that
660/// column, unmatch the partial index, and make an account-wide status change
661/// republish nothing.
662///
663/// Both rules say "one writer per fact". They differ on which copy survives
664/// because the generation needs a BIGINT comparison the JSON cannot give.
665#[derive(Serialize)]
666struct StoredSnapshotRef<'a> {
667    account_id: StoredId,
668    key_id: Option<StoredId>,
669    status: &'a AccountStatus,
670    /// The account-owned execution-capacity class (GL-99). Rides the JSONB
671    /// document; omitting it here would drop it on every publish.
672    capacity_class: &'a CapacityClass,
673    enforcement_mode: &'a EnforcementMode,
674    valid_until: &'a Timestamp,
675    permissions: &'a PermissionBits,
676    limits: &'a ResolvedLimits,
677    cost_table: &'a Arc<CostTable>,
678    /// Stored so a *pulled* snapshot carries what a pushed one does. Without
679    /// it an instance that refreshed instead of receiving a push would report
680    /// no budget at all, and the two would disagree about the same account.
681    budget: Option<&'a BudgetView>,
682    /// The consuming application's policy identity (GL-94). Rides the JSONB
683    /// column, so it needs no schema change of its own — but it does need to
684    /// be here: this DTO is the whole of what storage writes, and a field
685    /// omitted from it is dropped on every publish without a word.
686    policy_revision: &'a PolicyRevision,
687}
688
689impl<'a> From<&'a AccountSnapshot> for StoredSnapshotRef<'a> {
690    fn from(snapshot: &'a AccountSnapshot) -> Self {
691        StoredSnapshotRef {
692            account_id: StoredId(snapshot.account_id.0),
693            key_id: snapshot.key_id.map(|id| StoredId(id.0)),
694            status: &snapshot.status,
695            capacity_class: &snapshot.capacity_class,
696            enforcement_mode: &snapshot.enforcement_mode,
697            valid_until: &snapshot.valid_until,
698            permissions: &snapshot.permissions,
699            limits: &snapshot.limits,
700            cost_table: &snapshot.cost_table,
701            budget: snapshot.budget.as_ref(),
702            policy_revision: &snapshot.policy_revision,
703        }
704    }
705}
706
707/// Rows written before GL-54 still carry a `generation` key. Serde ignores
708/// unknown fields, so those rows decode unchanged and the vestigial key is
709/// simply not read — which is why this needed no backfill.
710#[derive(Deserialize)]
711struct StoredSnapshot {
712    account_id: StoredId,
713    key_id: Option<StoredId>,
714    status: AccountStatus,
715    /// Absent from every row written before GL-99, and `default` for the reason
716    /// the fields below are: a required field would make every pre-existing
717    /// row fail to decode and deny every principal until the whole catalogue
718    /// was republished. `Assured` is the safe default — the availability every
719    /// account already has.
720    #[serde(default)]
721    capacity_class: CapacityClass,
722    /// Absent from every row written before elastic mode existed, and those
723    /// rows are the overwhelming majority the first time this ships.
724    ///
725    /// `default` rather than a required field, because the alternative is not
726    /// a loud failure but a silent outage: a required field makes every
727    /// pre-existing row fail to decode, `SnapshotSource::snapshot` returns a
728    /// store error for each, and the request path denies every principal until
729    /// someone republishes the whole catalogue. Defaulting to
730    /// [`EnforcementMode::Strict`] is also the safe direction — an
731    /// undecided account enforces, it does not extend credit.
732    #[serde(default)]
733    enforcement_mode: EnforcementMode,
734    valid_until: Timestamp,
735    permissions: PermissionBits,
736    limits: ResolvedLimits,
737    cost_table: Arc<CostTable>,
738    /// Absent from every row written before periodic budgets, and `default`
739    /// for the reason `enforcement_mode` is: a required field would make every
740    /// pre-existing row fail to decode and deny every principal until the
741    /// whole catalogue was republished. `None` is "the control plane said
742    /// nothing", which readers report as such rather than as a zero balance.
743    #[serde(default)]
744    budget: Option<BudgetView>,
745    /// Absent from every row written before GL-94, and `default` for the reason
746    /// the two fields above are. The unstated revision is a value rather than
747    /// an absence, so a pre-existing row decodes to "this account's publisher
748    /// stated no policy identity" — which is true, and is what a consumer
749    /// reading it back should be told.
750    #[serde(default)]
751    policy_revision: PolicyRevision,
752}
753
754impl StoredSnapshot {
755    /// Rebuild the snapshot around a generation the JSON no longer carries.
756    ///
757    /// Deliberately not a `From` impl: the generation has to come from the
758    /// row's column, and a conversion that could be written without it would
759    /// be a conversion someone writes without it.
760    fn into_snapshot(self, generation: Generation) -> AccountSnapshot {
761        let builder = AccountSnapshot::builder(
762            AccountId(self.account_id.0),
763            generation,
764            self.status,
765            self.valid_until,
766            self.permissions,
767            self.limits,
768            self.cost_table,
769        )
770        .enforcement_mode(self.enforcement_mode)
771        .capacity_class(self.capacity_class)
772        .policy_revision(self.policy_revision);
773        match self.key_id {
774            Some(key_id) => builder.key_id(tollgate_core::KeyId(key_id.0)).build(),
775            None => builder.build(),
776        }
777    }
778}
779
780/// The reconciliation query's active-lease sum, named because two callers must
781/// agree on it: `conservation` runs it, and `explain_active_lease_sum` asks the
782/// planner what it does with it. A test that copied the text would keep
783/// reporting an index scan after the real predicate had drifted away from
784/// migration 0005's index (GL-12).
785///
786/// `state = 0` is spelled out rather than bound, so the predicate is a literal
787/// the partial index can match.
788///
789/// Deliberately still its own statement. Folding it into the account read as a
790/// lateral subquery would make the reconciliation read atomic without a
791/// transaction, but it also destabilises the plan: measured over twelve runs
792/// of `the_account_filter_is_answered_by_an_index_not_by_discarding_rows`, the
793/// correlated form kept the index 4 times in 5 and the uncorrelated form 10
794/// times in 12, against 12 in 12 for this shape. Atomicity is bought with a
795/// transaction instead (GL-56), which leaves this predicate — and the plan GL-12
796/// pinned — untouched.
797const ACTIVE_LEASE_SUM_SQL: &str =
798    "SELECT COALESCE(SUM(granted), 0)::BIGINT, COALESCE(SUM(used), 0)::BIGINT
799     FROM tollgate_leases WHERE account_id = $1 AND state = 0";
800
801/// The expiry sweep's selection, named for the same reason as the sum above:
802/// `reclaim_expired_batch` runs it and `explain_reclaim_due_leases` asks the
803/// planner what it does with it. A test that copied the text would keep
804/// reporting an index-ordered walk after the real `ORDER BY` had drifted away
805/// from `tollgate_leases_expiry` — which is exactly the drift GL-65 found.
806///
807/// The `ORDER BY` is the index's own column order, and that is the whole
808/// point. It was `(account_id, lease_id)`, which no index answers, so the
809/// `LIMIT` could not stop an index walk: every batch read and sorted the
810/// entire remaining backlog to return 256 rows, making a drain quadratic in
811/// the backlog it exists to clear (GL-65). Ordering by the expiry pair makes
812/// the `LIMIT` a range stop, and settles oldest-due-first like the reference
813/// backend does.
814/// Takes the account row before an issuance counts against its bound.
815///
816/// Hoisted so the test that demonstrates the race can drive the *same*
817/// statements this method does, rather than a copy that could drift from them
818/// — the convention `RECLAIM_DUE_LEASES_SQL` established.
819pub(crate) const ACCOUNT_LOCK_SQL: &str =
820    "SELECT 1 FROM tollgate_accounts WHERE account_id = $1 FOR UPDATE";
821
822/// Credentials that can still authenticate at the given instant: not revoked,
823/// and not past `not_after`. The predicate the bound counts with.
824pub(crate) const LIVE_KEY_COUNT_SQL: &str = "SELECT count(*) FROM tollgate_credential_keys
825     WHERE account_id = $1
826       AND revoked_at_us IS NULL
827       AND (not_after_floor_us IS NULL OR not_after_submicro_ns IS NULL
828            OR (not_after_floor_us, not_after_submicro_ns) > ($2, $3))";
829
830const RECLAIM_DUE_LEASES_SQL: &str = "SELECT lease_id, account_id, granted, used
831     FROM tollgate_leases
832     WHERE state = 0 AND (expires_at_floor_us, expires_at_submicro_ns) <= ($1, $2)
833     ORDER BY expires_at_floor_us, expires_at_submicro_ns
834     LIMIT $3 FOR UPDATE SKIP LOCKED";
835
836/// The rollover sweep's selection, hoisted and fixed for the same reason
837/// (GL-65's sibling) — but it needed two changes, not one.
838///
839/// `ORDER BY account_id` sorted every due account to return one bounded page,
840/// and served the lowest ids rather than the most overdue boundaries.
841/// `budget_period = $1` is an equality, so within that prefix of
842/// `tollgate_accounts_due_rollover (budget_period, period_start_us)` the scan
843/// is already ordered by `period_start_us`: ordering by it is index-native and
844/// crosses the oldest boundary first.
845///
846/// That alone changed nothing, because the index is *partial* on
847/// `budget_allowance IS NOT NULL` and this query never said so — the planner
848/// cannot apply a partial index it cannot prove applies, so the sweep had
849/// never once used the index added for it. Restating the predicate is free and
850/// selects exactly the same rows: `tollgate_accounts_budget_all_or_nothing` already
851/// requires the allowance, period and rollover columns to be all null or all
852/// non-null, and `budget_period = $1` has ruled out all-null. Measured on 400
853/// accounts with five due: seq scan and sort at cost 25.02, versus an index
854/// scan with no sort at 5.92.
855/// The crossed-from boundary is returned so the caller can order the batch.
856///
857/// `RETURNING` has no defined row order: PostgreSQL emits rows in whatever
858/// order the update's join produced, which is physical order, not the boundary
859/// order the CTE selected by. Ordering it *here* would mean a trailing
860/// `ORDER BY` and a `Sort` node, which is the cost
861/// `the_rollover_sweep_reaches_its_index_instead_of_sorting_the_due_set`
862/// forbids (GL-65) — so the page is ordered in Rust instead, where it is bounded
863/// by the batch limit rather than by how many accounts are due.
864const DUE_PERIODS_SQL: &str = "WITH due AS (
865         SELECT account_id, allowance_balance AS prior, budget_allowance AS allowance,
866                period_start_us AS crossed_from
867         FROM tollgate_accounts
868         WHERE budget_allowance IS NOT NULL
869           AND budget_period = $1 AND period_start_us < $2
870         ORDER BY period_start_us LIMIT $3 FOR UPDATE SKIP LOCKED
871     )
872     UPDATE tollgate_accounts AS account SET
873         deposited = account.deposited + due.allowance,
874         expired = account.expired + due.prior,
875         balance = account.balance - due.prior + due.allowance,
876         allowance_balance = due.allowance,
877         period_start_us = $2
878     FROM due
879     WHERE account.account_id = due.account_id
880     RETURNING due.account_id, due.allowance, due.prior, due.crossed_from";
881
882fn id_bytes(id: u128) -> Vec<u8> {
883    id.to_be_bytes().to_vec()
884}
885
886fn id_from(bytes: &[u8]) -> u128 {
887    let mut buf = [0u8; 16];
888    buf.copy_from_slice(bytes);
889    u128::from_be_bytes(buf)
890}
891
892fn to_i64(units: CostUnits, what: &str) -> Result<i64, StoreError> {
893    i64::try_from(units.get()).map_err(|_| StoreError(format!("{what} exceeds i64 range")))
894}
895
896/// Persisted capabilities are minted from counters seeded at one. Invalid
897/// storage must not become a caller mismatch or mint a zero capability.
898fn stored_fence(value: i64) -> Result<FencingToken, StoreError> {
899    u64::try_from(value)
900        .ok()
901        .filter(|value| *value > 0)
902        .map(FencingToken)
903        .ok_or_else(|| StoreError(format!("stored fencing token is not positive: {value}")))
904}
905
906/// A negative unit column is ledger corruption, never a value to normalise:
907/// clamping it to zero would let the conservation equation pass over exactly
908/// the discrepancy it exists to expose.
909fn to_units(value: i64, what: &str) -> Result<CostUnits, StoreError> {
910    u64::try_from(value)
911        .map(CostUnits)
912        .map_err(|_| StoreError(format!("{what} is negative in storage: {value}")))
913}
914
915fn ts_micros(ts: Timestamp) -> i64 {
916    ts.as_microsecond()
917}
918
919/// The inverse of [`ts_micros`]. A stored value outside the representable
920/// range is corruption to report, never an instant to clamp to.
921fn micros_ts(value: i64, what: &str) -> Result<Timestamp, StoreError> {
922    tollgate_store::clock::timestamp_from_micros(value).map_err(|e| {
923        StoreError(format!(
924            "{what} is not a representable instant: {value} ({e})"
925        ))
926    })
927}
928
929fn storage(e: sqlx::Error) -> StoreError {
930    StoreError(format!("postgres: {e}"))
931}
932
933fn alloc_storage(e: sqlx::Error) -> AllocateError {
934    AllocateError::Storage(storage(e))
935}
936
937/// Finish a transaction before making its result observable.
938///
939/// `sqlx::Transaction` only queues a rollback when it is dropped. A caller can
940/// therefore start a competing transaction after this future returns but
941/// before PostgreSQL has released the first transaction's row locks. That is
942/// especially visible to reclaim's `SKIP LOCKED` query, which may otherwise
943/// miss a lease owned by an operation that has already reported failure.
944///
945/// Generic over the error type rather than duplicated per error: this crate
946/// now finishes transactions that fail with `StoreError`, `AllocateError` and
947/// `SetStatusError`, and three copies of the rollback-and-report logic would
948/// be three places for it to drift.
949async fn finish_transaction<T, E>(
950    tx: Transaction<'_, Postgres>,
951    result: Result<T, E>,
952) -> Result<T, E>
953where
954    E: From<StoreError> + std::fmt::Display,
955{
956    match result {
957        Ok(value) => {
958            tx.commit().await.map_err(|e| E::from(storage(e)))?;
959            Ok(value)
960        }
961        Err(error) => {
962            if let Err(rollback_error) = tx.rollback().await {
963                return Err(E::from(StoreError(format!(
964                    "operation failed ({error}); transaction rollback failed ({})",
965                    storage(rollback_error)
966                ))));
967            }
968            Err(error)
969        }
970    }
971}
972
973/// Fixture reset is available only in the opt-in `test_support` module,
974/// never as a method on a normal store handle.
975///
976/// ```compile_fail,E0599
977/// # use tollgate_store_postgres::PostgresStore;
978/// fn reset(store: &PostgresStore) {
979///     let _ = store.truncate_all();
980/// }
981/// ```
982///
983/// Query-plan inspection, which refreshes database statistics, is likewise
984/// absent from the operational handle:
985///
986/// ```compile_fail,E0599
987/// # use tollgate_store_postgres::PostgresStore;
988/// # use tollgate_core::AccountId;
989/// fn explain(store: &PostgresStore) {
990///     let _ = store.explain_active_lease_sum(AccountId(1));
991/// }
992/// ```
993///
994/// The same holds for the two sweep plans (GL-65):
995///
996/// ```compile_fail,E0599
997/// # use tollgate_store_postgres::PostgresStore;
998/// fn explain_sweep(store: &PostgresStore) {
999///     let _ = store.explain_reclaim_due_leases(jiff::Timestamp::UNIX_EPOCH, 256);
1000/// }
1001/// ```
1002///
1003/// ```compile_fail,E0599
1004/// # use tollgate_store_postgres::PostgresStore;
1005/// fn explain_rollover(store: &PostgresStore) {
1006///     let _ = store.explain_due_periods("daily", 0, 256);
1007/// }
1008/// ```
1009pub struct PostgresStore {
1010    pool: PgPool,
1011    policy: GrantPolicy,
1012    push: broadcast::Sender<SnapshotPush>,
1013}
1014
1015/// Connection-pool bounds. Callers that hold background tasks open against
1016/// this store rely on `acquire_timeout`: when every connection is checked out
1017/// by a stalled query, it is the only thing that turns "wait forever" into an
1018/// error the caller can report (INVARIANTS.md GL-18).
1019#[derive(Debug, Clone, Copy)]
1020pub struct PoolConfig {
1021    /// The most connections the pool opens at once. Must be positive;
1022    /// defaults to 16.
1023    pub max_connections: u32,
1024    /// How long a caller may wait for a free connection before the call
1025    /// fails. Must be positive.
1026    pub acquire_timeout: std::time::Duration,
1027}
1028
1029impl Default for PoolConfig {
1030    fn default() -> Self {
1031        PoolConfig {
1032            max_connections: 16,
1033            acquire_timeout: std::time::Duration::from_secs(5),
1034        }
1035    }
1036}
1037
1038impl PoolConfig {
1039    /// Rejects a zero `max_connections` or a zero `acquire_timeout`.
1040    ///
1041    /// [`PostgresStore::connect_with`] calls this before it opens any
1042    /// connection, so an unsafe bound never reaches the database
1043    /// (INVARIANTS.md 16).
1044    pub fn validate(&self) -> Result<(), StoreError> {
1045        if self.max_connections == 0 {
1046            return Err(StoreError("max_connections must be positive".into()));
1047        }
1048        if self.acquire_timeout.is_zero() {
1049            return Err(StoreError("acquire_timeout must be positive".into()));
1050        }
1051        Ok(())
1052    }
1053}
1054
1055impl PostgresStore {
1056    /// Connect with default pool bounds and run pending migrations.
1057    pub async fn connect(url: &str, policy: GrantPolicy) -> Result<Arc<Self>, StoreError> {
1058        Self::connect_with(url, policy, PoolConfig::default()).await
1059    }
1060
1061    /// Connect and run pending migrations (versioned under ./migrations,
1062    /// tracked by sqlx's _sqlx_migrations table — review finding GL-11).
1063    ///
1064    /// Note the limits of what a pool bound can promise: `acquire_timeout`
1065    /// covers waiting for a connection, including establishing one, but a
1066    /// query already in flight on a healthy connection is bounded only by a
1067    /// server-side `statement_timeout`. Callers must still bound their own
1068    /// calls (INVARIANTS.md GL-18).
1069    pub async fn connect_with(
1070        url: &str,
1071        policy: GrantPolicy,
1072        pool_config: PoolConfig,
1073    ) -> Result<Arc<Self>, StoreError> {
1074        policy
1075            .validate()
1076            .map_err(|e| StoreError(format!("invalid grant policy: {e}")))?;
1077        pool_config.validate()?;
1078        let pool = PgPoolOptions::new()
1079            .max_connections(pool_config.max_connections)
1080            .acquire_timeout(pool_config.acquire_timeout)
1081            .connect(url)
1082            .await
1083            .map_err(storage)?;
1084        sqlx::migrate!("./migrations")
1085            .run(&pool)
1086            .await
1087            .map_err(|e| StoreError(format!("migrate: {e}")))?;
1088        let (push, _) = broadcast::channel(PUSH_CHANNEL_CAPACITY);
1089        Ok(Arc::new(PostgresStore { pool, policy, push }))
1090    }
1091
1092    /// Broadcast a control-plane change to in-process subscribers. No
1093    /// receivers is not a failure — a pull still observes the change — but
1094    /// how many instances the push reached is the difference between
1095    /// propagating in milliseconds and propagating at the next refresh, so
1096    /// it is reported rather than discarded.
1097    fn push_to_subscribers(&self, push: SnapshotPush) {
1098        let principal = push.principal;
1099        let subscribers = self.push.send(push).unwrap_or(0);
1100        tracing::debug!(
1101            %principal,
1102            subscribers,
1103            "snapshot pushed to subscribers"
1104        );
1105    }
1106
1107    /// Push a tombstone only when a committed removal changed something.
1108    fn announce_removal(&self, principal: Principal, receipt: &AdminReceipt<()>) {
1109        if receipt.before != receipt.after
1110            && let AdminState::Snapshot { generation, .. } = receipt.after
1111        {
1112            self.push_to_subscribers(SnapshotPush {
1113                principal,
1114                resolution: SnapshotResolution::Revoked { generation },
1115            });
1116        }
1117    }
1118
1119    // ---- reconciliation / test surface (mirrors MemoryStore) ---------
1120
1121    /// The account's unspent balance, read in one statement outside any
1122    /// transaction. An unknown account reads as zero, as in `MemoryStore`.
1123    ///
1124    /// A negative stored value is reported as a [`StoreError`], never
1125    /// clamped (INVARIANTS.md 11). For figures that must agree with each
1126    /// other, use [`conservation`](Self::conservation), which reads them
1127    /// from one snapshot.
1128    pub async fn balance(&self, account: AccountId) -> Result<CostUnits, StoreError> {
1129        let row = sqlx::query("SELECT balance FROM tollgate_accounts WHERE account_id = $1")
1130            .bind(id_bytes(account.0))
1131            .fetch_optional(&self.pool)
1132            .await
1133            .map_err(storage)?;
1134        row.map(|r| to_units(r.get::<i64, _>(0), "account balance"))
1135            .transpose()
1136            .map(|units| units.unwrap_or(CostUnits::ZERO))
1137    }
1138
1139    /// Total usage accepted into the account's billing ledger, including
1140    /// overage usage, read in one statement outside any transaction. An
1141    /// unknown account reads as zero, as in `MemoryStore`.
1142    ///
1143    /// A negative stored value is reported as a [`StoreError`], never
1144    /// clamped (INVARIANTS.md 11).
1145    pub async fn usage_recorded(&self, account: AccountId) -> Result<CostUnits, StoreError> {
1146        let row = sqlx::query("SELECT usage_recorded FROM tollgate_accounts WHERE account_id = $1")
1147            .bind(id_bytes(account.0))
1148            .fetch_optional(&self.pool)
1149            .await
1150            .map_err(storage)?;
1151        row.map(|r| to_units(r.get::<i64, _>(0), "account usage_recorded"))
1152            .transpose()
1153            .map(|units| units.unwrap_or(CostUnits::ZERO))
1154    }
1155
1156    /// The terms of the account's ledger equation, for reconciliation
1157    /// against the ledger contract in `INVARIANTS.md`. `None` when the account
1158    /// does not exist.
1159    ///
1160    /// The account totals and the sums over its active leases are read in one
1161    /// `REPEATABLE READ, READ ONLY` transaction, so both come from the same
1162    /// snapshot and a concurrent commit cannot land between them. The
1163    /// transaction writes nothing.
1164    ///
1165    /// Stored state that cannot be a valid ledger is reported as a
1166    /// [`StoreError`], never clamped or panicked on: a negative unit column
1167    /// (INVARIANTS.md 11), or active-lease usage exceeding recorded usage.
1168    /// Whether the returned terms balance is for the caller to check with
1169    /// [`Conservation::holds`](tollgate_store::Conservation::holds).
1170    pub async fn conservation(
1171        &self,
1172        account: AccountId,
1173    ) -> Result<Option<Conservation>, StoreError> {
1174        // One snapshot for both halves. The account's stored totals and the
1175        // sums over its live leases move together — `ingest` raises
1176        // `usage_recorded` and the lease's `used` in one transaction — so
1177        // reading them on two pooled connections let a commit land between
1178        // them: reconciliation then paired a pre-write total with a post-write
1179        // sum and reported corruption on a correct ledger, or underflowed the
1180        // subtraction outright (GL-56).
1181        //
1182        // `READ COMMITTED` is not enough, because it takes a fresh snapshot
1183        // per statement; `REPEATABLE READ` takes one at the first read and
1184        // holds it for the rest of the transaction. Read-only, so it cannot
1185        // hit the serialization failures that make `SERIALIZABLE` a retry
1186        // contract rather than a snapshot.
1187        let mut tx = self.pool.begin().await.map_err(storage)?;
1188        let result = async {
1189            sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
1190                .execute(&mut *tx)
1191                .await
1192                .map_err(storage)?;
1193            let account_row = sqlx::query(
1194                "SELECT deposited, balance, usage_recorded, settlement_loss, overage_recorded,
1195                        expired
1196                 FROM tollgate_accounts WHERE account_id = $1",
1197            )
1198            .bind(id_bytes(account.0))
1199            .fetch_optional(&mut *tx)
1200            .await
1201            .map_err(storage)?;
1202            let lease_row = sqlx::query(ACTIVE_LEASE_SUM_SQL)
1203                .bind(id_bytes(account.0))
1204                .fetch_one(&mut *tx)
1205                .await
1206                .map_err(storage)?;
1207            Ok::<_, StoreError>((account_row, lease_row))
1208        }
1209        .await;
1210        let (account_row, lease_row) = finish_transaction(tx, result).await?;
1211
1212        let Some(row) = account_row else {
1213            return Ok(None);
1214        };
1215        let active_grants = to_units(lease_row.get::<i64, _>(0), "active lease grants")?;
1216        let active_used = to_units(lease_row.get::<i64, _>(1), "active lease usage")?;
1217        let recorded = to_units(row.get::<i64, _>(2), "usage_recorded")?;
1218        Ok(Some(Conservation {
1219            deposited: to_units(row.get::<i64, _>(0), "deposited")?,
1220            overage_recorded: to_units(row.get::<i64, _>(4), "overage_recorded")?,
1221            balance: to_units(row.get::<i64, _>(1), "balance")?,
1222            active_lease_grants: active_grants,
1223            // Surfaced, never panicked on. These are two stored columns read
1224            // from a database this process does not exclusively own, so
1225            // `recorded < active_used` is corruption to report — the same
1226            // class as a negative unit column (INVARIANTS.md GL-11) — not an
1227            // internal invariant a caller could not violate. `MemoryStore`
1228            // keeps its `expect` because there the counters are maintained by
1229            // one process under one lock and the state is unrepresentable.
1230            settled_usage: recorded.checked_sub(active_used).ok_or_else(|| {
1231                StoreError(format!(
1232                    "active lease usage {} exceeds recorded usage {} for account {account}",
1233                    active_used.get(),
1234                    recorded.get()
1235                ))
1236            })?,
1237            settlement_loss: to_units(row.get::<i64, _>(3), "settlement_loss")?,
1238            expired: to_units(row.get::<i64, _>(5), "expired")?,
1239        }))
1240    }
1241}
1242
1243/// One locked lease row, named so its accounting fields cannot be confused at
1244/// the call site and tooling does not manufacture a Cartesian product of
1245/// arbitrary replacements for an anonymous seven-field tuple.
1246struct LockedLeaseRow {
1247    account_id: Vec<u8>,
1248    fencing_token: i64,
1249    granted: i64,
1250    used: i64,
1251    credited: i64,
1252    expires_at: StoredInstant,
1253    state: i16,
1254    /// The half of `granted` drawn from the account's periodic allowance, and
1255    /// the period that funded it. Settlement needs both: the split says which
1256    /// bucket each unspent unit belongs to, the period says whether the
1257    /// allowance half still exists (GL-97).
1258    from_allowance: i64,
1259    period_start_us: i64,
1260}
1261
1262/// Settlement evidence from the account update inside the current transaction.
1263/// Consolidation uses the credit that survived the period boundary as its floor.
1264struct ReleasedCredit {
1265    account: AccountId,
1266    restored: CostUnits,
1267    preserves_funding: bool,
1268}
1269
1270/// Lock one lease row for a settlement transition.
1271async fn lock_lease(
1272    tx: &mut Transaction<'_, Postgres>,
1273    lease_id: LeaseId,
1274) -> Result<Option<LockedLeaseRow>, sqlx::Error> {
1275    let row = sqlx::query(
1276        "SELECT account_id, fencing_token, granted, used, credited, expires_at_floor_us, state,
1277                from_allowance, period_start_us, expires_at_submicro_ns
1278         FROM tollgate_leases WHERE lease_id = $1 FOR UPDATE",
1279    )
1280    .bind(id_bytes(lease_id.0))
1281    .fetch_optional(&mut **tx)
1282    .await?;
1283    Ok(row.map(|row| LockedLeaseRow {
1284        account_id: row.get(0),
1285        fencing_token: row.get(1),
1286        granted: row.get(2),
1287        used: row.get(3),
1288        credited: row.get(4),
1289        expires_at: StoredInstant {
1290            micros: row.get(5),
1291            submicro_nanos: row.get(9),
1292        },
1293        state: row.get(6),
1294        from_allowance: row.get(7),
1295        period_start_us: row.get(8),
1296    }))
1297}
1298
1299/// What a consolidation's release half hands the grant half of the same
1300/// transaction. A plain acquire returns nothing: [`Exchange::ACQUIRE`].
1301struct Exchange {
1302    /// The credit this transaction restored to the account, which the grant
1303    /// policy's shrink cap may not size the result below (see
1304    /// [`LeaseAllocator::consolidate`]).
1305    floor: CostUnits,
1306    /// The largest quote the returned lease refused, which the grant may grow
1307    /// to when the restored balance funds it (GL-131).
1308    needed: CostUnits,
1309    /// Whether the settlement left total funding unchanged, so a refusal may
1310    /// attest the ledger it reads (GL-130).
1311    preserves_funding: bool,
1312}
1313
1314impl Exchange {
1315    const ACQUIRE: Exchange = Exchange {
1316        floor: CostUnits::ZERO,
1317        needed: CostUnits::ZERO,
1318        preserves_funding: true,
1319    };
1320}
1321
1322impl PostgresStore {
1323    /// One grant, inside a caller-owned transaction.
1324    async fn acquire_in_tx(
1325        &self,
1326        tx: &mut Transaction<'_, Postgres>,
1327        account: AccountId,
1328        requested: CostUnits,
1329        expires_at: Timestamp,
1330        exchange: Exchange,
1331    ) -> Result<Allocation, AllocateError> {
1332        let Exchange {
1333            floor,
1334            needed,
1335            preserves_funding: settlement_preserves_funding,
1336        } = exchange;
1337        let row = sqlx::query(
1338            "SELECT balance, status, next_fence, allowance_balance, period_start_us
1339             FROM tollgate_accounts WHERE account_id = $1 FOR UPDATE",
1340        )
1341        .bind(id_bytes(account.0))
1342        .fetch_optional(&mut **tx)
1343        .await
1344        .map_err(alloc_storage)?
1345        .ok_or(AllocateError::UnknownAccount)?;
1346
1347        // Suspended and Closed both refuse, under one deny reason: no
1348        // client acts on the distinction, and splitting it would widen
1349        // `AllocateError`'s per-reason tally for nothing (GL-51).
1350        if decode_status(row.get::<String, _>(1)).map_err(AllocateError::Storage)?
1351            != AccountStatus::Active
1352        {
1353            return Err(AllocateError::AccountInactive);
1354        }
1355        let balance =
1356            to_units(row.get::<i64, _>(0), "account balance").map_err(AllocateError::Storage)?;
1357        // The floor and demand are applied to the policy's answer, not to
1358        // the balance test: the units the caller returned rejoined `balance`
1359        // earlier in this same transaction, and both are capped by it, so
1360        // they can only re-select capacity the account demonstrably has, and
1361        // an account with nothing left still refuses.
1362        let granted = match self
1363            .policy
1364            .consolidation_grant(requested, balance, floor, needed)
1365        {
1366            Some(granted) => granted,
1367            None => {
1368                // A refused consolidation rolls its settlement back. Only
1369                // attest when that settlement did not remove funding itself.
1370                if settlement_preserves_funding && !requested.is_zero() {
1371                    let budget = sqlx::query(BUDGET_VIEW_SQL)
1372                        .bind(id_bytes(account.0))
1373                        .fetch_one(&mut **tx)
1374                        .await
1375                        .map_err(alloc_storage)?;
1376                    let evidence = budget_view(&budget)
1377                        .map_err(AllocateError::Storage)?
1378                        .shortfall();
1379                    return Err(match evidence.exhaustion() {
1380                        Some(exhausted) => AllocateError::BalanceExhausted(exhausted),
1381                        None => AllocateError::BalanceInsufficient(evidence),
1382                    });
1383                }
1384                return Err(AllocateError::InsufficientBalance);
1385            }
1386        };
1387        // Allowance first: the units with an expiry date are spent
1388        // before the manual credits sitting beside them, and the lease
1389        // remembers the split so settlement can return each half to where
1390        // it came from (GL-97).
1391        let granted_i = to_i64(granted, "grant").map_err(AllocateError::Storage)?;
1392        let allowance_balance = row.get::<i64, _>(3);
1393        let from_allowance = granted_i.min(allowance_balance);
1394        let period_start_us = row.get::<i64, _>(4);
1395        let fence = row.get::<i64, _>(2);
1396        let fence_token = stored_fence(fence).map_err(AllocateError::Storage)?;
1397
1398        let lease_id = LeaseId(uuid::Uuid::new_v4().as_u128());
1399
1400        // The grant commits with this transaction, and a consolidation's
1401        // settlement has already applied inside it, so the row this returns is
1402        // the ledger the grant lands in: loss and expiry included.
1403        let budget = sqlx::query(GRANT_DEBIT_SQL)
1404            .bind(id_bytes(account.0))
1405            .bind(granted_i)
1406            .bind(from_allowance)
1407            .fetch_one(&mut **tx)
1408            .await
1409            .map_err(alloc_storage)?;
1410        let funding = budget_view(&budget)
1411            .map_err(AllocateError::Storage)?
1412            .shortfall();
1413        sqlx::query(
1414            "INSERT INTO tollgate_leases
1415             (lease_id, account_id, fencing_token, granted, used, credited, expires_at_floor_us,
1416              state, from_allowance, period_start_us, expires_at_submicro_ns, expiry_is_upper_bound)
1417             VALUES ($1, $2, $3, $4, 0, 0, $5, 0, $6, $7, $8, FALSE)",
1418        )
1419        .bind(id_bytes(lease_id.0))
1420        .bind(id_bytes(account.0))
1421        .bind(fence)
1422        .bind(granted_i)
1423        .bind(StoredInstant::from(expires_at).micros)
1424        .bind(from_allowance)
1425        .bind(period_start_us)
1426        .bind(StoredInstant::from(expires_at).submicro_nanos)
1427        .execute(&mut **tx)
1428        .await
1429        .map_err(alloc_storage)?;
1430
1431        Ok(Allocation {
1432            grant: LeaseGrant {
1433                lease_id,
1434                account_id: account,
1435                fencing_token: fence_token,
1436                units: granted,
1437                expires_at,
1438            },
1439            funding: Some(funding),
1440        })
1441    }
1442
1443    /// One release, returning its account and the credit the account update
1444    /// actually restored, inside the caller-owned transaction.
1445    async fn release_in_tx(
1446        &self,
1447        tx: &mut Transaction<'_, Postgres>,
1448        lease_id: LeaseId,
1449        fencing_token: FencingToken,
1450        unspent: CostUnits,
1451        now: Timestamp,
1452    ) -> Result<ReleasedCredit, AllocateError> {
1453        let LockedLeaseRow {
1454            account_id,
1455            fencing_token: fence,
1456            granted,
1457            used,
1458            credited: _credited,
1459            expires_at,
1460            state,
1461            from_allowance,
1462            period_start_us,
1463        } = lock_lease(tx, lease_id)
1464            .await
1465            .map_err(alloc_storage)?
1466            .ok_or(AllocateError::UnknownLease)?;
1467
1468        if stored_fence(fence).map_err(AllocateError::Storage)? != fencing_token {
1469            return Err(AllocateError::Fenced);
1470        }
1471        // Releases are accepted through the grace window (see GrantPolicy::
1472        // reclaim_grace) — only a settled or grace-exhausted lease refuses.
1473        let expires_at = expires_at.timestamp().map_err(AllocateError::Storage)?;
1474        if state != STATE_ACTIVE
1475            || self
1476                .policy
1477                .reclaim_cutoff(now)
1478                .is_some_and(|cutoff| expires_at <= cutoff)
1479        {
1480            return Err(AllocateError::LeaseNotActive);
1481        }
1482        // Decoded before the byte form is consumed by the credit below, and
1483        // returned so a consolidation draws its replacement grant from the
1484        // account this lease named rather than one a caller supplied.
1485        let account = AccountId(id_from(&account_id));
1486        let unspent_i = to_i64(unspent, "unspent").map_err(AllocateError::Storage)?;
1487        let loss = granted
1488            .checked_sub(
1489                used.checked_add(unspent_i)
1490                    .ok_or(AllocateError::InvalidRelease)?,
1491            )
1492            .filter(|l| *l >= 0)
1493            .ok_or(AllocateError::InvalidRelease)?;
1494
1495        sqlx::query("UPDATE tollgate_leases SET state = $2, credited = $3 WHERE lease_id = $1")
1496            .bind(id_bytes(lease_id.0))
1497            .bind(STATE_RELEASED)
1498            .bind(unspent_i)
1499            .execute(&mut **tx)
1500            .await
1501            .map_err(alloc_storage)?;
1502        // The lease's usage is charged in the order the account spends,
1503        // allowance first, so the top-up half is what survives a partly
1504        // spent lease. Crediting the allowance half back first would close
1505        // the equation just as well while moving durable credits into the
1506        // bucket that expires at the next boundary (GL-97).
1507        let from_topup = granted
1508            .checked_sub(from_allowance)
1509            .filter(|t| *t >= 0)
1510            .ok_or_else(|| {
1511                AllocateError::Storage(StoreError(format!(
1512                    "lease allowance funding {from_allowance} exceeds its grant {granted}"
1513                )))
1514            })?;
1515        let to_topup = from_topup.min(unspent_i);
1516        let to_allowance = unspent_i - to_topup;
1517
1518        // One statement, and the `CASE` is the whole boundary rule: if the
1519        // account has moved on to a later period, the allowance half is
1520        // expired instead of returned. Deciding it in SQL against the
1521        // row's own `period_start_us` keeps the read and the write in one
1522        // atomic step, so a rollover committing between them cannot make
1523        // this credit an allowance the account no longer has.
1524        let restored: i64 = sqlx::query_scalar(
1525            "UPDATE tollgate_accounts SET
1526                 balance = balance + $2
1527                     + CASE WHEN period_start_us > $5 THEN 0 ELSE $3 END,
1528                 allowance_balance = allowance_balance
1529                     + CASE WHEN period_start_us > $5 THEN 0 ELSE $3 END,
1530                 expired = expired + CASE WHEN period_start_us > $5 THEN $3 ELSE 0 END,
1531                 settlement_loss = settlement_loss + $4
1532             WHERE account_id = $1
1533             RETURNING $2 + CASE WHEN period_start_us > $5 THEN 0 ELSE $3 END",
1534        )
1535        .bind(account_id)
1536        .bind(to_topup)
1537        .bind(to_allowance)
1538        .bind(loss)
1539        .bind(period_start_us)
1540        .fetch_one(&mut **tx)
1541        .await
1542        .map_err(alloc_storage)?;
1543        Ok(ReleasedCredit {
1544            account,
1545            preserves_funding: loss == 0 && restored == unspent_i,
1546            restored: to_units(restored, "restored release credit")
1547                .map_err(AllocateError::Storage)?,
1548        })
1549    }
1550
1551    /// The expiry a grant issued now would carry, clamping to the policy's
1552    /// `max_ttl`. Fallible before any transaction opens, so a bad TTL never
1553    /// settles a lease it cannot replace.
1554    fn grant_expiry(
1555        &self,
1556        ttl: SignedDuration,
1557        now: Timestamp,
1558    ) -> Result<Timestamp, AllocateError> {
1559        if ttl <= SignedDuration::ZERO {
1560            return Err(AllocateError::InvalidTtl);
1561        }
1562        now.checked_add(ttl.min(self.policy.max_ttl))
1563            .map_err(|e| AllocateError::Storage(StoreError(format!("ttl overflow: {e}"))))
1564    }
1565}
1566
1567#[async_trait]
1568impl LeaseAllocator for PostgresStore {
1569    async fn acquire(
1570        &self,
1571        account: AccountId,
1572        requested: CostUnits,
1573        ttl: SignedDuration,
1574        now: Timestamp,
1575    ) -> Result<Allocation, AllocateError> {
1576        let expires_at = self.grant_expiry(ttl, now)?;
1577        let mut tx = self.pool.begin().await.map_err(alloc_storage)?;
1578        let result = self
1579            .acquire_in_tx(&mut tx, account, requested, expires_at, Exchange::ACQUIRE)
1580            .await;
1581        finish_transaction(tx, result).await
1582    }
1583
1584    async fn release(
1585        &self,
1586        lease_id: LeaseId,
1587        fencing_token: FencingToken,
1588        unspent: CostUnits,
1589        now: Timestamp,
1590    ) -> Result<(), AllocateError> {
1591        let mut tx = self.pool.begin().await.map_err(alloc_storage)?;
1592        let result = self
1593            .release_in_tx(&mut tx, lease_id, fencing_token, unspent, now)
1594            .await
1595            .map(|_| ());
1596        finish_transaction(tx, result).await
1597    }
1598
1599    async fn consolidate(
1600        &self,
1601        lease_id: LeaseId,
1602        fencing_token: FencingToken,
1603        unspent: CostUnits,
1604        requested: CostUnits,
1605        needed: CostUnits,
1606        ttl: SignedDuration,
1607        now: Timestamp,
1608    ) -> Result<Allocation, AllocateError> {
1609        let expires_at = self.grant_expiry(ttl, now)?;
1610        let mut tx = self.pool.begin().await.map_err(alloc_storage)?;
1611        // One transaction for both halves is the whole point. Lease row first
1612        // and account row second, which is the order `release` already takes,
1613        // so a consolidation cannot invert the lock order against a concurrent
1614        // release of a sibling lease on the same account.
1615        //
1616        // The new grant is drawn from the account the released lease named,
1617        // never one the caller supplied, so the two halves cannot disagree
1618        // about whose balance moved.
1619        let result = async {
1620            let released = self
1621                .release_in_tx(&mut tx, lease_id, fencing_token, unspent, now)
1622                .await?;
1623            self.acquire_in_tx(
1624                &mut tx,
1625                released.account,
1626                requested,
1627                expires_at,
1628                Exchange {
1629                    floor: released.restored,
1630                    needed,
1631                    preserves_funding: released.preserves_funding,
1632                },
1633            )
1634            .await
1635        }
1636        .await;
1637        finish_transaction(tx, result).await
1638    }
1639
1640    async fn reclaim_expired_batch(
1641        &self,
1642        now: Timestamp,
1643        limit: NonZeroUsize,
1644    ) -> Result<ReclaimBatch, StoreError> {
1645        let limit_i = i64::try_from(limit.get())
1646            .map_err(|_| StoreError(format!("reclaim batch limit exceeds i64 range: {limit}")))?;
1647        let Some(cutoff) = self.policy.reclaim_cutoff(now) else {
1648            return ReclaimBatch::try_new(Vec::new(), limit);
1649        };
1650        let cutoff = StoredInstant::from(cutoff);
1651        let mut tx = self.pool.begin().await.map_err(storage)?;
1652        let result = async {
1653            // Reclaim only once the grace window past expiry has fully lapsed:
1654            // expires_at + grace <= now  iff  expires_at <= now - grace.
1655            // Both tuple components preserve the exact timestamp ordering.
1656            // SKIP LOCKED lets concurrent sweepers cooperate; the limit keeps
1657            // both the lease locks and the transaction's row work bounded.
1658            let rows = sqlx::query(RECLAIM_DUE_LEASES_SQL)
1659                .bind(cutoff.micros)
1660                .bind(cutoff.submicro_nanos)
1661                .bind(limit_i)
1662                .fetch_all(&mut *tx)
1663                .await
1664                .map_err(storage)?;
1665
1666            // A holder that never released cannot prove any unit unspent, so
1667            // nothing is credited: each remainder becomes provisional
1668            // settlement loss, which later usage for the lease converts into
1669            // billed usage (GL-136). `credited` stays zero so that usage fits.
1670            let mut reclaimed = Vec::with_capacity(rows.len());
1671            let mut lease_ids = Vec::with_capacity(rows.len());
1672            let mut forfeits: std::collections::BTreeMap<Vec<u8>, i64> =
1673                std::collections::BTreeMap::new();
1674            for row in rows {
1675                let lease_bytes: Vec<u8> = row.get(0);
1676                let account_bytes: Vec<u8> = row.get(1);
1677                let granted = row.get::<i64, _>(2);
1678                let used = row.get::<i64, _>(3);
1679                // Validate the whole batch before either set-wise UPDATE: a
1680                // negative remainder is corruption, and recording it would
1681                // shrink the account's loss rather than account for a lease.
1682                let forfeited = granted.checked_sub(used).ok_or_else(|| {
1683                    StoreError(format!(
1684                        "reclaim remainder overflow: granted {granted}, used {used}"
1685                    ))
1686                })?;
1687                let forfeited_units = to_units(forfeited, "reclaim remainder")?;
1688                let total = forfeits.entry(account_bytes.clone()).or_default();
1689                *total = total
1690                    .checked_add(forfeited)
1691                    .ok_or_else(|| StoreError("reclaim loss sum overflow".into()))?;
1692                lease_ids.push(lease_bytes.clone());
1693                reclaimed.push(ReclaimedLease {
1694                    lease_id: LeaseId(id_from(&lease_bytes)),
1695                    account_id: AccountId(id_from(&account_bytes)),
1696                    forfeited: forfeited_units,
1697                });
1698            }
1699
1700            let batch = ReclaimBatch::try_new(reclaimed, limit)?;
1701            if batch.is_empty() {
1702                return Ok(batch);
1703            }
1704
1705            let (account_ids, account_forfeits): (Vec<_>, Vec<_>) = forfeits.into_iter().unzip();
1706            let expected_lease_rows = u64::try_from(lease_ids.len())
1707                .map_err(|_| StoreError("reclaim lease row count exceeds u64 range".into()))?;
1708            let expected_account_rows = u64::try_from(account_ids.len())
1709                .map_err(|_| StoreError("reclaim account row count exceeds u64 range".into()))?;
1710
1711            // A set-wise UPDATE does not promise row-lock order. Lock every
1712            // affected account explicitly in byte-sorted order first (the
1713            // BTreeMap above), matching release and ingest's lease-then-account
1714            // order and preventing concurrent multi-account sweeps from
1715            // forming a deadlock cycle.
1716            let locked_accounts = sqlx::query(
1717                "SELECT account_id FROM tollgate_accounts
1718                 WHERE account_id = ANY($1) ORDER BY account_id FOR UPDATE",
1719            )
1720            .bind(&account_ids)
1721            .fetch_all(&mut *tx)
1722            .await
1723            .map_err(storage)?;
1724            if locked_accounts.len() != account_ids.len() {
1725                return Err(StoreError(format!(
1726                    "reclaim locked {} of {} referenced account rows",
1727                    locked_accounts.len(),
1728                    account_ids.len()
1729                )));
1730            }
1731
1732            let updated_leases = sqlx::query(
1733                "UPDATE tollgate_leases
1734                 SET state = $2, credited = 0
1735                 WHERE lease_id = ANY($1) AND state = $3",
1736            )
1737            .bind(&lease_ids)
1738            .bind(STATE_EXPIRED)
1739            .bind(STATE_ACTIVE)
1740            .execute(&mut *tx)
1741            .await
1742            .map_err(storage)?;
1743            if updated_leases.rows_affected() != expected_lease_rows {
1744                return Err(StoreError(format!(
1745                    "reclaim updated {} of {} locked lease rows",
1746                    updated_leases.rows_affected(),
1747                    lease_ids.len()
1748                )));
1749            }
1750
1751            let updated_accounts = sqlx::query(
1752                "UPDATE tollgate_accounts AS account
1753                 SET settlement_loss = account.settlement_loss + delta.forfeited
1754                 FROM UNNEST($1::bytea[], $2::bigint[]) AS delta(account_id, forfeited)
1755                 WHERE account.account_id = delta.account_id",
1756            )
1757            .bind(&account_ids)
1758            .bind(&account_forfeits)
1759            .execute(&mut *tx)
1760            .await
1761            .map_err(storage)?;
1762            if updated_accounts.rows_affected() != expected_account_rows {
1763                return Err(StoreError(format!(
1764                    "reclaim updated {} of {} locked account rows",
1765                    updated_accounts.rows_affected(),
1766                    account_ids.len()
1767                )));
1768            }
1769
1770            Ok(batch)
1771        }
1772        .await;
1773        finish_transaction(tx, result).await
1774    }
1775}
1776
1777#[async_trait]
1778impl UsageSink for PostgresStore {
1779    async fn ingest(
1780        &self,
1781        events: &[UsageEvent],
1782        _now: Timestamp,
1783    ) -> Result<IngestReport, IngestError> {
1784        // One transaction per *batch* (review finding GL-8): leases are locked
1785        // in a single sorted ANY() query (sorted to keep concurrent batches
1786        // deadlock-free), duplicates are detected with one lookup, events are
1787        // classified in memory against the locked rows, and the accepted set
1788        // lands via one bulk insert plus set-wise per-lease/per-account
1789        // updates. Classification in application code preserves the partial
1790        // acceptance contract without savepoints.
1791        let mut report = IngestReport {
1792            unattributed: Some(0),
1793            ..IngestReport::default()
1794        };
1795        if events.is_empty() {
1796            return Ok(report);
1797        }
1798
1799        // Encode each ID and timestamp exactly once before opening the
1800        // transaction. The prepared values are reused by the lock, dedup,
1801        // classification, insert, and aggregate-update phases. Units are
1802        // checked once later, after duplicate and capability classification, to
1803        // preserve the partial-acceptance ordering.
1804        struct PreparedEvent<'a> {
1805            event: &'a UsageEvent,
1806            request_id: Vec<u8>,
1807            account_id: Vec<u8>,
1808            key_id: Option<Vec<u8>>,
1809            /// `None` for overage, which names no lease. Every phase below
1810            /// keys off this rather than re-matching on the source, so a lease
1811            /// id can never be conjured for an event that has none.
1812            lease_id: Option<Vec<u8>>,
1813            occurred_at_us: i64,
1814        }
1815        let prepared: Vec<PreparedEvent<'_>> = events
1816            .iter()
1817            .map(|event| {
1818                Ok(PreparedEvent {
1819                    event,
1820                    request_id: id_bytes(event.request_id.0),
1821                    account_id: id_bytes(event.account_id.0),
1822                    key_id: event.key_id.map(|id| id_bytes(id.0)),
1823                    lease_id: event.source.lease_id().map(|id| id_bytes(id.0)),
1824                    occurred_at_us: ts_micros(event.occurred_at),
1825                })
1826            })
1827            .collect::<Result<_, StoreError>>()?;
1828
1829        let mut tx = self.pool.begin().await.map_err(storage)?;
1830        let result: Result<_, IngestError> = async {
1831            // Lock every referenced lease in one global account/lease order.
1832            // A cycle needs two waiters, and this is the path that waits:
1833            // release touches a single lease, and reclaim takes its leases
1834            // with SKIP LOCKED, so it abandons a contended row instead of
1835            // queueing behind it. Concurrent ingests are therefore what this
1836            // order is for -- reclaim selects in expiry order (GL-65) and is
1837            // still safe, because the lease-then-account phase order below is
1838            // what keeps the two from crossing.
1839            let lease_ids: Vec<Vec<u8>> = prepared
1840                .iter()
1841                .filter_map(|event| event.lease_id.clone())
1842                .collect::<std::collections::BTreeSet<_>>()
1843                .into_iter()
1844                .collect();
1845            struct LeaseRow {
1846                account_id: Vec<u8>,
1847                fence: i64,
1848                granted: i64,
1849                used: i64,
1850                used_delta: Option<NonZeroI64>,
1851                credited: i64,
1852                settled: bool,
1853            }
1854            let rows = sqlx::query(
1855                "SELECT lease_id, account_id, fencing_token, granted, used, credited, state
1856                 FROM tollgate_leases
1857                 WHERE lease_id = ANY($1)
1858                 ORDER BY account_id, lease_id FOR UPDATE",
1859            )
1860            .bind(&lease_ids)
1861            .fetch_all(&mut *tx)
1862            .await
1863            .map_err(storage)?;
1864            let mut leases: std::collections::BTreeMap<Vec<u8>, LeaseRow> =
1865                std::collections::BTreeMap::new();
1866            for row in rows {
1867                let lease_id: Vec<u8> = row.get(0);
1868                let fence: i64 = row.get(2);
1869                let granted: i64 = row.get(3);
1870                let used: i64 = row.get(4);
1871                let credited: i64 = row.get(5);
1872                leases.insert(
1873                    lease_id,
1874                    LeaseRow {
1875                        account_id: row.get(1),
1876                        fence,
1877                        granted,
1878                        used,
1879                        used_delta: None,
1880                        credited,
1881                        settled: row.get::<i16, _>(6) != STATE_ACTIVE,
1882                    },
1883                );
1884            }
1885
1886            // Existing request ids in one lookup.
1887            let request_ids: Vec<Vec<u8>> = prepared
1888                .iter()
1889                .map(|event| event.request_id.clone())
1890                .collect();
1891            let mut seen: std::collections::HashSet<Vec<u8>> = sqlx::query(
1892                "SELECT request_id FROM tollgate_usage_events WHERE request_id = ANY($1)",
1893            )
1894            .bind(&request_ids)
1895            .fetch_all(&mut *tx)
1896            .await
1897            .map_err(storage)?
1898            .into_iter()
1899            .map(|row| row.get::<Vec<u8>, _>(0))
1900            .collect();
1901
1902            // Which accounts referenced by *overage* events actually exist.
1903            //
1904            // A leased event proves its account exists by resolving a lease
1905            // row, whose `account_id` is a foreign key. An overage event has
1906            // no such proof, and the memory backend rejects one naming an
1907            // unknown account -- so this backend must too, or the two
1908            // classify the same batch differently.
1909            //
1910            // Read without `FOR UPDATE`, deliberately. The accounts these
1911            // events touch are locked in account order further down, after
1912            // the leases, and taking that lock here instead would invert the
1913            // lease-then-account order every writing transaction agrees on.
1914            // A row that vanished between this probe and that lock would fail
1915            // the lock's own count check and roll the batch back, which is the
1916            // fail-closed outcome; nothing in this store deletes accounts.
1917            let overage_account_ids: Vec<Vec<u8>> = prepared
1918                .iter()
1919                .filter(|event| event.lease_id.is_none())
1920                .map(|event| event.account_id.clone())
1921                .collect::<std::collections::BTreeSet<_>>()
1922                .into_iter()
1923                .collect();
1924            let known_overage_accounts: std::collections::HashSet<Vec<u8>> =
1925                if overage_account_ids.is_empty() {
1926                    std::collections::HashSet::new()
1927                } else {
1928                    sqlx::query(
1929                        "SELECT account_id FROM tollgate_accounts WHERE account_id = ANY($1)",
1930                    )
1931                    .bind(&overage_account_ids)
1932                    .fetch_all(&mut *tx)
1933                    .await
1934                    .map_err(storage)?
1935                    .into_iter()
1936                    .map(|row| row.get::<Vec<u8>, _>(0))
1937                    .collect()
1938                };
1939
1940            // Classify in memory against the locked rows (identical rules to
1941            // MemoryStore: capability triple, then the conservation fit that
1942            // also converts a released lease's provisional loss into billed usage).
1943            struct Accepted {
1944                event_index: usize,
1945                settled: bool,
1946                /// The lease's stored fence, already validated positive;
1947                /// acceptance required the event's token to equal it, so this
1948                /// is the event's fence in storage form with no reconversion.
1949                /// `None` for overage, whose row carries no capability.
1950                fence: Option<i64>,
1951                /// True when these units were extended as unfunded credit, so
1952                /// they fund the ledger as well as bill it.
1953                overage: bool,
1954                /// Checked once during classification and reused by every
1955                /// aggregate and insert array.
1956                units: i64,
1957            }
1958            let mut accepted: Vec<Accepted> = Vec::with_capacity(prepared.len());
1959            for (event_index, event) in prepared.iter().enumerate() {
1960                if seen.contains(event.request_id.as_slice()) {
1961                    report.duplicate += 1;
1962                    continue;
1963                }
1964                let Some(lease_key) = event.lease_id.as_deref() else {
1965                    // Overage: no capability to verify and no lease capacity
1966                    // to fit inside, so the only question is whether the
1967                    // account exists. It is accepted regardless of the
1968                    // account's current enforcement mode -- the ledger does
1969                    // not record which mode a request was admitted under, and
1970                    // discarding a charge because the account was switched
1971                    // back to `Strict` after the work ran would be fail-open
1972                    // on accounting.
1973                    if !known_overage_accounts.contains(event.account_id.as_slice()) {
1974                        report.rejected += 1;
1975                        continue;
1976                    }
1977                    let Ok(units) = i64::try_from(event.event.units.get()) else {
1978                        // Caller data outside this backend's storage domain
1979                        // rejects this event, not the valid neighboring work.
1980                        report.rejected += 1;
1981                        continue;
1982                    };
1983                    seen.insert(event.request_id.clone());
1984                    accepted.push(Accepted {
1985                        event_index,
1986                        // Overage belongs to no lease, so there is no
1987                        // provisional settlement loss for it to convert and
1988                        // it is settled the moment it is recorded.
1989                        settled: false,
1990                        fence: None,
1991                        overage: true,
1992                        units,
1993                    });
1994                    report.accepted += 1;
1995                    continue;
1996                };
1997                let Some(lease) = leases.get_mut(lease_key) else {
1998                    report.rejected += 1;
1999                    continue;
2000                };
2001                if Some(stored_fence(lease.fence)?) != event.event.source.fencing_token()
2002                    || lease.account_id.as_slice() != event.account_id.as_slice()
2003                {
2004                    report.rejected += 1;
2005                    continue;
2006                }
2007                to_units(lease.granted, "lease granted")?;
2008                to_units(lease.used, "lease used")?;
2009                to_units(lease.credited, "lease credited")?;
2010                let Ok(units) = i64::try_from(event.event.units.get()) else {
2011                    report.rejected += 1;
2012                    continue;
2013                };
2014                let committed = lease
2015                    .used
2016                    .checked_add(lease.used_delta.map_or(0, NonZeroI64::get))
2017                    .and_then(|used| used.checked_add(lease.credited))
2018                    .ok_or_else(|| {
2019                        StoreError(format!(
2020                            "lease accounting overflow for {:#034x}",
2021                            id_from(lease_key)
2022                        ))
2023                    })?;
2024                let remaining = lease.granted.checked_sub(committed).ok_or_else(|| {
2025                    StoreError(format!(
2026                        "lease accounting exceeds grant for {:#034x}: granted {}, committed {committed}",
2027                        id_from(lease_key), lease.granted
2028                    ))
2029                })?;
2030                if units > remaining {
2031                    report.rejected += 1;
2032                    continue;
2033                }
2034                let used_delta = lease
2035                    .used_delta
2036                    .map_or(0, NonZeroI64::get)
2037                    .checked_add(units)
2038                    .ok_or_else(|| {
2039                        StoreError(format!(
2040                            "lease usage delta overflow for {:#034x}",
2041                            id_from(lease_key)
2042                        ))
2043                    })?;
2044                lease.used_delta = NonZeroI64::new(used_delta);
2045                seen.insert(event.request_id.clone());
2046                accepted.push(Accepted {
2047                    event_index,
2048                    settled: lease.settled,
2049                    fence: Some(lease.fence),
2050                    overage: false,
2051                    units,
2052                });
2053                report.accepted += 1;
2054            }
2055
2056            if accepted.is_empty() {
2057                return Ok(report);
2058            }
2059
2060            #[derive(Default)]
2061            struct AccountDelta {
2062                usage: i64,
2063                loss: i64,
2064                /// Units that fund themselves: overage bills and funds in the
2065                /// same transaction, so the equation closes by construction
2066                /// rather than by a later reconciliation step.
2067                overage: i64,
2068            }
2069
2070            // Bulk insert the accepted events.
2071            let (mut rid, mut acct, mut lease, mut fence, mut units, mut at, mut revision, mut keys) = (
2072                Vec::with_capacity(accepted.len()),
2073                Vec::with_capacity(accepted.len()),
2074                Vec::with_capacity(accepted.len()),
2075                Vec::with_capacity(accepted.len()),
2076                Vec::with_capacity(accepted.len()),
2077                Vec::with_capacity(accepted.len()),
2078                Vec::with_capacity(accepted.len()),
2079                Vec::with_capacity(accepted.len()),
2080            );
2081            let mut account_deltas: std::collections::BTreeMap<Vec<u8>, AccountDelta> =
2082                std::collections::BTreeMap::new();
2083            for accepted_event in &accepted {
2084                let event = &prepared[accepted_event.event_index];
2085                rid.push(event.request_id.clone());
2086                acct.push(event.account_id.clone());
2087                keys.push(event.key_id.clone());
2088                lease.push(event.lease_id.clone());
2089                fence.push(accepted_event.fence);
2090                units.push(accepted_event.units);
2091                at.push(event.occurred_at_us);
2092                // Carried verbatim: 32 bytes in, 32 bytes out, and the
2093                // schema's length CHECK says so. Unlike the identifiers
2094                // above this needs no width conversion — the Rust type is
2095                // already the stored representation.
2096                revision.push(event.event.policy_revision.as_bytes().to_vec());
2097                if accepted_event.overage {
2098                    debug_assert!(
2099                        event.lease_id.is_none() && accepted_event.fence.is_none(),
2100                        "an overage row must carry neither half of a capability"
2101                    );
2102                }
2103
2104                let entry = account_deltas.entry(event.account_id.clone()).or_default();
2105                entry.usage = entry.usage.checked_add(accepted_event.units).ok_or_else(|| {
2106                    IngestError::Refused(StoreError(format!(
2107                        "usage delta overflow for account {:#034x}",
2108                        event.event.account_id.0
2109                    )))
2110                })?;
2111                if accepted_event.settled {
2112                    entry.loss = entry.loss.checked_add(accepted_event.units).ok_or_else(|| {
2113                        IngestError::Refused(StoreError(format!(
2114                            "settlement loss delta overflow for account {:#034x}",
2115                            event.event.account_id.0
2116                        )))
2117                    })?;
2118                }
2119                if accepted_event.overage {
2120                    entry.overage =
2121                        entry.overage.checked_add(accepted_event.units).ok_or_else(|| {
2122                            IngestError::Refused(StoreError(format!(
2123                                "overage delta overflow for account {:#034x}",
2124                                event.event.account_id.0
2125                            )))
2126                        })?;
2127                }
2128            }
2129            let inserted = sqlx::query(
2130                "INSERT INTO tollgate_usage_events
2131                 (request_id, account_id, lease_id, fencing_token, units, occurred_at_us, policy_revision, key_id)
2132                 SELECT * FROM UNNEST($1::bytea[], $2::bytea[], $3::bytea[], $4::bigint[], $5::bigint[], $6::bigint[], $7::bytea[], $8::bytea[])",
2133            )
2134            .bind(&rid)
2135            .bind(&acct)
2136            .bind(&lease)
2137            .bind(&fence)
2138            .bind(&units)
2139            .bind(&at)
2140            .bind(&revision)
2141            .bind(&keys)
2142            .execute(&mut *tx)
2143            .await
2144            .map_err(storage)?;
2145            let expected_event_rows = u64::try_from(accepted.len())
2146                .map_err(|_| StoreError("accepted event count exceeds u64 range".into()))?;
2147            if inserted.rows_affected() != expected_event_rows {
2148                return Err(IngestError::Unavailable(StoreError(format!(
2149                    "ingest inserted {} of {} accepted usage rows",
2150                    inserted.rows_affected(),
2151                    accepted.len()
2152                ))));
2153            }
2154
2155            // Every lease row was already locked by the sorted SELECT above,
2156            // so one set-wise statement can apply the grouped usage deltas
2157            // without adding a round trip per distinct lease.
2158            let (lease_update_ids, lease_used_deltas): (Vec<Vec<u8>>, Vec<i64>) = leases
2159                .iter()
2160                .filter_map(|(lease_id, row)| {
2161                    row.used_delta
2162                        .map(|used_delta| (lease_id.clone(), used_delta.get()))
2163                })
2164                .unzip();
2165            if !lease_update_ids.is_empty() {
2166                let updated_leases = sqlx::query(
2167                    "UPDATE tollgate_leases AS lease
2168                     SET used = lease.used + delta.used
2169                     FROM UNNEST($1::bytea[], $2::bigint[]) AS delta(lease_id, used)
2170                     WHERE lease.lease_id = delta.lease_id",
2171                )
2172                .bind(&lease_update_ids)
2173                .bind(&lease_used_deltas)
2174                .execute(&mut *tx)
2175                .await
2176                .map_err(storage)?;
2177                let expected_lease_rows = u64::try_from(lease_update_ids.len())
2178                    .map_err(|_| StoreError("ingest lease row count exceeds u64 range".into()))?;
2179                if updated_leases.rows_affected() != expected_lease_rows {
2180                    return Err(IngestError::Unavailable(StoreError(format!(
2181                        "ingest updated {} of {} locked lease rows",
2182                        updated_leases.rows_affected(),
2183                        lease_update_ids.len()
2184                    ))));
2185                }
2186            }
2187
2188            let mut account_ids = Vec::with_capacity(account_deltas.len());
2189            let mut account_usage_deltas = Vec::with_capacity(account_deltas.len());
2190            let mut account_loss_deltas = Vec::with_capacity(account_deltas.len());
2191            let mut account_overage_deltas = Vec::with_capacity(account_deltas.len());
2192            for (account_id, delta) in &account_deltas {
2193                account_ids.push(account_id.clone());
2194                account_usage_deltas.push(delta.usage);
2195                account_loss_deltas.push(delta.loss);
2196                account_overage_deltas.push(delta.overage);
2197            }
2198
2199            // A set-wise UPDATE does not promise row-lock order. Lock every
2200            // affected account explicitly in byte order first. BTreeMap made
2201            // account_ids sorted, preserving ingest's lease-then-account
2202            // order and preventing concurrent multi-account batches from
2203            // forming a deadlock cycle.
2204            let locked_accounts = sqlx::query(
2205                "SELECT account_id, usage_recorded, settlement_loss, overage_recorded
2206                 FROM tollgate_accounts
2207                 WHERE account_id = ANY($1)
2208                 ORDER BY account_id FOR UPDATE",
2209            )
2210            .bind(&account_ids)
2211            .fetch_all(&mut *tx)
2212            .await
2213            .map_err(storage)?;
2214            if locked_accounts.len() != account_ids.len() {
2215                return Err(IngestError::Unavailable(StoreError(format!(
2216                    "ingest locked {} of {} referenced account rows",
2217                    locked_accounts.len(),
2218                    account_ids.len()
2219                ))));
2220            }
2221
2222            // Validate every account before the set-wise mutation. The
2223            // per-lease fit check bounds each straggler by the provisional
2224            // loss its own release recorded, so a short account loss means
2225            // ledger corruption and the entire transaction must roll back.
2226            for row in locked_accounts {
2227                let account_id: Vec<u8> = row.get(0);
2228                let usage_recorded: i64 = row.get(1);
2229                let settlement_loss: i64 = row.get(2);
2230                let overage_recorded: i64 = row.get(3);
2231                to_units(usage_recorded, "account usage_recorded")?;
2232                to_units(settlement_loss, "account settlement_loss")?;
2233                to_units(overage_recorded, "account overage_recorded")?;
2234                let delta = account_deltas.get(&account_id).ok_or_else(|| {
2235                    StoreError(format!(
2236                        "ingest locked unexpected account {:#034x}",
2237                        id_from(&account_id)
2238                    ))
2239                })?;
2240                usage_recorded.checked_add(delta.usage).ok_or_else(|| {
2241                    IngestError::Refused(StoreError(format!(
2242                        "usage_recorded overflow for account {:#034x}",
2243                        id_from(&account_id)
2244                    )))
2245                })?;
2246                // The funding half of the same units. Both terms must be
2247                // representable or neither may move, or the batch would bill
2248                // overage it did not fund and leave the equation open.
2249                overage_recorded.checked_add(delta.overage).ok_or_else(|| {
2250                    IngestError::Refused(StoreError(format!(
2251                        "overage_recorded overflow for account {:#034x}",
2252                        id_from(&account_id)
2253                    )))
2254                })?;
2255                if settlement_loss < delta.loss {
2256                    return Err(IngestError::Unavailable(StoreError(format!(
2257                        "settlement_loss underflow for account {:#034x}: settled straggler \
2258                         usage {} exceeds recorded loss",
2259                        id_from(&account_id),
2260                        delta.loss
2261                    ))));
2262                }
2263            }
2264
2265            let updated_accounts = sqlx::query(
2266                "UPDATE tollgate_accounts AS account
2267                 SET usage_recorded = account.usage_recorded + delta.usage,
2268                     settlement_loss = account.settlement_loss - delta.loss,
2269                     overage_recorded = account.overage_recorded + delta.overage
2270                 FROM UNNEST($1::bytea[], $2::bigint[], $3::bigint[], $4::bigint[])
2271                      AS delta(account_id, usage, loss, overage)
2272                 WHERE account.account_id = delta.account_id
2273                   AND account.settlement_loss >= delta.loss",
2274            )
2275            .bind(&account_ids)
2276            .bind(&account_usage_deltas)
2277            .bind(&account_loss_deltas)
2278            .bind(&account_overage_deltas)
2279            .execute(&mut *tx)
2280            .await
2281            .map_err(storage)?;
2282            let expected_account_rows = u64::try_from(account_ids.len())
2283                .map_err(|_| StoreError("ingest account row count exceeds u64 range".into()))?;
2284            if updated_accounts.rows_affected() != expected_account_rows {
2285                return Err(IngestError::Unavailable(StoreError(format!(
2286                    "ingest updated {} of {} locked account rows",
2287                    updated_accounts.rows_affected(),
2288                    account_ids.len()
2289                ))));
2290            }
2291
2292            // Legacy/unscoped batches cannot name activity. Their complete
2293            // attribution result is already known, so avoid an empty database
2294            // round trip while the ledger transaction still commits normally.
2295            if keys.iter().all(Option::is_none) {
2296                report.unattributed = Some(expected_event_rows);
2297                return Ok(report);
2298            }
2299
2300            // Only newly accepted rows supply evidence. Count events before
2301            // grouping: two events for one key are both attributable, even
2302            // if neither advances an already-newer maximum. Sort writes to
2303            // share one activity-lock order across concurrent transactions.
2304            let attributed: i64 = sqlx::query_scalar(
2305                "WITH matched AS MATERIALIZED (
2306                    SELECT k.key_id, b.occurred_at_us
2307                    FROM UNNEST($1::bytea[], $2::bytea[], $3::bigint[])
2308                        AS b(key_id, account_id, occurred_at_us)
2309                    JOIN LATERAL (
2310                        SELECT key_id FROM tollgate_credential_keys
2311                        WHERE key_id = b.key_id AND account_id = b.account_id LIMIT 1
2312                    ) k ON true
2313                 ), updated AS (
2314                    INSERT INTO tollgate_credential_activity AS activity (key_id, last_committed_at_us)
2315                    SELECT key_id, MAX(occurred_at_us) FROM matched GROUP BY key_id ORDER BY key_id
2316                    ON CONFLICT (key_id) DO UPDATE
2317                    SET last_committed_at_us = EXCLUDED.last_committed_at_us
2318                    WHERE activity.last_committed_at_us < EXCLUDED.last_committed_at_us
2319                    RETURNING key_id
2320                 ) SELECT COUNT(*) FROM matched"
2321            ).bind(&keys).bind(&acct).bind(&at).fetch_one(&mut *tx).await.map_err(storage)?;
2322            report.unattributed = Some(expected_event_rows.checked_sub(
2323                u64::try_from(attributed).map_err(|_| StoreError("negative attribution count".into()))?
2324            ).ok_or_else(|| StoreError("attribution count exceeds accepted events".into()))?);
2325
2326            Ok(report)
2327        }
2328        .await;
2329        // Monotonic accounting overflow is a permanent batch refusal. Database
2330        // and stored-corruption errors remain retryable, and rollback must
2331        // complete before either outcome becomes observable.
2332        finish_transaction(tx, result).await
2333    }
2334}
2335
2336/// Decode a stored status column into the enum, refusing anything the
2337/// vocabulary does not contain.
2338///
2339/// Never defaults to `Active`. A value the `CHECK` constraint should have made
2340/// impossible means the row was written outside this code, and admitting it as
2341/// "active" would turn corruption into service (the rule INVARIANTS.md GL-11
2342/// applies to the ledger's numbers, applied to its status).
2343fn decode_status(stored: String) -> Result<AccountStatus, StoreError> {
2344    match stored.as_str() {
2345        s if s == AccountStatus::Active.as_str() => Ok(AccountStatus::Active),
2346        s if s == AccountStatus::Suspended.as_str() => Ok(AccountStatus::Suspended),
2347        s if s == AccountStatus::Closed.as_str() => Ok(AccountStatus::Closed),
2348        other => Err(StoreError(format!("unrecognized account status {other:?}"))),
2349    }
2350}
2351
2352/// The ledger's execution-capacity class, or a refusal for a spelling the
2353/// vocabulary does not contain (GL-99).
2354///
2355/// Never defaults to `Assured`, for the reason [`decode_status`] never defaults
2356/// to `Active`. A value the `CHECK` constraint should have made impossible
2357/// means the row was written outside this code, and admitting it as `Assured`
2358/// would turn corruption into unconditional capacity — the permissive answer,
2359/// arrived at by accident.
2360fn decode_capacity_class(stored: String) -> Result<CapacityClass, StoreError> {
2361    match stored.as_str() {
2362        s if s == CapacityClass::Assured.as_str() => Ok(CapacityClass::Assured),
2363        s if s == CapacityClass::BestEffort.as_str() => Ok(CapacityClass::BestEffort),
2364        other => Err(StoreError(format!("unrecognized capacity class {other:?}"))),
2365    }
2366}
2367
2368/// Rebuild an account's schedule from its three stored columns.
2369///
2370/// `None` is "no schedule", and it is only reachable when all three are NULL:
2371/// the row's all-or-nothing CHECK makes half a schedule unstorable, so a
2372/// half-populated row read here is corruption to report rather than a shape to
2373/// interpret. Unknown names are refused for the same reason `decode_status`
2374/// refuses them — a period this binary cannot evaluate must not be presented
2375/// as if it were monthly.
2376fn decode_schedule(
2377    allowance: Option<i64>,
2378    period: Option<String>,
2379    rollover: Option<String>,
2380) -> Result<Option<BudgetSchedule>, StoreError> {
2381    let populated = [allowance.is_some(), period.is_some(), rollover.is_some()];
2382    let (Some(allowance), Some(period), Some(rollover)) = (allowance, period, rollover) else {
2383        if populated.iter().any(|present| *present) {
2384            return Err(StoreError(
2385                "stored budget schedule is partially populated".into(),
2386            ));
2387        }
2388        return Ok(None);
2389    };
2390    let period = match period.as_str() {
2391        s if s == Period::UtcCalendarMonth.as_str() => Period::UtcCalendarMonth,
2392        other => return Err(StoreError(format!("unrecognized budget period {other:?}"))),
2393    };
2394    let rollover = match rollover.as_str() {
2395        s if s == Rollover::None.as_str() => Rollover::None,
2396        other => {
2397            return Err(StoreError(format!(
2398                "unrecognized budget rollover {other:?}"
2399            )));
2400        }
2401    };
2402    Ok(Some(BudgetSchedule {
2403        allowance: to_units(allowance, "budget allowance")?,
2404        period,
2405        rollover,
2406    }))
2407}
2408
2409/// The ledger row [`budget_view`] decodes, for allocator evidence read inside
2410/// the grant's own transaction.
2411const BUDGET_VIEW_SQL: &str = "SELECT account_id, deposited, overage_recorded, usage_recorded,
2412            settlement_loss, expired, budget_allowance, budget_period,
2413            budget_rollover, period_start_us
2414     FROM tollgate_accounts WHERE account_id = $1";
2415
2416/// Debit a grant and return the committed-to-be row [`budget_view`] decodes,
2417/// in the same column order as [`BUDGET_VIEW_SQL`].
2418const GRANT_DEBIT_SQL: &str = "UPDATE tollgate_accounts
2419     SET balance = balance - $2,
2420         allowance_balance = allowance_balance - $3,
2421         next_fence = next_fence + 1
2422     WHERE account_id = $1
2423     RETURNING account_id, deposited, overage_recorded, usage_recorded,
2424               settlement_loss, expired, budget_allowance, budget_period,
2425               budget_rollover, period_start_us";
2426
2427/// What the account could still spend, for a snapshot's budget view (GL-97).
2428///
2429/// Balance *plus* the unspent remainder of every active lease, because units
2430/// out on lease are still the account's. Derived from the account row alone:
2431/// the conservation equation makes `balance + active grants` equal to what the
2432/// account was funded with minus what it consumed, so no join over
2433/// `tollgate_leases` is needed at every publication to answer the same
2434/// number.
2435///
2436/// Reads the columns of the `FOR SHARE` row in `publish_snapshot`, positions 1
2437/// through 9. Corruption is reported rather than saturated: unlike
2438/// `MemoryStore`, this backend does not exclusively own the ledger it reads,
2439/// so an underflow here is the same class of event as a negative unit column
2440/// (INVARIANTS.md GL-11) — see the note in `PostgresStore::conservation`.
2441fn budget_view(row: &sqlx::postgres::PgRow) -> Result<BudgetView, StoreError> {
2442    let deposited = to_units(row.get::<i64, _>(1), "deposited")?;
2443    let overage = to_units(row.get::<i64, _>(2), "overage_recorded")?;
2444    let usage = to_units(row.get::<i64, _>(3), "usage_recorded")?;
2445    let loss = to_units(row.get::<i64, _>(4), "settlement_loss")?;
2446    let expired = to_units(row.get::<i64, _>(5), "expired")?;
2447    let schedule = decode_schedule(row.get(6), row.get(7), row.get(8))?;
2448    let period_start = micros_ts(row.get::<i64, _>(9), "period_start_us")?;
2449
2450    let funded = deposited
2451        .checked_add(overage)
2452        .ok_or_else(|| StoreError("account funding total overflows".into()))?;
2453    let consumed = usage
2454        .checked_add(loss)
2455        .and_then(|spent| spent.checked_add(expired))
2456        .ok_or_else(|| StoreError("account consumption total overflows".into()))?;
2457    Ok(BudgetView {
2458        balance_at_publish: funded.checked_sub(consumed).ok_or_else(|| {
2459            StoreError(format!(
2460                "consumption {} exceeds funding {}",
2461                consumed.get(),
2462                funded.get()
2463            ))
2464        })?,
2465        // The period the account is *in*, which is what it can spend against.
2466        // An account whose schedule was set but whose first rollover has not
2467        // run yet is still in its previous period, and says so, until the next
2468        // sweep tick moves it.
2469        period_end: schedule.map(|schedule| schedule.period.end_after(period_start)),
2470    })
2471}
2472
2473/// Decode one stored snapshot row into the publication proof.
2474///
2475/// Shared by `SnapshotSource::snapshot` and the status republish so the latter
2476/// reads back what it wrote through exactly the reader's path -- a row this
2477/// refuses is one the request path would refuse too, and finding that out at
2478/// write time is the point.
2479///
2480/// Since GL-54 that is literally true of the generation as well: both callers
2481/// hand this the row's `generation` column, so a live read and a republish
2482/// resolve it identically. Before, the JSON carried a second copy that only the
2483/// live path consulted.
2484fn decode_publishable(
2485    principal: Principal,
2486    generation: i64,
2487    value: serde_json::Value,
2488) -> Result<PublishableSnapshot, StoreError> {
2489    let generation = generation_from(generation)?;
2490    let snapshot: StoredSnapshot =
2491        serde_json::from_value(value).map_err(|e| StoreError(format!("snapshot decode: {e}")))?;
2492    // Carried across the rebuild rather than through the builder: the builder
2493    // has no setter for a budget on purpose, so that a publisher cannot supply
2494    // one. Re-attaching what this store itself wrote is the store writing it
2495    // again, which is the same rule and not an exception to it.
2496    let budget = snapshot.budget;
2497    let publishable = PublishableSnapshot::try_new(Arc::new(snapshot.into_snapshot(generation)))
2498        .map_err(|error| {
2499            StoreError(format!(
2500                "invalid stored snapshot for principal {:#034x}: {error}",
2501                principal.0
2502            ))
2503        })?;
2504    Ok(match budget {
2505        Some(budget) => publishable.with_budget(Some(budget)),
2506        None => publishable,
2507    })
2508}
2509
2510/// Read a stored generation column.
2511///
2512/// Takes the raw `i64` and converts here rather than at each call site, so
2513/// every caller shares one refusal. That matters for the republish
2514/// specifically: it must not fail an account-wide status change over one
2515/// unreadable row, and a caller-side conversion is one `?` away from doing
2516/// exactly that.
2517///
2518/// A negative value is corruption to surface, never to clamp (INVARIANTS GL-11).
2519/// Migration 0007's CHECK is what keeps it from being written in the first
2520/// place; this is the read-side backstop.
2521fn generation_from(column: i64) -> Result<Generation, StoreError> {
2522    u64::try_from(column)
2523        .map(Generation)
2524        .map_err(|_| StoreError("stored snapshot generation is negative".into()))
2525}
2526
2527#[async_trait]
2528impl SnapshotSource for PostgresStore {
2529    async fn snapshot(&self, principal: Principal) -> Result<SnapshotResolution, StoreError> {
2530        let row = sqlx::query(
2531            "SELECT generation, snapshot, deleted FROM tollgate_snapshots WHERE principal = $1",
2532        )
2533        .bind(id_bytes(principal.0))
2534        .fetch_optional(&self.pool)
2535        .await
2536        .map_err(storage)?;
2537        match row {
2538            // Both branches now resolve the generation from the same column.
2539            // They did not before: the tombstone read it here while a live read
2540            // decoded a second copy out of the JSON, so one row could answer
2541            // two different generations depending on which branch you reached
2542            // (GL-54).
2543            Some(row) if row.get::<bool, _>(2) => Ok(SnapshotResolution::Revoked {
2544                generation: generation_from(row.get::<i64, _>(0))?,
2545            }),
2546            Some(row) => Ok(SnapshotResolution::Present(decode_publishable(
2547                principal,
2548                row.get::<i64, _>(0),
2549                row.get(1),
2550            )?)),
2551            None => Ok(SnapshotResolution::Unknown),
2552        }
2553    }
2554
2555    fn subscribe(&self) -> broadcast::Receiver<SnapshotPush> {
2556        self.push.subscribe()
2557    }
2558
2559    /// Tombstones included, for the reason `MemoryStore::principals` gives:
2560    /// forgetting a revoked principal is how one gets resurrected.
2561    ///
2562    /// A primary-key scan, so no new index — the table holds one row per
2563    /// principal, not per request, and `ORDER BY` makes the result stable so
2564    /// a caller diffing two enumerations sees real changes rather than
2565    /// storage order.
2566    async fn principals(&self) -> Result<Option<Vec<Principal>>, StoreError> {
2567        let rows = sqlx::query("SELECT principal FROM tollgate_snapshots ORDER BY principal")
2568            .fetch_all(&self.pool)
2569            .await
2570            .map_err(storage)?;
2571        Ok(Some(
2572            rows.iter()
2573                .map(|row| Principal(id_from(row.get::<Vec<u8>, _>(0).as_slice())))
2574                .collect(),
2575        ))
2576    }
2577}
2578
2579#[async_trait]
2580impl StoreHealth for PostgresStore {
2581    async fn ping(&self) -> Result<(), StoreError> {
2582        sqlx::query("SELECT 1")
2583            .execute(&self.pool)
2584            .await
2585            .map_err(storage)?;
2586        Ok(())
2587    }
2588}
2589
2590/// Re-stamp every live snapshot of `account`, patching one JSON key, and
2591/// report which principals to push and how many rows could not be decoded.
2592///
2593/// Shared by the two account-owned facts that republish — status (GL-51) and
2594/// execution-capacity class (GL-99). They differ in their precondition and their
2595/// ledger column; everything below is identical, and it is the part where the
2596/// subtlety lives, so it is written once.
2597///
2598/// `jsonb_set` rather than read-modify-write in Rust, for three reasons any
2599/// one of which decides it:
2600///
2601/// 1. RMW reintroduces GL-51's own bug. A concurrent `publish_snapshot` landing
2602///    between the read and the write makes `generation + 1` no longer greater
2603///    than stored, and the monotonic guard then *silently drops the change*
2604///    for that principal.
2605/// 2. RMW loses fields. `StoredSnapshot` has no `flatten`, so decoding and
2606///    re-serialising a row written by a newer binary discards what this one
2607///    does not know about.
2608/// 3. RMW fails whole on one bad row. A safety operation must not be blockable
2609///    by one unrelated corrupt credential.
2610///
2611/// Only the named key is patched. The generation lives in the column alone
2612/// (GL-54), and `RETURNING generation` carries the new value out to the push.
2613///
2614/// `deleted = FALSE` leaves tombstones alone: republishing one would resurrect
2615/// a revoked principal (INVARIANTS.md GL-15), and revocation stays its own
2616/// per-credential mechanism. `IS DISTINCT FROM` makes a repeat converge,
2617/// bumping nothing — and note that a document predating the key has SQL NULL
2618/// there, so the first change of a newly added fact rewrites every row once.
2619async fn republish_patched_snapshots(
2620    tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
2621    account: AccountId,
2622    json_path: &'static str,
2623    value: &str,
2624) -> Result<(Vec<(Principal, PublishableSnapshot)>, usize), SetStatusError> {
2625    let rows = sqlx::query(
2626        "UPDATE tollgate_snapshots
2627            SET generation = generation + 1,
2628                snapshot   = jsonb_set(snapshot, $3::text[], to_jsonb($2::text))
2629          WHERE account_id = $1
2630            AND deleted = FALSE
2631            AND snapshot #>> $3::text[] IS DISTINCT FROM $2::text
2632        RETURNING principal, generation, snapshot",
2633    )
2634    .bind(id_bytes(account.0))
2635    .bind(value)
2636    .bind(json_path)
2637    .fetch_all(&mut **tx)
2638    .await
2639    .map_err(storage)?;
2640
2641    let mut republished = Vec::with_capacity(rows.len());
2642    let mut unreadable = 0usize;
2643    for row in rows {
2644        let principal = Principal(id_from(row.get::<Vec<u8>, _>(0).as_slice()));
2645        match decode_publishable(
2646            principal,
2647            row.get::<i64, _>(1),
2648            row.get::<serde_json::Value, _>(2),
2649        ) {
2650            Ok(snapshot) => republished.push((principal, snapshot)),
2651            // Skipped for push, not fatal: this row was already unreadable
2652            // before the change touched it, and refusing to suspend or
2653            // reclassify an account because one of its credentials is corrupt
2654            // is the worse outcome. Reported, never silent (INVARIANTS.md
2655            // GL-19). Counted as well as logged: the row changed durably but
2656            // will not be pushed, so those principals converge only at their
2657            // next refresh.
2658            Err(error) => {
2659                unreadable += 1;
2660                tracing::warn!(
2661                    %principal,
2662                    %error,
2663                    "restamped snapshot could not be decoded for push"
2664                );
2665            }
2666        }
2667    }
2668    // Ordered, so both backends emit the same sequence and a mirrored test
2669    // need not assert on incidental ordering.
2670    republished.sort_unstable_by_key(|(principal, _)| *principal);
2671    Ok((republished, unreadable))
2672}
2673#[async_trait]
2674impl AdminStore for PostgresStore {
2675    async fn create_account(
2676        &self,
2677        config: AccountConfig,
2678    ) -> Result<tollgate_store::AdminReceipt<()>, CreateAccountError> {
2679        let result = sqlx::query(
2680            "INSERT INTO tollgate_accounts
2681             (account_id, balance, deposited, status, capacity_class, next_fence,
2682              usage_recorded, settlement_loss, overage_recorded)
2683             VALUES ($1, $2, $2, $3, $4, 1, 0, 0, 0)
2684             ON CONFLICT (account_id) DO NOTHING",
2685        )
2686        .bind(id_bytes(config.account_id.0))
2687        .bind(to_i64(config.initial_balance, "balance").map_err(CreateAccountError::Storage)?)
2688        .bind(config.status.as_str())
2689        .bind(config.capacity_class.as_str())
2690        .execute(&self.pool)
2691        .await
2692        .map_err(|e| CreateAccountError::Storage(storage(e)))?;
2693        if result.rows_affected() == 0 {
2694            return Err(CreateAccountError::AlreadyExists);
2695        }
2696        Ok(AdminReceipt::new(
2697            (),
2698            AdminState::Absent,
2699            AdminState::AccountCreated {
2700                initial_balance: config.initial_balance,
2701                status: config.status,
2702                capacity_class: config.capacity_class,
2703            },
2704        ))
2705    }
2706
2707    async fn deposit(
2708        &self,
2709        account: AccountId,
2710        units: CostUnits,
2711    ) -> Result<AdminReceipt<()>, AllocateError> {
2712        let row = sqlx::query(
2713            "UPDATE tollgate_accounts
2714             SET balance = balance + $2, deposited = deposited + $2
2715             WHERE account_id = $1
2716             RETURNING balance - $2 AS old_topup, deposited - $2 AS old_deposited,
2717                       balance AS new_topup, deposited AS new_deposited",
2718        )
2719        .bind(id_bytes(account.0))
2720        .bind(i64::try_from(units.get()).map_err(|_| AllocateError::BalanceOverflow)?)
2721        .fetch_optional(&self.pool)
2722        .await
2723        .map_err(|error| {
2724            // Only this deposit's arithmetic refusal is permanent. Connection,
2725            // constraint and other database failures retain Storage semantics.
2726            // The single UPDATE rolls back both counters on numeric overflow.
2727            if error
2728                .as_database_error()
2729                .and_then(|db| db.code())
2730                .as_deref()
2731                == Some("22003")
2732            {
2733                AllocateError::BalanceOverflow
2734            } else {
2735                alloc_storage(error)
2736            }
2737        })?
2738        .ok_or(AllocateError::UnknownAccount)?;
2739        let state = |topup: &str, deposited: &str| -> Result<AdminState, StoreError> {
2740            Ok(AdminState::Funding {
2741                topup: to_units(row.get(topup), "audit topup")?,
2742                deposited: to_units(row.get(deposited), "audit deposited")?,
2743            })
2744        };
2745        Ok(AdminReceipt::new(
2746            (),
2747            state("old_topup", "old_deposited")?,
2748            state("new_topup", "new_deposited")?,
2749        ))
2750    }
2751
2752    async fn set_budget_schedule(
2753        &self,
2754        account: AccountId,
2755        schedule: Option<BudgetSchedule>,
2756    ) -> Result<AdminReceipt<()>, BudgetError> {
2757        // No deposit here. Were setting a schedule also a funding operation,
2758        // an operator correcting a mistyped allowance would fund the account
2759        // twice, and there would be no way to describe next month's budget
2760        // without paying it today. The first allowance arrives at the first
2761        // `roll_period` after this lands.
2762        //
2763        // All three columns are written together, NULL together, which is what
2764        // the row's all-or-nothing CHECK enforces: half a schedule satisfies
2765        // neither branch of the rollover.
2766        let allowance = schedule
2767            .map(|s| to_i64(s.allowance, "budget allowance"))
2768            .transpose()
2769            .map_err(BudgetError::Storage)?;
2770        let mut tx = self.pool.begin().await.map_err(storage)?;
2771        let result = async {
2772            // A self-join's `previous` alias retains its statement snapshot
2773            // after waiting for a writer. Lock first so the receipt observes
2774            // the committed predecessor, then update under that same lock.
2775            let row = sqlx::query(
2776                "SELECT budget_allowance, budget_period, budget_rollover
2777                 FROM tollgate_accounts WHERE account_id = $1 FOR UPDATE",
2778            )
2779            .bind(id_bytes(account.0))
2780            .fetch_optional(&mut *tx)
2781            .await
2782            .map_err(storage)?
2783            .ok_or(BudgetError::UnknownAccount)?;
2784            let before = decode_schedule(row.get(0), row.get(1), row.get(2))?;
2785            sqlx::query(
2786                "UPDATE tollgate_accounts
2787                 SET budget_allowance = $2, budget_period = $3, budget_rollover = $4
2788                 WHERE account_id = $1",
2789            )
2790            .bind(id_bytes(account.0))
2791            .bind(allowance)
2792            .bind(schedule.map(|s| s.period.as_str()))
2793            .bind(schedule.map(|s| s.rollover.as_str()))
2794            .execute(&mut *tx)
2795            .await
2796            .map_err(storage)?;
2797            Ok(AdminReceipt::new(
2798                (),
2799                AdminState::Budget { schedule: before },
2800                AdminState::Budget { schedule },
2801            ))
2802        }
2803        .await;
2804        finish_transaction(tx, result).await
2805    }
2806
2807    async fn account_view(&self, account: AccountId) -> Result<Option<AccountView>, StoreError> {
2808        // One `REPEATABLE READ, READ ONLY` snapshot over the account row and
2809        // its live leases, for the reason `conservation` takes one (GL-56): the
2810        // stored totals and the sums over active leases move together in a
2811        // single `ingest` transaction, so reading them under separate
2812        // snapshots can pair a pre-write total with a post-write sum and
2813        // report corruption on a correct ledger. Status and schedule join that
2814        // same snapshot here, so a view cannot straddle a suspension or a
2815        // rollover and describe a state the account was never in.
2816        let mut tx = self.pool.begin().await.map_err(storage)?;
2817        let result = async {
2818            sqlx::query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
2819                .execute(&mut *tx)
2820                .await
2821                .map_err(storage)?;
2822            let account_row = sqlx::query(
2823                "SELECT deposited, balance, usage_recorded, settlement_loss, overage_recorded,
2824                        expired, status, capacity_class, budget_allowance, budget_period,
2825                        budget_rollover, period_start_us
2826                 FROM tollgate_accounts WHERE account_id = $1",
2827            )
2828            .bind(id_bytes(account.0))
2829            .fetch_optional(&mut *tx)
2830            .await
2831            .map_err(storage)?;
2832            let lease_row = sqlx::query(ACTIVE_LEASE_SUM_SQL)
2833                .bind(id_bytes(account.0))
2834                .fetch_one(&mut *tx)
2835                .await
2836                .map_err(storage)?;
2837            Ok::<_, StoreError>((account_row, lease_row))
2838        }
2839        .await;
2840        let (account_row, lease_row) = finish_transaction(tx, result).await?;
2841
2842        let Some(row) = account_row else {
2843            return Ok(None);
2844        };
2845        let active_grants = to_units(lease_row.get::<i64, _>(0), "active lease grants")?;
2846        let active_used = to_units(lease_row.get::<i64, _>(1), "active lease usage")?;
2847        let recorded = to_units(row.get::<i64, _>(2), "usage_recorded")?;
2848        Ok(Some(AccountView {
2849            account_id: account,
2850            status: decode_status(row.get::<String, _>(6))?,
2851            capacity_class: decode_capacity_class(row.get::<String, _>(7))?,
2852            schedule: decode_schedule(
2853                row.get::<Option<i64>, _>(8),
2854                row.get::<Option<String>, _>(9),
2855                row.get::<Option<String>, _>(10),
2856            )?,
2857            period_start: StoredInstant {
2858                micros: row.get::<i64, _>(11),
2859                submicro_nanos: 0,
2860            }
2861            .timestamp()?,
2862            conservation: Conservation {
2863                deposited: to_units(row.get::<i64, _>(0), "deposited")?,
2864                overage_recorded: to_units(row.get::<i64, _>(4), "overage_recorded")?,
2865                balance: to_units(row.get::<i64, _>(1), "balance")?,
2866                active_lease_grants: active_grants,
2867                // Surfaced rather than panicked on, as `conservation` does:
2868                // two stored columns from a database this process does not
2869                // exclusively own, so this is corruption to report.
2870                settled_usage: recorded.checked_sub(active_used).ok_or_else(|| {
2871                    StoreError(format!(
2872                        "active lease usage {} exceeds recorded usage {} for account {account}",
2873                        active_used.get(),
2874                        recorded.get()
2875                    ))
2876                })?,
2877                settlement_loss: to_units(row.get::<i64, _>(3), "settlement_loss")?,
2878                expired: to_units(row.get::<i64, _>(5), "expired")?,
2879            },
2880        }))
2881    }
2882
2883    async fn roll_due_periods(
2884        &self,
2885        now: Timestamp,
2886        limit: NonZeroUsize,
2887    ) -> Result<RolloverBatch, StoreError> {
2888        let limit_i = i64::try_from(limit.get())
2889            .map_err(|_| StoreError(format!("rollover batch limit exceeds i64 range: {limit}")))?;
2890        let mut rolled = Vec::new();
2891        // One statement per period kind, because the boundary is a property of
2892        // the period: a weekly schedule and a monthly one are due at different
2893        // instants, and a single comparison could only be right for one of
2894        // them. The exhaustive match in `every_period_is_swept` is what makes
2895        // a new variant a compile error here rather than a silently unswept
2896        // schedule.
2897        for period in Period::ALL {
2898            let boundary_us = ts_micros(period.start_of(now));
2899            let remaining = limit_i
2900                - i64::try_from(rolled.len())
2901                    .map_err(|_| StoreError("rollover batch row count exceeds i64 range".into()))?;
2902            if remaining <= 0 {
2903                break;
2904            }
2905            // One statement, so the selection and the crossing cannot be
2906            // separated. `FOR UPDATE SKIP LOCKED` is what lets replicas
2907            // cooperate: a concurrent pass skips the rows this one holds, and
2908            // once this commits their `period_start_us < boundary` test is
2909            // false — so a boundary is crossed exactly once however many
2910            // passes race it.
2911            //
2912            // The prior `allowance_balance` is read in the CTE because the
2913            // UPDATE cannot return it: PostgreSQL's RETURNING sees the new row
2914            // only, and `expired` has to be reported as the delta it is.
2915            let rows = sqlx::query(DUE_PERIODS_SQL)
2916                .bind(period.as_str())
2917                .bind(boundary_us)
2918                .bind(remaining)
2919                .fetch_all(&self.pool)
2920                .await
2921                .map_err(storage)?;
2922
2923            for row in rows {
2924                rolled.push((
2925                    row.get::<i64, _>(3),
2926                    RolledAccount {
2927                        account_id: AccountId(id_from(&row.get::<Vec<u8>, _>(0))),
2928                        deposited: to_units(row.get::<i64, _>(1), "budget allowance")?,
2929                        expired: to_units(row.get::<i64, _>(2), "expiring allowance")?,
2930                    },
2931                ));
2932            }
2933        }
2934        // Oldest boundary first, account id breaking ties, which is the order
2935        // `MemoryStore` reports and therefore the one this backend owes. The
2936        // rows arrive unordered from `RETURNING` and, across more than one
2937        // period kind, in per-period groups; sorting the assembled page is what
2938        // makes the report a function of the stored state rather than of the
2939        // planner. Bounded by `limit`, not by how many accounts are due.
2940        rolled
2941            .sort_unstable_by_key(|(crossed_from, account)| (*crossed_from, account.account_id.0));
2942        RolloverBatch::try_new(
2943            rolled.into_iter().map(|(_, account)| account).collect(),
2944            limit,
2945        )
2946    }
2947
2948    async fn set_account_status(
2949        &self,
2950        account: AccountId,
2951        status: AccountStatus,
2952    ) -> Result<tollgate_store::AdminReceipt<StatusChange>, SetStatusError> {
2953        let mut tx = self.pool.begin().await.map_err(storage)?;
2954        let result = async {
2955            // The account row first, and its lock is the serialization point:
2956            // two concurrent status changes cannot interleave their snapshot
2957            // updates, so "ledger says Active, snapshots say Suspended" is
2958            // unrepresentable rather than merely unlikely.
2959            let row = sqlx::query(
2960                "SELECT status FROM tollgate_accounts WHERE account_id = $1 FOR UPDATE",
2961            )
2962            .bind(id_bytes(account.0))
2963            .fetch_optional(&mut *tx)
2964            .await
2965            .map_err(storage)?
2966            .ok_or(SetStatusError::UnknownAccount)?;
2967
2968            let before = decode_status(row.get::<String, _>(0))?;
2969            if before == AccountStatus::Closed && status != AccountStatus::Closed {
2970                // Terminal, and the refusal changes nothing: no ledger write,
2971                // no generation bump. Rolling back here is what makes that so.
2972                return Err(SetStatusError::AccountClosed);
2973            }
2974
2975            sqlx::query("UPDATE tollgate_accounts SET status = $2 WHERE account_id = $1")
2976                .bind(id_bytes(account.0))
2977                .bind(status.as_str())
2978                .execute(&mut *tx)
2979                .await
2980                .map_err(storage)?;
2981
2982            let (republished, unreadable) =
2983                republish_patched_snapshots(&mut tx, account, "{status}", status.as_str()).await?;
2984            Ok((before, republished, unreadable))
2985        }
2986        .await;
2987
2988        let (before, republished, unreadable) = finish_transaction(tx, result).await?;
2989        // Publish only after commit, as with snapshot publication itself.
2990        if pushes_exceed_capacity(republished.len()) {
2991            tracing::warn!(
2992                %account,
2993                principals = republished.len(),
2994                capacity = PUSH_CHANNEL_CAPACITY,
2995                "status change emitted more pushes than the channel holds; subscribers will resync"
2996            );
2997        }
2998        let count = republished.len();
2999        for (principal, snapshot) in republished {
3000            self.push_to_subscribers(SnapshotPush {
3001                principal,
3002                resolution: SnapshotResolution::Present(snapshot),
3003            });
3004        }
3005        Ok(AdminReceipt::new(
3006            StatusChange {
3007                // Rows that changed durably: the ones pushed, plus any that could
3008                // not be decoded to push. The caller is told both numbers.
3009                republished: count + unreadable,
3010                unreadable,
3011            },
3012            AdminState::Status { status: before },
3013            AdminState::Status { status },
3014        ))
3015    }
3016
3017    async fn set_capacity_class(
3018        &self,
3019        account: AccountId,
3020        class: CapacityClass,
3021    ) -> Result<tollgate_store::AdminReceipt<StatusChange>, SetStatusError> {
3022        let mut tx = self.pool.begin().await.map_err(storage)?;
3023        let result = async {
3024            // `FOR UPDATE` for the reason `set_account_status` takes it: the
3025            // account row's lock is the serialization point, so two concurrent
3026            // class changes cannot interleave their snapshot updates and leave
3027            // the ledger saying one thing while some snapshots say another.
3028            // It also serialises against a concurrent status change, which is
3029            // what keeps the two account-owned facts from racing each other's
3030            // republications.
3031            let row = sqlx::query(
3032                "SELECT status, capacity_class FROM tollgate_accounts WHERE account_id = $1 FOR UPDATE",
3033            )
3034            .bind(id_bytes(account.0))
3035            .fetch_optional(&mut *tx)
3036            .await
3037            .map_err(storage)?
3038            .ok_or(SetStatusError::UnknownAccount)?;
3039
3040            let before = decode_capacity_class(row.get::<String, _>(1))?;
3041            // Closed is terminal, so reclassifying is meaningless. Unlike a
3042            // status change there is no "already at the target" escape: every
3043            // class is equally meaningless on a closed account. Rolling back
3044            // here is what makes the refusal change nothing.
3045            if decode_status(row.get::<String, _>(0))? == AccountStatus::Closed {
3046                return Err(SetStatusError::AccountClosed);
3047            }
3048
3049            sqlx::query("UPDATE tollgate_accounts SET capacity_class = $2 WHERE account_id = $1")
3050                .bind(id_bytes(account.0))
3051                .bind(class.as_str())
3052                .execute(&mut *tx)
3053                .await
3054                .map_err(storage)?;
3055
3056            let (republished, unreadable) =
3057                republish_patched_snapshots(&mut tx, account, "{capacity_class}", class.as_str()).await?;
3058            Ok((before, republished, unreadable))
3059        }
3060        .await;
3061
3062        let (before, republished, unreadable) = finish_transaction(tx, result).await?;
3063        if pushes_exceed_capacity(republished.len()) {
3064            tracing::warn!(
3065                %account,
3066                principals = republished.len(),
3067                capacity = PUSH_CHANNEL_CAPACITY,
3068                "capacity class change emitted more pushes than the channel holds; \
3069                 subscribers will resync"
3070            );
3071        }
3072        let count = republished.len();
3073        for (principal, snapshot) in republished {
3074            self.push_to_subscribers(SnapshotPush {
3075                principal,
3076                resolution: SnapshotResolution::Present(snapshot),
3077            });
3078        }
3079        Ok(AdminReceipt::new(
3080            StatusChange {
3081                // Rows that changed durably, whether or not they could be decoded
3082                // for a push — the same accounting `set_account_status` reports.
3083                republished: count + unreadable,
3084                unreadable,
3085            },
3086            AdminState::CapacityClass {
3087                capacity_class: before,
3088            },
3089            AdminState::CapacityClass {
3090                capacity_class: class,
3091            },
3092        ))
3093    }
3094
3095    async fn publish_snapshot(
3096        &self,
3097        principal: Principal,
3098        snapshot: PublishableSnapshot,
3099    ) -> Result<tollgate_store::AdminReceipt<()>, PublishSnapshotError> {
3100        let generation = i64::try_from(snapshot.generation.0).map_err(|_| {
3101            StoreError("snapshot generation exceeds PostgreSQL BIGINT range".into())
3102        })?;
3103
3104        let mut tx = self.pool.begin().await.map_err(storage)?;
3105        let result = publish_in_tx(&mut tx, principal, generation, snapshot).await;
3106
3107        let (written, published, before, after) = finish_transaction(tx, result).await?;
3108        if written {
3109            // The stamped snapshot, not the submitted one: a subscriber must
3110            // receive exactly what was stored, or a pushed instance and a
3111            // pulling one would disagree about the account's balance.
3112            self.push_to_subscribers(SnapshotPush {
3113                principal,
3114                resolution: SnapshotResolution::Present(published),
3115            });
3116        }
3117        Ok(AdminReceipt::new((), before, after))
3118    }
3119
3120    async fn remove_snapshot(&self, principal: Principal) -> Result<AdminReceipt<()>, StoreError> {
3121        let mut tx = self.pool.begin().await.map_err(storage)?;
3122        let result = remove_in_tx(&mut tx, principal).await;
3123        let receipt = finish_transaction(tx, result).await?;
3124        self.announce_removal(principal, &receipt);
3125        Ok(receipt)
3126    }
3127}
3128
3129/// Resolve `account`'s credential `key` to its principal and retirement,
3130/// holding the credential row `FOR SHARE` until the transaction ends (GL-143).
3131///
3132/// `FOR SHARE` conflicts with `revoke_key_audited`'s `FOR UPDATE`, so a
3133/// key-bound publication and a revocation serialize: the publish either sees
3134/// the retirement and is refused, or commits first.
3135///
3136/// **Lock order: credential, then account, then snapshot.** A cycle needs a
3137/// path that holds an account lock and then waits on a credential row lock.
3138/// None does:
3139/// - `revoke_key_audited` locks only the credential row, never an account.
3140/// - `insert_credential` takes the account `FOR UPDATE`, then inserts a *new*
3141///   credential row; it never locks an existing one.
3142/// - `ingest` locks accounts `FOR UPDATE`, then reads credentials with a plain
3143///   `SELECT` and writes activity rows whose foreign key takes only
3144///   `FOR KEY SHARE`, which `FOR SHARE` does not block.
3145/// - Every other account-locking path — `set_account_status`,
3146///   `set_capacity_class`, `set_budget_schedule`, `deposit`, rollover,
3147///   reclaim and lease acquisition — touches account, lease, usage and
3148///   snapshot rows only; `tollgate_credential_keys` appears in none of them.
3149///
3150/// After the credential row, this path takes the account and snapshot rows in
3151/// the same order `publish_snapshot` and `set_account_status` do.
3152async fn lock_account_key(
3153    tx: &mut Transaction<'_, Postgres>,
3154    account: AccountId,
3155    key: KeyId,
3156) -> Result<(Principal, bool), KeySnapshotError> {
3157    let row = sqlx::query(
3158        "SELECT principal, revoked_at_us IS NOT NULL
3159         FROM tollgate_credential_keys WHERE key_id = $1 AND account_id = $2 FOR SHARE",
3160    )
3161    .bind(id_bytes(key.0))
3162    .bind(id_bytes(account.0))
3163    .fetch_optional(&mut **tx)
3164    .await
3165    .map_err(storage)?
3166    .ok_or(KeySnapshotError::UnknownCredential)?;
3167    let principal: Vec<u8> = row.get(0);
3168    let principal: [u8; 16] = principal
3169        .try_into()
3170        .map_err(|_| StoreError("credential principal is not 16 bytes".into()))?;
3171    Ok((Principal(u128::from_be_bytes(principal)), row.get(1)))
3172}
3173
3174/// A publication's checks and write inside the caller's transaction: the
3175/// stated credential binding (GL-35), the ledger status and capacity-class
3176/// guards (GL-51, GL-99), the budget stamp (GL-97), and the generation-ordered
3177/// write. Returns whether a row was written, the stamped snapshot to push
3178/// after commit, and the audited predecessor and successor.
3179async fn publish_in_tx(
3180    tx: &mut Transaction<'_, Postgres>,
3181    principal: Principal,
3182    generation: i64,
3183    snapshot: PublishableSnapshot,
3184) -> Result<(bool, PublishableSnapshot, AdminState, AdminState), PublishSnapshotError> {
3185    if let Some(key_id) = snapshot.key_id {
3186        let matches: bool = sqlx::query_scalar(
3187            "SELECT EXISTS(SELECT 1 FROM tollgate_credential_keys
3188                 WHERE key_id = $1 AND principal = $2 AND account_id = $3)",
3189        )
3190        .bind(id_bytes(key_id.0))
3191        .bind(id_bytes(principal.0))
3192        .bind(id_bytes(snapshot.account_id.0))
3193        .fetch_one(&mut **tx)
3194        .await
3195        .map_err(storage)?;
3196        if !matches {
3197            return Err(PublishSnapshotError::CredentialMismatch { key_id });
3198        }
3199    }
3200    // The ledger decides an account's status; a publish may carry it
3201    // but not change it, or the two records `set_account_status`
3202    // unified could be pulled apart again one principal at a time
3203    // (GL-51). FOR SHARE, not FOR UPDATE: this only has to hold the
3204    // status still, and a status change takes FOR UPDATE on the same
3205    // row, so the two serialize without publishes blocking each other.
3206    //
3207    // The budget columns ride along on the read that was already being
3208    // taken, under the same lock, so the view stamped below is the
3209    // ledger as of this publication rather than a second read that
3210    // could straddle a lease or a rollover (GL-97).
3211    let ledger = sqlx::query(
3212        "SELECT status, deposited, overage_recorded, usage_recorded, settlement_loss,
3213                    expired, budget_allowance, budget_period, budget_rollover, period_start_us,
3214                    capacity_class
3215             FROM tollgate_accounts WHERE account_id = $1 FOR SHARE",
3216    )
3217    .bind(id_bytes(snapshot.account_id.0))
3218    .fetch_optional(&mut **tx)
3219    .await
3220    .map_err(storage)?;
3221
3222    // An account the ledger does not hold publishes unchanged: this
3223    // adds no account-existence requirement to publication, and an
3224    // unstamped snapshot reports no budget rather than a zero.
3225    // Unconditional, including the `None` arm: a snapshot arrives here
3226    // having crossed a wire, where nothing stops a publisher putting a
3227    // balance in the JSON. Overwriting always is what makes the store
3228    // the field's only writer, rather than only usually.
3229    let view = match &ledger {
3230        Some(row) => {
3231            let ledger = decode_status(row.get::<String, _>(0))?;
3232            if ledger != snapshot.status {
3233                return Err(PublishSnapshotError::StatusMismatch {
3234                    ledger,
3235                    submitted: snapshot.status,
3236                });
3237            }
3238            // The capacity class is the same kind of fact and gets the
3239            // same guard (GL-99): the ledger owns it, a publish may
3240            // carry it, and only `set_capacity_class` may change it.
3241            let ledger_class = decode_capacity_class(row.get::<String, _>(10))?;
3242            if ledger_class != snapshot.capacity_class {
3243                return Err(PublishSnapshotError::CapacityClassMismatch {
3244                    ledger: ledger_class,
3245                    submitted: snapshot.capacity_class,
3246                });
3247            }
3248            Some(budget_view(row)?)
3249        }
3250        None => None,
3251    };
3252    let published = snapshot.with_budget(view);
3253    let value = serde_json::to_value(StoredSnapshotRef::from(published.as_snapshot()))
3254        .map_err(|e| StoreError(format!("snapshot encode: {e}")))?;
3255
3256    let (written, before, after) = write_snapshot_audited(tx, principal, generation, value).await?;
3257    Ok((written, published, before, after))
3258}
3259
3260/// Tombstone a live snapshot inside the caller's transaction, keeping its
3261/// generation as the watermark.
3262async fn remove_in_tx(
3263    tx: &mut Transaction<'_, Postgres>,
3264    principal: Principal,
3265) -> Result<AdminReceipt<()>, StoreError> {
3266    let before = snapshot_audit_row(tx, principal).await?;
3267    let after = match before {
3268        AdminState::Snapshot {
3269            generation,
3270            revoked: false,
3271        } => {
3272            sqlx::query("UPDATE tollgate_snapshots SET deleted = TRUE WHERE principal = $1")
3273                .bind(id_bytes(principal.0))
3274                .execute(&mut **tx)
3275                .await
3276                .map_err(storage)?;
3277            AdminState::Snapshot {
3278                generation,
3279                revoked: true,
3280            }
3281        }
3282        state => state,
3283    };
3284    Ok(AdminReceipt::new((), before, after))
3285}
3286
3287/// Holding the snapshot row while both observing and replacing it makes the
3288/// receipt exact even when another server publishes or revokes concurrently.
3289async fn snapshot_audit_row(
3290    tx: &mut Transaction<'_, Postgres>,
3291    principal: Principal,
3292) -> Result<AdminState, StoreError> {
3293    let row = sqlx::query(
3294        "SELECT generation, deleted FROM tollgate_snapshots WHERE principal = $1 FOR UPDATE",
3295    )
3296    .bind(id_bytes(principal.0))
3297    .fetch_optional(&mut **tx)
3298    .await
3299    .map_err(storage)?;
3300    row.map(|row| {
3301        Ok(AdminState::Snapshot {
3302            generation: generation_from(row.get(0))?,
3303            revoked: row.get(1),
3304        })
3305    })
3306    .unwrap_or(Ok(AdminState::Absent))
3307}
3308
3309async fn write_snapshot_audited(
3310    tx: &mut Transaction<'_, Postgres>,
3311    principal: Principal,
3312    generation: i64,
3313    value: serde_json::Value,
3314) -> Result<(bool, AdminState, AdminState), StoreError> {
3315    let mut before = snapshot_audit_row(tx, principal).await?;
3316    let after = AdminState::Snapshot {
3317        generation: generation_from(generation)?,
3318        revoked: false,
3319    };
3320    if before == AdminState::Absent {
3321        let inserted = sqlx::query(
3322            "INSERT INTO tollgate_snapshots (principal, generation, snapshot, deleted)
3323            VALUES ($1, $2, $3, FALSE) ON CONFLICT (principal) DO NOTHING",
3324        )
3325        .bind(id_bytes(principal.0))
3326        .bind(generation)
3327        .bind(&value)
3328        .execute(&mut **tx)
3329        .await
3330        .map_err(storage)?;
3331        if inserted.rows_affected() == 1 {
3332            return Ok((true, before, after));
3333        }
3334        // Another creator won the unique constraint. READ COMMITTED sees its
3335        // row here; lock and observe that actual predecessor before replacing.
3336        before = snapshot_audit_row(tx, principal).await?;
3337    }
3338    let AdminState::Snapshot {
3339        generation: previous,
3340        ..
3341    } = before
3342    else {
3343        return Err(StoreError("snapshot disappeared during publication".into()));
3344    };
3345    if previous >= generation_from(generation)? {
3346        return Ok((false, before, before));
3347    }
3348    sqlx::query("UPDATE tollgate_snapshots SET generation = $2, snapshot = $3, deleted = FALSE WHERE principal = $1")
3349        .bind(id_bytes(principal.0)).bind(generation).bind(value)
3350        .execute(&mut **tx).await.map_err(storage)?;
3351    Ok((true, before, after))
3352}
3353
3354#[cfg(test)]
3355mod tests {
3356    use super::*;
3357
3358    #[test]
3359    fn stored_fences_use_the_exact_positive_bigint_domain() {
3360        for invalid in [i64::MIN, -1, 0] {
3361            assert!(stored_fence(invalid).is_err());
3362        }
3363        for valid in [1, 2, i64::MAX] {
3364            assert_eq!(stored_fence(valid).unwrap(), FencingToken(valid as u64));
3365        }
3366    }
3367
3368    /// A readiness probe is evidence that PostgreSQL answered, not merely that
3369    /// a store object exists (INVARIANTS.md GL-19).
3370    #[tokio::test]
3371    async fn ping_surfaces_a_closed_pool() {
3372        let pool = PgPoolOptions::new()
3373            .connect_lazy("postgres://localhost/tollgate")
3374            .unwrap();
3375        pool.close().await;
3376        let (push, _) = broadcast::channel(1);
3377        let store = PostgresStore {
3378            pool,
3379            policy: GrantPolicy::default(),
3380            push,
3381        };
3382
3383        assert!(store.ping().await.is_err());
3384        assert!(matches!(
3385            AdminStore::deposit(&store, AccountId(1), CostUnits(1)).await,
3386            Err(AllocateError::Storage(_))
3387        ));
3388    }
3389}