agentplane 0.38.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
//! Per-tenant quota accounting on `PostgreSQL`.
//!
//! This is the backend the guarantee is actually about. On a single node a
//! ceiling can be held up by almost anything; the moment two instances admit
//! concurrently, only the database can arbitrate — which is why the reservation
//! below is **one statement**, not a count followed by an insert.

use async_trait::async_trait;

use crate::core::{RunId, Spend, StoreError, Timestamp};
use crate::quota::{Halt, HaltScope, QuotaError, QuotaSettlement, QuotaStore};

use super::postgres::{PostgresStore, amount_of, be, pool_err, sql_amount};

#[async_trait]
impl QuotaStore for PostgresStore {
    fn tenant(&self) -> &str {
        crate::journal::JournalStore::tenant(self)
    }

    async fn reserve(
        &self,
        run: RunId,
        limit: Option<u32>,
        at: Timestamp,
    ) -> Result<(), QuotaError> {
        let mut client = self
            .pool_ref()
            .get()
            .await
            .map_err(|e| QuotaError::Unavailable(pool_err(&e).to_string()))?;
        let tenant = self.tenant_name();
        let run = run.to_string();

        let Some(limit) = limit else {
            // No ceiling: still record the run, so `running()` answers honestly
            // and a ceiling added later starts from the truth. No lock either —
            // there is no decision here for two admissions to disagree about.
            client
                .execute(
                    "INSERT INTO quota_running (tenant, run_id, admitted_at)
                     VALUES ($1, $2, $3)
                     ON CONFLICT (tenant, run_id) DO NOTHING",
                    &[&tenant, &run, &at.unix_timestamp()],
                )
                .await
                .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
            return Ok(());
        };

        // The count and the insert decide together **under a per-tenant
        // advisory lock**, because nothing weaker serialises them. This
        // statement's earlier form ran without the lock, on a comment claiming
        // the decision happened "inside the row lock the write takes" — and no
        // such lock exists: two INSERTs of *different* rows lock nothing in
        // common, each count subquery reads its own statement snapshot under
        // READ COMMITTED, and two admissions racing for one remaining slot
        // both passed `count < limit` and both landed. A ceiling that admits
        // limit+k exactly under concurrent load is the catalogued
        // yields-under-load shape, on the one control whose whole promise is
        // surviving a second instance.
        //
        // The lock is transaction-scoped and per tenant — admissions for one
        // tenant serialise, which is the semantic the ceiling requires; other
        // tenants' admissions do not wait. The length prefix keeps
        // `("acme", …)` from colliding with a tenant literally named
        // `"acme…"` under concatenation.
        //
        // `ON CONFLICT DO NOTHING` keeps a retried admission idempotent: the
        // run already holds its slot, so re-reserving must neither take a
        // second nor be refused against a ceiling it is already counted in.
        let tx = client
            .transaction()
            .await
            .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
        tx.query_one(
            "SELECT pg_advisory_xact_lock(hashtextextended($1, 0))",
            &[&format!("quota-admission:{}:{tenant}", tenant.len())],
        )
        .await
        .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
        let inserted = tx
            .execute(
                "INSERT INTO quota_running (tenant, run_id, admitted_at)
                 SELECT $1, $2, $3
                  WHERE (SELECT COUNT(*) FROM quota_running WHERE tenant = $1) < $4
                     OR EXISTS (SELECT 1 FROM quota_running
                                 WHERE tenant = $1 AND run_id = $2)
                 ON CONFLICT (tenant, run_id) DO NOTHING",
                &[&tenant, &run, &at.unix_timestamp(), &i64::from(limit)],
            )
            .await
            .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;

        if inserted == 1 {
            tx.commit()
                .await
                .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
            return Ok(());
        }

        // Nothing was written. Either the tenant is at its ceiling, or the run
        // already held a slot and `DO NOTHING` fired — and those are opposite
        // answers, so it is read back rather than assumed. Still inside the
        // transaction, so the counts are the ones the decision was made from.
        let held: i64 = tx
            .query_one(
                "SELECT COUNT(*) FROM quota_running WHERE tenant = $1 AND run_id = $2",
                &[&tenant, &run],
            )
            .await
            .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?
            .get(0);
        if held > 0 {
            tx.commit()
                .await
                .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
            return Ok(());
        }

        let running: i64 = tx
            .query_one(
                "SELECT COUNT(*) FROM quota_running WHERE tenant = $1",
                &[&tenant],
            )
            .await
            .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?
            .get(0);
        tx.commit()
            .await
            .map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
        Err(QuotaError::TooManyRuns {
            tenant,
            running: u32::try_from(running).unwrap_or(u32::MAX),
        })
    }

