audit-core/audit_core/postgres_backend.py
tegwick bd274f6269
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 1s
Add the PostgreSQL audit backend and a shared conformance suite
AUDIT-WP-0005-T01, built and verified against PostgreSQL 16 locally in
Docker; the Railiance cluster was not needed.

tests/test_backend_conformance.py is one suite run against every backend, so
"the Postgres backend is done" means it satisfies the same contract SQLite
already does rather than having its own green tests. It skips cleanly with no
server reachable; make pg-test-up and make test-pg run it. Suite 50 -> 71.

RetentionPolicy declares immutable=True and earns it: migration 0002 installs
a trigger rejecting UPDATE and DELETE on the events table, so a leaked runtime
credential can append but cannot rewrite or erase the trail. That materially
narrows the residual risk ADR-0001 section 5 called out. tamper_evidence stays
False because nothing here would prove a database owner had dropped the
trigger - hash-chaining or external anchoring would be needed and is not
implemented.

Idempotency is one statement (INSERT ... ON CONFLICT DO NOTHING RETURNING),
verified to behave identically to the SQLite backend under 12 concurrent
submissions of the same event. Migrations are ordered, recorded and
idempotent. Replay reconciles rather than duplicating - the piece deferred out
of WP-0004-T05 - and is tested to leave exactly one custody record.

Backend selection is by AUDIT_CORE_DATABASE_URL; the SQLite fallback logs a
warning so a deployment that lost its URL is visible rather than quietly
running on the wrong store.

Also fixed: ingestion had no __main__ guard, so python -m audit_core.ingestion
silently did nothing. Found during end-to-end smoke.

Counting semantics documented: occurrences counts transmissions, not stored
events, so a retry of a secret-shaped field increments it again. That is the
sender behaviour being optimized away.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 17:09:46 +02:00

426 lines
16 KiB
Python

