3505 lines
168 KiB
Python
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())
|