use std::time::Duration;
use async_trait::async_trait;
use deadpool_postgres::{Config, Pool, Runtime as PoolRuntime};
use tokio_postgres::NoTls;
use tokio_postgres::error::SqlState;
use crate::core::{Digest, EffectKey, Epoch, RunId, Seq, StoreError};
use std::sync::Arc;
use serde_json::Value;
use crate::journal::{
Append, AtomicJournal, AtomicTx, AtomicWork, Cancellation, Head, JournalStore, Lease, Record,
SqlValue,
};
const SCHEMA: &str = "
CREATE TABLE IF NOT EXISTS journal (
-- The tenant leads every key in this schema. A run id is unique, so a
-- filter would do; a key component is what makes a query that forgets the
-- predicate return nothing rather than another tenant's row, and what keeps
-- every index physically clustered per tenant so a scan cannot walk into
-- somebody else's traffic.
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
-- Every counter and instant in this schema carries a CHECK for the reason
-- `amount_of` gives: the types are unsigned, PostgreSQL has no unsigned
-- integer, and a row edited around the application — negative by typo or
-- by malice — must fail the edit rather than read back as a value that
-- reverses an accumulation. Chain positions and epochs start at one, so
-- their floor is 1, not 0.
seq BIGINT NOT NULL CHECK (seq >= 1),
epoch BIGINT NOT NULL CHECK (epoch >= 1),
kind TEXT NOT NULL,
effect_key TEXT,
prev_hash BYTEA NOT NULL,
hash BYTEA NOT NULL,
raw BYTEA NOT NULL,
-- Who wrote it. Null when no signer is configured: a hash chain still
-- detects edits, it just cannot say who made them.
key_id TEXT,
signature BYTEA,
-- Which matter this record belongs to, when it belongs to one.
--
-- A column and not only a field inside `raw`: *show me everything about
-- this matter* is answerable by a range scan here, and by a join over the
-- case's runs otherwise — a join that also **misses** every record written
-- by a run the case does not own, which is exactly what a sweep is.
case_id TEXT,
PRIMARY KEY (tenant, run_id, seq)
);
-- One matter's history, in order, without touching another tenant's.
CREATE INDEX IF NOT EXISTS journal_by_case
ON journal (tenant, case_id, run_id, seq)
WHERE case_id IS NOT NULL;
CREATE TABLE IF NOT EXISTS run_activity (
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
updated_at BIGINT NOT NULL CHECK (updated_at >= 0),
PRIMARY KEY (tenant, run_id)
);
CREATE INDEX IF NOT EXISTS run_activity_recent
ON run_activity (tenant, updated_at DESC, run_id DESC);
-- Exactly-once, as a constraint rather than a code path. A second
-- `EffectStarted` for one effect key in one run cannot be written, whichever
-- instance is writing and whatever it believes about the journal.
CREATE UNIQUE INDEX IF NOT EXISTS journal_effect_started
ON journal (tenant, run_id, effect_key)
WHERE kind = 'EffectStarted' AND effect_key IS NOT NULL;
CREATE TABLE IF NOT EXISTS run_lease (
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
owner TEXT NOT NULL,
epoch BIGINT NOT NULL CHECK (epoch >= 1),
expires_at BIGINT NOT NULL CHECK (expires_at >= 0),
PRIMARY KEY (tenant, run_id)
);
-- The recovery sweep's read: leases that expired without being released. A
-- partial index over held rows only, because released rows blank the owner and
-- dominate the table — every run ever admitted leaves one, since the epoch
-- lives in the row and deleting it would un-fence the run. Without the WHERE
-- clause this index would be mostly rows the query filters out.
CREATE INDEX IF NOT EXISTS run_lease_abandoned
ON run_lease (tenant, expires_at)
WHERE owner <> '';
-- An operator's stop request. Beside the chain rather than in it, and
-- deliberately not fenced: whoever wants a run stopped is not its owner, holds
-- no epoch, and is usually asking because the owner is busy. The primary key
-- makes the request idempotent, so a retry cannot overwrite the original asker.
-- The plane's Merkle log: one row per sealed run, positioned by a sequence.
--
-- A sequence rather than `MAX + 1`, and that is what makes a log possible on
-- this backend at all. Several instances seal concurrently here — that is the
-- topology this backend exists for — and `MAX + 1` computed by two transactions
-- at once hands both the same position. A sequence is monotonic under
-- concurrency and never reissues a number after a delete: exactly the two
-- properties the log needs.
CREATE SEQUENCE IF NOT EXISTS run_log_position;
-- One log per tenant, so a checkpoint commits to that tenant's runs and no
-- others. `log_index` is unique *within* a tenant: a shared sequence hands out
-- distinct numbers globally, and each tenant simply sees gaps, which the dense
-- rank in `inclusion_proof` removes before anything is proved.
CREATE TABLE IF NOT EXISTS run_seal (
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
chain_head BYTEA NOT NULL,
log_index BIGINT NOT NULL CHECK (log_index >= 1),
sealed_at BIGINT NOT NULL CHECK (sealed_at >= 0),
-- How the run ended *when it sealed*. Descriptive: the outcome index that
-- answers queries lives in run_outcome, fed by the chain's `RunSealed`
-- records, where the last conclusion wins.
outcome TEXT NOT NULL,
PRIMARY KEY (tenant, run_id),
UNIQUE (tenant, log_index)
);
-- Conclusion ordering under concurrency, for the same reason as
-- run_log_position: several instances conclude runs at once here.
CREATE SEQUENCE IF NOT EXISTS run_conclusion_ordinal;
-- Concluded runs by how they ended. **Derived**, not authoritative: fed by the
-- `RunSealed` record inside `append`, in the same transaction, so the outcome's
-- home stays the chain and this table can be rebuilt from it. The last
-- conclusion wins — a run that failed, was resumed and then succeeded moves
-- from `failed` to `succeeded`; an index fed by the seal would report the
-- first conclusion forever, and a wrong answer reads exactly like a right one.
CREATE TABLE IF NOT EXISTS run_outcome (
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
outcome TEXT NOT NULL,
ordinal BIGINT NOT NULL CHECK (ordinal >= 1),
PRIMARY KEY (tenant, run_id)
);
-- One outcome's backlog without touching another tenant's. Scanned **newest
-- first** — see `JournalStore::runs_by_outcome` for why the direction is part
-- of the contract rather than a detail of this index.
CREATE INDEX IF NOT EXISTS run_outcome_by_outcome
ON run_outcome (tenant, outcome, ordinal);
-- The admission keys this tenant has issued. **Derived**, not authoritative:
-- fed by the `RunAdmitted` record inside `append`, in the same transaction, so
-- the key's home stays the chain and this table can be rebuilt from it.
--
-- The primary key is the mechanism, not a hint. Two instances admitting one
-- inbound message concurrently both find nothing and both insert; exactly one
-- commits, and the loser is told which run won.
--
-- Tenant-scoped for the reason `case_correlation` is: an admission key is a
-- business value, and two tenants using the same one is ordinary.
CREATE TABLE IF NOT EXISTS run_admission (
tenant TEXT NOT NULL,
key TEXT NOT NULL,
run_id TEXT NOT NULL,
admitted_at BIGINT NOT NULL CHECK (admitted_at >= 0),
PRIMARY KEY (tenant, key)
);
CREATE TABLE IF NOT EXISTS run_cancel (
tenant TEXT NOT NULL,
run_id TEXT NOT NULL,
actor TEXT NOT NULL,
reason TEXT NOT NULL,
requested_at BIGINT NOT NULL CHECK (requested_at >= 0),
PRIMARY KEY (tenant, run_id)
);
-- A2A push uses the journal itself as the outbox. `next_seq` is the first task
-- record this receiver has not acknowledged; advancing only after a 2xx makes a
-- crash produce a duplicate rather than a lost notification.
CREATE TABLE IF NOT EXISTS push_delivery (
tenant TEXT NOT NULL,
task_id TEXT NOT NULL,
config_id TEXT NOT NULL,
url TEXT NOT NULL,
token TEXT,
auth_scheme TEXT,
auth_credentials TEXT,
next_seq BIGINT NOT NULL CHECK (next_seq >= 0),
attempts INTEGER NOT NULL DEFAULT 0 CHECK (attempts >= 0),
next_attempt_at BIGINT NOT NULL DEFAULT 0 CHECK (next_attempt_at >= 0),
last_error TEXT,
-- Stopped, with the cursor kept: a receiver that answered permanently or
-- outlasted the retry ceiling. Excluded from the due index so a parked row
-- costs nothing per sweep, and still names how far its receiver got.
parked BOOLEAN NOT NULL DEFAULT FALSE,
PRIMARY KEY (tenant, task_id, config_id)
);
CREATE INDEX IF NOT EXISTS push_delivery_due
ON push_delivery (tenant, next_attempt_at, task_id, config_id)
WHERE NOT parked;
";
#[derive(Debug, Clone)]
pub struct PostgresStore {
pool: Pool,
signer: Option<Arc<dyn crate::core::Signer>>,
origin: String,
tenant: crate::core::TenantId,
}
impl PostgresStore {
#[must_use]
pub fn signing_as(mut self, signer: Arc<dyn crate::core::Signer>) -> Self {
self.signer = Some(signer);
self
}
#[must_use]
pub fn origin(mut self, origin: impl Into<String>) -> Self {
self.origin = origin.into();
self
}
#[must_use]
pub fn for_tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.tenant = tenant;
self
}
fn log_origin(&self) -> String {
if self.tenant.as_str() == crate::core::TenantId::DEFAULT {
self.origin.clone()
} else {
format!("{}/{}", self.origin, self.tenant)
}
}
async fn log_leaves(&self) -> Result<Vec<crate::core::merkle::LeafHash>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT chain_head FROM run_seal WHERE tenant = $1 ORDER BY log_index ASC",
&[&self.tenant_name()],
)
.await
.map_err(|e| be(&e))?;
rows.iter()
.map(|r| {
digest_from(&r.get::<_, Vec<u8>>(0)).map(|d| crate::core::merkle::leaf_hash(&d))
})
.collect::<Result<Vec<_>, _>>()
}
}
pub(super) fn detail(error: &tokio_postgres::Error) -> String {
if let Some(db) = error.as_db_error() {
let extra = db.detail().map_or(String::new(), |d| format!(": {d}"));
return format!("{} ({}{})", db.message(), db.code().code(), extra);
}
let mut out = error.to_string();
let mut cause: Option<&(dyn std::error::Error + 'static)> = std::error::Error::source(error);
while let Some(next) = cause {
use std::fmt::Write as _;
let _ = write!(out, ": {next}");
cause = next.source();
}
out
}
pub(super) fn be(error: &tokio_postgres::Error) -> StoreError {
StoreError::Backend(detail(error))
}
fn commit_refused_or_in_doubt(e: &tokio_postgres::Error) -> StoreError {
if e.as_db_error().is_some() {
be(e)
} else {
StoreError::CommitUnknown {
detail: e.to_string(),
}
}
}
pub(super) fn pool_err(e: &impl std::fmt::Display) -> StoreError {
StoreError::Backend(e.to_string())
}
pub(super) const fn amount_of(v: i64) -> u64 {
if v < 0 { 0 } else { v.cast_unsigned() }
}
#[expect(
clippy::cast_possible_wrap,
reason = "guarded above: values past i64::MAX clamp rather than wrapping"
)]
pub(super) const fn sql_amount(v: u64) -> i64 {
if v > i64::MAX as u64 {
i64::MAX
} else {
v as i64
}
}
#[allow(clippy::disallowed_methods)]
fn now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs())
}
fn lease_expiry(now: u64, ttl: Duration) -> Result<u64, StoreError> {
let secs = ttl.as_secs();
if secs == 0 {
return Err(StoreError::Backend(format!(
"a lease TTL of {ttl:?} is below this store's whole-second \
granularity and would round to zero — pass at least one second, \
or use the runtime's lease_ttl which enforces its own minimum"
)));
}
now.checked_add(secs).ok_or_else(|| {
StoreError::Backend(format!(
"a lease TTL of {ttl:?} overflows the expiry instant — the lease \
would wrap into the past and read as already expired"
))
})
}
impl PostgresStore {
pub(super) fn pool_ref(&self) -> &Pool {
&self.pool
}
#[cfg(feature = "keyring")]
pub async fn erasure_probe(&self, scope: &str) -> Result<bool, crate::core::StoreError> {
let key = crate::keyring::coordinator::PostgresCoordinator::scope_key(scope);
let client = self
.pool
.get()
.await
.map_err(|e| crate::core::StoreError::Backend(e.to_string()))?;
let row = client
.query_one("SELECT pg_try_advisory_lock($1)", &[&key])
.await
.map_err(|e| crate::core::StoreError::Backend(e.to_string()))?;
let got: bool = row.get(0);
if got {
client
.execute("SELECT pg_advisory_unlock($1)", &[&key])
.await
.map_err(|e| crate::core::StoreError::Backend(e.to_string()))?;
}
Ok(!got)
}
#[cfg(feature = "keyring")]
#[must_use]
pub fn erasure_coordinator(&self) -> crate::keyring::coordinator::PostgresCoordinator {
crate::keyring::coordinator::PostgresCoordinator::new(self.pool.clone())
}
pub(super) fn tenant_name(&self) -> String {
self.tenant.to_string()
}
pub(super) fn tenant_str(&self) -> &str {
self.tenant.as_str()
}
pub async fn connect(url: &str) -> Result<Self, StoreError> {
Self::connect_sized(url, None).await
}
pub async fn connect_sized(
url: &str,
max_connections: Option<usize>,
) -> Result<Self, StoreError> {
let pg: tokio_postgres::Config = url
.parse()
.map_err(|e: tokio_postgres::Error| StoreError::Backend(e.to_string()))?;
match pg.get_ssl_mode() {
tokio_postgres::config::SslMode::Disable | tokio_postgres::config::SslMode::Prefer => {}
demanded => {
return Err(StoreError::Backend(format!(
"the connection URL demands TLS (sslmode {demanded:?}) and this \
build connects without it — connecting anyway would silently \
send the journal in plaintext. Use sslmode=disable (or prefer) \
against a trusted network, or terminate TLS in a proxy this \
store connects to locally"
)));
}
}
let mut cfg = Config::new();
cfg.host = pg.get_hosts().first().map(|h| match h {
tokio_postgres::config::Host::Tcp(s) => s.clone(),
#[cfg(unix)]
tokio_postgres::config::Host::Unix(p) => p.to_string_lossy().into_owned(),
});
cfg.port = pg.get_ports().first().copied();
cfg.user = pg.get_user().map(ToOwned::to_owned);
cfg.password = pg
.get_password()
.map(|p| String::from_utf8_lossy(p).into_owned());
cfg.dbname = pg.get_dbname().map(ToOwned::to_owned);
if let Some(max) = max_connections {
cfg.pool = Some(deadpool_postgres::PoolConfig::new(max));
}
let pool = cfg
.create_pool(Some(PoolRuntime::Tokio1), NoTls)
.map_err(|e| pool_err(&e))?;
let client = pool.get().await.map_err(|e| pool_err(&e))?;
client.batch_execute(SCHEMA).await.map_err(|e| be(&e))?;
client
.batch_execute(super::postgres_cases::CASE_SCHEMA)
.await
.map_err(|e| be(&e))?;
client
.batch_execute(super::postgres_authority::AUTHORITY_SCHEMA)
.await
.map_err(|e| be(&e))?;
client
.batch_execute(super::postgres_memory::MEMORY_SCHEMA)
.await
.map_err(|e| be(&e))?;
Ok(Self {
pool,
signer: None,
origin: "agentplane".to_owned(),
tenant: crate::core::TenantId::default(),
})
}
}
impl PostgresStore {
async fn fence(
&self,
tx: &deadpool_postgres::tokio_postgres::Transaction<'_>,
run: RunId,
epoch: Epoch,
) -> Result<(), StoreError> {
let lease = tx
.query_opt(
"SELECT epoch FROM run_lease WHERE tenant = $1 AND run_id = $2 FOR UPDATE",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
if let Some(row) = lease {
let current: i64 = row.get(0);
let current = current.cast_unsigned();
if epoch != current {
return Err(StoreError::Fenced {
run: run.to_string(),
held: epoch,
current,
});
}
}
Ok(())
}
async fn append_within(
&self,
tx: &deadpool_postgres::tokio_postgres::Transaction<'_>,
run: RunId,
epoch: Epoch,
batch: Vec<Append>,
) -> Result<Vec<Record>, StoreError> {
if let Some(a) = batch.iter().find(|a| a.run != run) {
return Err(StoreError::Backend(format!(
"batch spans runs {run} and {} — a batch is one run's atomic unit",
a.run
)));
}
let sealed_row = tx
.query_opt(
"SELECT outcome FROM run_seal WHERE tenant = $1 AND run_id = $2",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
if let Some(row) = sealed_row {
return Err(StoreError::RunSealed {
run: run.to_string(),
outcome: row.get(0),
});
}
let head = tx
.query_opt(
"SELECT seq, hash FROM journal
WHERE tenant = $1 AND run_id = $2 ORDER BY seq DESC LIMIT 1",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
let (mut seq, mut prev) = match head {
Some(row) => {
let s: i64 = row.get(0);
let h: Vec<u8> = row.get(1);
(s.cast_unsigned(), digest_from(&h)?)
}
None => (0, Digest::ZERO),
};
let mut sealed = Vec::with_capacity(batch.len());
let mut conclusion: Option<String> = None;
let mut claimed: Option<String> = None;
for append in batch {
seq += 1;
let body = append.into_body(seq, epoch);
match &body.kind {
crate::journal::RecordKind::RunConcluded { outcome, .. } => {
conclusion = Some(outcome.clone());
}
crate::journal::RecordKind::RunAdmitted {
idempotency_key: Some(key),
..
} => claimed = Some(key.clone()),
_ => {}
}
let record = Record::seal_signed(body, prev, self.signer.as_deref())?;
prev = record.hash;
let effect = record.effect_key().map(EffectKey::to_hex);
let result = tx
.execute(
"INSERT INTO journal
(tenant, run_id, seq, epoch, kind, effect_key, prev_hash, hash,
raw, key_id, signature, case_id)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12)",
&[
&self.tenant_name(),
&run.to_string(),
&seq.cast_signed(),
&epoch.cast_signed(),
&record.kind().kind_str(),
&effect,
&record.prev_hash.as_bytes().to_vec(),
&record.hash.as_bytes().to_vec(),
&record.raw().to_vec(),
&record.attestation.as_ref().map(|a| a.key_id.clone()),
&record.attestation.as_ref().map(|a| a.signature.clone()),
&record.body.case.map(|c| c.to_string()),
],
)
.await;
if let Err(e) = result {
return Err(refused_insert(&e, &record, run, seq));
}
sealed.push(record);
}
tx.execute(
"INSERT INTO run_activity (tenant, run_id, updated_at) VALUES ($1, $2, $3)
ON CONFLICT (tenant, run_id) DO UPDATE SET updated_at = EXCLUDED.updated_at",
&[
&self.tenant_name(),
&run.to_string(),
&now_secs().cast_signed(),
],
)
.await
.map_err(|error| be(&error))?;
self.derive_indexes(tx, run, claimed, conclusion).await?;
Ok(sealed)
}
async fn derive_indexes(
&self,
tx: &deadpool_postgres::tokio_postgres::Transaction<'_>,
run: RunId,
claimed: Option<String>,
conclusion: Option<String>,
) -> Result<(), StoreError> {
if let Some(key) = claimed {
let taken = tx
.execute(
"INSERT INTO run_admission (tenant, key, run_id, admitted_at)
VALUES ($1, $2, $3, $4)
ON CONFLICT (tenant, key) DO NOTHING",
&[
&self.tenant_name(),
&key,
&run.to_string(),
&now_secs().cast_signed(),
],
)
.await
.map_err(|error| be(&error))?;
if taken == 0 {
let holder = tx
.query_opt(
"SELECT run_id FROM run_admission WHERE tenant = $1 AND key = $2",
&[&self.tenant_name(), &key],
)
.await
.map_err(|error| be(&error))?
.map_or_else(String::new, |row| row.get(0));
return Err(StoreError::DuplicateAdmission { key, run: holder });
}
}
if let Some(outcome) = conclusion {
tx.execute(
"INSERT INTO run_outcome (tenant, run_id, outcome, ordinal)
VALUES ($1, $2, $3, nextval('run_conclusion_ordinal'))
ON CONFLICT (tenant, run_id) DO UPDATE
SET outcome = EXCLUDED.outcome, ordinal = EXCLUDED.ordinal",
&[&self.tenant_name(), &run.to_string(), &outcome],
)
.await
.map_err(|error| be(&error))?;
}
Ok(())
}
}
fn refused_insert(e: &tokio_postgres::Error, record: &Record, run: RunId, seq: u64) -> StoreError {
if e.code() == Some(&SqlState::UNIQUE_VIOLATION) {
let constraint = e
.as_db_error()
.and_then(|d| d.constraint())
.unwrap_or_default();
if constraint == "journal_effect_started"
&& let Some(k) = record.effect_key()
{
return StoreError::DuplicateEffect(k);
}
return StoreError::Backend(format!(
"append raced another writer on run {run} at seq {seq} \
(unique violation on '{constraint}') — two connections \
extended one chain concurrently, which the lease fence \
should have serialised"
));
}
be(e)
}
#[async_trait]
impl JournalStore for PostgresStore {
fn is_shared(&self) -> bool {
true
}
fn tenant(&self) -> &str {
self.tenant.as_str()
}
fn atomic(&self) -> Option<&dyn AtomicJournal> {
Some(self)
}
async fn append(&self, epoch: Epoch, batch: Vec<Append>) -> Result<Vec<Record>, StoreError> {
if batch.is_empty() {
return Ok(Vec::new());
}
let run = batch[0].run;
let mut client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let tx = client.transaction().await.map_err(|e| be(&e))?;
self.fence(&tx, run, epoch).await?;
let sealed = self.append_within(&tx, run, epoch, batch).await?;
tx.commit()
.await
.map_err(|e| commit_refused_or_in_doubt(&e))?;
Ok(sealed)
}
async fn read(&self, run: RunId, from: Seq) -> Result<Vec<Record>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT seq, prev_hash, hash, raw, key_id, signature FROM journal
WHERE tenant = $1 AND run_id = $2 AND seq >= $3 ORDER BY seq ASC",
&[&self.tenant_name(), &run.to_string(), &from.cast_signed()],
)
.await
.map_err(|e| be(&e))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let seq: i64 = row.get(0);
let prev: Vec<u8> = row.get(1);
let hash: Vec<u8> = row.get(2);
let raw: Vec<u8> = row.get(3);
let key_id: Option<String> = row.get(4);
let signature: Option<Vec<u8>> = row.get(5);
let _ = seq;
let attestation = key_id
.zip(signature)
.map(|(key_id, signature)| crate::core::Attestation { key_id, signature });
out.push(Record::from_stored_attested(
raw,
digest_from(&prev)?,
digest_from(&hash)?,
attestation,
)?);
}
Ok(out)
}
async fn admitted_as(&self, key: &str) -> Result<Option<RunId>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let row = client
.query_opt(
"SELECT run_id FROM run_admission WHERE tenant = $1 AND key = $2",
&[&self.tenant_name(), &key],
)
.await
.map_err(|e| be(&e))?;
row.map(|row| {
let id: &str = row.get(0);
RunId::parse(id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_admission holds an unparseable run id '{id}': {e}"),
})
})
.transpose()
}
async fn forget_admissions(
&self,
older_than: crate::core::Timestamp,
) -> Result<usize, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let removed = client
.execute(
"DELETE FROM run_admission WHERE tenant = $1 AND admitted_at < $2",
&[&self.tenant_name(), &older_than.unix_timestamp()],
)
.await
.map_err(|e| be(&e))?;
Ok(usize::try_from(removed).unwrap_or(usize::MAX))
}
async fn count_by_outcome(&self, outcome: &str) -> Result<u64, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let row = client
.query_one(
"SELECT count(*) FROM run_outcome WHERE tenant = $1 AND outcome = $2",
&[&self.tenant_name(), &outcome],
)
.await
.map_err(|e| be(&e))?;
let n: i64 = row.get(0);
u64::try_from(n).map_err(|_| StoreError::Corrupt {
seq: 0,
detail: format!("run_outcome counted {n} rows, which is not a count"),
})
}
async fn runs_by_outcome(&self, outcome: &str, limit: usize) -> Result<Vec<RunId>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT run_id FROM run_outcome
WHERE tenant = $1 AND outcome = $2
ORDER BY ordinal DESC
LIMIT $3",
&[
&self.tenant_name(),
&outcome,
&i64::try_from(limit).unwrap_or(i64::MAX),
],
)
.await
.map_err(|e| be(&e))?;
rows.iter()
.map(|r| {
let id: &str = r.get(0);
RunId::parse(id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_outcome holds an unparsable run id '{id}': {e}"),
})
})
.collect()
}
async fn recent_runs(
&self,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let cap = i64::try_from(limit).unwrap_or(i64::MAX);
let rows = match after {
Some((updated, run)) => {
client
.query(
"SELECT run_id, updated_at FROM run_activity
WHERE tenant = $1 AND (updated_at, run_id) < ($2, $3)
ORDER BY updated_at DESC, run_id DESC LIMIT $4",
&[
&self.tenant_name(),
&i64::try_from(updated).unwrap_or(i64::MAX),
&run.to_string(),
&cap,
],
)
.await
}
None => {
client
.query(
"SELECT run_id, updated_at FROM run_activity
WHERE tenant = $1
ORDER BY updated_at DESC, run_id DESC LIMIT $2",
&[&self.tenant_name(), &cap],
)
.await
}
}
.map_err(|error| be(&error))?;
rows.into_iter()
.map(|row| {
let id: String = row.get(0);
let updated: i64 = row.get(1);
Ok((
RunId::parse(&id).map_err(|error| StoreError::Backend(error.to_string()))?,
updated.cast_unsigned(),
))
})
.collect()
}
async fn case_history(
&self,
case: crate::core::CaseId,
limit: usize,
) -> Result<Vec<Record>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT raw, prev_hash, hash, key_id, signature FROM journal
WHERE tenant = $1 AND case_id = $2
ORDER BY run_id, seq ASC
LIMIT $3",
&[
&self.tenant_name(),
&case.to_string(),
&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: Vec<u8> = row.get(0);
let prev: Vec<u8> = row.get(1);
let hash: Vec<u8> = row.get(2);
let key_id: Option<String> = row.get(3);
let signature: Option<Vec<u8>> = row.get(4);
let attestation = key_id
.zip(signature)
.map(|(key_id, signature)| crate::core::Attestation { key_id, signature });
out.push(Record::from_stored_attested(
raw,
digest_from(&prev)?,
digest_from(&hash)?,
attestation,
)?);
}
Ok(out)
}
async fn head(&self, run: RunId) -> Result<Head, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let row = client
.query_opt(
"SELECT seq, hash FROM journal
WHERE tenant = $1 AND run_id = $2 ORDER BY seq DESC LIMIT 1",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
match row {
None => Ok(Head::genesis()),
Some(row) => {
let seq: i64 = row.get(0);
let hash: Vec<u8> = row.get(1);
Ok(Head {
seq: seq.cast_unsigned(),
hash: digest_from(&hash)?,
})
}
}
}
async fn acquire(&self, run: RunId, owner: &str, ttl: Duration) -> Result<Lease, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let now = now_secs();
let expires = lease_expiry(now, ttl)?;
let key = run.to_string();
let row = client
.query_opt(
"INSERT INTO run_lease (tenant, run_id, owner, epoch, expires_at)
VALUES ($1, $2, $3, 1, $4)
ON CONFLICT (tenant, run_id) DO UPDATE SET
owner = EXCLUDED.owner,
epoch = run_lease.epoch + 1,
expires_at = EXCLUDED.expires_at
WHERE run_lease.expires_at <= $5
RETURNING epoch",
&[
&self.tenant_name(),
&key,
&owner.to_owned(),
&expires.cast_signed(),
&now.cast_signed(),
],
)
.await
.map_err(|e| be(&e))?;
if let Some(row) = row {
let epoch: i64 = row.get(0);
return Ok(Lease {
run,
owner: owner.to_owned(),
epoch: epoch.cast_unsigned(),
});
}
let held = client
.query_opt(
"SELECT owner, epoch, expires_at FROM run_lease
WHERE tenant = $1 AND run_id = $2",
&[&self.tenant_name(), &key],
)
.await
.map_err(|e| be(&e))?;
match held {
Some(row) => Err(StoreError::LeaseHeld {
run: key,
owner: row.get(0),
epoch: row.get::<_, i64>(1).cast_unsigned(),
remaining_secs: row.get::<_, i64>(2).cast_unsigned().saturating_sub(now),
}),
None => Err(StoreError::Backend(format!(
"acquire for run {key} was refused but no lease row exists — \
lease rows are never deleted, so this store's state is not \
one this backend can have written"
))),
}
}
async fn renew(
&self,
run: RunId,
owner: &str,
epoch: Epoch,
ttl: Duration,
) -> Result<Lease, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let now = now_secs();
let expires = lease_expiry(now, ttl)?;
let n = client
.execute(
"UPDATE run_lease SET expires_at = $5
WHERE tenant = $1 AND run_id = $2 AND owner = $3 AND epoch = $4
AND expires_at > $6",
&[
&self.tenant_name(),
&run.to_string(),
&owner.to_owned(),
&epoch.cast_signed(),
&expires.cast_signed(),
&now.cast_signed(),
],
)
.await
.map_err(|e| be(&e))?;
if n == 0 {
return Err(StoreError::LeaseNotHeld {
run: run.to_string(),
epoch,
});
}
Ok(Lease {
run,
owner: owner.to_owned(),
epoch,
})
}
async fn abandoned_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT run_id FROM run_lease
WHERE tenant = $1 AND owner <> '' AND expires_at <= $2
ORDER BY expires_at ASC
LIMIT $3",
&[
&self.tenant_name(),
&now_secs().cast_signed(),
&i64::try_from(limit).unwrap_or(i64::MAX),
],
)
.await
.map_err(|e| be(&e))?;
rows.iter()
.map(|row| {
let id: String = row.get(0);
RunId::parse(&id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_lease holds an unparsable run id '{id}': {e}"),
})
})
.collect()
}
async fn release_lease(&self, run: RunId, epoch: Epoch) -> Result<(), StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
client
.execute(
"UPDATE run_lease SET owner = '', expires_at = 0 \
WHERE tenant = $1 AND run_id = $2 AND epoch = $3",
&[&self.tenant_name(), &run.to_string(), &epoch.cast_signed()],
)
.await
.map_err(|e| be(&e))?;
Ok(())
}
async fn seal(&self, run: RunId, epoch: Epoch, outcome: &str) -> Result<Digest, StoreError> {
let mut client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let tx = client.transaction().await.map_err(|e| be(&e))?;
self.fence(&tx, run, epoch).await?;
let head_row = tx
.query_opt(
"SELECT hash FROM journal
WHERE tenant = $1 AND run_id = $2 ORDER BY seq DESC LIMIT 1",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
let head = match head_row {
Some(row) => digest_from(&row.get::<_, Vec<u8>>(0))?,
None => Digest::ZERO,
};
tx.execute(
"INSERT INTO run_seal (tenant, run_id, chain_head, log_index, sealed_at, outcome)
VALUES ($1, $2, $3, nextval('run_log_position'), $4, $5)
ON CONFLICT (tenant, run_id) DO NOTHING",
&[
&self.tenant_name(),
&run.to_string(),
&head.as_bytes().to_vec(),
&now_secs().cast_signed(),
&outcome,
],
)
.await
.map_err(|e| be(&e))?;
tx.commit()
.await
.map_err(|e| commit_refused_or_in_doubt(&e))?;
Ok(head)
}
async fn checkpoint(&self) -> Result<crate::journal::Checkpoint, StoreError> {
let leaves = self.log_leaves().await?;
Ok(crate::journal::Checkpoint {
origin: self.log_origin(),
size: leaves.len() as u64,
root: crate::core::merkle::root(&leaves),
})
}
async fn consistency_proof(&self, old_size: u64) -> Result<Vec<Digest>, StoreError> {
let leaves = self.log_leaves().await?;
let old = usize::try_from(old_size).unwrap_or(usize::MAX);
if old > leaves.len() {
return Err(StoreError::Backend(format!(
"a checkpoint of size {old_size} is larger than this log ({}) — \
either it belongs to another plane, or runs were removed",
leaves.len()
)));
}
Ok(crate::core::merkle::consistency_proof(&leaves, old))
}
async fn inclusion_proof(
&self,
run: RunId,
) -> Result<Option<crate::journal::Inclusion>, StoreError> {
let leaves = self.log_leaves().await?;
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let Some(row) = client
.query_opt(
"SELECT rank, chain_head FROM (
SELECT run_id, chain_head,
ROW_NUMBER() OVER (ORDER BY log_index) - 1 AS rank
FROM run_seal WHERE tenant = $1
) ranked WHERE run_id = $2",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?
else {
return Ok(None);
};
let index: i64 = row.get(0);
let seal = digest_from(&row.get::<_, Vec<u8>>(1))?;
let index = usize::try_from(index).unwrap_or(0);
Ok(Some(crate::journal::Inclusion {
index: index as u64,
size: leaves.len() as u64,
seal,
proof: crate::core::merkle::inclusion_proof(&leaves, index),
}))
}
async fn request_cancel(
&self,
run: RunId,
actor: &str,
reason: &str,
) -> Result<bool, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let n = client
.execute(
"INSERT INTO run_cancel (tenant, run_id, actor, reason, requested_at)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (tenant, run_id) DO NOTHING",
&[
&self.tenant_name(),
&run.to_string(),
&actor.to_owned(),
&reason.to_owned(),
&now_secs().cast_signed(),
],
)
.await
.map_err(|e| be(&e))?;
Ok(n == 1)
}
async fn cancellation(&self, run: RunId) -> Result<Option<Cancellation>, StoreError> {
let client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let row = client
.query_opt(
"SELECT actor, reason FROM run_cancel WHERE tenant = $1 AND run_id = $2",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
Ok(row.map(|r| Cancellation {
actor: r.get(0),
reason: r.get(1),
}))
}
}
fn digest_from(bytes: &[u8]) -> Result<Digest, StoreError> {
let arr: [u8; 32] = bytes.try_into().map_err(|_| StoreError::Corrupt {
seq: 0,
detail: format!("a hash column holds {} bytes, not 32", bytes.len()),
})?;
Ok(Digest::from_bytes(arr))
}
struct PgAtomicTx<'a>(&'a deadpool_postgres::tokio_postgres::Transaction<'a>);
fn bind(params: &[SqlValue]) -> Vec<Box<dyn tokio_postgres::types::ToSql + Sync + Send>> {
params
.iter()
.map(|v| -> Box<dyn tokio_postgres::types::ToSql + Sync + Send> {
match v {
SqlValue::Null => Box::new(Option::<i64>::None),
SqlValue::Bool(b) => Box::new(*b),
SqlValue::Int(i) => Box::new(*i),
SqlValue::Float(f) => Box::new(*f),
SqlValue::Text(s) => Box::new(s.clone()),
SqlValue::Bytes(b) => Box::new(b.clone()),
SqlValue::Json(j) => Box::new(j.clone()),
}
})
.collect()
}
fn as_refs(
owned: &[Box<dyn tokio_postgres::types::ToSql + Sync + Send>],
) -> Vec<&(dyn tokio_postgres::types::ToSql + Sync)> {
owned.iter().map(|b| &**b as _).collect()
}
fn refuse_transaction_control(sql: &str) -> Result<(), StoreError> {
let mut words = sql
.split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
.filter(|w| !w.is_empty());
let first = words.next().unwrap_or("").to_ascii_uppercase();
let second = words.next().unwrap_or("").to_ascii_uppercase();
let forbidden = match first.as_str() {
"COMMIT" | "END" | "ROLLBACK" | "ABORT" | "BEGIN" | "START" | "SAVEPOINT" | "RELEASE" => {
true
}
"PREPARE" | "SET" => second == "TRANSACTION",
_ => false,
};
if forbidden {
return Err(StoreError::Backend(format!(
"the statement '{first}' would end or control the journal's \
transaction — a co-located resource writes *inside* the \
transaction that records it; it does not own that transaction, \
and dissolving it from within would break the atomicity the \
seam exists to provide"
)));
}
Ok(())
}
#[async_trait]
impl AtomicTx for PgAtomicTx<'_> {
async fn execute(&self, sql: &str, params: &[SqlValue]) -> Result<u64, StoreError> {
refuse_transaction_control(sql)?;
let owned = bind(params);
self.0
.execute(sql, &as_refs(&owned))
.await
.map_err(|e| be(&e))
}
async fn query(&self, sql: &str, params: &[SqlValue]) -> Result<Vec<Value>, StoreError> {
refuse_transaction_control(sql)?;
let owned = bind(params);
let rows = self
.0
.query(sql, &as_refs(&owned))
.await
.map_err(|e| be(&e))?;
let mut out = Vec::with_capacity(rows.len());
for row in &rows {
let mut obj = serde_json::Map::new();
for (i, col) in row.columns().iter().enumerate() {
obj.insert(col.name().to_owned(), column_json(row, i, col.type_())?);
}
out.push(Value::Object(obj));
}
Ok(out)
}
}
fn column_json(
row: &deadpool_postgres::tokio_postgres::Row,
i: usize,
ty: &tokio_postgres::types::Type,
) -> Result<Value, StoreError> {
use tokio_postgres::types::Type;
macro_rules! get {
($t:ty) => {
row.try_get::<_, Option<$t>>(i)
.map_err(|e| StoreError::Backend(e.to_string()))?
.map_or(Value::Null, Value::from)
};
}
Ok(match *ty {
Type::BOOL => get!(bool),
Type::INT2 => get!(i16),
Type::INT4 => get!(i32),
Type::INT8 => get!(i64),
Type::FLOAT4 => get!(f32),
Type::FLOAT8 => get!(f64),
Type::TEXT | Type::VARCHAR | Type::BPCHAR | Type::NAME => get!(String),
Type::JSON | Type::JSONB => row
.try_get::<_, Option<Value>>(i)
.map_err(|e| StoreError::Backend(e.to_string()))?
.unwrap_or(Value::Null),
Type::BYTEA => row
.try_get::<_, Option<Vec<u8>>>(i)
.map_err(|e| StoreError::Backend(e.to_string()))?
.map_or(Value::Null, |b| Value::String(hex(&b))),
ref other => {
return Err(StoreError::Backend(format!(
"column '{}' has type {other}, which this seam does not convert — \
select it as text or jsonb rather than being handed a null that \
reads as a zero",
row.columns()[i].name()
)));
}
})
}
fn hex(bytes: &[u8]) -> String {
use std::fmt::Write as _;
bytes.iter().fold(String::new(), |mut s, b| {
let _ = write!(s, "{b:02x}");
s
})
}
#[async_trait]
impl AtomicJournal for PostgresStore {
async fn append_atomic(
&self,
run: RunId,
epoch: Epoch,
work: &dyn AtomicWork,
) -> Result<Vec<Record>, StoreError> {
let mut client = self.pool.get().await.map_err(|e| pool_err(&e))?;
let tx = client.transaction().await.map_err(|e| be(&e))?;
self.fence(&tx, run, epoch).await?;
let batch = work
.run(&PgAtomicTx(&tx))
.await
.map_err(|e| StoreError::Backend(e.to_string()))?;
let sealed = self.append_within(&tx, run, epoch, batch).await?;
tx.commit()
.await
.map_err(|e| commit_refused_or_in_doubt(&e))?;
Ok(sealed)
}
}