Files
truf-server/app/container_projection_recovery.py
2026-09-30 20:30:56 +03:00

417 lines
28 KiB
Python

"""Explicit maintenance-only loss acknowledgement for the reviewed copied output.
No CLI, startup hook, connection creation, PostgreSQL lifecycle, or output rebuild.
The caller must retain initialize.lock and ClusterAuthorityLock through this call
AND subsequent positively verified maintenance stop, including every exception or
uncertain commit. It must keep the target isolated with no workers/network clients.
Use an idle, writable, autocommit=True psycopg connection with dict_row and quiet server log
settings. The supplied runtime is the prepared container_runtime module, not its
main()/initialize() entrypoint. Never call this during normal initialized startup.
The immutable PREPARED journal describes intent, not a fabricated completed rename.
An exact before state can be applied; an exact after state is a read-only retry.
Partial journals and later output require review, never automatic cleanup/rewind.
"""
import hashlib
import json
import os
from pathlib import Path
import re
import time
from datetime import datetime, timezone
DATA = Path('/data')
RUN = Path('/run/truf')
APPROVED_MANIFEST_SHA256 = '08344147133c37d4b6f404cf4fac3e59d58f94917f1fa58a77cbb68c36db7e8a'
JOURNAL_NAME = 'found-secrets-loss-g13-g14.prepared.json'
FORMAT = 'truf-found-secrets-loss-g13-g14-v1'
LOSS = 'Previously published copied output intentionally lost; all PostgreSQL history retained. No rotation or rename occurred.'
OLD_GENERATION, NEW_GENERATION, OLD_OFFSET = 13, 14, 97783145
APPEND_COUNT, ROTATION_COUNT = 38024, 13
MAX_JOURNAL = 1024 * 1024
STREAM_COLUMNS = {'stream_name', 'base_relative_path', 'current_generation', 'rotation_bytes',
'max_generations', 'created_at', 'updated_at'}
CURSOR_COLUMNS = {'stream_name', 'generation', 'committed_offset', 'last_append_id', 'last_job_id',
'last_event_id', 'last_event_hash', 'updated_at'}
class ProjectionRecoveryError(RuntimeError):
"""Safe diagnostic only; caller still owns maintenance/stop authority."""
def _encoded(value):
return (json.dumps(value, ensure_ascii=True, sort_keys=True,
separators=(',', ':'), allow_nan=False) + '\n').encode('ascii')
def recover_found_secrets_projection(runtime, connection, *, system_identifier, manifest_sha256,
initialize_lock, authority_lock):
"""Apply only the pinned g13 loss transition, or recognize its exact retry.
Caller-owned locks must be acquired runtime_security lock objects for the fixed
target. This function never closes the connection or releases those locks.
Its own file lock and transaction-scoped advisory/table locks exclude writers.
Returned 'committed'/'already-committed' is not target readiness or stop proof.
journal_sha256 hashes the complete immutable file, not just its inner record.
"""
stage = 'preflight'
deadline = time.monotonic() + 10800
try:
from container_import import _fsync_dir, _identifier, _input, _integer, _json, _manifest, _regular, _write, MAX_MANIFEST
from runtime_security import ClusterAuthorityLock, PrivateFileLock
runtime.require_container()
if (runtime.DATA != DATA or runtime.RUN != RUN
or runtime.INITIALIZED != DATA / 'initialized.json'
or runtime.INITIALIZE_LOCK != DATA / 'initialize.lock'
or manifest_sha256 != APPROVED_MANIFEST_SHA256
or not isinstance(system_identifier, str)
or re.fullmatch(r'[1-9][0-9]{0,19}', system_identifier) is None
or connection.closed or connection.broken or connection.autocommit is not True
or int(connection.info.transaction_status) != 0):
raise ValueError()
config_dir, results = DATA / 'config', DATA / 'runtime-linux/results'
identity_path = DATA / 'runtime-linux/postgres/cluster_identity.json'
manifest_path = config_dir / 'windows-import-manifest.json'
journal_path = config_dir / JOURNAL_NAME
projector_lock = results / '.jsonl-projector.lock'
endpoint = _encoded({'database': 'truf', 'host': '127.0.0.1', 'port': 5432,
'schema': 'public'}).decode('ascii').strip()
data_hash = hashlib.sha256(str(DATA / 'postgres-linux').encode('utf-8')).hexdigest()
endpoint_hash = hashlib.sha256(endpoint.encode('ascii')).hexdigest()
def files_stopped():
if (time.monotonic() >= deadline or getattr(runtime, '_shutdown_requested', True) is not False
or not isinstance(initialize_lock, PrivateFileLock) or initialize_lock.acquired is not True
or initialize_lock.path != os.path.normcase(str(runtime.INITIALIZE_LOCK))
or not isinstance(authority_lock, ClusterAuthorityLock) or authority_lock.acquired is not True
or authority_lock.data_directory != str(DATA / 'postgres-linux')
or authority_lock.endpoint_identity != endpoint
or authority_lock.path != str(DATA / 'runtime-linux/postgres' / f'.cluster-authority-{data_hash}.lock')
or authority_lock.endpoint_path != str(RUN / 'authority' / f'endpoint-{endpoint_hash}.lock')):
raise ValueError()
for directory in (DATA, config_dir, results, identity_path.parent,
DATA / 'runtime-linux/logs', RUN, RUN / 'control'):
runtime.private_path(directory, directory=True)
_regular(runtime.INITIALIZE_LOCK, runtime)
for path in (runtime.INITIALIZED, RUN / 'control/supervisor.instance.json',
RUN / 'control/supervisor.pid', DATA / 'runtime-linux/logs/supervisor.instance.json',
DATA / 'runtime-linux/logs/supervisor.pid'):
try:
path.lstat()
except FileNotFoundError:
continue
raise ValueError()
with os.scandir(results) as entries:
for index, entry in enumerate(entries):
if index >= 100000 or entry.name.casefold() == 'found_secrets' or entry.name.casefold().startswith('found_secrets.'):
raise ValueError()
def read_private(path, limit):
with _input(path, runtime) as (handle, before):
if not 0 < before[4] <= limit:
raise ValueError()
raw = handle.read(limit + 1)
if len(raw) != before[4]:
raise ValueError()
return raw, _json(raw)
files_stopped()
identity_raw, identity = read_private(identity_path, MAX_JOURNAL)
manifest_raw, _ = read_private(manifest_path, MAX_MANIFEST)
manifest, _ = _manifest(manifest_raw, manifest_sha256)
expected_identity = {'pg_major': 16, 'system_identifier': system_identifier,
'data_directory': str(DATA / 'postgres-linux'), 'database': 'truf',
'user': 'truf', 'port': 5432}
if (not isinstance(identity, dict) or any(identity.get(k) != v for k, v in expected_identity.items())
or manifest['database']['system_identifier'] == system_identifier
or 'sequence_states' not in manifest['database']):
raise ValueError()
binding = {'system_identifier': system_identifier, 'manifest_sha256': manifest_sha256,
'identity_sha256': hashlib.sha256(identity_raw).hexdigest()}
def online():
if time.monotonic() >= deadline or runtime._shutdown_requested:
raise ValueError()
connection.execute('SELECT pg_catalog.pg_stat_clear_snapshot()')
row = connection.execute("""SELECT pg_catalog.current_database() AS database,
current_user AS user_name, pg_catalog.current_setting('data_directory') AS data_directory,
pg_catalog.current_setting('port')::int AS port,
pg_catalog.current_setting('server_version_num')::int AS version_num,
pg_catalog.pg_is_in_recovery() AS in_recovery,
(SELECT system_identifier::text FROM pg_catalog.pg_control_system()) AS system_identifier,
(SELECT rolsuper FROM pg_catalog.pg_roles WHERE rolname = current_user) AS superuser,
(SELECT count(*) FROM pg_catalog.pg_stat_activity WHERE backend_type = 'client backend'
AND pid <> pg_catalog.pg_backend_pid()) AS other_clients,
pg_catalog.current_schema() AS schema_name,
pg_catalog.current_setting('search_path') AS search_path,
pg_catalog.current_setting('transaction_read_only') AS read_only,
pg_catalog.current_setting('fsync') AS fsync,
pg_catalog.current_setting('full_page_writes') AS full_page_writes,
EXISTS (SELECT 1 FROM pg_catalog.pg_namespace n CROSS JOIN LATERAL pg_catalog.aclexplode(
COALESCE(n.nspacl, pg_catalog.acldefault('n', n.nspowner))) acl
WHERE n.nspname = 'public' AND acl.grantee = 0 AND acl.privilege_type = 'CREATE') AS public_create""").fetchone()
expected = {'database': 'truf', 'user_name': 'truf', 'data_directory': str(DATA / 'postgres-linux'),
'port': 5432, 'system_identifier': system_identifier, 'in_recovery': False,
'superuser': True, 'other_clients': 0, 'schema_name': 'public',
'search_path': 'public', 'read_only': 'off', 'public_create': False,
'fsync': 'on', 'full_page_writes': 'on'}
if (not isinstance(row, dict) or any(row.get(k) != v for k, v in expected.items())
or _integer(row.get('version_num')) // 10000 != 16):
raise ValueError()
def metadata():
rows = connection.execute("""SELECT s.stream_name, c.stream_name AS cursor_stream_name,
pg_catalog.row_to_json(s) AS stream, pg_catalog.row_to_json(c) AS cursor
FROM public.projection_streams s FULL JOIN public.projection_cursors c
ON c.stream_name = s.stream_name""").fetchall()
seen, found = set(), None
scans = {'scan_results': 'scan_results.jsonl', 'found_secrets': 'found_secrets.jsonl',
'scan_errors': 'scan_errors.log'}
for row in rows:
name, stream, cursor = row['stream_name'], row['stream'], row['cursor']
if (not isinstance(name, str) or name in seen or row['cursor_stream_name'] != name
or not isinstance(stream, dict) or set(stream) != STREAM_COLUMNS
or not isinstance(cursor, dict) or set(cursor) != CURSOR_COLUMNS
or stream['stream_name'] != name or cursor['stream_name'] != name
or _integer(stream['current_generation']) != _integer(cursor['generation'])):
raise ValueError()
seen.add(name)
_integer(cursor['committed_offset'])
_integer(stream['rotation_bytes'], 1)
_integer(stream['max_generations'])
expected_path = scans.get(name)
if expected_path is None:
match = re.fullmatch(r'keycheck:([a-z0-9][a-z0-9_.-]{0,63}):(results|status)', name)
if not match:
raise ValueError()
suffix = 'Results.jsonl' if match[2] == 'results' else 'Checked.txt'
expected_path = f'{match[1]}/{match[1]}{suffix}'
if stream['base_relative_path'] != expected_path:
raise ValueError()
if name == 'found_secrets':
found = {'stream': stream, 'cursor': cursor}
if len(seen) != 34 or not scans.keys() <= seen:
raise ValueError()
return found
tables = sorted(manifest['database']['table_counts'])
sequences = manifest['database']['sequence_states']
if set(sequences) != {'public'}:
raise ValueError()
def preserved():
proof = {}
for table in tables:
if time.monotonic() >= deadline or runtime._shutdown_requested:
raise ValueError()
where = " WHERE t.stream_name <> 'found_secrets'" if table in ('projection_streams', 'projection_cursors') else ''
digest, count = hashlib.sha256(), 0
with connection.cursor(name='found_loss_digest') as cursor:
cursor.itersize = 1000
cursor.execute("SELECT pg_catalog.encode(pg_catalog.sha256(pg_catalog.convert_to("
"pg_catalog.row_to_json(t)::text, 'UTF8')), 'hex') COLLATE \"C\" AS digest FROM public."
+ _identifier(table) + ' AS t' + where + ' ORDER BY digest')
for row in cursor:
value = row['digest']
if not isinstance(value, str) or re.fullmatch(r'[0-9a-f]{64}', value) is None:
raise ValueError()
digest.update(value.encode('ascii'))
count += 1
if count % 1000 == 0 and (time.monotonic() >= deadline or runtime._shutdown_requested):
raise ValueError()
if count != manifest['database']['table_counts'][table] - int(bool(where)):
raise ValueError()
proof[table] = {'rows': count, 'sha256': digest.hexdigest()}
rows = connection.execute("""SELECT c.relname AS name FROM pg_catalog.pg_class c
JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = 'public' AND c.relkind = 'S' ORDER BY c.relname""").fetchall()
if {row['name'] for row in rows} != set(sequences['public']):
raise ValueError()
for row in rows:
state = connection.execute('SELECT last_value, is_called FROM public.' + _identifier(row['name'])).fetchone()
if _encoded(state) != _encoded(sequences['public'][row['name']]):
raise ValueError()
return {'tables': proof, 'sequences': {'count': len(rows),
'sha256': hashlib.sha256(_encoded(sequences)).hexdigest()}}
with PrivateFileLock(str(projector_lock)):
_regular(projector_lock, runtime)
with connection.transaction():
stage = 'database-fences'
connection.execute('SET TRANSACTION ISOLATION LEVEL READ COMMITTED')
online()
connection.execute("SET LOCAL lock_timeout = '5s'")
connection.execute("SET LOCAL statement_timeout = '10800s'")
connection.execute("SET LOCAL temp_file_limit = '4GB'")
connection.execute("SET LOCAL work_mem = '128MB'")
connection.execute('SET LOCAL max_parallel_workers_per_gather = 0')
connection.execute("SET LOCAL row_security = off")
connection.execute('SET LOCAL synchronous_commit = on')
for setting, value in (('log_min_error_statement', 'panic'), ('log_min_messages', 'panic'),
('log_statement', 'none'), ('log_min_duration_statement', '-1')):
connection.execute('SET LOCAL ' + setting + " = '" + value + "'")
for arguments in ((1414681926, 1785753445), (1414681926, 1768842867), (781273968142991337,)):
lock = connection.execute('SELECT pg_catalog.pg_try_advisory_xact_lock('
+ ','.join(['%s'] * len(arguments)) + ') AS locked', arguments).fetchone()
if not lock or lock['locked'] is not True:
raise ValueError()
catalog_sql = """SELECT c.relname AS name, c.relkind AS kind FROM pg_catalog.pg_class c
JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = 'public' AND c.relkind IN ('r','p','f') ORDER BY c.relname"""
catalog = connection.execute(catalog_sql).fetchall()
if len(catalog) != len(tables) or {row['name'] for row in catalog} != set(tables) or any(row['kind'] != 'r' for row in catalog):
raise ValueError()
connection.execute('LOCK TABLE ' + ','.join('public.' + _identifier(table) for table in tables)
+ ' IN EXCLUSIVE MODE NOWAIT')
online()
stage = 'history-gates'
gates = connection.execute("""WITH a AS (
SELECT count(*) AS appends, max(a.generation) AS append_highwater,
count(*) FILTER (WHERE a.state IS DISTINCT FROM 'appended' OR j.id IS NULL
OR j.status IS DISTINCT FROM 'completed' OR j.capacity_released IS DISTINCT FROM 1
OR j.job_kind IS DISTINCT FROM 'scan_event' OR (j.required_stream_mask & 2) IS DISTINCT FROM 2
OR a.event_id IS DISTINCT FROM j.event_id OR a.event_hash IS DISTINCT FROM j.event_hash
OR a.generation IS NULL OR a.generation NOT BETWEEN 0 AND 13
OR a.byte_offset IS NULL OR a.byte_offset < 0 OR a.byte_length IS NULL OR a.byte_length < 0
OR a.record_count IS NULL OR a.record_count < 0
OR (a.generation = 13 AND a.byte_offset::numeric + a.byte_length::numeric > 97783145)) AS bad_appends
FROM public.projection_appends a LEFT JOIN public.projection_jobs j ON j.id = a.job_id
WHERE a.stream_name = 'found_secrets'), r AS (
SELECT count(*) AS rotations, max(to_generation) AS rotation_highwater,
count(*) FILTER (WHERE state IS DISTINCT FROM 'completed' OR from_generation IS NULL OR from_generation < 0
OR to_generation IS DISTINCT FROM from_generation + 1 OR to_generation > 13
OR source_bytes IS NULL OR source_bytes < 0
OR segment_relative_path IS DISTINCT FROM
'found_secrets.g' || lpad(from_generation::text, 6, '0') || '.jsonl') AS bad_rotations
FROM public.projection_rotations WHERE stream_name = 'found_secrets')
SELECT a.*, r.*, (SELECT count(*) FROM public.projection_append_audit) AS audits,
(SELECT count(*) FROM public.projection_jobs WHERE (required_stream_mask & 2) <> 0
AND status <> 'completed') AS pending,
(SELECT count(*) FROM public.result_reservations
WHERE state IN ('scanning','ready','ingesting','db_committed')) AS reservations,
(SELECT count(*) FROM public.target_queue q LEFT JOIN public.result_reservations r
ON r.id = q.current_result_reservation_id WHERE q.status = 'in_progress'
OR r.state IN ('scanning','ready','ingesting','db_committed')) AS queue_leases,
(SELECT count(*) FROM public.docker_content_blobs
WHERE state IN ('leased','submitted') OR lease_reservation_id IS NOT NULL) AS blob_leases,
(SELECT count(*) FROM pg_catalog.pg_trigger WHERE NOT tgisinternal
AND tgrelid IN ('public.projection_streams'::regclass,'public.projection_cursors'::regclass)) AS triggers,
(SELECT count(*) FROM pg_catalog.pg_rewrite
WHERE ev_class IN ('public.projection_streams'::regclass,'public.projection_cursors'::regclass)) AS rules
FROM a CROSS JOIN r""").fetchone()
expected_gates = {'appends': APPEND_COUNT, 'append_highwater': 13, 'bad_appends': 0,
'rotations': ROTATION_COUNT, 'rotation_highwater': 13, 'bad_rotations': 0,
'audits': 0, 'pending': 0, 'reservations': 0, 'queue_leases': 0,
'blob_leases': 0, 'triggers': 0, 'rules': 0}
if _encoded(gates) != _encoded(expected_gates):
raise ValueError()
current = metadata()
proof = preserved()
stage = 'journal'
try:
journal_path.lstat()
except FileNotFoundError:
journal = None
else:
raw, journal = read_private(journal_path, MAX_JOURNAL)
if (not isinstance(journal, dict) or set(journal) != {'record', 'sha256'}
or raw != _encoded(journal)
or journal['sha256'] != hashlib.sha256(_encoded(journal['record'])).hexdigest()):
raise ValueError()
if journal is None:
stamp = datetime.now(timezone.utc).isoformat(timespec='seconds')
before = current
record = {'format': FORMAT, 'state': 'PREPARED', 'loss': LOSS, 'binding': binding,
'prepared_at': stamp, 'before': before, 'preserved': proof,
'digest_algorithm': 'sha256-concatenated-sorted-pg-row-sha256-hex-v1'}
else:
record = journal['record']
if (not isinstance(record, dict) or set(record) != {'format', 'state', 'loss', 'binding',
'prepared_at', 'before', 'after', 'preserved', 'digest_algorithm'}
or record['format'] != FORMAT or record['state'] != 'PREPARED' or record['loss'] != LOSS
or _encoded(record['binding']) != _encoded(binding)
or _encoded(record['preserved']) != _encoded(proof)
or record['digest_algorithm'] != 'sha256-concatenated-sorted-pg-row-sha256-hex-v1'):
raise ValueError()
stamp, before = record['prepared_at'], record['before']
if (not isinstance(stamp, str) or datetime.fromisoformat(stamp).isoformat(timespec='seconds') != stamp
or not stamp.endswith('+00:00') or set(before) != {'stream', 'cursor'}
or set(before['stream']) != STREAM_COLUMNS or set(before['cursor']) != CURSOR_COLUMNS
or before['stream']['stream_name'] != 'found_secrets'
or before['stream']['base_relative_path'] != 'found_secrets.jsonl'
or before['cursor']['stream_name'] != 'found_secrets'
or _integer(before['stream']['current_generation']) != OLD_GENERATION
or _integer(before['cursor']['generation']) != OLD_GENERATION
or _integer(before['cursor']['committed_offset']) != OLD_OFFSET):
raise ValueError()
_integer(before['cursor']['last_append_id'], 1)
_integer(before['cursor']['last_job_id'], 1)
last = connection.execute("""SELECT id, job_id, stream_name, generation, byte_offset,
byte_length, event_id, event_hash, state FROM public.projection_appends WHERE id = %s""",
(before['cursor']['last_append_id'],)).fetchone()
if (not last or last['id'] != before['cursor']['last_append_id']
or last['job_id'] != before['cursor']['last_job_id'] or last['stream_name'] != 'found_secrets'
or last['generation'] != OLD_GENERATION or last['state'] != 'appended'
or _integer(last['byte_offset']) + _integer(last['byte_length']) != OLD_OFFSET
or last['event_id'] != before['cursor']['last_event_id']
or last['event_hash'] != before['cursor']['last_event_hash']
or not isinstance(last['event_id'], str) or not last['event_id']
or re.fullmatch(r'[0-9a-f]{64}', last['event_hash'] or '') is None):
raise ValueError()
after = {key: dict(value) for key, value in before.items()}
after['stream'].update(current_generation=NEW_GENERATION, updated_at=stamp)
after['cursor'].update(generation=NEW_GENERATION, committed_offset=0, last_append_id=None, updated_at=stamp)
if journal is not None and _encoded(record['after']) != _encoded(after):
raise ValueError()
already = _encoded(current) == _encoded(after)
if (not already and _encoded(current) != _encoded(before)) or (already and journal is None):
raise ValueError()
if journal is None:
record['after'] = after
journal = {'record': record, 'sha256': hashlib.sha256(_encoded(record)).hexdigest()}
encoded = _encoded(journal)
if len(encoded) > MAX_JOURNAL:
raise ValueError()
_write(runtime, journal_path, encoded)
if read_private(journal_path, MAX_JOURNAL)[0] != encoded:
raise ValueError()
# A prior failure may have left complete bytes without a confirmed fsync.
with _input(journal_path, runtime) as (handle, _):
if handle.read(MAX_JOURNAL + 1) != _encoded(journal):
raise ValueError()
os.fsync(handle.fileno())
_fsync_dir(config_dir)
if not already:
stage = 'compare-and-swap'
result = connection.execute("""UPDATE public.projection_streams AS s
SET current_generation = 14, updated_at = %s
WHERE s.stream_name = 'found_secrets' AND pg_catalog.to_jsonb(s) = %s::jsonb""",
(stamp, _encoded(before['stream']).decode('ascii')))
if result.rowcount != 1:
raise ValueError()
result = connection.execute("""UPDATE public.projection_cursors AS c
SET generation = 14, committed_offset = 0, last_append_id = NULL, updated_at = %s
WHERE c.stream_name = 'found_secrets' AND pg_catalog.to_jsonb(c) = %s::jsonb""",
(stamp, _encoded(before['cursor']).decode('ascii')))
if result.rowcount != 1:
raise ValueError()
if _encoded(metadata()) != _encoded(after) or _encoded(preserved()) != _encoded(proof):
raise ValueError()
stage = 'precommit'
files_stopped()
if (read_private(identity_path, MAX_JOURNAL)[0] != identity_raw
or read_private(manifest_path, MAX_MANIFEST)[0] != manifest_raw
or read_private(journal_path, MAX_JOURNAL)[0] != _encoded(journal)
or _encoded(connection.execute(catalog_sql).fetchall()) != _encoded(catalog)):
raise ValueError()
online()
stage = 'commit'
return {'status': 'already-committed' if already else 'committed', 'journal_path': str(journal_path),
'journal_sha256': hashlib.sha256(_encoded(journal)).hexdigest(), **binding, 'generation': NEW_GENERATION}
except BaseException:
raise ProjectionRecoveryError('Projection recovery refused at ' + stage
+ '; retain caller maintenance authority and verify stop.') from None