Add audit maintenance, verified recovery and reproducible verification
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e6f1-443f-7783-9920-a16b2ffc467f
This commit is contained in:
parent
820f0ff7b5
commit
3166da1d2e
13 changed files with 413 additions and 5 deletions
|
|
@ -214,6 +214,19 @@ class Store:
|
|||
return {row['status']: row['n'] for row in db.execute(
|
||||
"SELECT status,COUNT(*) AS n FROM outbox WHERE status!='delivered' GROUP BY status")}
|
||||
|
||||
def audit_metrics(self, last_drain):
|
||||
"""Private aggregate metrics: no event, actor, or alert labels."""
|
||||
with self.connect() as db:
|
||||
counts = dict(db.execute('SELECT status,COUNT(*) FROM outbox GROUP BY status'))
|
||||
lines = ['# HELP railiance_ack_audit_events Stored audit events by delivery state.',
|
||||
'# TYPE railiance_ack_audit_events gauge']
|
||||
for status in ('pending', 'blocked', 'delivered'):
|
||||
lines.append(f'railiance_ack_audit_events{{status="{status}"}} {counts.get(status, 0)}')
|
||||
lines += ['# HELP railiance_ack_audit_last_drain_timestamp_seconds Last completed local outbox pass; not receiver reachability.',
|
||||
'# TYPE railiance_ack_audit_last_drain_timestamp_seconds gauge',
|
||||
f'railiance_ack_audit_last_drain_timestamp_seconds {float(last_drain)}']
|
||||
return '\n'.join(lines) + '\n'
|
||||
|
||||
|
||||
class Application:
|
||||
"""WSGI adapter. authenticate and authorize are required, never default-allow.
|
||||
|
|
|
|||
130
scripts/alert_admin.py
Normal file
130
scripts/alert_admin.py
Normal file
|
|
@ -0,0 +1,130 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Local maintenance of the private acknowledgment store; no HTTP admin route."""
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import sqlite3
|
||||
import stat
|
||||
|
||||
from alert_ack import Store
|
||||
|
||||
TRIGGERS = {'ack_immutable_update', 'ack_immutable_delete', 'outbox_body_immutable',
|
||||
'outbox_immutable_delete', 'alert_immutable_update', 'alert_immutable_delete'}
|
||||
|
||||
|
||||
def private_path(path, existing=True):
|
||||
path = Path(path).absolute()
|
||||
parent = path.parent.lstat()
|
||||
if not stat.S_ISDIR(parent.st_mode) or parent.st_uid != os.getuid() or parent.st_mode & 0o077:
|
||||
raise ValueError('private owned directory required')
|
||||
if existing:
|
||||
info = path.lstat()
|
||||
if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid()
|
||||
or info.st_nlink != 1 or info.st_mode & 0o077):
|
||||
raise ValueError('private owned database required')
|
||||
return path
|
||||
|
||||
|
||||
def readonly(path):
|
||||
return sqlite3.connect(private_path(path).as_uri() + '?mode=ro', uri=True)
|
||||
|
||||
|
||||
def validate(db):
|
||||
if db.execute('PRAGMA integrity_check').fetchall() != [('ok',)] or db.execute('PRAGMA foreign_key_check').fetchall():
|
||||
raise ValueError('database integrity check failed')
|
||||
triggers = {r[0] for r in db.execute("SELECT name FROM sqlite_master WHERE type='trigger'")}
|
||||
if not TRIGGERS <= triggers:
|
||||
raise ValueError('acknowledgment schema required')
|
||||
# Every durable acknowledgment has exactly its original outbox event.
|
||||
if db.execute('''SELECT COUNT(*) FROM acknowledgments a LEFT JOIN outbox o
|
||||
ON a.event_id=o.event_id WHERE o.event_id IS NULL''').fetchone()[0]:
|
||||
raise ValueError('incomplete audit outbox')
|
||||
if db.execute('''SELECT COUNT(*) FROM outbox o LEFT JOIN acknowledgments a
|
||||
ON a.event_id=o.event_id WHERE a.event_id IS NULL
|
||||
OR o.status NOT IN ('pending','blocked','delivered')''').fetchone()[0]:
|
||||
raise ValueError('invalid audit outbox')
|
||||
|
||||
|
||||
def snapshot(source, destination):
|
||||
"""SQLite online backup, never overwrite; restore uses the same fresh-target path.
|
||||
|
||||
Source may be live for backup. Restore to a NEW path, then switch the stopped
|
||||
service to it. Neither operation rewinds or replaces an existing database.
|
||||
"""
|
||||
target = private_path(destination, existing=False)
|
||||
source_db = readonly(source)
|
||||
created = False
|
||||
try:
|
||||
fd = os.open(target, os.O_CREAT | os.O_EXCL | os.O_WRONLY | os.O_NOFOLLOW, 0o600)
|
||||
os.close(fd)
|
||||
created = True
|
||||
target_db = sqlite3.connect(target)
|
||||
try:
|
||||
source_db.backup(target_db)
|
||||
validate(target_db)
|
||||
finally:
|
||||
target_db.close()
|
||||
with target.open('rb') as saved:
|
||||
os.fsync(saved.fileno())
|
||||
directory = os.open(target.parent, os.O_RDONLY | os.O_DIRECTORY)
|
||||
try:
|
||||
os.fsync(directory)
|
||||
finally:
|
||||
os.close(directory)
|
||||
except Exception:
|
||||
if created:
|
||||
target.unlink()
|
||||
raise
|
||||
finally:
|
||||
source_db.close()
|
||||
|
||||
|
||||
def main(argv=None):
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('--database', required=True, type=Path)
|
||||
commands = parser.add_subparsers(dest='action', required=True)
|
||||
commands.add_parser('status')
|
||||
listing = commands.add_parser('list-blocked')
|
||||
listing.add_argument('--limit', type=int, default=100)
|
||||
retry = commands.add_parser('requeue')
|
||||
retry.add_argument('--event-id', required=True)
|
||||
for action in ('backup', 'restore'):
|
||||
command = commands.add_parser(action)
|
||||
command.add_argument('--destination', type=Path, required=True)
|
||||
args = parser.parse_args(argv)
|
||||
try:
|
||||
if args.action in ('backup', 'restore'):
|
||||
snapshot(args.database, args.destination)
|
||||
result = {'status': 'verified', 'action': args.action}
|
||||
else:
|
||||
db = readonly(args.database)
|
||||
try:
|
||||
validate(db)
|
||||
counts = dict(db.execute('SELECT status,COUNT(*) FROM outbox GROUP BY status'))
|
||||
if args.action == 'status':
|
||||
result = {'outbox': {s: counts.get(s, 0) for s in ('pending','blocked','delivered')}}
|
||||
elif args.action == 'list-blocked':
|
||||
if not 1 <= args.limit <= 1000:
|
||||
raise ValueError('limit must be 1..1000')
|
||||
result = {'event_ids': [r[0] for r in db.execute(
|
||||
"SELECT event_id FROM outbox WHERE status='blocked' ORDER BY event_id LIMIT ?", (args.limit,))]}
|
||||
else:
|
||||
# Requeue only, never manufacture delivery or change evidence.
|
||||
store = Store(args.database)
|
||||
changed = store.requeue(args.event_id)
|
||||
result = {'event_id': args.event_id, 'requeued': changed}
|
||||
if not changed:
|
||||
print(json.dumps(result))
|
||||
return 2
|
||||
finally:
|
||||
db.close()
|
||||
print(json.dumps(result, sort_keys=True))
|
||||
return 0
|
||||
except (OSError, ValueError, sqlite3.Error):
|
||||
print(json.dumps({'status': 'refused', 'reason': 'maintenance failed; check paths, permissions and database integrity'}))
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
raise SystemExit(main())
|
||||
|
|
@ -7,6 +7,7 @@ import os
|
|||
from pathlib import Path
|
||||
import secrets
|
||||
import signal
|
||||
import sqlite3
|
||||
import threading
|
||||
import time
|
||||
from urllib.parse import parse_qs
|
||||
|
|
@ -61,6 +62,11 @@ class Router:
|
|||
return ('Set-Cookie', f'{name}={value}; Path=/; Secure; HttpOnly; SameSite=Lax; Max-Age={age}')
|
||||
path, method = env.get('PATH_INFO'), env.get('REQUEST_METHOD')
|
||||
try:
|
||||
if method == 'GET' and path == '/metrics':
|
||||
body = self.store.audit_metrics(self.last_drain)
|
||||
start('200 OK', [('Content-Type', 'text/plain; version=0.0.4; charset=utf-8'),
|
||||
('Cache-Control', 'no-store')])
|
||||
return [body.encode()]
|
||||
if method == 'GET' and path in ('/healthz', '/readyz'):
|
||||
ready = path == '/healthz' or time.time() - self.last_drain < 90 and not self.store.audit_debt()
|
||||
return reply('200 OK' if ready else '503 Service Unavailable', 'ready' if ready else 'audit delivery pending')
|
||||
|
|
@ -87,7 +93,7 @@ class Router:
|
|||
if path == '/webhook':
|
||||
self.application.webhook_token = credential(self.config['webhook_token_file'])
|
||||
return self.application(env, start)
|
||||
except (ValueError, OSError, KeyError, TypeError):
|
||||
except (ValueError, OSError, KeyError, TypeError, sqlite3.Error):
|
||||
return reply('503 Service Unavailable', 'Identity or service unavailable. Retry later.')
|
||||
|
||||
|
||||
|
|
@ -109,7 +115,7 @@ def main():
|
|||
try:
|
||||
store.drain(transport)
|
||||
router.last_drain = time.time()
|
||||
except (OSError, ValueError):
|
||||
except (OSError, ValueError, sqlite3.Error):
|
||||
router.last_drain = 0
|
||||
stop.wait(30)
|
||||
worker = threading.Thread(target=drain, daemon=True)
|
||||
|
|
|
|||
26
scripts/verify.sh
Normal file
26
scripts/verify.sh
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
#!/usr/bin/env bash
|
||||
# Full Python and native owner-contract suite; fail if native dependencies are absent.
|
||||
set -euo pipefail
|
||||
cd -- "$(dirname -- "${BASH_SOURCE[0]}")/.."
|
||||
repo_root="$PWD"
|
||||
audit_source="${RTEL_AUDIT_CORE_SOURCE:-$repo_root/../audit-core}"
|
||||
policy_source="${RTEL_FLEX_AUTH_SOURCE:-$repo_root/../flex-auth}"
|
||||
if [[ ! -f "$audit_source/audit_core/ingestion.py" || ! -f "$policy_source/go.mod" ]]; then
|
||||
echo 'Required: sibling audit-core and flex-auth checkouts, or RTEL_AUDIT_CORE_SOURCE / RTEL_FLEX_AUTH_SOURCE.' >&2
|
||||
exit 1
|
||||
fi
|
||||
mkdir -p .cache/verify
|
||||
export UV_CACHE_DIR="$repo_root/.cache/verify/uv"
|
||||
export GOCACHE="$repo_root/.cache/verify/go-build"
|
||||
if [[ ! -x .venv-verify/bin/python ]]; then
|
||||
uv venv --python python3 .venv-verify
|
||||
fi
|
||||
uv pip sync --python .venv-verify/bin/python --require-hashes requirements-runtime.lock
|
||||
go -C "$policy_source" build -mod=readonly -o "$repo_root/.cache/verify/flex-auth" ./cmd/flex-auth
|
||||
export RTEL_AUDIT_CORE_SOURCE="$audit_source"
|
||||
export RTEL_FLEX_AUTH_BINARY="$repo_root/.cache/verify/flex-auth"
|
||||
echo 'Native owner source revisions:'
|
||||
git -C "$audit_source" rev-parse HEAD
|
||||
git -C "$policy_source" rev-parse HEAD
|
||||
.venv-verify/bin/python -m unittest discover -s tests -v
|
||||
.venv-verify/bin/python -m unittest discover -s tests_runtime -v
|
||||
Loading…
Add table
Add a link
Reference in a new issue