    async fn set_halt(&self, scope: &HaltScope, reason: Option<&str>) -> Result<(), StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let scope = scope.key();
        match reason {
            Some(reason) => {
                client
                    .execute(
                        "INSERT INTO quota_halted (tenant, scope, reason) VALUES ($1, $2, $3)
                         ON CONFLICT (tenant, scope) DO UPDATE SET reason = EXCLUDED.reason",
                        &[&self.tenant_name(), &scope, &reason],
                    )
                    .await
                    .map_err(|e| be(&e))?;
            }
            None => {
                client
                    .execute(
                        "DELETE FROM quota_halted WHERE tenant = $1 AND scope = $2",
                        &[&self.tenant_name(), &scope],
                    )
                    .await
                    .map_err(|e| be(&e))?;
            }
        }
        Ok(())
    }

    async fn halts(&self) -> Result<Vec<Halt>, StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let rows = client
            .query(
                "SELECT scope, reason FROM quota_halted WHERE tenant = $1 ORDER BY scope",
                &[&self.tenant_name()],
            )
            .await
            .map_err(|e| be(&e))?;
        rows.into_iter()
            .map(|row| {
                let stored: String = row.get(0);
                // Corruption, not a row to skip: a halt this build cannot read
                // is one it would run straight through, and from the outside
                // that is indistinguishable from a halt that was lifted.
                let scope = HaltScope::parse(&stored).ok_or_else(|| StoreError::Corrupt {
                    seq: 0,
                    detail: format!(
                        "quota_halted holds the scope '{stored}', which this build cannot \
                         read — refusing rather than admitting work an operator stopped"
                    ),
                })?;
                Ok(Halt {
                    scope,
                    reason: row.get(1),
                })
            })
            .collect()
    }

    async fn release(&self, run: RunId) -> Result<(), StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        client
            .execute(
                "DELETE FROM quota_running WHERE tenant = $1 AND run_id = $2",
                &[&self.tenant_name(), &run.to_string()],
            )
            .await
            .map_err(|e| be(&e))?;
        Ok(())
    }

    async fn settle(&self, settlement: &QuotaSettlement) -> Result<(), StoreError> {
        let mut client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let tx = client.transaction().await.map_err(|e| be(&e))?;
        let tenant = self.tenant_name();
        let run = settlement.run.to_string();
        let epoch = settlement.epoch.cast_signed();
        let tokens = sql_amount(settlement.spend.tokens);
        let minor_units = sql_amount(settlement.spend.minor_units);

        let inserted = tx
            .execute(
                "INSERT INTO quota_settled
                   (tenant, run_id, epoch, period, tokens, minor_units, release_slot)
                 VALUES ($1, $2, $3, $4, $5, $6, $7)
                 ON CONFLICT (tenant, run_id, epoch) DO NOTHING",
                &[
                    &tenant,
                    &run,
                    &epoch,
                    &settlement.period,
                    &tokens,
                    &minor_units,
                    &settlement.release_slot,
                ],
            )
            .await
            .map_err(|e| be(&e))?;

        if inserted == 0 {
            let stored = tx
                .query_one(
                    "SELECT period, tokens, minor_units, release_slot
                       FROM quota_settled
                      WHERE tenant = $1 AND run_id = $2 AND epoch = $3",
                    &[&tenant, &run, &epoch],
                )
                .await
                .map_err(|e| be(&e))?;
            let same = stored.get::<_, Option<String>>(0) == settlement.period
                && stored.get::<_, i64>(1) == tokens
                && stored.get::<_, i64>(2) == minor_units
                && stored.get::<_, bool>(3) == settlement.release_slot;
            if !same {
                return Err(StoreError::Corrupt {
                    seq: 0,
                    detail: format!(
                        "quota pass {run}/{} was settled twice with different payloads",
                        settlement.epoch
                    ),
                });
            }
            tx.commit().await.map_err(|e| be(&e))?;
            return Ok(());
        }

        if let Some(period) = settlement.period.as_deref() {
            // Cast before addition: BIGINT addition can overflow before LEAST
            // sees it. The numeric intermediate is exact and the stored total
            // saturates at the largest representable non-negative amount.
            tx.execute(
                "INSERT INTO quota_spent (tenant, period, tokens, minor_units)
                 VALUES ($1, $2, $3, $4)
                 ON CONFLICT (tenant, period) DO UPDATE SET
                   tokens = LEAST(9223372036854775807::numeric,
                                  quota_spent.tokens::numeric + EXCLUDED.tokens::numeric)::bigint,
                   minor_units = LEAST(9223372036854775807::numeric,
                                       quota_spent.minor_units::numeric + EXCLUDED.minor_units::numeric)::bigint",
                &[&tenant, &period, &tokens, &minor_units],
            )
            .await
            .map_err(|e| be(&e))?;
        }
        if settlement.release_slot {
            tx.execute(
                "DELETE FROM quota_running WHERE tenant = $1 AND run_id = $2",
                &[&tenant, &run],
            )
            .await
            .map_err(|e| be(&e))?;
        }
        tx.commit().await.map_err(|e| be(&e))?;
        Ok(())
    }

    async fn spent(&self, period: &str) -> Result<Spend, StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let row = client
            .query_opt(
                "SELECT tokens, minor_units FROM quota_spent
                  WHERE tenant = $1 AND period = $2",
                &[&self.tenant_name(), &period.to_owned()],
            )
            .await
            .map_err(|e| be(&e))?;
        Ok(row.map_or_else(Spend::default, |r| Spend {
            tokens: amount_of(r.get::<_, i64>(0)),
            minor_units: amount_of(r.get::<_, i64>(1)),
        }))
    }

    async fn running(&self) -> Result<u32, StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let n: i64 = client
            .query_one(
                "SELECT COUNT(*) FROM quota_running WHERE tenant = $1",
                &[&self.tenant_name()],
            )
            .await
            .map_err(|e| be(&e))?
            .get(0);
        Ok(u32::try_from(n).unwrap_or(u32::MAX))
    }

    async fn running_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError> {
        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
        let rows = client
            .query(
                "SELECT run_id FROM quota_running WHERE tenant = $1 \
                 ORDER BY run_id ASC LIMIT $2",
                &[
                    &self.tenant_name(),
                    &i64::try_from(limit).unwrap_or(i64::MAX),
                ],
            )
            .await
            .map_err(|e| be(&e))?;
        let mut out = Vec::with_capacity(rows.len());
        for row in rows {
            let raw: String = row.get(0);
            // Damage rather than absence — see the embedded backend for why a
            // skipped row is the worst available answer here.
            out.push(RunId::parse(&raw).map_err(|e| StoreError::Corrupt {
                seq: 0,
                detail: format!("bad run id '{raw}' in the quota slot table: {e}"),
            })?);
        }
        Ok(out)
    }
}