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