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>
203 lines
7.8 KiB
Python
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
|