audit-core/audit_core/migrate_store.py
tegwick 88d16847ff
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 1s
Establish pre-production record disposition and add the migration tool
AUDIT-WP-0005-T04.

Disposition: there are no pre-production records. No audit-core SQLite store
on this host, no mock-file-backend output, and no audit-core pod, deployment
or PVC on railiance01 - the only audit-* PVC there is OpenBao's own audit
device. Consistent with the history: WP-0003-T03 was cancelled before the
receiver was ever deployed, so every SQLite store that has existed was a test
fixture. Nothing is being discarded because nothing was ever accepted outside
tests.

The tool is built anyway because the SQLite path stays reachable - the
entrypoint falls back to it when AUDIT_CORE_DATABASE_URL is unset. If that
fallback is ever used in anger the records are audit records, and writing the
migration afterwards under pressure is the wrong time.

audit_core.migrate_store and `python -m audit_core migrate-store` transfer
events, dead letters and secret-finding counters. Records keep their original
event_id, payload_hash and accepted_at, which is why this bypasses accept():
that stamps acceptance with the current time, and a migration that rewrote
acceptance times would destroy the evidence it exists to preserve.

Idempotent, and verification reads back from the destination rather than
trusting the write path. A destination record with a differing payload hash is
reported as a conflict and left untouched - silently overwriting a stored audit
record is the same class of failure as losing it. Conflicts and failed
verification exit non-zero; a partial migration is not a success.

Tests 71 -> 77.

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

203 lines
7.8 KiB
Python

