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

3505 lines
168 KiB
Python

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'<sha256:{hashlib.sha256(payload).hexdigest()[:24]} bytes:{len(payload)}>'
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 '<unnamed>')[: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())