import sys sys.dont_write_bytecode = True import argparse import contextlib import hashlib import json import os import re import sqlite3 import stat import time from datetime import datetime, timezone from paths import apply_path_config from query_policy import validate_rejected_query_policy from db_backend import ( DatabaseUrlError, POSTGRES_APPLICATION_SCHEMA, canonical_postgres_url, database_url_from_env, is_postgres_url, ) from postgres_runtime import ( _listener_present, configured_cluster_values, load_postgres_environment, postgres_runtime_paths, verify_cluster_identity, ) from scanner_db import ( PIPELINE_MIGRATION_VERSIONS, ScanEventConflictError, ScannerDB, json_dumps, migrate_runtime_safety_schema, normalize_target, utc_now_iso, ) from process_identity import current_process_identity from result_ingester import ResultIngester from result_spool import ResultSpool from result_bundle import ResultBundleReader from runtime_security import ( atomic_write_private_json, ClusterAuthorityLock, PrivatePathState, PrivateFileError, canonical_cluster_data_directory, canonical_path, harden_private_directory, harden_private_file, harden_private_tree, durable_publish, durable_unlink, ensure_private_directory, fsync_directory, is_reparse_point, private_directory_ready, private_file_ready, inspect_private_relative_path, read_private_json, require_private_directory, require_private_file, require_trusted_native_executable, reject_reparse_components, sha256_file, preflight_lifecycle_paths, ) RECONCILIATION_ABSOLUTE_ROW_BYTES = 16 * 1024 * 1024 RECONCILIATION_MUTATION_RETRIES = 3 TARGET_QUEUE_POLICY_MANIFEST_MAX_BYTES = 16 * 1024 * 1024 REVIEWED_LEGACY_RESULT_SPOOL = os.path.normcase( os.path.abspath(r'D:\truf\runtime\result_spool') ) class ReconciliationFileChanged(RuntimeError): pass def reviewed_legacy_result_spool(config): global_config = (config or {}).get('global') or {} value = global_config.get('legacy_result_spool_dir') or global_config.get('result_spool_dir') if not value or '{' in str(value) or '}' in str(value): raise RuntimeError('legacy result spool path is unresolved') if os.name != 'nt': runtime = global_config.get('runtime_dir') for label, path in (('runtime', runtime), ('legacy result spool', value)): if ( not isinstance(path, str) or not os.path.isabs(path) or path != os.path.normpath(path) or any(char in path for char in ('{', '}', '\\', '\x00')) ): raise RuntimeError(f'{label} path must be resolved, exact and absolute') reject_reparse_components(path) expected = os.path.join(runtime, 'result_spool') if value != expected: raise RuntimeError('legacy result spool must be exactly runtime_dir/result_spool') for path in (runtime, value): require_private_directory(path, create=False) resolved = canonical_path(value) if resolved != os.path.join(canonical_path(runtime), 'result_spool'): raise RuntimeError('legacy result spool resolved outside its runtime authority') for name in ('reservations', 'quarantine'): path = os.path.join(value, name) if os.path.lexists(path): require_private_directory(path, create=False) return resolved resolved = os.path.normcase(os.path.abspath(value)) if resolved != REVIEWED_LEGACY_RESULT_SPOOL: raise RuntimeError( 'legacy result spool path does not match reviewed D:\\truf\\runtime\\result_spool authority' ) return os.path.abspath(value) def _mtime_ns(stat_result): return int(getattr(stat_result, 'st_mtime_ns', int(stat_result.st_mtime * 1_000_000_000))) def _stat_identity(stat_result): inode = int(getattr(stat_result, 'st_ino', 0) or 0) device = int(getattr(stat_result, 'st_dev', 0) or 0) return device, inode, int(stat_result.st_size), _mtime_ns(stat_result) def todo_file_identity(path, stat_result=None): stat_result = stat_result or os.stat(path, follow_symlinks=False) device, inode, size, mtime_ns = _stat_identity(stat_result) return f'{device}:{inode}:{size}:{mtime_ns}:{os.path.realpath(os.path.abspath(path))}' def _safe_target_preview(value): if isinstance(value, bytes): payload = value else: payload = str(value or '').encode('utf-8', errors='replace') return f'' def _same_file_snapshot(path, handle_stat): try: current_handle = os.fstat(handle_stat[0]) if isinstance(handle_stat, tuple) else handle_stat current_path = os.stat(path, follow_symlinks=False) except OSError: return False expected = handle_stat[1] if isinstance(handle_stat, tuple) else _stat_identity(handle_stat) return _stat_identity(current_handle) == expected and _stat_identity(current_path) == expected def docker_target_is_bare(target): text = str(target or '').strip() if not text or '@' in text: return False return ':' not in text.rsplit('/', 1)[-1] def _contains_nul(value): if isinstance(value, str): return '\x00' in value if isinstance(value, dict): return any(_contains_nul(key) or _contains_nul(item) for key, item in value.items()) if isinstance(value, (list, tuple)): return any(_contains_nul(item) for item in value) return False def reconciliation_target_error(target, platform): if not target: return 'empty row' if '\x00' in target: return 'NUL byte in row' if platform == 'docker' and docker_target_is_bare(target): return 'unresolved bare Docker repository' if platform in ('npm', 'pypi', 'package_git', 'postman', 'github_actions', 'gitlab_ci'): try: value = json.loads(target) except (TypeError, ValueError, json.JSONDecodeError): return f'malformed {platform} JSON target' if not isinstance(value, dict): return f'malformed {platform} JSON target' if _contains_nul(value): return 'NUL byte in decoded JSON row' if not normalize_target(target, platform): return 'target normalizes to an empty value' return '' def require_legacy_cutover_clear(db, config, max_entries=10000): """Refuse v2 activation while any legacy publication state is unresolved.""" conn = getattr(db, 'conn', None) if not conn: raise RuntimeError('database connection is unavailable') outbox_rows = 0 if conn.table_exists('scan_publication_outbox'): row = conn.execute( 'SELECT COUNT(*) AS count FROM scan_publication_outbox' ).fetchone() outbox_rows = int(row['count'] or 0) if outbox_rows: raise RuntimeError( f'final cutover refused: scan_publication_outbox has {outbox_rows} unresolved row(s)' ) legacy_raw_rows = 0 if conn.table_exists('target_scans'): row = conn.execute( '''SELECT COUNT(*) AS count FROM target_scans WHERE raw_result_storage != 'normalized_v2' AND raw_result_json IS NOT NULL''' ).fetchone() legacy_raw_rows = int(row['count'] or 0) if legacy_raw_rows: raise RuntimeError( f'final cutover refused: target_scans has {legacy_raw_rows} legacy raw result row(s); ' 'run --backfill-normalized-results in bounded offline batches' ) global_config = (config or {}).get('global') or {} spool_dir = ( reviewed_legacy_result_spool(config) if getattr(conn, 'is_postgres', False) else ( global_config.get('legacy_result_spool_dir') or global_config.get('result_spool_dir') or os.path.join(global_config.get('runtime_dir') or '', 'result_spool') ) ) inspected = 0 unresolved = [] spool_present = False if spool_dir: try: spool_details = os.stat(spool_dir, follow_symlinks=False) except FileNotFoundError: spool_details = None except OSError as exc: raise RuntimeError('final cutover refused: legacy spool state is unknown') from exc if spool_details is not None and not stat.S_ISDIR(spool_details.st_mode): raise RuntimeError('final cutover refused: legacy spool path is not a directory') spool_present = spool_details is not None if spool_present: for relative in ('', 'reservations', 'quarantine'): parent = os.path.join(spool_dir, relative) if relative else spool_dir try: parent_details = os.stat(parent, follow_symlinks=False) except FileNotFoundError: continue except OSError as exc: raise RuntimeError('final cutover refused: legacy spool shard state is unknown') from exc if not stat.S_ISDIR(parent_details.st_mode): raise RuntimeError('final cutover refused: legacy spool shard is not a directory') if os.name != 'nt': require_private_directory(parent, create=False) with os.scandir(parent) as entries: for entry in entries: inspected += 1 if inspected > max(1, int(max_entries)): raise RuntimeError('final cutover refused: legacy spool inspection exceeded its entry bound') if os.name != 'nt' and (entry.is_symlink() or is_reparse_point(entry.path)): raise RuntimeError('final cutover refused: legacy spool contains a link') if entry.is_file(follow_symlinks=False) and entry.name.lower().endswith('.json'): unresolved.append(os.path.join(relative, entry.name)) if len(unresolved) >= 10: break if len(unresolved) >= 10: break if unresolved: raise RuntimeError( 'final cutover refused: legacy result spool has unresolved objects: ' + ', '.join(unresolved) ) results_dir = global_config.get('results_dir') ledgers_checked = 0 if results_dir: for name in ('scan_results', 'found_secrets', 'scan_errors'): ledger_path = os.path.join(results_dir, f'{name}.publication-ledger.sqlite3') try: ledger_details = os.stat(ledger_path, follow_symlinks=False) except FileNotFoundError: continue except OSError as exc: raise RuntimeError(f'final cutover refused: {name} ledger state is unknown') from exc if not stat.S_ISREG(ledger_details.st_mode): raise RuntimeError(f'final cutover refused: {name} ledger is not a regular file') ledgers_checked += 1 uri = 'file:' + os.path.abspath(ledger_path).replace('\\', '/') + '?mode=ro' ledger = sqlite3.connect(uri, uri=True, timeout=5) try: table = ledger.execute( "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'publication_identity'" ).fetchone() if table: pending = ledger.execute( "SELECT 1 FROM publication_identity WHERE state = 'prepared' LIMIT 1" ).fetchone() if pending: raise RuntimeError( f'final cutover refused: {name} publication ledger has a prepared append' ) finally: ledger.close() if conn.is_postgres: conn.commit() return { 'legacy_outbox_rows': outbox_rows, 'legacy_raw_result_rows': legacy_raw_rows, 'legacy_spool_present': spool_present, 'legacy_spool_canonical': os.path.normcase(os.path.abspath(spool_dir)), 'legacy_spool_entries_inspected': inspected, 'legacy_projection_ledgers_checked': ledgers_checked, 'legacy_prepared_appends': 0, } def drain_legacy_scan_outbox(db, config, max_rows=1000): max_rows = min(10000, max(1, int(max_rows))) from console_runner import apply_global_config, drain_scan_publication_outbox from scanner import initialize_scanner_runtime global_config = (config or {}).get('global') or {} apply_global_config(global_config) initialize_scanner_runtime(preflight_complete=True, register_cleanup=False) drained = drain_scan_publication_outbox( db, max_rows, require_v2_schema=False, ) remaining = db.conn.execute( 'SELECT COUNT(*) AS count FROM scan_publication_outbox' ).fetchone() db.conn.commit() return {'drained': int(drained), 'remaining': int(remaining['count'] or 0)} def import_legacy_outbox_to_projection(db, max_rows=1000): conn = getattr(db, 'conn', None) if not conn or not conn.is_postgres: raise RuntimeError('legacy outbox projection transfer requires PostgreSQL') max_rows = min(10000, max(1, int(max_rows))) now = utc_now_iso() transferred = 0 try: rows = conn.execute( '''SELECT o.id, o.target_scan_id, ts.scan_event_id, ts.scan_event_hash, COALESCE(ts.findings_count, 0) AS findings_count, COALESCE(ts.error_count, 0) AS error_count, CASE WHEN ts.raw_result_json IS NULL THEN 0 ELSE octet_length(ts.raw_result_json) END AS raw_result_bytes FROM scan_publication_outbox o JOIN target_scans ts ON ts.id = o.target_scan_id WHERE NOT EXISTS ( SELECT 1 FROM pipeline_quarantine q WHERE q.subsystem = 'migration' AND q.object_type = 'legacy_outbox' AND q.object_id = o.id AND q.review_status = 'pending' ) ORDER BY o.id LIMIT ? FOR UPDATE OF o, ts''', (max_rows,), ).fetchall() for row in rows: raw_bytes = int(row['raw_result_bytes'] or 0) event_id = str(row['scan_event_id'] or f'legacy-target-scan-{row["target_scan_id"]}') event_hash = str(row['scan_event_hash'] or '') if not event_hash: if raw_bytes <= 0 or raw_bytes > 192 * 1024 * 1024: conn.execute( '''INSERT INTO pipeline_quarantine( subsystem, object_type, object_id, reason_code, reason_detail, byte_count, capacity_credit_applied, review_status, detected_at ) VALUES ('migration','legacy_outbox',?, 'legacy_payload_oversized',?,?,0,'pending',?) ON CONFLICT DO NOTHING''', ( row['id'], f'legacy outbox payload is {raw_bytes} bytes and has no event hash', max(0, raw_bytes), utc_now_iso(), ), ) continue payload = conn.execute( '''SELECT raw_result_json FROM target_scans WHERE id = ? AND raw_result_json IS NOT NULL AND octet_length(raw_result_json) = ?''', (row['target_scan_id'], raw_bytes), ).fetchone() if not payload: raise RuntimeError('legacy outbox payload changed during bounded identity recovery') event_hash = hashlib.sha256( str(payload['raw_result_json']).encode('utf-8') ).hexdigest() stream_mask = ( 1 | (2 if int(row['findings_count']) else 0) | (4 if int(row['error_count']) else 0) ) job = conn.execute( '''SELECT id, event_hash FROM projection_jobs WHERE job_kind = 'scan_event' AND event_id = ? FOR UPDATE''', (event_id,), ).fetchone() if job and str(job['event_hash']) != event_hash: raise RuntimeError(f'legacy outbox event hash conflicts with projection job: {event_id}') if not job: capacity_bytes = max(1, raw_bytes) * 2 conn.execute('SELECT id FROM pipeline_capacity WHERE id = 1 FOR UPDATE') conn.insert_returning_id( '''INSERT INTO projection_jobs( job_kind, event_id, event_hash, target_scan_id, status, required_stream_mask, capacity_items, capacity_bytes, created_at, updated_at ) VALUES ('scan_event', ?, ?, ?, 'pending', ?, 1, ?, ?, ?)''', ( event_id, event_hash, row['target_scan_id'], stream_mask, capacity_bytes, now, now, ), ) conn.execute( '''UPDATE pipeline_capacity SET projection_items = projection_items + 1, projection_bytes = projection_bytes + ?, updated_at = ? WHERE id = 1''', (capacity_bytes, now), ) deleted = conn.execute( 'DELETE FROM scan_publication_outbox WHERE id = ? AND target_scan_id = ?', (row['id'], row['target_scan_id']), ) if int(deleted.rowcount or 0) != 1: raise RuntimeError(f'legacy outbox row changed during projection transfer: {row["id"]}') transferred += 1 conn.commit() remaining = conn.execute( 'SELECT COUNT(*) AS count FROM scan_publication_outbox' ).fetchone() conn.commit() return {'transferred': transferred, 'remaining': int(remaining['count'] or 0)} except Exception: conn.rollback() raise def backfill_normalized_raw_results( db, max_rows=1000, max_bytes=192 * 1024 * 1024, max_seconds=30, max_object_bytes=192 * 1024 * 1024, ): conn = getattr(db, 'conn', None) if not conn or not conn.is_postgres: raise RuntimeError('normalized raw-result backfill requires PostgreSQL') max_rows = min(10000, max(1, int(max_rows))) max_bytes = max(1, int(max_bytes)) deadline = time.monotonic() + max(0.1, float(max_seconds)) processed = 0 consumed = 0 blocked_id = None reviewed = 0 while processed < max_rows and time.monotonic() < deadline: try: candidate = conn.execute( '''SELECT s.id, octet_length(s.raw_result_json) AS payload_bytes FROM target_scans s WHERE s.raw_result_storage != 'normalized_v2' AND s.raw_result_json IS NOT NULL AND NOT EXISTS ( SELECT 1 FROM pipeline_quarantine q WHERE q.subsystem = 'migration' AND q.object_type = 'legacy_target_scan' AND q.object_id = s.id AND q.review_status = 'pending' ) ORDER BY s.id LIMIT 1''' ).fetchone() if not candidate: conn.commit() break candidate_id = int(candidate['id']) payload_size = int(candidate['payload_bytes'] or 0) if payload_size <= 0 or payload_size > max(1, int(max_object_bytes)): conn.execute( '''INSERT INTO pipeline_quarantine( subsystem, object_type, object_id, reason_code, reason_detail, byte_count, capacity_credit_applied, review_status, detected_at ) VALUES ('migration','legacy_target_scan',?, 'legacy_payload_oversized',?,?,0,'pending',?) ON CONFLICT DO NOTHING''', ( candidate_id, f'legacy raw result is {payload_size} bytes; reviewed bounded import required', max(0, payload_size), utc_now_iso(), ), ) conn.commit() reviewed += 1 continue if consumed + payload_size > max_bytes: blocked_id = candidate_id conn.commit() break row = conn.execute( '''SELECT id, raw_result_json FROM target_scans WHERE id = ? AND raw_result_storage != 'normalized_v2' AND raw_result_json IS NOT NULL AND octet_length(raw_result_json) = ? FOR UPDATE''', (candidate_id, payload_size), ).fetchone() if not row: conn.rollback() continue prepared = [] payload = str(row['raw_result_json'] or '') if len(payload.encode('utf-8')) != payload_size: raise RuntimeError(f'legacy target scan {candidate_id} changed during bounded fetch') try: result = json.loads(payload) except (TypeError, ValueError, json.JSONDecodeError) as exc: raise RuntimeError(f'legacy target scan {candidate_id} contains invalid JSON') from exc if not isinstance(result, dict): raise RuntimeError(f'legacy target scan {candidate_id} is not a JSON object') findings = result.get('findings') or [] errors = result.get('errors') or [] if not isinstance(findings, list) or not isinstance(errors, list): raise RuntimeError(f'legacy target scan {candidate_id} has invalid finding/error arrays') metadata = dict(result) metadata.pop('findings', None) metadata.pop('errors', None) encoded_metadata = json.dumps( metadata, ensure_ascii=True, sort_keys=True, separators=(',', ':'), default=str, ).encode('utf-8') if len(encoded_metadata) > 16 * 1024 * 1024: raise RuntimeError( f'legacy target scan {candidate_id} metadata exceeds its normalized bound' ) prepared.append({ 'id': candidate_id, 'payload': payload, 'payload_bytes': payload_size, 'findings': findings, 'errors': errors, 'metadata': encoded_metadata, 'metadata_sha256': hashlib.sha256(encoded_metadata).hexdigest(), }) ids = [item['id'] for item in prepared] placeholders = ','.join('?' for _ in ids) finding_rows = conn.execute( f'''SELECT id, target_scan_id, finding_uid FROM findings WHERE target_scan_id IN ({placeholders}) ORDER BY target_scan_id, id''', ids, ).fetchall() findings_by_scan = {item['id']: [] for item in prepared} for finding_row in finding_rows: findings_by_scan[int(finding_row['target_scan_id'])].append(finding_row) error_rows = conn.execute( f'''SELECT target_scan_id, COUNT(*) AS count FROM errors WHERE target_scan_id IN ({placeholders}) GROUP BY target_scan_id''', ids, ).fetchall() errors_by_scan = {int(row['target_scan_id']): int(row['count']) for row in error_rows} existing_metadata_rows = conn.execute( f'''SELECT target_scan_id, metadata_sha256, metadata_bytes FROM scan_result_compat WHERE target_scan_id IN ({placeholders})''', ids, ).fetchall() existing_metadata = {int(row['target_scan_id']): row for row in existing_metadata_rows} now = utc_now_iso() metadata_inserts = [] finding_inserts = [] finding_updates = [] for item in prepared: finding_group = findings_by_scan[item['id']] if ( len(finding_group) != len(item['findings']) or errors_by_scan.get(item['id'], 0) != len(item['errors']) ): raise RuntimeError( f'legacy target scan {item["id"]} normalized child counts do not match raw history' ) existing = existing_metadata.get(item['id']) if existing and ( existing['metadata_sha256'] != item['metadata_sha256'] or int(existing['metadata_bytes']) != len(item['metadata']) ): raise RuntimeError( f'legacy target scan {item["id"]} has conflicting compatibility metadata' ) if not existing: metadata_inserts.append(( item['id'], item['metadata'].decode('utf-8'), item['metadata_sha256'], len(item['metadata']), now, )) for finding_row, finding in zip(finding_group, item['findings']): if not isinstance(finding, dict): raise RuntimeError( f'legacy target scan {item["id"]} contains a non-object finding' ) finding_uid = str(finding.get('finding_uid') or '') if finding_uid and finding_uid != str(finding_row['finding_uid'] or ''): raise RuntimeError( f'legacy target scan {item["id"]} finding order/identity changed' ) compat = db._compat_payload(finding) finding_inserts.append(( finding_row['id'], compat['raw_value'], compat['raw_v2_value'], compat['structured_data_json'], compat['extra_data_json'], compat['analysis_info_json'], compat['extension_json'], compat['payload_sha256'], compat['payload_bytes'], compat['payload_omitted'], now, )) finding_updates.append((finding_row['id'], compat)) if metadata_inserts: values = ','.join('(?, 2, ?, ?, ?, \'bounded\', ?)' for _ in metadata_inserts) conn.execute( '''INSERT INTO scan_result_compat( target_scan_id, schema_version, metadata_json, metadata_sha256, metadata_bytes, reconstruction_status, created_at ) VALUES ''' + values, tuple(value for row in metadata_inserts for value in row), ) if finding_inserts: finding_ids = [int(row[0]) for row in finding_inserts] existing_payloads = {} for start in range(0, len(finding_ids), 10000): id_batch = finding_ids[start:start + 10000] finding_placeholders = ','.join('?' for _ in id_batch) existing_payload_rows = conn.execute( f'''SELECT finding_id, payload_sha256, payload_bytes, payload_omitted FROM finding_compat_payloads WHERE finding_id IN ({finding_placeholders})''', id_batch, ).fetchall() existing_payloads.update({ int(row['finding_id']): row for row in existing_payload_rows }) missing_payloads = [] for values_row in finding_inserts: existing = existing_payloads.get(int(values_row[0])) if existing and ( existing['payload_sha256'] != values_row[7] or int(existing['payload_bytes']) != int(values_row[8]) or bool(existing['payload_omitted']) != bool(values_row[9]) ): raise RuntimeError('legacy finding has conflicting compatibility data') if not existing: missing_payloads.append(values_row) if missing_payloads: for start in range(0, len(missing_payloads), 5000): payload_batch = missing_payloads[start:start + 5000] values = ','.join( '(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)' for _ in payload_batch ) conn.execute( '''INSERT INTO finding_compat_payloads( finding_id, raw_value, raw_v2_value, structured_data_json, extra_data_json, analysis_info_json, extension_json, payload_sha256, payload_bytes, payload_omitted, created_at ) VALUES ''' + values, tuple(value for row in payload_batch for value in row), ) for start in range(0, len(finding_updates), 10000): update_batch = finding_updates[start:start + 10000] values = ','.join( '(?::bigint, ?::text, ?::bigint, ?::integer)' for _ in update_batch ) conn.execute( '''UPDATE findings AS f SET raw_finding_json = NULL, raw_payload_sha256 = v.payload_sha256, raw_payload_bytes = v.payload_bytes, raw_payload_omitted = v.payload_omitted FROM (VALUES ''' + values + ''') AS v( id, payload_sha256, payload_bytes, payload_omitted ) WHERE f.id = v.id''', tuple( value for finding_id, compat in update_batch for value in ( finding_id, compat['payload_sha256'], compat['payload_bytes'], compat['payload_omitted'], ) ), ) updated = conn.execute( '''UPDATE target_scans SET raw_result_json = NULL, compat_schema_version = 2, raw_result_storage = 'normalized_v2' WHERE id IN (''' + placeholders + ''') AND raw_result_storage != 'normalized_v2' AND raw_result_json IS NOT NULL''', ids, ) if int(updated.rowcount or 0) != len(prepared): raise RuntimeError('legacy target scan normalization lost its row fence') conn.commit() except Exception: conn.rollback() raise processed += len(prepared) consumed += sum(item['payload_bytes'] for item in prepared) if blocked_id is not None: break remaining = conn.execute( '''SELECT COUNT(*) AS count FROM target_scans WHERE raw_result_storage != 'normalized_v2' AND raw_result_json IS NOT NULL''' ).fetchone() conn.commit() return { 'processed': processed, 'bytes': consumed, 'remaining': int(remaining['count'] or 0), 'blocked_target_scan_id': blocked_id, 'reviewed_oversized': reviewed, } def import_legacy_result_spool(db, config, max_rows=1000): conn = getattr(db, 'conn', None) if not conn or not conn.is_postgres: raise RuntimeError('legacy result-spool import requires PostgreSQL') global_config = (config or {}).get('global') or {} spool_dir = reviewed_legacy_result_spool(config) spool = ResultSpool( spool_dir, max_event_bytes=int(global_config.get( 'legacy_result_spool_max_event_bytes', global_config.get('result_spool_max_event_bytes', 192 * 1024 * 1024), )), max_events=int(global_config.get( 'legacy_result_spool_max_events', global_config.get('result_spool_max_events', 10000), )), max_total_bytes=int(global_config.get( 'legacy_result_spool_max_total_bytes', global_config.get('result_spool_max_total_bytes', 3 * 1024 * 1024 * 1024), )), min_free_bytes=0, ) max_rows = min(10000, max(1, int(max_rows))) if not db.try_acquire_result_spool_publisher(): raise RuntimeError('legacy result-spool publisher advisory lock is unavailable') imported = 0 try: while imported < max_rows: record = spool.next_pending_event() if record is None: break try: outcome = db.ingest_scan_event(record.envelope) except ScanEventConflictError as exc: spool.quarantine_event( record.event_id, f'database scan-event hash conflict during offline import: {exc}', ) raise RuntimeError( f'legacy scan event {record.event_id} was quarantined after a hash conflict' ) from exc if ( not outcome or not outcome.get('ingested') or str(outcome.get('scan_event_id') or '') != record.event_id or str(outcome.get('scan_event_hash') or '') != record.event_hash ): raise RuntimeError(f'legacy scan event {record.event_id} was not confirmed by PostgreSQL') if spool.acknowledge(record.event_id, record.event_hash) is not True: raise RuntimeError(f'legacy scan event {record.event_id} acknowledgement was not confirmed') imported += 1 remaining = spool.next_pending_event() is not None return {'imported': imported, 'remaining': int(remaining)} finally: if db.release_result_spool_publisher() is not True: raise RuntimeError('legacy result-spool publisher advisory lock release was not confirmed') def review_pipeline_quarantine_manifest(db, manifest_path, max_rows=1000, config=None): require_private_file(manifest_path) manifest = read_private_json(manifest_path, max_bytes=1024 * 1024) if not isinstance(manifest, dict) or manifest.get('type') != 'truf-pipeline-quarantine-review-v1': raise RuntimeError('pipeline quarantine review manifest type is invalid') entries = manifest.get('entries') if not isinstance(entries, list) or len(entries) > min(10000, max(1, int(max_rows))): raise RuntimeError('pipeline quarantine review manifest exceeds its row bound') manifest_bytes = json.dumps( manifest, ensure_ascii=True, sort_keys=True, separators=(',', ':'), default=str, ).encode('utf-8') manifest_hash = hashlib.sha256(manifest_bytes).hexdigest() reviewed = 0 duplicates = 0 for entry in entries: if not isinstance(entry, dict): raise RuntimeError('pipeline quarantine review entry is not an object') required = ('id', 'reason_code', 'payload_sha256', 'action') if any(name not in entry for name in required): raise RuntimeError('pipeline quarantine review entry is incomplete') canonical = json.dumps( entry, ensure_ascii=True, sort_keys=True, separators=(',', ':'), default=str, ).encode('utf-8') audit_hash = hashlib.sha256( b'truf-pipeline-quarantine-review-v1|' + manifest_hash.encode('ascii') + b'|' + canonical ).hexdigest() row = db.pipeline_quarantine_for_review(entry['id']) if not row: raise RuntimeError('pipeline quarantine review row is absent') global_config = (config or {}).get('global') or {} if row['object_type'] == 'result_bundle': bundle_root = global_config.get('result_bundle_dir') if not bundle_root: raise RuntimeError('bundle quarantine review root is unavailable') reservation = db.conn.execute( '''SELECT id, state, ready_relative_path, docker_layer_plan_json FROM result_reservations WHERE id = ?''', (row['reservation_id'],), ).fetchone() db.conn.commit() if not reservation: raise RuntimeError('bundle quarantine reservation is absent') if reservation['state'] != 'quarantined': ready = inspect_private_relative_path( bundle_root, reservation['ready_relative_path'], ) quarantined = inspect_private_relative_path( bundle_root, row['source_relative_path'], ) if ( ready.state == PrivatePathState.UNKNOWN or quarantined.state == PrivatePathState.UNKNOWN ): raise RuntimeError('prepared bundle quarantine physical state is unknown') if ( ready.state == PrivatePathState.PRESENT and quarantined.state == PrivatePathState.ABSENT ): durable_publish(ready.path, quarantined.path) quarantined = inspect_private_relative_path( bundle_root, row['source_relative_path'], ) ready = inspect_private_relative_path( bundle_root, reservation['ready_relative_path'], ) if not ( ready.state == PrivatePathState.ABSENT and quarantined.state == PrivatePathState.PRESENT ): raise RuntimeError( 'prepared bundle quarantine cannot be finalized from conflicting physical state' ) db.quarantine_result_bundle( reservation['id'], row['reason_code'], row['reason_detail'] or '', row['source_relative_path'], payload_sha256=row['payload_sha256'] or '', byte_count=row['byte_count'], physical_confirmed=True, ) row = db.pipeline_quarantine_for_review(entry['id']) if str(entry['action']).lower() == 'rescan' and ( row['object_type'] != 'result_bundle' or reservation['docker_layer_plan_json'] is None ): raise RuntimeError('quarantine rescan requires a Docker layer result bundle') if str(entry['action']).lower() == 'retry' and row['object_type'] == 'result_bundle': global_config = (config or {}).get('global') or {} bundle_root = global_config.get('result_bundle_dir') if not bundle_root: raise RuntimeError('bundle quarantine retry root is unavailable') reservation = db.conn.execute( '''SELECT id, bundle_id, scan_event_id, ready_relative_path, docker_layer_plan_json FROM result_reservations WHERE id = ?''', (row['reservation_id'],), ).fetchone() db.conn.commit() if not reservation: raise RuntimeError('bundle quarantine retry reservation is absent') if reservation['docker_layer_plan_json'] is not None: raise RuntimeError( 'Docker layer bundle quarantine requires a fresh parent claim' ) quarantined = inspect_private_relative_path(bundle_root, row['source_relative_path']) ready = inspect_private_relative_path(bundle_root, reservation['ready_relative_path']) if quarantined.state == PrivatePathState.UNKNOWN or ready.state == PrivatePathState.UNKNOWN: raise RuntimeError('bundle quarantine retry physical state is unknown') if quarantined.state == PrivatePathState.PRESENT and ready.state == PrivatePathState.ABSENT: candidate_path = quarantined.path elif quarantined.state == PrivatePathState.ABSENT and ready.state == PrivatePathState.PRESENT: candidate_path = ready.path else: raise RuntimeError('bundle quarantine retry requires exactly one physical artifact') if not private_file_ready(candidate_path): raise RuntimeError('bundle quarantine retry artifact is not an exact private file') if os.path.getsize(candidate_path) != int(row['byte_count']): raise RuntimeError('bundle quarantine retry byte count conflicts with reviewed evidence') if sha256_file(candidate_path) != str(row['payload_sha256'] or ''): raise RuntimeError('bundle quarantine retry payload hash conflicts with reviewed evidence') validated = ResultBundleReader( candidate_path, max_event_bytes=int(global_config.get('result_bundle_max_event_bytes', 192 * 1024 * 1024)), ).validate() if ( int(validated.reservation_id) != int(reservation['id']) or str(validated.bundle_id) != str(reservation['bundle_id']) or str(validated.scan_event_id) != str(reservation['scan_event_id']) ): raise RuntimeError('bundle quarantine retry identity conflicts with its reservation') if quarantined.state == PrivatePathState.PRESENT: ensure_private_directory(os.path.dirname(ready.path), reject_reparse=True) durable_publish(quarantined.path, ready.path) quarantined = inspect_private_relative_path(bundle_root, row['source_relative_path']) ready = inspect_private_relative_path(bundle_root, reservation['ready_relative_path']) if not ( quarantined.state == PrivatePathState.ABSENT and ready.state == PrivatePathState.PRESENT and private_file_ready(ready.path) ): raise RuntimeError('bundle quarantine retry publication was not confirmed') if str(entry['action']).lower() in ('discard', 'rescan') and row.get('source_relative_path'): root = ( global_config.get('result_bundle_dir') if row['subsystem'] == 'result_ingester' else global_config.get('results_dir') ) if not root: raise RuntimeError('physical quarantine review root is unavailable') inspection = inspect_private_relative_path(root, row['source_relative_path']) if inspection.state == PrivatePathState.UNKNOWN: raise RuntimeError('physical quarantine artifact state is unknown') if inspection.state == PrivatePathState.PRESENT: durable_unlink(inspection.path) inspection = inspect_private_relative_path(root, row['source_relative_path']) if inspection.state != PrivatePathState.ABSENT: raise RuntimeError('physical quarantine artifact unlink was not confirmed') if row['object_type'] == 'projection_tail': db.mark_pipeline_artifact_deleted(row['object_id']) elif row['object_type'] == 'result_bundle': artifact = db.conn.execute( '''SELECT id FROM pipeline_artifacts WHERE subsystem = 'result_bundle' AND artifact_kind = 'bundle_quarantine' AND owner_id = ?''', (row['reservation_id'],), ).fetchone() db.conn.commit() if not artifact: raise RuntimeError('bundle quarantine artifact index is absent') db.mark_pipeline_artifact_deleted(artifact['id']) outcome = db.review_pipeline_quarantine( entry['id'], entry['reason_code'], entry['payload_sha256'], entry['action'], audit_hash, ) reviewed += int(bool(outcome and outcome.get('reviewed'))) duplicates += int(bool(outcome and outcome.get('duplicate'))) return { 'manifest_sha256': manifest_hash, 'reviewed': reviewed, 'duplicates': duplicates, } def rebuild_jsonl_projections( db, output_root, max_rows=1000, max_bytes=192 * 1024 * 1024, max_seconds=30, ): from jsonl_projector import JsonlProjector if not db.conn or not db.conn.is_postgres: raise RuntimeError('full JSONL rebuild requires PostgreSQL') db.require_runtime_safety_schema() db.require_final_cutover() output_root = os.path.abspath(output_root) existed = os.path.lexists(output_root) output_root = ensure_private_directory(output_root, reject_reparse=True) state_path = os.path.join(output_root, '.rebuild-state.json') identity = db.conn.execute( '''SELECT current_database() AS database_name, current_user AS database_user, inet_server_port() AS database_port''' ).fetchone() db.conn.commit() cutover = db.final_cutover_status() authority = hashlib.sha256(json.dumps({ 'database': dict(identity), 'cutover_evidence_sha256': cutover['evidence_sha256'], }, ensure_ascii=True, sort_keys=True, separators=(',', ':')).encode('utf-8')).hexdigest() if not os.path.lexists(state_path): if existed: with os.scandir(output_root) as entries: if next(entries, None) is not None: raise RuntimeError('new JSONL rebuild output directory must be empty') state = { 'schema': 1, 'authority_sha256': authority, 'phase': 'scans', 'last_scan_id': 0, 'last_keycheck_result_id': 0, 'offsets': {}, 'prepared': None, 'completed': False, } atomic_write_private_json(state_path, state) state = read_private_json(state_path, max_bytes=1024 * 1024) if ( not isinstance(state, dict) or state.get('schema') != 1 or state.get('authority_sha256') != authority ): raise RuntimeError('JSONL rebuild state is invalid') results_dir = ensure_private_directory( os.path.join(output_root, 'results'), reject_reparse=True, ) keycheck_dir = ensure_private_directory( os.path.join(output_root, 'keychecks'), reject_reparse=True, ) projector = JsonlProjector( db, results_dir, 'offline-rebuild', keycheck_dir=keycheck_dir, artifact_tracking=False, ) def output_path(relative): path = os.path.abspath(os.path.join(output_root, str(relative).replace('/', os.sep))) if path == output_root or os.path.commonpath((output_root, path)) != output_root: raise RuntimeError('JSONL rebuild path escapes its output root') ensure_private_directory(os.path.dirname(path), reject_reparse=True) if os.path.lexists(path): require_private_file(path) else: descriptor = os.open( path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) os.close(descriptor) harden_private_file(path) fsync_directory(os.path.dirname(path)) return path prepared = state.get('prepared') if prepared: if not isinstance(prepared, dict) or not isinstance(prepared.get('offsets'), dict): raise RuntimeError('JSONL rebuild prepared state is invalid') for relative, offset in prepared['offsets'].items(): path = output_path(relative) size = os.path.getsize(path) offset = int(offset) if size < offset: raise RuntimeError('JSONL rebuild output is shorter than its committed offset') if size != offset: with open(path, 'r+b', buffering=0) as handle: handle.truncate(offset) handle.flush() os.fsync(handle.fileno()) state['prepared'] = None atomic_write_private_json(state_path, state) max_rows = min(100000, max(1, int(max_rows))) max_bytes = max(1, int(max_bytes)) deadline = time.monotonic() + max(0.1, float(max_seconds)) rows_written = 0 bytes_written = 0 blocked = None def stream_relative(stream_name, service=''): if stream_name == 'scan_results': return 'results/scan_results.jsonl' if stream_name == 'found_secrets': return 'results/found_secrets.jsonl' if stream_name == 'scan_errors': return 'results/scan_errors.log' if stream_name.endswith(':results'): return f'keychecks/{service}/{service}Results.jsonl' return f'keychecks/{service}/{service}Checked.txt' def append_serialized(job, streams, service=''): nonlocal bytes_written, blocked total = sum(int(item.byte_length) for item in streams) if bytes_written + total > max_bytes: blocked = {'kind': job['job_kind'], 'id': int(job['id']), 'bytes': total} return False offsets = {} destinations = [] for item in streams: relative = stream_relative(item.stream_name, service) path = output_path(relative) expected = int(state['offsets'].get(relative, 0)) if os.path.getsize(path) != expected: raise RuntimeError('JSONL rebuild output/cursor mismatch') offsets[relative] = expected destinations.append((item, relative, path)) state['prepared'] = { 'kind': job['job_kind'], 'id': int(job['id']), 'offsets': offsets, } atomic_write_private_json(state_path, state) for item, relative, path in destinations: with open(path, 'ab', buffering=0) as target, open(item.path, 'rb', buffering=0) as source: while True: block = source.read(1024 * 1024) if not block: break target.write(block) target.flush() os.fsync(target.fileno()) state['offsets'][relative] = offsets[relative] + int(item.byte_length) bytes_written += total state['prepared'] = None return True try: while rows_written < max_rows and time.monotonic() < deadline and not state['completed']: streams = [] if state['phase'] == 'scans': row = db.conn.execute( '''SELECT id, scan_event_id, scan_event_hash FROM target_scans WHERE id > ? ORDER BY id LIMIT 1''', (int(state['last_scan_id']),), ).fetchone() db.conn.commit() if not row: state['phase'] = 'keychecks' atomic_write_private_json(state_path, state) continue event_id = str(row['scan_event_id'] or f'legacy-target-scan-{row["id"]}') event_hash = str(row['scan_event_hash'] or hashlib.sha256( f'legacy-target-scan|{row["id"]}'.encode('utf-8') ).hexdigest()) job = { 'id': int(row['id']), 'job_kind': 'scan_event', 'target_scan_id': int(row['id']), 'event_id': event_id, 'event_hash': event_hash, 'required_stream_mask': 7, 'capacity_bytes': 192 * 1024 * 1024, } streams = projector._serialize(job) if not append_serialized(job, streams): break state['last_scan_id'] = int(row['id']) else: row = db.conn.execute( '''SELECT id, event_id, service FROM keycheck_results WHERE id > ? ORDER BY id LIMIT 1''', (int(state['last_keycheck_result_id']),), ).fetchone() db.conn.commit() if not row: state['completed'] = True atomic_write_private_json(state_path, state) break event_id = str(row['event_id'] or f'legacy-keycheck-result-{row["id"]}') event_hash = hashlib.sha256( f'keycheck-rebuild|{row["id"]}|{event_id}'.encode('utf-8') ).hexdigest() job = { 'id': int(row['id']), 'job_kind': 'keycheck_event', 'keycheck_result_id': int(row['id']), 'event_id': event_id, 'event_hash': event_hash, 'required_stream_mask': 8, 'capacity_bytes': 192 * 1024 * 1024, } streams = projector._serialize(job) if not append_serialized(job, streams, str(row['service'])): break state['last_keycheck_result_id'] = int(row['id']) rows_written += 1 atomic_write_private_json(state_path, state) for item in streams: try: os.remove(item.path) except OSError: pass finally: for item in locals().get('streams', []): try: os.remove(item.path) except OSError: pass return { 'rows_written': rows_written, 'bytes_written': bytes_written, 'phase': state['phase'], 'completed': bool(state['completed']), 'blocked': blocked, 'output_root': output_root, } def initialize_projection_cursors_from_existing_files(db, config): conn = getattr(db, 'conn', None) if not conn: raise RuntimeError('projection cursor initialization requires a database') global_config = ((config or {}).get('global') or {}) results_value = global_config.get('results_dir') if not results_value and global_config.get('runtime_dir'): results_value = os.path.join(global_config['runtime_dir'], 'results') if not results_value: raise RuntimeError('projection cursor initialization requires results_dir') results_dir = os.path.abspath(results_value) initialized = {} try: for stream_name in ('scan_results', 'found_secrets', 'scan_errors'): suffix = ' FOR UPDATE' if conn.is_postgres else '' row = conn.execute( f'''SELECT s.base_relative_path, c.committed_offset, c.last_append_id, c.last_job_id FROM projection_streams s JOIN projection_cursors c ON c.stream_name = s.stream_name WHERE s.stream_name = ?{suffix}''', (stream_name,), ).fetchone() if not row: raise RuntimeError(f'projection cursor is absent: {stream_name}') path = os.path.abspath(os.path.join(results_dir, str(row['base_relative_path']))) if os.path.commonpath((results_dir, path)) != results_dir or path == results_dir: raise RuntimeError(f'projection stream path escapes results_dir: {stream_name}') size = 0 try: details = os.stat(path, follow_symlinks=False) except FileNotFoundError: details = None except OSError as exc: raise RuntimeError( f'projection cursor file state is unknown: {stream_name}' ) from exc if details is not None: if not stat.S_ISREG(details.st_mode): raise RuntimeError(f'projection stream is not a regular file: {stream_name}') require_private_file(path) size = int(details.st_size) offset = int(row['committed_offset'] or 0) if offset == size: initialized[stream_name] = size continue append_count = conn.execute( 'SELECT COUNT(*) AS count FROM projection_appends WHERE stream_name = ?', (stream_name,), ).fetchone() if ( offset != 0 or row['last_append_id'] is not None or row['last_job_id'] is not None or int(append_count['count'] or 0) != 0 ): raise RuntimeError( f'projection cursor/file mismatch requires reviewed recovery: {stream_name}' ) conn.execute( '''UPDATE projection_cursors SET committed_offset = ?, updated_at = ? WHERE stream_name = ? AND committed_offset = 0 AND last_append_id IS NULL AND last_job_id IS NULL''', (size, utc_now_iso(), stream_name), ) initialized[stream_name] = size conn.commit() return initialized except Exception: conn.rollback() raise def reconcile_todo_file( db, todo_path, source, platform, query='reconciled', max_rows=1000, max_bytes=4 * 1024 * 1024, max_seconds=5.0, ): if not db or not getattr(db, 'conn', None): raise RuntimeError('database connection is unavailable') db.require_runtime_safety_schema() todo_path = os.path.normcase(os.path.abspath(todo_path)) reject_reparse_components(todo_path) if 'checked' in os.path.basename(todo_path).lower(): raise ValueError('checked files are not accepted as reconciliation input') if is_reparse_point(todo_path) or not os.path.isfile(todo_path): raise FileNotFoundError(todo_path) max_rows = max(1, int(max_rows)) max_bytes = max(1, int(max_bytes)) max_seconds = max(0.001, float(max_seconds)) conn = db.conn last_change = None for mutation_attempt in range(RECONCILIATION_MUTATION_RETRIES): started = time.monotonic() descriptor = None try: flags = os.O_RDONLY | getattr(os, 'O_BINARY', 0) | getattr(os, 'O_NOFOLLOW', 0) descriptor = os.open(todo_path, flags) initial_stat = os.fstat(descriptor) if not stat.S_ISREG(initial_stat.st_mode): raise ValueError('reconciliation input must be a regular file') snapshot = _stat_identity(initial_stat) if not _same_file_snapshot(todo_path, (descriptor, snapshot)): raise ReconciliationFileChanged('reconciliation input changed while opening') identity = todo_file_identity(todo_path, initial_stat) file_size = snapshot[2] file_mtime_ns = snapshot[3] if conn.is_sqlite: conn.execute('BEGIN IMMEDIATE') cursor = conn.execute( '''SELECT file_identity, source, platform, byte_offset, line_number, discarding_oversized, oversized_line_start, cumulative_rows, cumulative_bytes, cumulative_inserted, cumulative_rejected, completed_at FROM target_queue_reconciliation_cursors WHERE source_file = ?''', (todo_path,), ).fetchone() if cursor and (str(cursor['source']) != str(source) or str(cursor['platform']) != str(platform)): raise ValueError( f'reconciliation cursor is already bound to {cursor["source"]}/{cursor["platform"]}; ' f'it cannot be reused for {source}/{platform}' ) same_identity = bool(cursor and cursor['file_identity'] == identity) offset = int(cursor['byte_offset'] or 0) if same_identity else 0 line_number = int(cursor['line_number'] or 0) if same_identity else 0 discarding = bool(cursor['discarding_oversized']) if same_identity else False oversized_line_start = int(cursor['oversized_line_start'] or 0) if same_identity and cursor['oversized_line_start'] is not None else None cumulative_rows = int(cursor['cumulative_rows'] or 0) if same_identity else 0 cumulative_bytes = int(cursor['cumulative_bytes'] or 0) if same_identity else 0 cumulative_inserted = int(cursor['cumulative_inserted'] or 0) if same_identity else 0 cumulative_rejected = int(cursor['cumulative_rejected'] or 0) if same_identity else 0 previous_completed_at = cursor['completed_at'] if same_identity else None if offset < 0 or offset > file_size: offset = line_number = 0 discarding = False oversized_line_start = None cumulative_rows = cumulative_bytes = cumulative_inserted = cumulative_rejected = 0 rows_read = 0 bytes_read = 0 rejected = 0 eligible = [] resolver_rows = [] seen = set() malformed = [] unresolved_docker = [] bounded_row = None issues = [] next_offset = offset next_line = line_number with os.fdopen(descriptor, 'rb', closefd=False) as handle: handle.seek(offset) while time.monotonic() - started < max_seconds: if not discarding and rows_read >= max_rows: break if bytes_read >= max_bytes and (rows_read or discarding): break row_start = handle.tell() if discarding: read_limit = min( RECONCILIATION_ABSOLUTE_ROW_BYTES + 1, max(1, max_bytes - bytes_read), ) else: read_limit = RECONCILIATION_ABSOLUTE_ROW_BYTES + 1 raw = handle.readline(read_limit) if not raw: break candidate_offset = handle.tell() reached_eof = candidate_offset >= file_size if discarding: bytes_read += len(raw) next_offset = candidate_offset if raw.endswith(b'\n') or reached_eof: discarding = False oversized_line_start = None next_line += 1 continue oversized = len(raw) > RECONCILIATION_ABSOLUTE_ROW_BYTES and not raw.endswith(b'\n') if oversized: if rows_read and bytes_read + len(raw) > max_bytes: handle.seek(row_start) break bytes_read += len(raw) next_offset = candidate_offset issue_line = next_line + 1 rows_read += 1 rejected += 1 oversized_line_start = row_start discarding = not reached_eof if not discarding: next_line += 1 oversized_line_start = None reason = f'row exceeds absolute {RECONCILIATION_ABSOLUTE_ROW_BYTES}-byte reconciliation row limit' bounded_row = {'line': issue_line, 'offset': row_start, 'reason': reason} issue = { 'line': issue_line, 'offset': row_start, 'reason': reason, 'target': _safe_target_preview(raw), } malformed.append(issue) issues.append(issue) continue if rows_read and bytes_read + len(raw) > max_bytes: handle.seek(row_start) break decode_error = None target = '' try: target = raw.decode('utf-8').strip().lstrip('\ufeff') except UnicodeDecodeError as exc: decode_error = exc target_error = reconciliation_target_error(target, platform) if target and not decode_error else '' bytes_read += len(raw) next_offset = candidate_offset rows_read += 1 next_line += 1 if decode_error is not None: rejected += 1 issue = { 'line': next_line, 'offset': row_start, 'reason': f'invalid UTF-8 at byte {decode_error.start}', 'target': _safe_target_preview(raw), } malformed.append(issue) issues.append(issue) continue if not target: continue error = target_error if error: rejected += 1 item = { 'line': next_line, 'offset': row_start, 'target': _safe_target_preview(target), 'reason': error, } issues.append(item) if error == 'unresolved bare Docker repository': unresolved_docker.append(item) normalized = normalize_target(target, platform) if normalized and normalized not in seen: seen.add(normalized) resolver_rows.append((target, normalized)) else: malformed.append(item) continue normalized = normalize_target(target, platform) if normalized in seen: continue seen.add(normalized) eligible.append((target, normalized)) now = utc_now_iso() for issue in issues: conn.execute( '''INSERT INTO target_queue_reconciliation_issues ( source_file, file_identity, source, platform, line_number, byte_offset, reason, target_preview, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(source_file, file_identity, line_number, byte_offset, reason) DO NOTHING''', ( todo_path, identity, source, platform, issue['line'], issue['offset'], issue['reason'], issue.get('target'), now, ), ) inserted = 0 for target, normalized in eligible: cur = conn.execute( '''INSERT INTO target_queue ( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES (?, ?, ?, ?, ?, 'pending', ?, ?) ON CONFLICT(source, normalized_target) DO NOTHING''', (source, platform, query, target, normalized, now, now), ) inserted += max(0, int(getattr(cur, 'rowcount', 0) or 0)) resolver_due = datetime.fromtimestamp(time.time() + 3600, timezone.utc).isoformat(timespec='seconds') resolver_inserted = 0 for target, normalized in resolver_rows: cur = conn.execute( '''INSERT INTO target_queue ( source, platform, query, target, normalized_target, status, available_after, resolver_state, resolver_due_at, last_error, created_at, updated_at ) VALUES (?, ?, ?, ?, ?, 'deferred', ?, 'pending', ?, 'Docker tag resolution unresolved', ?, ?) ON CONFLICT(source, normalized_target) DO NOTHING''', (source, platform, query, target, normalized, resolver_due, resolver_due, now, now), ) resolver_inserted += max(0, int(getattr(cur, 'rowcount', 0) or 0)) inserted += resolver_inserted if not _same_file_snapshot(todo_path, (descriptor, snapshot)): raise ReconciliationFileChanged('reconciliation input changed before cursor update') at_eof = next_offset >= file_size and not discarding completed_at = (previous_completed_at or now) if at_eof else None cumulative_rows += rows_read cumulative_bytes += bytes_read cumulative_inserted += inserted cumulative_rejected += rejected unresolved_row = conn.execute( '''SELECT COUNT(*) AS count FROM target_queue_reconciliation_issues WHERE source_file = ? AND source = ? AND platform = ? AND resolved_at IS NULL''', (todo_path, source, platform), ).fetchone() unresolved_count = int(unresolved_row['count'] or 0) report = { 'source_file': todo_path, 'file_identity': identity, 'file_size': file_size, 'file_mtime_ns': file_mtime_ns, 'source': source, 'platform': platform, 'start_offset': offset, 'byte_offset': next_offset, 'line_number': next_line, 'rows_read': rows_read, 'bytes_read': bytes_read, 'eligible_rows': len(eligible), 'inserted_rows': inserted, 'resolver_rows_inserted': resolver_inserted, 'rejected_rows': rejected, 'malformed_rows': malformed, 'unresolved_docker_rows': unresolved_docker, 'unresolved_issues': unresolved_count, 'bounded_row': bounded_row, 'discarding_oversized': discarding, 'at_eof': at_eof, 'completed_at': completed_at, 'cumulative_rows': cumulative_rows, 'cumulative_bytes': cumulative_bytes, 'cumulative_inserted': cumulative_inserted, 'cumulative_rejected': cumulative_rejected, 'elapsed_sec': max(0.0, time.monotonic() - started), } conn.execute( '''INSERT INTO target_queue_reconciliation_cursors ( source_file, file_identity, file_size, file_mtime_ns, source, platform, byte_offset, line_number, discarding_oversized, oversized_line_start, cumulative_rows, cumulative_bytes, cumulative_inserted, cumulative_rejected, completed_at, last_report_json, updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(source_file) DO UPDATE SET file_identity = excluded.file_identity, file_size = excluded.file_size, file_mtime_ns = excluded.file_mtime_ns, byte_offset = excluded.byte_offset, line_number = excluded.line_number, discarding_oversized = excluded.discarding_oversized, oversized_line_start = excluded.oversized_line_start, cumulative_rows = excluded.cumulative_rows, cumulative_bytes = excluded.cumulative_bytes, cumulative_inserted = excluded.cumulative_inserted, cumulative_rejected = excluded.cumulative_rejected, completed_at = excluded.completed_at, last_report_json = excluded.last_report_json, updated_at = excluded.updated_at''', ( todo_path, identity, file_size, file_mtime_ns, source, platform, next_offset, next_line, 1 if discarding else 0, oversized_line_start, cumulative_rows, cumulative_bytes, cumulative_inserted, cumulative_rejected, completed_at, json_dumps(report), now, ), ) if not _same_file_snapshot(todo_path, (descriptor, snapshot)): raise ReconciliationFileChanged('reconciliation input changed before commit') conn.commit() return report except ReconciliationFileChanged as exc: last_change = exc try: conn.rollback() except Exception: pass if mutation_attempt + 1 >= RECONCILIATION_MUTATION_RETRIES: raise time.sleep(0.01) except Exception: try: conn.rollback() except Exception: pass raise finally: if descriptor is not None: try: os.close(descriptor) except OSError: pass raise last_change or ReconciliationFileChanged('reconciliation input remained unstable') def resolve_reconciliation_issues(db, issue_ids): requested = sorted({int(value) for value in issue_ids or []}) if not requested: return 0 conn = getattr(db, 'conn', None) if not conn: raise RuntimeError('database connection is unavailable') placeholders = ','.join('?' for _ in requested) try: if conn.is_sqlite: conn.execute('BEGIN IMMEDIATE') lock_suffix = ' FOR UPDATE' if conn.is_postgres else '' rows = conn.execute( f'''SELECT id, resolved_at FROM target_queue_reconciliation_issues WHERE id IN ({placeholders}){lock_suffix}''', requested, ).fetchall() found = {int(row['id']): row['resolved_at'] for row in rows} if set(found) != set(requested) or any(found[issue_id] is not None for issue_id in requested): raise RuntimeError('one or more requested reconciliation issues were absent or already resolved') cursor = conn.execute( f'''UPDATE target_queue_reconciliation_issues SET resolved_at = ? WHERE id IN ({placeholders}) AND resolved_at IS NULL''', (utc_now_iso(), *requested), ) if int(getattr(cursor, 'rowcount', 0) or 0) != len(requested): raise RuntimeError('reconciliation issue set changed before resolution commit') conn.commit() return len(requested) except Exception: try: conn.rollback() except Exception: pass raise def _canonical_json_sha256(value): encoded = json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), default=str, ).encode('utf-8') return hashlib.sha256(encoded).hexdigest() def _target_queue_platform_for_source(source): source = str(source or '').strip() return 'docker' if source in ('docker', 'dockerhub') else source def configured_target_queue_query_policy(config): policy = {} document = [] rejected_entries = validate_rejected_query_policy(config) rejected = {} sources = (config or {}).get('sources') or {} if not isinstance(sources, dict): raise RuntimeError('configured source policy is not a mapping') for source in sorted(sources): source_config = sources[source] if not isinstance(source_config, dict) or 'queries' not in source_config: continue raw_queries = source_config.get('queries') if isinstance(raw_queries, str): raw_queries = raw_queries.split(',') if not isinstance(raw_queries, (list, tuple)): raise RuntimeError(f'configured query policy is invalid for source {source}') queries = [] for raw_query in raw_queries: query = str(raw_query or '').strip() if not query: raise RuntimeError(f'configured query policy contains an empty query for source {source}') if query in queries: raise RuntimeError(f'configured query policy contains a duplicate query for source {source}') queries.append(query) platform = _target_queue_platform_for_source(source) key = (str(source), platform) policy[key] = tuple(queries) document.append({ 'source': str(source), 'platform': platform, 'queries': sorted(queries), }) for entry in rejected_entries: source = entry['source'] platform = _target_queue_platform_for_source(source) rejected.setdefault((source, platform), []).append(entry['query']) rejected = { key: tuple(sorted(queries)) for key, queries in sorted(rejected.items()) } registry_present = (config or {}).get('query_policy') is not None policy_document = document if registry_present: policy_document = { 'active': document, 'rejected': [ {**entry, 'platform': _target_queue_platform_for_source(entry['source'])} for entry in rejected_entries ], } return { 'queries': policy, 'rejected': rejected, 'rejected_registry_present': registry_present, 'document': policy_document, 'policy_sha256': _canonical_json_sha256(policy_document), } def _stale_target_queue_rows(db, policy, source, platform, max_rows): source = str(source or '').strip() platform = str(platform or '').strip() queries = policy['queries'].get((source, platform)) if queries is None or not queries: raise RuntimeError('target queue policy scope has no configured query allowlist') maximum = min(100000, max(1, int(max_rows))) placeholders = ','.join('?' for _ in queries) rejected = policy.get('rejected', {}).get((source, platform), ()) registry_present = bool(policy.get('rejected_registry_present')) if registry_present and not rejected: return [], [] if registry_present: rejected_placeholders = ','.join('?' for _ in rejected) query_filter = ( f'AND queue.query IN ({rejected_placeholders}) ' f'AND queue.query NOT IN ({placeholders})' ) query_params = (*rejected, *queries) else: query_filter = f'AND queue.query NOT IN ({placeholders})' query_params = tuple(queries) rows = db.conn.execute( f'''SELECT queue.*, EXISTS ( SELECT 1 FROM result_reservations reservation WHERE reservation.queue_id = queue.id AND reservation.state IN ('scanning','ready','ingesting','db_committed') ) AS active_reservation, EXISTS ( SELECT 1 FROM docker_image_blob_coverage coverage JOIN docker_content_blobs blob ON blob.digest = coverage.blob_digest AND blob.coverage_policy_sha256 = coverage.coverage_policy_sha256 WHERE coverage.queue_id = queue.id AND blob.state IN ('leased','submitted') ) AS active_docker_blob FROM target_queue queue WHERE queue.source = ? AND queue.platform = ? AND queue.status IN ('pending','deferred','in_progress') AND queue.query IS NOT NULL AND BTRIM(queue.query) <> '' {query_filter} ORDER BY queue.id LIMIT ?''', (source, platform, *query_params, maximum + 1), ).fetchall() db.conn.commit() if len(rows) > maximum: raise RuntimeError(f'stale target queue selection exceeds its {maximum} row bound') eligible = [] blockers = [] fence_fields = ( 'lease_owner', 'lease_token', 'claim_batch', 'leased_at', 'lease_expires_at', 'current_result_reservation_id', 'claim_event_id', 'resolver_token', ) for row in rows: blocked = ( row['status'] not in ('pending', 'deferred') or any(row[field] is not None for field in fence_fields) or str(row['resolver_state'] or '') == 'resolving' or bool(row['active_reservation']) or bool(row['active_docker_blob']) ) if blocked: blockers.append(int(row['id'])) continue eligible.append({ 'queue_id': int(row['id']), 'source': str(row['source']), 'platform': str(row['platform']), 'query': str(row['query']), 'prior_status': str(row['status']), 'prior_updated_at': str(row['updated_at']), }) return eligible, blockers def plan_stale_target_queue_cold( db, config, *, source, platform, config_sha256, max_rows=10000, ): policy = configured_target_queue_query_policy(config) entries, blockers = _stale_target_queue_rows( db, policy, source, platform, max_rows, ) if blockers: raise RuntimeError( f'stale target queue policy selection has {len(blockers)} fenced row(s)' ) manifest = { 'schema': 1, 'type': 'truf-target-queue-cold-review-v1', 'config_sha256': str(config_sha256), 'policy_sha256': policy['policy_sha256'], 'source': str(source), 'platform': str(platform), 'reason_code': ( 'rejected_zero_alive' if policy.get('rejected_registry_present') else 'query_not_in_canonical_policy' ), 'selection_sha256': _canonical_json_sha256(entries), 'generated_at': utc_now_iso(), 'entries': entries, } return manifest def _cold_reactivation_rows(db, source, platform, query, max_rows): maximum = min(100000, max(1, int(max_rows))) rows = db.conn.execute( '''SELECT queue.*, cold_event.id AS cold_event_id, cold_event.prior_status AS restore_status, reverse_event.id AS reverse_event_id, EXISTS ( SELECT 1 FROM result_reservations reservation WHERE reservation.queue_id = queue.id AND reservation.state IN ('scanning','ready','ingesting','db_committed') ) AS active_reservation, EXISTS ( SELECT 1 FROM docker_image_blob_coverage coverage JOIN docker_content_blobs blob ON blob.digest = coverage.blob_digest AND blob.coverage_policy_sha256 = coverage.coverage_policy_sha256 WHERE coverage.queue_id = queue.id AND blob.state IN ('leased','submitted') ) AS active_docker_blob FROM target_queue queue JOIN target_queue_policy_events cold_event ON cold_event.queue_id = queue.id AND cold_event.action = 'cold' AND cold_event.experiment_id IS NULL LEFT JOIN target_queue_policy_events reverse_event ON reverse_event.reverses_event_id = cold_event.id WHERE queue.source = ? AND queue.platform = ? AND queue.query = ? AND queue.status = 'cold' AND reverse_event.id IS NULL ORDER BY queue.id LIMIT ?''', (str(source), str(platform), str(query), maximum + 1), ).fetchall() db.conn.commit() if len(rows) > maximum: raise RuntimeError(f'cold target queue selection exceeds its {maximum} row bound') entries = [] fence_fields = ( 'lease_owner', 'lease_token', 'claim_batch', 'leased_at', 'lease_expires_at', 'current_result_reservation_id', 'claim_event_id', 'resolver_token', ) for row in rows: if ( any(row[field] is not None for field in fence_fields) or str(row['resolver_state'] or '') == 'resolving' or bool(row['active_reservation']) or bool(row['active_docker_blob']) ): raise RuntimeError('cold target queue reactivation has a fenced row') entries.append({ 'queue_id': int(row['id']), 'source': str(row['source']), 'platform': str(row['platform']), 'query': str(row['query']), 'cold_event_id': int(row['cold_event_id']), 'restore_status': str(row['restore_status']), 'prior_updated_at': str(row['updated_at']), }) return entries def plan_cold_target_queue_reactivation( db, config, *, source, platform, query, config_sha256, max_rows=10000, ): policy = configured_target_queue_query_policy(config) entries = _cold_reactivation_rows(db, source, platform, query, max_rows) return { 'schema': 1, 'type': 'truf-target-queue-reactivation-review-v1', 'config_sha256': str(config_sha256), 'policy_sha256': policy['policy_sha256'], 'source': str(source), 'platform': str(platform), 'query': str(query), 'reason_code': 'reviewed_policy_reactivation', 'selection_sha256': _canonical_json_sha256(entries), 'generated_at': utc_now_iso(), 'entries': entries, } def load_target_queue_policy_manifest(path, expected_type, max_rows=10000): require_private_file(path) manifest = read_private_json(path, max_bytes=TARGET_QUEUE_POLICY_MANIFEST_MAX_BYTES) if not isinstance(manifest, dict) or manifest.get('type') != expected_type: raise RuntimeError('target queue policy manifest type is invalid') expected_keys = { 'schema', 'type', 'config_sha256', 'policy_sha256', 'source', 'platform', 'reason_code', 'selection_sha256', 'generated_at', 'entries', } if expected_type == 'truf-target-queue-reactivation-review-v1': expected_keys.add('query') if set(manifest) != expected_keys or manifest.get('schema') != 1: raise RuntimeError('target queue policy manifest shape is invalid') entries = manifest.get('entries') maximum = min(100000, max(1, int(max_rows))) if not isinstance(entries, list) or not entries or len(entries) > maximum: raise RuntimeError('target queue policy manifest entry count is outside its bound') action = 'cold' if expected_type == 'truf-target-queue-cold-review-v1' else 'reactivate' entries = ScannerDB._normalize_target_queue_policy_entries(entries, action, maximum) if entries != manifest['entries']: raise RuntimeError('target queue policy manifest entries are not canonical') for name in ('config_sha256', 'policy_sha256', 'selection_sha256'): if not re.fullmatch(r'[a-f0-9]{64}', str(manifest.get(name) or '')): raise RuntimeError(f'target queue policy manifest {name} is invalid') if manifest['selection_sha256'] != _canonical_json_sha256(entries): raise RuntimeError('target queue policy selection hash is invalid') return manifest, _canonical_json_sha256(manifest) def _manifest_already_applied(db, manifest, manifest_sha256, action): rows = db.conn.execute( '''SELECT id, queue_id, action FROM target_queue_policy_events WHERE manifest_sha256 = ? ORDER BY queue_id''', (manifest_sha256,), ).fetchall() db.conn.commit() if not rows: return False expected_ids = [int(entry['queue_id']) for entry in manifest['entries']] actual_ids = [int(row['queue_id']) for row in rows] if actual_ids != expected_ids or any(row['action'] != action for row in rows): raise RuntimeError('target queue policy manifest was only partially or differently applied') return True def apply_stale_target_queue_cold( db, config, manifest, manifest_sha256, *, config_sha256, max_rows=10000, ): policy = configured_target_queue_query_policy(config) if manifest['config_sha256'] != config_sha256: raise RuntimeError('target queue policy manifest config identity drifted') if manifest['policy_sha256'] != policy['policy_sha256']: raise RuntimeError('target queue policy manifest query policy drifted') if not _manifest_already_applied(db, manifest, manifest_sha256, 'cold'): current = plan_stale_target_queue_cold( db, config, source=manifest['source'], platform=manifest['platform'], config_sha256=config_sha256, max_rows=max_rows, ) if ( current['selection_sha256'] != manifest['selection_sha256'] or current['entries'] != manifest['entries'] ): raise RuntimeError('target queue policy selection drifted after review') return db.cold_target_queue_rows( manifest['entries'], reason_code=manifest['reason_code'], config_sha256=config_sha256, policy_sha256=policy['policy_sha256'], manifest_sha256=manifest_sha256, max_rows=max_rows, ) def apply_cold_target_queue_reactivation( db, config, manifest, manifest_sha256, *, config_sha256, max_rows=10000, ): policy = configured_target_queue_query_policy(config) if manifest['config_sha256'] != config_sha256: raise RuntimeError('target queue reactivation manifest config identity drifted') if manifest['policy_sha256'] != policy['policy_sha256']: raise RuntimeError('target queue reactivation manifest query policy drifted') if not _manifest_already_applied(db, manifest, manifest_sha256, 'reactivate'): current = plan_cold_target_queue_reactivation( db, config, source=manifest['source'], platform=manifest['platform'], query=manifest['query'], config_sha256=config_sha256, max_rows=max_rows, ) if ( current['selection_sha256'] != manifest['selection_sha256'] or current['entries'] != manifest['entries'] ): raise RuntimeError('target queue reactivation selection drifted after review') return db.reactivate_cold_target_queue_rows( manifest['entries'], reason_code=manifest['reason_code'], config_sha256=config_sha256, policy_sha256=policy['policy_sha256'], manifest_sha256=manifest_sha256, max_rows=max_rows, ) def load_config(path): try: import yaml except ImportError as exc: raise SystemExit('PyYAML is required for runtime safety migration') from exc with open(path, 'r', encoding='utf-8') as handle: return apply_path_config(yaml.safe_load(handle) or {}, path) def _supervisor_metadata_paths(config): global_config = config.get('global') or {} supervisor = config.get('supervisor') or {} log_dir = supervisor.get('log_dir') or global_config.get('log_dir') or '' control_dir = supervisor.get('control_dir') or global_config.get('control_dir') or '' paths = [ supervisor.get('instance_file') or os.path.join(control_dir, 'supervisor.instance.json'), os.path.join(log_dir, 'supervisor.instance.json'), os.path.join(log_dir, 'supervisor.pid'), ] return list(dict.fromkeys(os.path.abspath(path) for path in paths if path)) def require_local_sources_stopped(config, inspect_scan_slots=True): for path in _supervisor_metadata_paths(config): if os.path.lexists(path): try: metadata = read_private_json(path) from process_identity import exact_process_identity_state state = exact_process_identity_state( metadata.get('pid'), metadata.get('process_creation_time'), metadata.get('executable'), ) except (OSError, ValueError): state = 'unknown' if state not in ('dead', 'reused'): raise RuntimeError( f'refusing migration while supervisor metadata identity is {state}: {path}' ) global_config = config.get('global') or {} limiter_path = global_config.get('scan_limiter_db') or os.path.join(global_config.get('state_dir') or '', 'scan_limiter.db') if inspect_scan_slots and limiter_path and os.path.isfile(limiter_path): try: uri = 'file:' + os.path.abspath(limiter_path).replace('\\', '/') + '?mode=ro' with sqlite3.connect(uri, uri=True, timeout=1) as local_db: row = local_db.execute("SELECT COUNT(*) FROM scan_slots").fetchone() if row and int(row[0] or 0) > 0: raise RuntimeError(f'refusing migration while {row[0]} local scan slot(s) are active') except sqlite3.OperationalError as exc: if 'no such table' not in str(exc).lower(): raise RuntimeError(f'unable to prove local scan slots are stopped: {exc}') from exc def recover_dead_scan_slots(config, identity_live=None, now=None, max_rows=10000): from scanner import exact_process_identity_live identity_live = identity_live or exact_process_identity_live now = time.time() if now is None else float(now) max_rows = min(10000, max(1, int(max_rows))) global_config = config.get('global') or {} path = global_config.get('scan_limiter_db') or os.path.join( global_config.get('state_dir') or '', 'scan_limiter.db', ) path = os.path.abspath(path) if not os.path.isfile(path): return {'path': path, 'scanned': 0, 'deleted': 0, 'remaining': 0, 'rows': []} require_private_file(path) connection = sqlite3.connect(path, timeout=30) connection.row_factory = sqlite3.Row try: connection.execute('PRAGMA busy_timeout=30000') connection.execute('BEGIN IMMEDIATE') columns = {row[1] for row in connection.execute('PRAGMA table_info(scan_slots)').fetchall()} required = { 'slot_id', 'owner_pid', 'owner_thread', 'owner_source', 'owner_creation_time', 'owner_executable', 'child_pid', 'child_creation_time', 'child_executable', 'acquired_at', 'updated_at', } if not required.issubset(columns): raise RuntimeError('scan-slot table lacks exact process identity columns') rows = connection.execute( '''SELECT slot_id, owner_pid, owner_thread, owner_source, owner_creation_time, owner_executable, child_pid, child_creation_time, child_executable, acquired_at, updated_at FROM scan_slots ORDER BY acquired_at LIMIT ?''', (max_rows + 1,), ).fetchall() if len(rows) > max_rows: raise RuntimeError(f'scan-slot recovery exceeds its {max_rows} row bound') report_rows = [] deleted = 0 for row in rows: checks = [] for prefix in ('owner', 'child'): try: live = identity_live( row[f'{prefix}_pid'], row[f'{prefix}_creation_time'], row[f'{prefix}_executable'], ) except Exception: live = None checks.append(live if live is None else bool(live)) owner_live, child_live = checks deleted_row = False if owner_live is False and child_live is False: cursor = connection.execute( '''DELETE FROM scan_slots WHERE slot_id = ? AND owner_pid = ? AND owner_thread = ? AND owner_source IS ? AND owner_creation_time IS ? AND owner_executable IS ? AND child_pid IS ? AND child_creation_time IS ? AND child_executable IS ? AND acquired_at = ? AND updated_at = ?''', ( row['slot_id'], row['owner_pid'], row['owner_thread'], row['owner_source'], row['owner_creation_time'], row['owner_executable'], row['child_pid'], row['child_creation_time'], row['child_executable'], row['acquired_at'], row['updated_at'], ), ) if int(cursor.rowcount or 0) != 1: raise RuntimeError(f'scan-slot identity changed during recovery: {row["slot_id"]}') deleted += 1 deleted_row = True report_rows.append({ 'slot_id': row['slot_id'], 'owner_pid': row['owner_pid'], 'owner_thread': row['owner_thread'], 'owner_source': row['owner_source'], 'owner_creation_time': row['owner_creation_time'], 'owner_executable': row['owner_executable'], 'owner_live': owner_live, 'child_pid': row['child_pid'], 'child_creation_time': row['child_creation_time'], 'child_executable': row['child_executable'], 'child_live': child_live, 'acquired_at': row['acquired_at'], 'updated_at': row['updated_at'], 'age_seconds': max(0.0, now - float(row['updated_at'] or row['acquired_at'])), 'deleted': deleted_row, }) remaining = int(connection.execute('SELECT COUNT(*) FROM scan_slots').fetchone()[0]) connection.commit() return { 'path': path, 'scanned': len(rows), 'deleted': deleted, 'remaining': remaining, 'rows': report_rows, } except BaseException: connection.rollback() raise finally: connection.close() harden_private_file(path) def _offline_result_pipeline_state(db): conn = getattr(db, 'conn', None) if not conn or not conn.is_postgres: raise RuntimeError('offline result-pipeline recovery requires PostgreSQL') return { 'worker_leases': int(conn.execute( "SELECT COUNT(*) AS count FROM pipeline_leases WHERE state NOT IN ('released','failed')" ).fetchone()['count'] or 0), 'result_reservations': int(conn.execute( """SELECT COUNT(*) AS count FROM result_reservations WHERE state IN ('scanning','ready','ingesting','db_committed')""" ).fetchone()['count'] or 0), 'queue_leases': int(conn.execute( """SELECT COUNT(*) AS count FROM target_queue q LEFT JOIN 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')""" ).fetchone()['count'] or 0), 'blob_leases': int(conn.execute( """SELECT COUNT(*) AS count FROM docker_content_blobs WHERE state IN ('leased','submitted') OR lease_reservation_id IS NOT NULL""" ).fetchone()['count'] or 0), } def _require_offline_result_pipeline_schema(db): conn = getattr(db, 'conn', None) if not conn or not conn.is_postgres: raise RuntimeError('offline result-pipeline recovery requires PostgreSQL') marker = '20260901_13_docker_layer_content_scanning' if not conn.table_exists('runtime_schema_migrations') or not conn.execute( 'SELECT 1 AS present FROM runtime_schema_migrations WHERE version = ?', (marker,), ).fetchone(): raise RuntimeError('offline result-pipeline recovery requires the marker-13 schema') required = { 'pipeline_leases': {'worker_name', 'generation', 'lease_token', 'state'}, 'result_reservations': { 'id', 'state', 'queue_id', 'claim_lease_token', 'scan_event_id', 'producer_pid', 'producer_creation_time', 'producer_executable', 'ready_relative_path', 'bundle_id', 'reservation_token', }, 'result_bundles': {'reservation_id', 'state', 'relative_path', 'scan_event_hash'}, 'target_queue': {'id', 'status', 'current_result_reservation_id', 'claim_event_id'}, 'docker_content_blobs': {'state', 'lease_reservation_id'}, 'keycheck_candidates': {'candidate_uid', 'credential_id', 'service', 'secret_hash'}, } for table, columns in required.items(): if not conn.table_exists(table) or not columns.issubset(conn.table_columns(table)): raise RuntimeError(f'offline result-pipeline recovery schema is incomplete: {table}') conn.commit() def recover_stale_result_pipeline(db, config, max_rows=1000, max_seconds=300.0): _require_offline_result_pipeline_schema(db) max_rows = min(10000, max(1, int(max_rows))) max_seconds = min(3600.0, max(1.0, float(max_seconds))) initial = _offline_result_pipeline_state(db) if initial['result_reservations'] > max_rows: raise RuntimeError('offline result-pipeline recovery exceeds its reservation bound') global_config = config.get('global') or {} worker_config = ((config.get('supervisor') or {}).get('result_ingester') or {}) lease_seconds = min(3600, max(60, int(worker_config.get('lease_seconds', 300)))) instance_id = f'offline-result-pipeline-{os.getpid()}' identity = current_process_identity() projector_lease = None worker = None processed = 0 recovery_passes = 0 deadline = time.monotonic() + max_seconds try: projector_lease = db.acquire_pipeline_lease( 'jsonl_projector', instance_id, identity, lease_seconds=lease_seconds, initial_state='recovering', ) if not projector_lease: raise RuntimeError('jsonl projector advisory lock is held') worker = ResultIngester( db, global_config['result_bundle_dir'], instance_id, lease_seconds=lease_seconds, quarantine_max_items=int(global_config.get('pipeline_quarantine_max_items', 10000)), quarantine_max_bytes=int( global_config.get('pipeline_quarantine_max_bytes', 1024 * 1024 * 1024) ), recover_expired_ready=True, ) worker.lease = db.acquire_pipeline_lease( 'result_ingester', instance_id, identity, lease_seconds=lease_seconds, initial_state='recovering', ) if not worker.lease: raise RuntimeError('result ingester advisory lock is held') previous_active = None while time.monotonic() < deadline: current = _offline_result_pipeline_state(db)['result_reservations'] if current == 0: break if previous_active is not None and current >= previous_active: break previous_active = current worker.recovery_after_id = 0 worker.recover( page_size=min(100, max_rows), max_pages=max(1, (max_rows + 99) // 100), ) recovery_passes += 1 if not worker.heartbeat('recovering'): raise RuntimeError('offline result ingester heartbeat fence was lost') while processed < max_rows and time.monotonic() < deadline: if not worker.process_one(): break processed += 1 if not worker.heartbeat('recovering'): raise RuntimeError('offline result ingester heartbeat fence was lost') if worker and not worker.stop(): raise RuntimeError('offline result ingester lease release was not fenced') worker = None if not db.release_pipeline_lease( 'jsonl_projector', projector_lease['generation'], projector_lease['lease_token'], ): raise RuntimeError('offline JSONL projector lease release was not fenced') projector_lease = None final = _offline_result_pipeline_state(db) return { 'initial': initial, 'processed_bundles': processed, 'recovery_passes': recovery_passes, 'final': final, 'quiescent': not any(final.values()), } except BaseException: if worker is not None: try: worker.stop('offline_result_pipeline_recovery_failed') except Exception: pass if projector_lease is not None: try: db.release_pipeline_lease( 'jsonl_projector', projector_lease['generation'], projector_lease['lease_token'], state='failed', error='offline_result_pipeline_recovery_failed', ) except Exception: pass raise def _pid_running(pid): try: pid = int(pid) if pid <= 0 or pid == os.getpid(): return False if os.name == 'nt': from process_identity import open_process process = open_process(pid) try: return process.is_running() finally: process.close() os.kill(pid, 0) return True except (OSError, TypeError, ValueError): return False def known_application_process_markers(config): global_config = config.get('global') or {} roots = [global_config.get('work_dir')] active = [] for root in roots: if not root or not os.path.isdir(root) or is_reparse_point(root): continue with os.scandir(root) as entries: for index, entry in enumerate(entries): if index >= 256: raise RuntimeError(f'unable to bound known application process markers under {root}') if entry.is_symlink() or is_reparse_point(entry.path) or not entry.is_dir(follow_symlinks=False): continue marker = os.path.join(entry.path, '.scanner-owner.json') if not os.path.isfile(marker) or not private_file_ready(marker): continue try: with open(marker, 'r', encoding='utf-8') as handle: value = json.load(handle) prefix = 'owner' if value.get('owner_pid') else 'parent' pid = int(value.get(f'{prefix}_pid') or 0) creation_time = value.get(f'{prefix}_creation_time') executable = value.get(f'{prefix}_executable') except (OSError, TypeError, ValueError, json.JSONDecodeError): continue from process_identity import exact_process_identity_state state = exact_process_identity_state(pid, creation_time, executable) if state == 'alive' or (state == 'unknown' and _pid_running(pid)): active.append((pid, marker)) return active def require_runtime_hardening_stopped(config): require_local_sources_stopped(config, inspect_scan_slots=False) paths = postgres_runtime_paths(config) postmaster_pid = os.path.join(paths['data_dir'], 'postmaster.pid') if os.path.lexists(postmaster_pid): raise RuntimeError(f'refusing hardening while postmaster.pid exists: {postmaster_pid}') port = configured_cluster_values()['port'] config_path = str((config.get('global') or {}).get('project_dir') or '') env_path = find_postgres_environment_path(os.path.join(config_path, 'config.yaml'), config) if config_path else None if env_path: try: with open(env_path, 'r', encoding='utf-8') as handle: for line in handle: key, separator, value = line.strip().partition('=') if separator and key.strip() == 'TRUF_POSTGRES_PORT': port = int(value.strip().strip('"').strip("'")) if not 0 < port <= 65535: raise ValueError('port outside valid range') break except (OSError, ValueError) as exc: raise RuntimeError(f'unable to determine the offline PostgreSQL listener port: {exc}') from exc if _listener_present(port): raise RuntimeError(f'refusing hardening while a listener is present on 127.0.0.1:{port}') active = known_application_process_markers(config) if active: pid, marker = active[0] raise RuntimeError(f'refusing hardening while known application process {pid} is active: {marker}') def find_postgres_environment_path(config_path, config): global_config = (config or {}).get('global') or {} candidates = [] if global_config.get('root_dir'): candidates.append(os.path.join(global_config['root_dir'], '.env.postgres')) candidates.extend(( os.path.join(os.path.dirname(config_path), '..', '.env.postgres'), os.path.join(os.path.dirname(config_path), '.env.postgres'), )) return next((os.path.abspath(path) for path in candidates if os.path.isfile(path)), None) def _online_postgres_identity(db, dsn, config, identity=None): identity = identity or verify_cluster_identity(config) canonical = canonical_postgres_url(dsn, identity['database'], identity['user'], identity['port']) if canonical != dsn: db.url = canonical row = db.conn.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')::integer AS port, (SELECT system_identifier::text FROM pg_catalog.pg_control_system()) AS system_identifier''' ).fetchone() checks = { 'database': (str(row['database']), identity['database']), 'user': (str(row['user_name']), identity['user']), 'data_directory': ( os.path.normcase(os.path.realpath(os.path.abspath(row['data_directory']))), os.path.normcase(os.path.realpath(os.path.abspath(identity['data_directory']))), ), 'port': (int(row['port']), int(identity['port'])), 'system_identifier': (str(row['system_identifier']), str(identity['system_identifier'])), } for name, (actual, expected) in checks.items(): if actual != expected: raise RuntimeError(f'online migration target {name} does not match cluster_identity.json') db.conn.commit() return identity, canonical def _require_no_postgres_application_sessions(db): rows = db.conn.execute( '''SELECT pid, application_name, state FROM pg_catalog.pg_stat_activity WHERE datname = pg_catalog.current_database() AND pid <> pg_catalog.pg_backend_pid() AND backend_type = 'client backend' LIMIT 20''' ).fetchall() db.conn.commit() if rows: names = ', '.join(str(row['application_name'] or '')[:80] for row in rows[:5]) raise RuntimeError(f'refusing migration while {len(rows)} other database application session(s) are active: {names}') @contextlib.contextmanager def postgres_migration_guard(db): db.conn.execute(f'SET search_path = {POSTGRES_APPLICATION_SCHEMA}') schema_row = db.conn.execute( """SELECT pg_catalog.current_schema() AS schema_name, pg_catalog.current_setting('search_path') AS search_path, 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() normalized_path = ''.join(str(schema_row['search_path'] or '').lower().replace('"', '').split()) if schema_row else '' if ( not schema_row or schema_row['schema_name'] != POSTGRES_APPLICATION_SCHEMA or normalized_path != 'public' or bool(schema_row.get('public_create', False)) ): raise RuntimeError('migration PostgreSQL schema/search_path validation failed') db.conn.execute("SET application_name = 'truf-offline-migration'") db.conn.execute("SET statement_timeout = 0") db.conn.execute("SET lock_timeout = '10s'") db.conn.execute("SET idle_in_transaction_session_timeout = 0") db.conn.commit() _require_no_postgres_application_sessions(db) row = db.conn.execute('SELECT pg_catalog.pg_try_advisory_lock(781273968142991337) AS locked').fetchone() db.conn.commit() if not row or not row['locked']: raise RuntimeError('another PostgreSQL runtime-safety migration holds the advisory lock') try: _require_no_postgres_application_sessions(db) yield finally: try: db.conn.execute('SELECT pg_catalog.pg_advisory_unlock(781273968142991337)') db.conn.commit() except Exception: try: db.conn.rollback() except Exception: pass def harden_runtime_paths(config, env_path=None, config_path=None): global_config = config.get('global') or {} supervisor_config = config.get('supervisor') or {} bundle_dir = global_config.get('result_bundle_dir') if bundle_dir and os.name == 'nt' and os.path.splitdrive(os.path.abspath(bundle_dir))[0].upper() != 'S:': raise RuntimeError('production result_bundle_dir must be on S:') directories = [] for key in ( 'root_dir', 'project_dir', 'runtime_dir', 'results_dir', 'result_spool_dir', 'result_bundle_dir', 'control_dir', 'queue_dir', 'state_dir', 'keycheck_dir', 'postman_cache_dir', 'gharchive_cache_dir', 'log_dir', 'work_dir', ): if global_config.get(key): directories.append(global_config[key]) for key in ('control_dir', 'log_dir', 'state_dir'): if supervisor_config.get(key): directories.append(supervisor_config[key]) if bundle_dir: directories.extend( os.path.join(bundle_dir, name) for name in ('tmp', 'ready', 'quarantine') ) postgres_paths = postgres_runtime_paths(config) postgres_paths['data_dir'] = canonical_cluster_data_directory(config) directories.append(postgres_paths['postgres_dir']) try: data_inside_postgres_dir = ( os.path.commonpath((postgres_paths['postgres_dir'], postgres_paths['data_dir'])) == postgres_paths['postgres_dir'] ) except ValueError: data_inside_postgres_dir = False if not data_inside_postgres_dir: directories.append(postgres_paths['data_dir']) normalized_directories = sorted( {os.path.normcase(os.path.abspath(path)) for path in directories if path}, key=lambda path: (path.count(os.sep), path), ) files = [ config_path, env_path, global_config.get('secrets_file'), global_config.get('proxy_file'), global_config.get('api_proxy_file'), global_config.get('download_proxy_file'), global_config.get('trufflehog_config'), ] from lifecycle_authority import ( GIT_MANIFEST_NAME, TRUFFLEHOG_MANIFEST_NAME, manifest_authority_paths, resolve_manifest_executable, ) authority_files = manifest_authority_paths( global_config.get('project_dir'), global_config.get('trufflehog_path'), policy_paths=[global_config.get('trufflehog_config')], existing_only=True, include_executables=os.name == 'nt', ) native_files = [] if os.name != 'nt': for name, value in ( (TRUFFLEHOG_MANIFEST_NAME, global_config.get('trufflehog_path')), (GIT_MANIFEST_NAME, None), ): native_files.append(require_trusted_native_executable(resolve_manifest_executable( value, name=name, app_dir=global_config.get('project_dir'), ))) for key in ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata', 'initdb', 'psql'): native_files.append(require_trusted_native_executable(postgres_paths[key])) files.extend(authority_files) required_authority_files = { os.path.normcase(os.path.abspath(path)) for path in authority_files } config_dir = os.path.dirname(os.path.abspath(config_path)) if config_path else None for parent in ( global_config.get('root_dir'), global_config.get('project_dir'), config_dir, os.path.dirname(config_dir) if config_dir else None, ): if parent: candidate = os.path.join(parent, '.env.postgres') if candidate not in files: files.append(candidate) # Private trees may not encompass immutable native code, including its # parent. Reject conflicting configuration rather than chmod system paths. private_parents = normalized_directories + [ os.path.dirname(os.path.abspath(path)) for path in files if path ] for native in native_files: parent = os.path.dirname(native) for directory in private_parents: if os.path.commonpath((directory, parent)) in (directory, parent): raise PrivateFileError(f'private runtime authority overlaps a native executable directory: {parent}') hardened = 0 for directory in normalized_directories: reject_reparse_components(os.path.dirname(directory)) os.makedirs(directory, mode=0o700, exist_ok=True) reject_reparse_components(directory) harden_private_directory(directory) hardened += 1 for path in files: if not path: continue absolute = os.path.abspath(path) if not os.path.lexists(absolute): if os.path.normcase(absolute) in required_authority_files: raise PrivateFileError(f'required code authority file is absent: {absolute}') if config_path and os.path.normcase(absolute) == os.path.normcase(os.path.abspath(config_path)): raise PrivateFileError(f'required config file is absent: {absolute}') continue reject_reparse_components(absolute) if not os.path.isfile(absolute): raise PrivateFileError(f'sensitive configured file is not regular: {absolute}') harden_private_directory(os.path.dirname(absolute)) harden_private_file(absolute) hardened += 1 # Preserve recursive hardening of runtime-owned content after every # configured parent has been created and hardened in a deterministic order. tree_roots = [] for key in ('runtime_dir', 'results_dir', 'result_spool_dir', 'result_bundle_dir', 'queue_dir', 'state_dir', 'keycheck_dir', 'postman_cache_dir', 'log_dir', 'work_dir'): if global_config.get(key): tree_roots.append(global_config[key]) tree_roots.append(postgres_paths['postgres_dir']) if not data_inside_postgres_dir: tree_roots.append(postgres_paths['data_dir']) seen_trees = set() for path in tree_roots: absolute = os.path.normcase(os.path.abspath(path)) if absolute in seen_trees: continue seen_trees.add(absolute) hardened += harden_private_tree(absolute) for directory in normalized_directories: if not private_directory_ready(directory): raise PrivateFileError(f'private directory verification failed after hardening: {directory}') for path in files: if path and os.path.isfile(path) and not private_file_ready(path): raise PrivateFileError(f'private file verification failed after hardening: {path}') for path in authority_files: if not private_file_ready(path): raise PrivateFileError(f'required code authority verification failed after hardening: {path}') for path in native_files: require_trusted_native_executable(path) return hardened def load_jsonl_issue_review_manifest(path): path = os.path.abspath(path) require_private_file(path) if os.path.getsize(path) > 16 * 1024 * 1024: raise RuntimeError('JSONL issue review manifest exceeds its byte bound') with open(path, 'r', encoding='utf-8') as handle: value = json.load(handle) if not isinstance(value, dict) or set(value) != {'schema', 'type', 'audit_sha256', 'issues'}: raise RuntimeError('JSONL issue review manifest root is invalid') if value.get('schema') != 1 or value.get('type') != 'truf-jsonl-projection-review': raise RuntimeError('JSONL issue review manifest schema is unsupported') if not re.fullmatch(r'[a-f0-9]{64}', str(value.get('audit_sha256') or '')): raise RuntimeError('JSONL issue review manifest audit identity is invalid') issues = value.get('issues') if not isinstance(issues, list) or not issues or len(issues) > 100000: raise RuntimeError('JSONL issue review manifest issue count is outside its bound') allowed_classifications = { 'invalid_json', 'legacy_numeric_prefix_corrupt_json', 'invalid_utf8', 'utf8_bom_prefix', 'unterminated_record', 'invalid_error_projection', 'oversized_record', } normalized = [] seen = set() expected_keys = { 'ledger_kind', 'file', 'offset', 'length', 'sha256', 'classification', 'safe_identity', } for issue in issues: if not isinstance(issue, dict) or set(issue) != expected_keys: raise RuntimeError('JSONL issue review manifest entry is invalid') ledger_kind = str(issue.get('ledger_kind') or '') file_name = str(issue.get('file') or '') offset = int(issue.get('offset', -1)) length = int(issue.get('length', 0)) digest = str(issue.get('sha256') or '').lower() classification = str(issue.get('classification') or '') if ( ledger_kind not in ('scan_results', 'found_secrets', 'scan_errors') or not file_name or os.path.basename(file_name) != file_name or offset < 0 or length <= 0 or not re.fullmatch(r'[a-f0-9]{64}', digest) or classification not in allowed_classifications or issue.get('safe_identity') is not False ): raise RuntimeError('JSONL issue review manifest metadata is invalid') identity = (ledger_kind, file_name, offset, digest) if identity in seen: raise RuntimeError('JSONL issue review manifest contains duplicate metadata') seen.add(identity) normalized.append({ 'ledger_kind': ledger_kind, 'file': file_name, 'offset': offset, 'length': length, 'sha256': digest, 'classification': classification, }) return normalized def load_jsonl_conflict_review_manifest(path): path = os.path.abspath(path) require_private_file(path) if os.path.getsize(path) > 16 * 1024 * 1024: raise RuntimeError('JSONL conflict review manifest exceeds its byte bound') with open(path, 'r', encoding='utf-8') as handle: value = json.load(handle) if not isinstance(value, dict) or set(value) != {'schema', 'type', 'audit_sha256', 'variants'}: raise RuntimeError('JSONL conflict review manifest root is invalid') if value.get('schema') != 1 or value.get('type') != 'truf-jsonl-projection-conflict-review': raise RuntimeError('JSONL conflict review manifest schema is unsupported') if not re.fullmatch(r'[a-f0-9]{64}', str(value.get('audit_sha256') or '')): raise RuntimeError('JSONL conflict review manifest audit identity is invalid') variants = value.get('variants') if not isinstance(variants, list) or not variants or len(variants) > 100000: raise RuntimeError('JSONL conflict review manifest variant count is outside its bound') expected_keys = { 'ledger_kind', 'file', 'offset', 'length', 'identity_sha256', 'payload_sha256', 'field_name_set_sha256', 'classification', 'occurrences', } normalized = [] seen = set() for variant in variants: if not isinstance(variant, dict) or set(variant) != expected_keys: raise RuntimeError('JSONL conflict review manifest entry is invalid') ledger_kind = str(variant.get('ledger_kind') or '') file_name = str(variant.get('file') or '') offset = int(variant.get('offset', -1)) length = int(variant.get('length', 0)) identity_sha256 = str(variant.get('identity_sha256') or '').lower() payload_sha256 = str(variant.get('payload_sha256') or '').lower() field_name_set_sha256 = str(variant.get('field_name_set_sha256') or '').lower() occurrences = int(variant.get('occurrences', 0)) if ( ledger_kind not in ('scan_results', 'found_secrets', 'scan_errors') or not file_name or os.path.basename(file_name) != file_name or offset < 0 or length <= 0 or occurrences <= 0 or not re.fullmatch(r'[a-f0-9]{64}', identity_sha256) or not re.fullmatch(r'[a-f0-9]{64}', payload_sha256) or not re.fullmatch(r'[a-f0-9]{64}', field_name_set_sha256) or variant.get('classification') != 'historical_payload_variant' ): raise RuntimeError('JSONL conflict review manifest metadata is invalid') identity = (ledger_kind, identity_sha256, payload_sha256) if identity in seen: raise RuntimeError('JSONL conflict review manifest contains duplicate variants') seen.add(identity) normalized.append({ 'ledger_kind': ledger_kind, 'file': file_name, 'offset': offset, 'length': length, 'identity_sha256': identity_sha256, 'payload_sha256': payload_sha256, 'field_name_set_sha256': field_name_set_sha256, 'classification': 'historical_payload_variant', 'occurrences': occurrences, }) return normalized def reconcile_jsonl_projection_ledgers(config, args): from scanner import ( approve_projection_conflict_variants, approve_projection_reconciliation_issues, reconcile_projection_ledger_batch, ) global_config = config.get('global') or {} results_dir = global_config.get('results_dir') if not results_dir: raise RuntimeError('configured results_dir is required for JSONL projection reconciliation') specs = ( ('scan_results.jsonl', 'scan_event_id'), ('found_secrets.jsonl', 'finding_uid'), ('scan_errors.log', 'error_row_id'), ) reviewed = {name: [] for name, _ in specs} for value in args.resolve_jsonl_issue: parts = str(value or '').rsplit(':', 2) if len(parts) != 3: raise SystemExit('--resolve-jsonl-issue requires FILE:OFFSET:SHA256') physical_file, raw_offset, digest = parts selected = None for name, identity_key in specs: stem, extension = os.path.splitext(name) if physical_file == name or ( physical_file.startswith(stem + '.') and physical_file.endswith(extension) ): selected = (name, identity_key) break if selected is None: raise SystemExit(f'unsupported JSONL projection issue file: {physical_file}') name, identity_key = selected reviewed[name].append({ 'file': physical_file, 'offset': int(raw_offset), 'sha256': digest, }) if args.resolve_jsonl_issue_manifest: for issue in load_jsonl_issue_review_manifest(args.resolve_jsonl_issue_manifest): reviewed[issue['ledger_kind'] + ('.log' if issue['ledger_kind'] == 'scan_errors' else '.jsonl')].append(issue) for name, identity_key in specs: if not reviewed[name]: continue report = approve_projection_reconciliation_issues( os.path.join(results_dir, name), identity_key, reviewed[name], max_record_bytes=int(global_config.get('result_spool_max_event_bytes', 192 * 1024 * 1024) or 192 * 1024 * 1024), ) print(json.dumps({ 'resolved_issue_review': {'ledger': name, **report}, }, ensure_ascii=True, sort_keys=True), flush=True) if args.resolve_jsonl_conflict_manifest: conflict_reviews = {name: [] for name, _ in specs} for variant in load_jsonl_conflict_review_manifest(args.resolve_jsonl_conflict_manifest): logical_name = variant['ledger_kind'] + ( '.log' if variant['ledger_kind'] == 'scan_errors' else '.jsonl' ) conflict_reviews[logical_name].append(variant) for name, identity_key in specs: if not conflict_reviews[name]: continue report = approve_projection_conflict_variants( os.path.join(results_dir, name), identity_key, conflict_reviews[name], max_record_bytes=int(global_config.get('result_spool_max_event_bytes', 192 * 1024 * 1024) or 192 * 1024 * 1024), progress_callback=lambda value, ledger_name=name: print(json.dumps({ 'conflict_review_progress': {'ledger': ledger_name, **value}, }, ensure_ascii=True, sort_keys=True), flush=True), ) print(json.dumps({ 'resolved_conflict_review': {'ledger': name, **report}, }, ensure_ascii=True, sort_keys=True), flush=True) all_complete = True for name, identity_key in specs: path = os.path.join(results_dir, name) report = None batch_limit = max(1, int(args.max_batches)) if args.until_complete else 1 for _ in range(batch_limit): report = reconcile_projection_ledger_batch( path, identity_key, max_rows=args.max_rows, max_bytes=args.max_bytes, max_seconds=args.max_seconds, row_limit=int(global_config.get('jsonl_ledger_max_rows', 1000000) or 1000000), ledger_byte_limit=int(global_config.get('jsonl_ledger_max_bytes', 512 * 1024 * 1024) or 512 * 1024 * 1024), max_record_bytes=int(global_config.get('result_spool_max_event_bytes', 192 * 1024 * 1024) or 192 * 1024 * 1024), ) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) if report.get('complete'): break if not report or not report.get('complete'): all_complete = False return 0 if all_complete else 2 def _metadata_is_docker_command_timeout(metadata_json): try: metadata = json.loads(metadata_json or '{}') except (TypeError, ValueError): return False if not isinstance(metadata, dict): return False scan_meta = metadata.get('scan_meta') return bool( isinstance(scan_meta, dict) and scan_meta.get('command_timed_out') or metadata.get('error_class') == 'timeout' ) def repair_docker_timeout_attempts(db, max_attempts=3, max_rows=1000): max_attempts = max(1, int(max_attempts or 3)) max_rows = max(1, int(max_rows or 1000)) now = utc_now_iso() report = {'examined': 0, 'repaired': 0, 'terminal': 0, 'unchanged': 0} candidates = db.conn.execute( '''SELECT id, attempts FROM target_queue WHERE source = 'dockerhub' AND platform = 'docker' AND status = 'deferred' AND normalized_target LIKE '%@sha256:%' AND lease_owner IS NULL AND lease_token IS NULL AND current_result_reservation_id IS NULL AND claim_event_id IS NULL AND resolver_token IS NULL ORDER BY id LIMIT ?''', (max_rows,), ).fetchall() try: for candidate in candidates: report['examined'] += 1 rows = db.conn.execute( '''SELECT s.id, c.metadata_json FROM target_scans s LEFT JOIN scan_result_compat c ON c.target_scan_id = s.id WHERE s.queue_id = ? ORDER BY s.id''', (candidate['id'],), ).fetchall() if not rows or not _metadata_is_docker_command_timeout(rows[-1]['metadata_json']): report['unchanged'] += 1 continue timeout_attempts = sum( 1 for row in rows if _metadata_is_docker_command_timeout(row['metadata_json']) ) repaired_attempts = min( max_attempts, max(int(candidate['attempts'] or 0), timeout_attempts), ) if repaired_attempts == int(candidate['attempts'] or 0): report['unchanged'] += 1 continue terminal = repaired_attempts >= max_attempts cursor = db.conn.execute( '''UPDATE target_queue SET attempts = ?, status = CASE WHEN ? != 0 THEN 'failed' ELSE status END, available_after = CASE WHEN ? != 0 THEN NULL ELSE available_after END, completed_at = CASE WHEN ? != 0 THEN COALESCE(completed_at, ?) ELSE completed_at END, last_error = CASE WHEN ? != 0 THEN 'Docker timeout attempts exhausted by guarded repair' ELSE last_error END, updated_at = ? WHERE id = ? AND source = 'dockerhub' AND platform = 'docker' AND status = 'deferred' AND lease_owner IS NULL AND lease_token IS NULL AND current_result_reservation_id IS NULL AND claim_event_id IS NULL AND resolver_token IS NULL''', ( repaired_attempts, 1 if terminal else 0, 1 if terminal else 0, 1 if terminal else 0, now, 1 if terminal else 0, now, candidate['id'], ), ) if cursor.rowcount != 1: raise RuntimeError('Docker timeout repair lost its queue fence') report['repaired'] += 1 report['terminal'] += 1 if terminal else 0 db.conn.commit() return report except Exception: db.conn.rollback() raise def run_target_queue_policy_action(config_path, config, args): cold_action = bool(args.cold_stale_backlog) reactivate_action = bool(args.reactivate_cold_backlog) if cold_action == reactivate_action: raise SystemExit('Select exactly one target queue policy action') if bool(args.dry_run) == bool(args.apply): raise SystemExit('Target queue policy action requires exactly one of --dry-run or --apply') if not args.sources_stopped: raise SystemExit('Target queue policy action requires --sources-stopped') if not args.policy_manifest: raise SystemExit('Target queue policy action requires --policy-manifest') policy_source = str(args.policy_source or '').strip() policy_platform = str(args.policy_platform or '').strip() policy_query = str(args.policy_query or '') if not policy_source or not policy_platform: raise SystemExit('Target queue policy action requires --policy-source and --policy-platform') if cold_action and policy_query: raise SystemExit('--policy-query is only valid for cold-backlog reactivation') if reactivate_action and not policy_query: raise SystemExit('Cold-backlog reactivation requires --policy-query') if any(( args.sqlite, args.initialize_base, args.todo, args.source, args.platform, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.recover_stale_result_pipeline, args.repair_docker_timeout_attempts, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, args.backfill_normalized_results, args.import_legacy_spool, args.review_pipeline_quarantine, args.rebuild_jsonl_output, args.until_complete, )): raise SystemExit('Target queue policy action cannot be combined with another migration action') global_config = config.get('global') or {} load_postgres_environment(config_path, config) requested_db_url = ( args.database_url or database_url_from_env() or global_config.get('database_url') ) if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for target queue policy review') preflight_lifecycle_paths(config_path, config, authority_profile='server') manifest_path = os.path.abspath(args.policy_manifest) with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('target queue policy PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) db.set_application_name('truf-offline-target-queue-policy') config_sha256 = sha256_file(config_path) with postgres_migration_guard(db): db.require_runtime_safety_schema() if args.dry_run: if cold_action: manifest = plan_stale_target_queue_cold( db, config, source=policy_source, platform=policy_platform, config_sha256=config_sha256, max_rows=args.max_rows, ) action = 'cold' else: manifest = plan_cold_target_queue_reactivation( db, config, source=policy_source, platform=policy_platform, query=policy_query, config_sha256=config_sha256, max_rows=args.max_rows, ) action = 'reactivate' atomic_write_private_json( manifest_path, manifest, max_bytes=TARGET_QUEUE_POLICY_MANIFEST_MAX_BYTES, ) require_private_file(manifest_path) report = { 'action': action, 'planned': len(manifest['entries']), 'config_sha256': manifest['config_sha256'], 'policy_sha256': manifest['policy_sha256'], 'selection_sha256': manifest['selection_sha256'], 'manifest_sha256': _canonical_json_sha256(manifest), } else: expected_type = ( 'truf-target-queue-cold-review-v1' if cold_action else 'truf-target-queue-reactivation-review-v1' ) manifest, manifest_sha256 = load_target_queue_policy_manifest( manifest_path, expected_type, max_rows=args.max_rows, ) if ( manifest['source'] != policy_source or manifest['platform'] != policy_platform or ( reactivate_action and manifest['query'] != policy_query ) ): raise RuntimeError('target queue policy manifest scope does not match CLI scope') if cold_action: report = apply_stale_target_queue_cold( db, config, manifest, manifest_sha256, config_sha256=config_sha256, max_rows=args.max_rows, ) else: report = apply_cold_target_queue_reactivation( db, config, manifest, manifest_sha256, config_sha256=config_sha256, max_rows=args.max_rows, ) report = { 'action': 'cold' if cold_action else 'reactivate', **report, } print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 finally: db.close() def parse_args(): parser = argparse.ArgumentParser(description='Offline/idempotent runtime safety schema migration and bounded todo reconciliation.') parser.add_argument('--config', default=os.path.join(os.path.dirname(__file__), 'config.yaml')) parser.add_argument('--database-url') parser.add_argument('--sqlite') parser.add_argument('--initialize-base', action='store_true', help='Create the fresh-install base schema before applying additive migration') parser.add_argument('--todo', help='Optional todo file to reconcile; checked files are intentionally unsupported') parser.add_argument('--source') parser.add_argument('--platform') parser.add_argument('--query', default='offline-reconciliation') parser.add_argument('--max-rows', type=int, default=1000) parser.add_argument('--max-bytes', type=int, default=4 * 1024 * 1024) parser.add_argument('--max-seconds', type=float, default=5.0) parser.add_argument('--until-complete', action='store_true', help='Run bounded batches until EOF or --max-batches') parser.add_argument('--max-batches', type=int, default=100) parser.add_argument('--resolve-issue', action='append', type=int, default=[], help='Mark a reviewed reconciliation issue resolved') parser.add_argument('--harden-runtime', action='store_true', help='Offline recursive hardening/verification of sensitive runtime trees') parser.add_argument('--reconcile-jsonl-projections', action='store_true', help='Offline bounded/resumable publication-ledger initialization') parser.add_argument('--resolve-jsonl-issue', action='append', default=[], help='Approve exact reviewed FILE:OFFSET:SHA256 corrupt record metadata') parser.add_argument('--resolve-jsonl-issue-manifest', help='Approve a private exact reviewed issue manifest') parser.add_argument('--resolve-jsonl-conflict-manifest', help='Approve a private exact historical payload-variant manifest') parser.add_argument('--recover-dead-scan-slots', action='store_true', help='Explicitly remove only exact-identity dead local scan slots') parser.add_argument('--recover-stale-result-pipeline', action='store_true', help='Explicitly reconcile fenced stale result reservations and singleton worker leases') parser.add_argument('--repair-docker-timeout-attempts', action='store_true', help='Repair only unfenced Docker targets whose latest durable result is a command timeout') parser.add_argument('--drain-legacy-outbox', action='store_true', help='Explicit bounded offline projection of legacy scan_publication_outbox rows') parser.add_argument('--import-legacy-outbox-to-projection', action='store_true', help='Transfer bounded legacy outbox references into PostgreSQL projection jobs') parser.add_argument('--backfill-normalized-results', action='store_true', help='Boundedly convert legacy raw PostgreSQL result blobs to normalized-v2 compatibility rows') parser.add_argument('--import-legacy-spool', action='store_true', help='Import a bounded batch from the retired durable result spool into PostgreSQL') parser.add_argument('--review-pipeline-quarantine', help='Apply an exact private audited quarantine review manifest') parser.add_argument('--rebuild-jsonl-output', help='Build a full resumable PostgreSQL-derived compatibility projection in this dedicated directory') parser.add_argument('--cold-stale-backlog', action='store_true', help='Plan or apply exact policy-stale target queue cold transitions') parser.add_argument('--reactivate-cold-backlog', action='store_true', help='Plan or apply exact reviewed cold target queue reactivation') parser.add_argument('--policy-manifest', help='Private exact target queue policy review manifest path') parser.add_argument('--policy-source', help='Exact source scope for target queue policy review') parser.add_argument('--policy-platform', help='Exact platform scope for target queue policy review') parser.add_argument('--policy-query', help='Exact query scope for cold target queue reactivation') parser.add_argument('--dry-run', action='store_true') parser.add_argument('--apply', action='store_true') parser.add_argument('--sources-stopped', action='store_true') return parser.parse_args() def main(): args = parse_args() for name in ( 'drain_legacy_outbox', 'import_legacy_outbox_to_projection', 'backfill_normalized_results', 'import_legacy_spool', 'review_pipeline_quarantine', 'rebuild_jsonl_output', 'recover_stale_result_pipeline', ): if not hasattr(args, name): setattr(args, name, False) config_path = os.path.abspath(args.config) config = load_config(config_path) if args.cold_stale_backlog or args.reactivate_cold_backlog: return run_target_queue_policy_action(config_path, config, args) if args.dry_run: raise SystemExit('--dry-run is only supported for a target queue policy action') if not args.apply or not args.sources_stopped: raise SystemExit('Refusing migration without both --apply and --sources-stopped') global_config = config.get('global', {}) if args.recover_stale_result_pipeline: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.repair_docker_timeout_attempts, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, args.backfill_normalized_results, args.import_legacy_spool, args.review_pipeline_quarantine, args.rebuild_jsonl_output, args.until_complete, )): raise SystemExit('--recover-stale-result-pipeline is a separate PostgreSQL recovery action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for result-pipeline recovery') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('result-pipeline recovery PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) db.set_application_name('truf-offline-result-recovery') with postgres_migration_guard(db): report = recover_stale_result_pipeline( db, config, max_rows=args.max_rows, max_seconds=args.max_seconds, ) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['quiescent'] else 2 except BaseException as exc: raise SystemExit( f'offline result-pipeline recovery failed: {type(exc).__name__}' ) from None finally: db.close() if args.rebuild_jsonl_output: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, args.backfill_normalized_results, args.import_legacy_spool, args.review_pipeline_quarantine, )): raise SystemExit('--rebuild-jsonl-output is a separate bounded compatibility action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for full JSONL rebuild') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('full JSONL rebuild PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): report = rebuild_jsonl_projections( db, args.rebuild_jsonl_output, args.max_rows, args.max_bytes, args.max_seconds, ) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['completed'] else 2 finally: db.close() if args.review_pipeline_quarantine: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, args.backfill_normalized_results, args.import_legacy_spool, )): raise SystemExit('--review-pipeline-quarantine is a separate audited action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for quarantine review') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('quarantine review PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): db.require_runtime_safety_schema() report = review_pipeline_quarantine_manifest( db, os.path.abspath(args.review_pipeline_quarantine), args.max_rows, config, ) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 finally: db.close() if args.import_legacy_spool: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, args.backfill_normalized_results, )): raise SystemExit('--import-legacy-spool is a separate bounded compatibility action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for legacy spool import') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('legacy spool import PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): db.require_runtime_safety_schema() db.revoke_final_cutover() report = import_legacy_result_spool(db, config, args.max_rows) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['remaining'] == 0 else 2 finally: db.close() if args.backfill_normalized_results: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.drain_legacy_outbox, args.import_legacy_outbox_to_projection, )): raise SystemExit('--backfill-normalized-results is a separate bounded compatibility action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for normalized result backfill') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('normalized result backfill PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): db.require_runtime_safety_schema() db.revoke_final_cutover() batch_limit = max(1, int(args.max_batches)) if args.until_complete else 1 report = None for batch_number in range(1, batch_limit + 1): report = backfill_normalized_raw_results( db, args.max_rows, args.max_bytes, args.max_seconds, ) print(json.dumps( {'batch': batch_number, **report}, ensure_ascii=True, sort_keys=True, ), flush=True) if report['remaining'] == 0 or report['processed'] == 0: break return 0 if report['remaining'] == 0 else 2 finally: db.close() if args.import_legacy_outbox_to_projection: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, args.drain_legacy_outbox, )): raise SystemExit('--import-legacy-outbox-to-projection is a separate bounded compatibility action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for legacy outbox transfer') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('legacy outbox transfer PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): db.require_runtime_safety_schema() db.revoke_final_cutover() report = import_legacy_outbox_to_projection(db, args.max_rows) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['remaining'] == 0 else 2 finally: db.close() if args.drain_legacy_outbox: if any(( args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.recover_dead_scan_slots, )): raise SystemExit('--drain-legacy-outbox is a separate bounded compatibility action') load_postgres_environment(config_path, config) requested_db_url = args.database_url or database_url_from_env() or global_config.get('database_url') if not is_postgres_url(requested_db_url): raise SystemExit('A canonical PostgreSQL DSN is required for legacy outbox drain') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, endpoint_dsn=requested_db_url): require_local_sources_stopped(config) identity = verify_cluster_identity(config) db_url = canonical_postgres_url( requested_db_url, identity['database'], identity['user'], identity['port'], ) db = ScannerDB(db_url=db_url, initialize=False) try: if not db.enabled: raise RuntimeError('legacy outbox PostgreSQL connection is unavailable') _online_postgres_identity(db, db_url, config, identity=identity) with postgres_migration_guard(db): db.require_runtime_safety_schema() db.revoke_final_cutover() report = drain_legacy_scan_outbox(db, config, args.max_rows) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['remaining'] == 0 else 2 finally: db.close() if args.recover_dead_scan_slots: if any(( args.database_url, args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.harden_runtime, args.reconcile_jsonl_projections, args.resolve_jsonl_issue, args.resolve_jsonl_issue_manifest, args.resolve_jsonl_conflict_manifest, args.until_complete, )): raise SystemExit('--recover-dead-scan-slots is a separate local recovery action') load_postgres_environment(config_path, config) requested_db_url = global_config.get('database_url') or database_url_from_env() if not is_postgres_url(requested_db_url): raise SystemExit('A caller-selected canonical PostgreSQL DSN is required for maintenance authority') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock(config, create_parent=True, endpoint_dsn=requested_db_url): require_local_sources_stopped(config, inspect_scan_slots=False) report = recover_dead_scan_slots(config) print(json.dumps(report, ensure_ascii=True, sort_keys=True), flush=True) return 0 if report['remaining'] == 0 else 2 if ( args.resolve_jsonl_issue or args.resolve_jsonl_issue_manifest or args.resolve_jsonl_conflict_manifest ) and not args.reconcile_jsonl_projections: raise SystemExit('JSONL review requires --reconcile-jsonl-projections') if args.harden_runtime and args.reconcile_jsonl_projections: raise SystemExit('--harden-runtime and --reconcile-jsonl-projections are separate actions') if args.harden_runtime: if any((args.database_url, args.sqlite, args.initialize_base, args.todo, args.resolve_issue, args.until_complete)): raise SystemExit('--harden-runtime is a separate filesystem-only action and cannot be combined with migration/reconciliation options') load_postgres_environment(config_path, config) requested_db_url = global_config.get('database_url') or database_url_from_env() if not is_postgres_url(requested_db_url): raise SystemExit('A caller-selected canonical PostgreSQL DSN is required for maintenance authority') with ClusterAuthorityLock(config, create_parent=True, endpoint_dsn=requested_db_url): require_runtime_hardening_stopped(config) hardened = harden_runtime_paths( config, find_postgres_environment_path(config_path, config), config_path=config_path, ) print(f'Runtime hardening verification: OK ({hardened} entries)') return 0 if args.reconcile_jsonl_projections: if any((args.database_url, args.sqlite, args.initialize_base, args.todo, args.resolve_issue)): raise SystemExit('--reconcile-jsonl-projections cannot be combined with database migration/reconciliation options') load_postgres_environment(config_path, config) requested_db_url = global_config.get('database_url') or database_url_from_env() if not is_postgres_url(requested_db_url): raise SystemExit('A caller-selected canonical PostgreSQL DSN is required for maintenance authority') with ClusterAuthorityLock(config, create_parent=True, endpoint_dsn=requested_db_url): require_runtime_hardening_stopped(config) return reconcile_jsonl_projection_ledgers(config, args) load_postgres_environment(config_path, config) requested_db_url = ( args.database_url or database_url_from_env() or global_config.get('database_url') ) if not args.sqlite and not is_postgres_url(requested_db_url): raise SystemExit('A PostgreSQL DSN is required unless --sqlite is explicit; refusing scanner_active.db fallback') if args.sqlite and not requested_db_url: raise SystemExit('A caller-selected canonical PostgreSQL DSN is required for maintenance authority') preflight_lifecycle_paths(config_path, config, authority_profile='server') with ClusterAuthorityLock( config, endpoint_dsn=requested_db_url, ): db_url = requested_db_url or database_url_from_env() if not args.sqlite and not is_postgres_url(db_url): raise SystemExit('A PostgreSQL DSN is required unless --sqlite is explicit; refusing scanner_active.db fallback') if args.sqlite and args.database_url: raise SystemExit('--sqlite and --database-url are mutually exclusive') if args.todo and (not args.source or not args.platform): raise SystemExit('--todo requires --source and --platform') if args.until_complete and not args.todo: raise SystemExit('--until-complete requires --todo') require_local_sources_stopped(config) db = None exit_code = 0 try: identity = None if not args.sqlite: try: identity = verify_cluster_identity(config) db_url = canonical_postgres_url( db_url, identity['database'], identity['user'], identity['port'], ) except (DatabaseUrlError, OSError, ValueError) as exc: raise RuntimeError(f'PostgreSQL migration authority validation failed: {exc}') from exc db = ScannerDB(db_path=args.sqlite, db_url='' if args.sqlite else db_url, initialize=False) if not db.enabled: raise RuntimeError(f'Unable to open migration database: {db.db_display}') guard = contextlib.nullcontext() if db.conn.is_postgres: _online_postgres_identity(db, db_url, config, identity=identity) guard = postgres_migration_guard(db) with guard: require_local_sources_stopped(config) if db.conn.is_postgres and db.conn.table_exists('runtime_final_cutover'): db.revoke_final_cutover() migrate_runtime_safety_schema(db, initialize_base=args.initialize_base) cursor_report = initialize_projection_cursors_from_existing_files(db, config) legacy_evidence = require_legacy_cutover_clear(db, config) db.require_runtime_safety_schema() if db.conn.is_postgres: migrations = db.conn.execute( 'SELECT version, code_sha256 FROM runtime_schema_migrations ORDER BY version' ).fetchall() db.conn.commit() marker = db.record_final_cutover({ 'legacy': legacy_evidence, 'projection_cursors': cursor_report, 'schema_migrations': [dict(row) for row in migrations], }) db.require_final_cutover() print( f'Final PostgreSQL cutover marker: {marker["evidence_sha256"]}', flush=True, ) print('Runtime safety schema migration: OK') if args.repair_docker_timeout_attempts: report = repair_docker_timeout_attempts( db, max_attempts=int(global_config.get('target_retry_max_attempts', 3) or 3), max_rows=args.max_rows, ) print('Docker timeout attempt repair: ' + json.dumps(report, sort_keys=True)) if args.resolve_issue: resolve_reconciliation_issues(db, args.resolve_issue) if args.todo: batch_limit = max(1, int(args.max_batches)) if args.until_complete else 1 report = None for _ in range(batch_limit): previous_offset = report.get('byte_offset') if report else None report = reconcile_todo_file( db, args.todo, args.source, args.platform, args.query, args.max_rows, args.max_bytes, args.max_seconds, ) print(json.dumps(report, indent=2, sort_keys=True)) if report['at_eof']: break if previous_offset is not None and report['byte_offset'] <= previous_offset: raise RuntimeError('bounded reconciliation made no forward progress') if not report or not report['at_eof'] or int(report.get('unresolved_issues') or 0) != 0: exit_code = 2 finally: if db is not None: db.close() return exit_code if __name__ == '__main__': raise SystemExit(main())