"""Move audit records from a SQLite store into PostgreSQL (AUDIT-WP-0005-T04).
As of 2026-08-10 there are no pre-production records to move: audit-core was
never deployed, so the only SQLite stores that ever existed were test
fixtures. See the workplan for that finding.
The tool exists anyway because the SQLite path is still reachable — the
entrypoint falls back to it when ``AUDIT_CORE_DATABASE_URL`` is unset. If that
fallback is ever used in anger, the records it accumulates are audit records
and cannot simply be dropped.
Fidelity is the whole point. Records are transferred with their original
``event_id``, ``payload_hash``, and ``accepted_at`` intact; a migration that
rewrote acceptance times would destroy the evidence it was meant to preserve.
That is why this bypasses ``accept()``, which stamps ``accepted_at`` with the
current time.
"""
from __future__ import annotations
import json
import sqlite3
from dataclasses import dataclass, field
from audit_core.interface import BackendUnavailableError
@dataclass
class MigrationReport:
"""What happened, in enough detail to be evidence."""
events_read: int = 0
events_written: int = 0
events_already_present: int = 0
dead_letters_written: int = 0
findings_written: int = 0
conflicts: list[str] = field(default_factory=list)
verified: bool = False
@property
def ok(self) -> bool:
return not self.conflicts and self.verified
def summary(self) -> str:
lines = [
f"events read: {self.events_read}",
f"events written: {self.events_written}",
f"already present: {self.events_already_present}",
f"dead letters written: {self.dead_letters_written}",
f"secret findings written:{self.findings_written}",
f"verified: {self.verified}",
]
if self.conflicts:
lines.append(f"CONFLICTS ({len(self.conflicts)}):")
lines += [f" - {c}" for c in self.conflicts]
return "\n".join(lines)
def migrate(sqlite_path: str, backend, *, verify: bool = True) -> MigrationReport:
"""Copy every record from ``sqlite_path`` into ``backend``.
Idempotent: re-running skips events already present with a matching payload
hash. An event present with a *different* hash is reported as a conflict
and never overwritten — the destination record wins, and a human decides.
"""
report = MigrationReport()
source = sqlite3.connect(f"file:{sqlite_path}?mode=ro", uri=True)
try:
rows = source.execute(
"SELECT event_id, payload_hash, accepted_at, correlation_id, tenant, record "
"FROM events ORDER BY accepted_at, event_id"
).fetchall()
report.events_read = len(rows)
for event_id, payload_hash, accepted_at, correlation_id, tenant, record in rows:
outcome = _import_event(
backend, event_id, payload_hash, accepted_at, correlation_id,
tenant, json.loads(record),
)
if outcome == "written":
report.events_written += 1
elif outcome == "present":
report.events_already_present += 1
else:
report.conflicts.append(outcome)
report.dead_letters_written = _copy_dead_letters(source, backend)
report.findings_written = _copy_findings(source, backend)
if verify:
report.verified = _verify(source, backend, report)
finally:
source.close()
return report
def _import_event(
backend, event_id, payload_hash, accepted_at, correlation_id, tenant, record
) -> str:
existing = backend._query(
f"SELECT payload_hash FROM {backend._events} WHERE event_id = %s", (event_id,)
)
if existing:
if existing[0][0] != payload_hash:
return (
f"{event_id}: destination holds a different payload "
f"(source {payload_hash[:12]}…, destination {existing[0][0][:12]}…) "
"— left untouched"
)
return "present"
backend._execute(
f"INSERT INTO {backend._events} "
"(event_id, payload_hash, accepted_at, observed_at, tenant, correlation_id, "
" source, action, record) "
"VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) "
"ON CONFLICT (event_id) DO NOTHING",
(
event_id, payload_hash, accepted_at, record.get("observed_at"),
tenant, correlation_id, record.get("source", ""), record.get("action", ""),
json.dumps(record, sort_keys=True),
),
)
return "written"
def _copy_dead_letters(source: sqlite3.Connection, backend) -> int:
try:
rows = source.execute(
"SELECT event_id, received_at, sender, reason, payload_hash, payload, "
"payload_withheld FROM dead_letters ORDER BY id"
).fetchall()
except sqlite3.Error:
return 0
for event_id, received_at, sender, reason, payload_hash, payload, withheld in rows:
backend._execute(
f'INSERT INTO "{backend.schema}".dead_letters '
"(event_id, received_at, sender, reason, payload_hash, payload, payload_withheld) "
"VALUES (%s, %s, %s, %s, %s, %s, %s)",
(event_id, received_at, sender, reason, payload_hash, payload, bool(withheld)),
)
return len(rows)
def _copy_findings(source: sqlite3.Connection, backend) -> int:
try:
rows = source.execute(
"SELECT sender, source, action, field_path, outcome, persisted, "
"occurrences, first_seen, last_seen FROM secret_findings"
).fetchall()
except sqlite3.Error:
return 0
for r in rows:
backend._execute(
f'INSERT INTO "{backend.schema}".secret_findings '
"(sender, source, action, field_path, outcome, persisted, occurrences, "
" first_seen, last_seen) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s) "
"ON CONFLICT (sender, source, action, field_path, outcome) DO UPDATE "
"SET occurrences = secret_findings.occurrences + excluded.occurrences",
(r[0], r[1], r[2], r[3], r[4], bool(r[5]), r[6], r[7], r[8]),
)
return len(rows)
def _verify(source: sqlite3.Connection, backend, report: MigrationReport) -> bool:
"""Confirm every source event is present with its hash and time intact.
Verification reads back rather than trusting the write path: a migration
that reports success without checking is the same as no migration at all.
"""
rows = source.execute("SELECT event_id, payload_hash, accepted_at FROM events").fetchall()
for event_id, payload_hash, accepted_at in rows:
found = backend._query(
f"SELECT payload_hash, accepted_at FROM {backend._events} WHERE event_id = %s",
(event_id,),
)
if not found:
report.conflicts.append(f"{event_id}: missing from destination after migration")
return False
if found[0][0] != payload_hash:
report.conflicts.append(f"{event_id}: payload hash differs after migration")
return False
if not _same_instant(found[0][1], accepted_at):
report.conflicts.append(
f"{event_id}: accepted_at not preserved "
f"(source {accepted_at}, destination {found[0][1]})"
)
return False
return True
def _same_instant(destination, source_iso: str) -> bool:
from datetime import datetime
try:
want = datetime.fromisoformat(str(source_iso).replace("Z", "+00:00"))
except ValueError:
return False
if not hasattr(destination, "timestamp"):
return False
if want.tzinfo is None or destination.tzinfo is None:
return False
return abs(destination.timestamp() - want.timestamp()) < 1.0