audit-core/audit_core/sqlite_backend.py
tegwick 5fd04e2095 Implement AUDIT-WP-0007 hash-chain integrity.
Accept now extends a single-schema chain. Verify walks it; a rewritten
payload_hash is a break. Tamper evidence is that detector plus an
external chain-head attestation, not WORM.
2026-08-16 01:18:30 +02:00

402 lines
15 KiB
Python

"""Durable SQLite audit backend.
Implements the idempotent backend contract for single-node deployments and for
development. Production custody moves to PostgreSQL under AUDIT-WP-0005; this
backend stays the development and test implementation and defines the
behaviour the Postgres backend must match.
"""
from __future__ import annotations
import json
import sqlite3
import threading
from datetime import datetime, timezone
from audit_core.interface import (
AcceptResult,
AuditEvent,
BackendUnavailableError,
EventConflictError,
EventValidationError,
RetentionPolicy,
validate_event,
)
from audit_core.integrity import (
GENESIS,
ChainRow,
attestation_from_report,
chain_link,
verify_rows,
)
_SCHEMA = """
CREATE TABLE IF NOT EXISTS events (
event_id TEXT PRIMARY KEY,
payload_hash TEXT NOT NULL,
accepted_at TEXT NOT NULL,
correlation_id TEXT,
tenant TEXT NOT NULL,
record TEXT NOT NULL,
chain_hash TEXT,
chain_prev TEXT
);
CREATE INDEX IF NOT EXISTS events_correlation_idx ON events (correlation_id);
CREATE INDEX IF NOT EXISTS events_tenant_idx ON events (tenant);
CREATE TABLE IF NOT EXISTS dead_letters (
id INTEGER PRIMARY KEY AUTOINCREMENT,
event_id TEXT,
received_at TEXT NOT NULL,
sender TEXT,
reason TEXT NOT NULL,
payload_hash TEXT NOT NULL,
payload TEXT,
payload_withheld INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS dead_letters_event_idx ON dead_letters (event_id);
-- Counted per field path, not merely per event: the point is to stop senders
-- emitting secret-shaped fields, and that needs the offending path named.
CREATE TABLE IF NOT EXISTS secret_findings (
sender TEXT NOT NULL,
source TEXT NOT NULL,
action TEXT NOT NULL,
field_path TEXT NOT NULL,
outcome TEXT NOT NULL,
persisted INTEGER NOT NULL DEFAULT 0,
occurrences INTEGER NOT NULL DEFAULT 0,
first_seen TEXT NOT NULL,
last_seen TEXT NOT NULL,
PRIMARY KEY (sender, source, action, field_path, outcome)
);
"""
# Rejection reasons whose payload must never be persisted. Storing the body of
# an event rejected *for containing secret-shaped material* would write that
# material into the audit store — the precise outcome the rejection prevents.
WITHHOLD_PAYLOAD_REASONS = frozenset({"secret_shaped_field"})
class SQLiteAuditBackend:
"""Store audit events in SQLite with idempotent accept semantics.
Configured for durability rather than speed: WAL journalling, ``synchronous
= FULL`` so an acknowledged write has reached disk, and a busy timeout so
concurrent writers wait instead of raising immediately.
"""
def __init__(self, path: str, retention_days: int | None = None, busy_timeout_ms: int = 5000) -> None:
self.path = path
self.retention_days = retention_days
self.busy_timeout_ms = int(busy_timeout_ms)
# One connection per thread. A shared connection lets concurrent
# statements interleave, which was observed to let two callers both
# believe they were the first to accept the same event.
self._local = threading.local()
with self._connect_raw() as setup:
setup.executescript(_SCHEMA)
self._ensure_chain_columns(setup)
def _connect_raw(self) -> sqlite3.Connection:
try:
db = sqlite3.connect(self.path, isolation_level=None)
db.execute("PRAGMA journal_mode = WAL")
db.execute("PRAGMA synchronous = FULL")
db.execute(f"PRAGMA busy_timeout = {self.busy_timeout_ms}")
return db
except sqlite3.Error as exc:
raise BackendUnavailableError(f"cannot open audit store: {exc}") from exc
@property
def db(self) -> sqlite3.Connection:
conn = getattr(self._local, "conn", None)
if conn is None:
conn = self._local.conn = self._connect_raw()
return conn
@property
def retention_policy(self) -> RetentionPolicy:
return RetentionPolicy(
custody_class="development",
retention_days=self.retention_days,
immutable=False,
tamper_evidence=False,
durable=True,
)
def emit(self, event: AuditEvent) -> str:
"""Persist an event, generating no idempotency guarantee."""
return self.accept(event, payload_hash=_record_hash(event)).reference
def accept(self, event: AuditEvent, payload_hash: str) -> AcceptResult:
try:
validate_event(event)
except EventValidationError:
raise
reference = f"audit:{event.event_id}"
details = event.details if isinstance(event.details, dict) else {}
db = self.db
try:
# BEGIN IMMEDIATE takes the write lock up front, so the insert and
# the follow-up read are one atomic pair. Without it, two callers
# racing on the same event id can both be told they were first.
db.execute("BEGIN IMMEDIATE")
except sqlite3.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
try:
head = db.execute(
"SELECT chain_hash FROM events "
"ORDER BY accepted_at DESC, event_id DESC LIMIT 1"
).fetchone()
previous = head[0] if head and head[0] else GENESIS
link = chain_link(previous, payload_hash, event.event_id)
inserted = db.execute(
"""
INSERT INTO events
(event_id, payload_hash, accepted_at, correlation_id, tenant, record,
chain_hash, chain_prev)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(event_id) DO NOTHING
RETURNING event_id
""",
(
event.event_id,
payload_hash,
datetime.now(timezone.utc).replace(microsecond=0).isoformat(),
str(details.get("correlation_id") or "") or None,
event.tenant,
json.dumps(event.as_record(), sort_keys=True),
link,
previous,
),
).fetchone()
existing = None
if inserted is None:
existing = db.execute(
"SELECT payload_hash FROM events WHERE event_id = ?", (event.event_id,)
).fetchone()
db.execute("COMMIT")
except sqlite3.Error as exc:
_rollback(db)
raise BackendUnavailableError(str(exc)) from exc
except BaseException:
_rollback(db)
raise
if inserted is not None:
return AcceptResult(duplicate=False, reference=reference)
if existing is None:
# The row vanished between the insert and the read inside one
# transaction, which should be impossible. Retryable rather than
# guessed at.
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 read surface (AUDIT-WP-0004-T05) --------------------------
def get(self, event_id: str) -> dict | None:
"""Return one stored event record, or None."""
row = self._query(
"SELECT record, accepted_at FROM events WHERE event_id = ?", (event_id,)
)
if not row:
return None
return {"accepted_at": row[0][1], **json.loads(row[0][0])}
def by_correlation(self, correlation_id: str, limit: int = 100) -> list[dict]:
"""Return every stored event carrying ``correlation_id``, oldest first."""
rows = self._query(
"SELECT record, accepted_at FROM events WHERE correlation_id = ? "
"ORDER BY accepted_at, event_id LIMIT ?",
(correlation_id, int(limit)),
)
return [{"accepted_at": at, **json.loads(rec)} for rec, at in rows]
def record_rejection(
self,
*,
event_id: str | None,
reason: str,
payload_hash: str,
sender: str | None = None,
payload: str | None = None,
) -> None:
"""Record a rejected event so it is visible to an operator.
A rejection is not a silent drop: the sender dead-letters the event and
somebody has to be able to see why. The payload is withheld when the
rejection reason implies it carries secret-shaped material.
"""
withheld = reason in WITHHOLD_PAYLOAD_REASONS
try:
self.db.execute(
"INSERT INTO dead_letters "
"(event_id, received_at, sender, reason, payload_hash, payload, payload_withheld) "
"VALUES (?, ?, ?, ?, ?, ?, ?)",
(
event_id,
datetime.now(timezone.utc).replace(microsecond=0).isoformat(),
sender,
reason,
payload_hash,
None if withheld else payload,
1 if withheld else 0,
),
)
except sqlite3.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def dead_letters(self, limit: int = 100) -> list[dict]:
"""Return recent rejections, newest first."""
rows = self._query(
"SELECT event_id, received_at, sender, reason, payload_hash, payload, "
"payload_withheld FROM dead_letters ORDER BY id DESC LIMIT ?",
(int(limit),),
)
return [
{
"event_id": r[0],
"received_at": 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:
"""Increment the per-path counter for each finding.
Written durably rather than held in memory: these counters exist to
drive a fix in the sending service, and that work outlives a pod
restart.
"""
now = datetime.now(timezone.utc).replace(microsecond=0).isoformat()
try:
for finding in findings:
self.db.execute(
"""
INSERT INTO secret_findings
(sender, source, action, field_path, outcome, persisted,
occurrences, first_seen, last_seen)
VALUES (?, ?, ?, ?, ?, ?, 1, ?, ?)
ON CONFLICT(sender, source, action, field_path, outcome)
DO UPDATE SET occurrences = occurrences + 1, last_seen = excluded.last_seen
""",
(
sender, source, action, finding.path, outcome,
1 if finding.in_persisted_data else 0, now, now,
),
)
except sqlite3.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def secret_findings(self, limit: int = 100) -> list[dict]:
"""Return secret-shaped field counters, most frequent first."""
rows = self._query(
"SELECT sender, source, action, field_path, outcome, persisted, "
"occurrences, first_seen, last_seen FROM secret_findings "
"ORDER BY occurrences DESC, last_seen DESC LIMIT ?",
(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": r[7], "last_seen": r[8],
}
for r in rows
]
def _query(self, sql: str, params: tuple) -> list:
try:
return self.db.execute(sql, params).fetchall()
except sqlite3.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def verify_chain(self, attestation: dict | None = None):
rows = self._query(
"SELECT event_id, payload_hash, chain_hash, chain_prev, accepted_at "
"FROM events ORDER BY accepted_at, event_id",
(),
)
return verify_rows(
[
ChainRow(
event_id=r[0],
payload_hash=r[1],
chain_hash=r[2] or "",
chain_prev=r[3] or "",
accepted_at=r[4],
)
for r in rows
],
attestation=attestation,
)
def attest_chain(self) -> dict:
return attestation_from_report(self.verify_chain())
def _ensure_chain_columns(self, db: sqlite3.Connection) -> None:
cols = {row[1] for row in db.execute("PRAGMA table_info(events)")}
if "chain_hash" not in cols:
db.execute("ALTER TABLE events ADD COLUMN chain_hash TEXT")
if "chain_prev" not in cols:
db.execute("ALTER TABLE events ADD COLUMN chain_prev TEXT")
missing = db.execute(
"SELECT event_id, payload_hash FROM events "
"WHERE chain_hash IS NULL OR chain_prev IS NULL "
"ORDER BY accepted_at, event_id"
).fetchall()
if not missing:
return
prev = GENESIS
# Recompute the whole chain so a partial backfill cannot fork.
for event_id, payload_hash in db.execute(
"SELECT event_id, payload_hash FROM events ORDER BY accepted_at, event_id"
):
link = chain_link(prev, payload_hash, event_id)
db.execute(
"UPDATE events SET chain_prev = ?, chain_hash = ? WHERE event_id = ?",
(prev, link, event_id),
)
prev = link
def health(self) -> None:
"""Raise :class:`BackendUnavailableError` if the store is unusable."""
try:
self.db.execute("SELECT 1 FROM events LIMIT 1").fetchone()
except sqlite3.Error as exc:
raise BackendUnavailableError(str(exc)) from exc
def close(self) -> None:
"""Close this thread's connection, if it has one."""
conn = getattr(self._local, "conn", None)
if conn is not None:
conn.close()
self._local.conn = None
def _rollback(db: sqlite3.Connection) -> None:
try:
db.execute("ROLLBACK")
except sqlite3.Error:
pass
def _record_hash(event: AuditEvent) -> str:
import hashlib
return hashlib.sha256(
json.dumps(event.as_record(), sort_keys=True).encode("utf-8")
).hexdigest()