"""Durable PostgreSQL audit backend (AUDIT-WP-0005-T01).
Production custody. Runs on the Railiance shared PostgreSQL platform
(`rapp-postgres`), inside the database-per-consumer boundary fixed by that
repo's ADR-0001.
Two properties are pushed into the database rather than the application:
* **Idempotency** — ``INSERT ... ON CONFLICT DO NOTHING RETURNING`` resolves
insert-or-detect in one statement, so two concurrent submissions of the same
event cannot both be told they were first. The SQLite backend learned this
the hard way; see its test.
* **Append-only custody** — a trigger rejects ``UPDATE`` and ``DELETE`` on the
events table, so a leaked runtime credential can add records but cannot
rewrite or erase existing ones. This is what lets ``RetentionPolicy`` declare
``immutable=True`` honestly; see :meth:`retention_policy` for its limits.
"""
from __future__ import annotations
import json
import os
from datetime import datetime, timezone
from typing import Any
from audit_core.interface import (
AcceptResult,
AuditEvent,
BackendUnavailableError,
EventConflictError,
EventValidationError,
RetentionPolicy,
validate_event,
)
from audit_core.redaction import Finding
try: # pragma: no cover - import guard
import psycopg
from psycopg_pool import ConnectionPool
except ImportError as exc: # pragma: no cover
raise ImportError(
"the postgres backend needs psycopg: pip install 'audit-core[postgres]'"
) from exc
DEFAULT_SCHEMA = "audit_core"
# Ordered, append-only. Each entry runs once and is recorded in
# schema_migrations. Never edit a released migration — add a new one.
MIGRATIONS: list[tuple[str, str]] = [
(
"0001-events",
"""
CREATE TABLE IF NOT EXISTS {schema}.events (
event_id text PRIMARY KEY,
payload_hash text NOT NULL,
accepted_at timestamptz NOT NULL DEFAULT now(),
observed_at timestamptz,
tenant text NOT NULL,
correlation_id text,
source text NOT NULL,
action text NOT NULL,
record jsonb NOT NULL
);
CREATE INDEX IF NOT EXISTS events_correlation_idx
ON {schema}.events (correlation_id);
-- Tenant keying is mandatory per business-app-service-contract 1.3;
-- the index makes per-tenant export and retention tractable.
CREATE INDEX IF NOT EXISTS events_tenant_idx ON {schema}.events (tenant);
""",
),
(
"0002-append-only",
"""
CREATE OR REPLACE FUNCTION {schema}.reject_mutation() RETURNS trigger AS $fn$
BEGIN
RAISE EXCEPTION 'audit events are append-only (attempted %)', TG_OP
USING ERRCODE = 'restrict_violation';
END;
$fn$ LANGUAGE plpgsql;
DROP TRIGGER IF EXISTS events_append_only ON {schema}.events;
CREATE TRIGGER events_append_only
BEFORE UPDATE OR DELETE ON {schema}.events
FOR EACH ROW EXECUTE FUNCTION {schema}.reject_mutation();
""",
),
(
"0003-dead-letters",
"""
CREATE TABLE IF NOT EXISTS {schema}.dead_letters (
id bigserial PRIMARY KEY,
event_id text,
received_at timestamptz NOT NULL DEFAULT now(),
sender text,
reason text NOT NULL,
payload_hash text NOT NULL,
payload text,
payload_withheld boolean NOT NULL DEFAULT false
);
CREATE INDEX IF NOT EXISTS dead_letters_event_idx
ON {schema}.dead_letters (event_id);
""",
),
(
"0004-secret-findings",
"""
CREATE TABLE IF NOT EXISTS {schema}.secret_findings (
sender text NOT NULL,
source text NOT NULL,
action text NOT NULL,
field_path text NOT NULL,
outcome text NOT NULL,
persisted boolean NOT NULL DEFAULT false,
occurrences bigint NOT NULL DEFAULT 0,
first_seen timestamptz NOT NULL DEFAULT now(),
last_seen timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (sender, source, action, field_path, outcome)
);
""",
),
]
# Rejection reasons whose payload must never be persisted — storing the body of
# an event rejected *for carrying secret-shaped material* would write that
# material into the audit store.
WITHHOLD_PAYLOAD_REASONS = frozenset({"secret_shaped_field"})
class PostgresAuditBackend:
"""Audit custody in PostgreSQL."""
def __init__(
self,
dsn: str | None = None,
*,
schema: str = DEFAULT_SCHEMA,
retention_days: int | None = None,
min_size: int = 1,
max_size: int = 8,
statement_timeout_ms: int = 30_000,
migrate: bool = True,
) -> None:
self.dsn = dsn or os.environ.get("AUDIT_CORE_DATABASE_URL") or ""
if not self.dsn:
raise ValueError("a DSN is required (AUDIT_CORE_DATABASE_URL)")
if not schema.isidentifier():
raise ValueError(f"unsafe schema name: {schema!r}")
self.schema = schema
self.retention_days = retention_days
try:
self.pool = ConnectionPool(
self.dsn,
min_size=min_size,
max_size=max_size,
kwargs={
"autocommit": True,
# A stalled write must surface as unavailable rather than
# hold a request open indefinitely.
"options": f"-c statement_timeout={int(statement_timeout_ms)}",
},
open=True,
timeout=10,
)
except Exception as exc: # psycopg raises a wide family here
raise BackendUnavailableError(f"cannot connect: {exc}") from exc
if migrate:
self.migrate()
# --- schema ------------------------------------------------------------
def migrate(self) -> list[str]:
"""Apply pending migrations. Returns the ids applied this call."""
applied: list[str] = []
try:
with self.pool.connection() as conn:
conn.execute(f'CREATE SCHEMA IF NOT EXISTS "{self.schema}"')
conn.execute(
f'CREATE TABLE IF NOT EXISTS "{self.schema}".schema_migrations ('
" id text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())"
)
done = {
row[0]
for row in conn.execute(
f'SELECT id FROM "{self.schema}".schema_migrations'
).fetchall()
}
for migration_id, body in MIGRATIONS:
if migration_id in done:
continue
conn.execute(body.format(schema=f'"{self.schema}"'))
conn.execute(
f'INSERT INTO "{self.schema}".schema_migrations (id) VALUES (%s)',
(migration_id,),
)
applied.append(migration_id)
except psycopg.Error as exc:
raise BackendUnavailableError(f"migration failed: {exc}") from exc
return applied
@property
def _events(self) -> str:
return f'"{self.schema}".events'
# --- contract ----------------------------------------------------------
@property
def retention_policy(self) -> RetentionPolicy:
"""Declared custody guarantees.
``immutable`` is True because migration 0002 installs a trigger that
rejects UPDATE and DELETE, so no consumer credential can alter a stored
record. It is not a claim against the database owner or a superuser,
who can drop the trigger; ``tamper_evidence`` is correspondingly False,
because nothing here would *prove* they had. Hash-chaining or external
anchoring would be needed for that, and is not implemented.
"""
return RetentionPolicy(
custody_class="archive",
retention_days=self.retention_days,
immutable=True,
tamper_evidence=False,
durable=True,
)
def emit(self, event: AuditEvent) -> str:
return self.accept(event, payload_hash=_record_hash(event)).reference
def accept(self, event: AuditEvent, payload_hash: str) -> AcceptResult:
validate_event(event)
reference = f"audit:{event.event_id}"
details = event.details if isinstance(event.details, dict) else {}
record = event.as_record()
try:
with self.pool.connection() as conn:
inserted = conn.execute(
f"""
INSERT INTO {self._events}
(event_id, payload_hash, observed_at, tenant, correlation_id,
source, action, record)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (event_id) DO NOTHING
RETURNING event_id
""",
(
event.event_id,
payload_hash,
_timestamp(event.observed_at),
event.tenant,
str(details.get("correlation_id") or "") or None,
event.source,
event.action,
json.dumps(record, sort_keys=True),
),
).fetchone()
if inserted is not None:
return AcceptResult(duplicate=False, reference=reference)
existing = conn.execute(
f"SELECT payload_hash FROM {self._events} WHERE event_id = %s",
(event.event_id,),
).fetchone()
except psycopg.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
if existing is None:
raise BackendUnavailableError("event disappeared during accept")
if existing[0] != payload_hash:
raise EventConflictError(
f"event_id {event.event_id} already held with a different payload"
)
return AcceptResult(duplicate=True, reference=reference)
# --- operator surface --------------------------------------------------
def get(self, event_id: str) -> dict | None:
rows = self._query(
f"SELECT record, accepted_at FROM {self._events} WHERE event_id = %s",
(event_id,),
)
if not rows:
return None
return {"accepted_at": _iso(rows[0][1]), **rows[0][0]}
def by_correlation(self, correlation_id: str, limit: int = 100) -> list[dict]:
rows = self._query(
f"SELECT record, accepted_at FROM {self._events} "
"WHERE correlation_id = %s ORDER BY accepted_at, event_id LIMIT %s",
(correlation_id, int(limit)),
)
return [{"accepted_at": _iso(at), **rec} for rec, at in rows]
def replay(self, event_id: str) -> AcceptResult:
"""Re-submit a stored event through :meth:`accept`.
Reconciliation, not re-creation: replaying an event already in custody
must return ``duplicate=True`` against the same record. If it were to
report a first acceptance, the store would be producing a second
custody record for one source event — the exact failure the whole
idempotency design exists to prevent.
"""
rows = self._query(
f"SELECT record, payload_hash FROM {self._events} WHERE event_id = %s",
(event_id,),
)
if not rows:
raise KeyError(event_id)
record, payload_hash = rows[0]
return self.accept(_event_from_record(record), payload_hash)
def record_rejection(
self,
*,
event_id: str | None,
reason: str,
payload_hash: str,
sender: str | None = None,
payload: str | None = None,
) -> None:
withheld = reason in WITHHOLD_PAYLOAD_REASONS
self._execute(
f'INSERT INTO "{self.schema}".dead_letters '
"(event_id, sender, reason, payload_hash, payload, payload_withheld) "
"VALUES (%s, %s, %s, %s, %s, %s)",
(event_id, sender, reason, payload_hash,
None if withheld else payload, withheld),
)
def dead_letters(self, limit: int = 100) -> list[dict]:
rows = self._query(
"SELECT event_id, received_at, sender, reason, payload_hash, payload, "
f'payload_withheld FROM "{self.schema}".dead_letters '
"ORDER BY id DESC LIMIT %s",
(int(limit),),
)
return [
{
"event_id": r[0], "received_at": _iso(r[1]), "sender": r[2],
"reason": r[3], "payload_hash": r[4], "payload": r[5],
"payload_withheld": bool(r[6]),
}
for r in rows
]
def count_secret_findings(
self, *, sender: str, source: str, action: str, outcome: str, findings
) -> None:
for finding in findings:
self._execute(
f'INSERT INTO "{self.schema}".secret_findings '
"(sender, source, action, field_path, outcome, persisted, occurrences) "
"VALUES (%s, %s, %s, %s, %s, %s, 1) "
"ON CONFLICT (sender, source, action, field_path, outcome) DO UPDATE "
"SET occurrences = secret_findings.occurrences + 1, last_seen = now()",
(sender, source, action, finding.path, outcome,
bool(getattr(finding, "in_persisted_data", False))),
)
def secret_findings(self, limit: int = 100) -> list[dict]:
rows = self._query(
"SELECT sender, source, action, field_path, outcome, persisted, "
f'occurrences, first_seen, last_seen FROM "{self.schema}".secret_findings '
"ORDER BY occurrences DESC, last_seen DESC LIMIT %s",
(int(limit),),
)
return [
{
"sender": r[0], "source": r[1], "action": r[2], "field_path": r[3],
"outcome": r[4], "persisted": bool(r[5]), "occurrences": r[6],
"first_seen": _iso(r[7]), "last_seen": _iso(r[8]),
}
for r in rows
]
# --- lifecycle ---------------------------------------------------------
def health(self) -> None:
self._query("SELECT 1", ())
def close(self) -> None:
self.pool.close()
def _query(self, sql: str, params: tuple) -> list:
try:
with self.pool.connection() as conn:
return conn.execute(sql, params).fetchall()
except psycopg.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def _execute(self, sql: str, params: tuple) -> None:
try:
with self.pool.connection() as conn:
conn.execute(sql, params)
except psycopg.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def _timestamp(value: str | None) -> datetime | None:
if not value:
return None
try:
return datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except ValueError:
return None
def _iso(value: Any) -> str | None:
if isinstance(value, datetime):
return value.astimezone(timezone.utc).replace(microsecond=0).isoformat()
return value
def _event_from_record(record: dict) -> AuditEvent:
return AuditEvent(
source=record["source"], action=record["action"], resource=record["resource"],
outcome=record["outcome"], tenant=record["tenant"], scope=record["scope"],
actor=record.get("actor"), reason=record.get("reason"),
details=record.get("details") or {}, event_id=record["event_id"],
observed_at=record["observed_at"], schema_version=record["schema_version"],
)
def _record_hash(event: AuditEvent) -> str:
import hashlib
return hashlib.sha256(
json.dumps(event.as_record(), sort_keys=True).encode("utf-8")
).hexdigest()