import sys import os if __name__ == '__main__': if sys.platform != 'linux' or not os.path.isfile('/.dockerenv') or os.path.abspath(__file__) != '/opt/truf/app/supervisor.py': raise SystemExit('Docker development copy: runtime control is disabled outside the prepared container. See DOCKER_MIGRATION.md.') import runpy runpy.run_path('/opt/truf/app/container_runtime.py')['require_container']() sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('supervisor could not disable bytecode writes') import argparse import json import queue import re import selectors import secrets import shlex import shutil import signal import socket import socketserver import sqlite3 import stat import subprocess import threading import time import urllib.request from concurrent.futures import ThreadPoolExecutor from contextlib import redirect_stdout from datetime import datetime, timezone from io import StringIO from urllib.parse import quote if os.name == 'nt': import ctypes from ctypes import wintypes _SUPERVISOR_KERNEL32 = ctypes.WinDLL('kernel32', use_last_error=True) _SUPERVISOR_OPEN_PROCESS = _SUPERVISOR_KERNEL32.OpenProcess _SUPERVISOR_OPEN_PROCESS.argtypes = [wintypes.DWORD, wintypes.BOOL, wintypes.DWORD] _SUPERVISOR_OPEN_PROCESS.restype = wintypes.HANDLE _SUPERVISOR_GET_EXIT_CODE_PROCESS = _SUPERVISOR_KERNEL32.GetExitCodeProcess _SUPERVISOR_GET_EXIT_CODE_PROCESS.argtypes = [ wintypes.HANDLE, ctypes.POINTER(wintypes.DWORD), ] _SUPERVISOR_GET_EXIT_CODE_PROCESS.restype = wintypes.BOOL _SUPERVISOR_CLOSE_HANDLE = _SUPERVISOR_KERNEL32.CloseHandle _SUPERVISOR_CLOSE_HANDLE.argtypes = [wintypes.HANDLE] _SUPERVISOR_CLOSE_HANDLE.restype = wintypes.BOOL else: _SUPERVISOR_KERNEL32 = None _SUPERVISOR_OPEN_PROCESS = None _SUPERVISOR_GET_EXIT_CODE_PROCESS = None _SUPERVISOR_CLOSE_HANDLE = None def _preimport_runtime_launch_requested(arguments): flags = {str(value).split('=', 1)[0] for value in arguments if str(value).startswith('--')} if flags.intersection({'--background', '--background-child'}): return True return not flags.intersection({'--dry-run', '--stop-background', '--background-status', '--attach', '--cmd'}) def _preimport_is_reparse_point(path): details = os.lstat(path) if stat.S_ISLNK(details.st_mode): return True attributes = getattr(details, 'st_file_attributes', 0) reparse_attribute = getattr(stat, 'FILE_ATTRIBUTE_REPARSE_POINT', 0) return bool(attributes & reparse_attribute) or getattr(os.path, 'isjunction', lambda _path: False)(path) def _preimport_reject_cached_bytecode(app_dir): def raise_walk_error(exc): raise SystemExit(f'Unable to inspect application root: {exc}') from exc try: root_details = os.lstat(app_dir) except OSError as exc: raise SystemExit(f'Application root is unavailable: {app_dir}') from exc if _preimport_is_reparse_point(app_dir): raise SystemExit(f'Application root reparse point is forbidden: {app_dir}') if not stat.S_ISDIR(root_details.st_mode): raise SystemExit(f'Application root is not a directory: {app_dir}') canonical_root = os.path.normcase(os.path.realpath(os.path.abspath(app_dir))) for current, directories, files in os.walk(app_dir, followlinks=False, onerror=raise_walk_error): for name in directories: candidate = os.path.join(current, name) if _preimport_is_reparse_point(candidate): if name.lower() == '__pycache__': raise SystemExit(f'Application bytecode cache link is forbidden: {candidate}') raise SystemExit(f'Application directory reparse point is forbidden: {candidate}') relative = os.path.relpath(current, app_dir) in_cache = any(part.lower() == '__pycache__' for part in relative.split(os.sep)) for name in files: candidate = os.path.join(current, name) if _preimport_is_reparse_point(candidate): raise SystemExit(f'Application file reparse point is forbidden: {candidate}') if name.lower().endswith(('.py', '.pyw', '.pyc', '.pyd')): path = os.path.normcase(os.path.realpath(os.path.abspath(candidate))) try: contained = os.path.commonpath((canonical_root, path)) == canonical_root except ValueError: contained = False if not contained: raise SystemExit(f'Application Python authority escapes its root: {candidate}') if in_cache and name.lower().endswith('.pyc'): raise SystemExit( f'Application __pycache__ bytecode is forbidden: ' f'{os.path.relpath(candidate, app_dir)}' ) if __name__ == '__main__' and _preimport_runtime_launch_requested(sys.argv[1:]): if '--with-postgres' not in sys.argv[1:]: raise SystemExit('Supervisor unmanaged PostgreSQL mutation is retired; use --with-postgres.') if not (sys.flags.isolated and sys.flags.no_site and sys.flags.dont_write_bytecode): raise SystemExit('Mutating supervisor runtime requires python -I -S -B via runtime_bootstrap.py.') if os.getenv('TRUF_RUNTIME_BOOTSTRAP') != '1': raise SystemExit('Mutating supervisor runtime requires the canonical runtime bootstrap.') _preimport_reject_cached_bytecode(os.path.dirname(os.path.abspath(__file__))) from paths import apply_path_config, resolve_optional_path from docker_depth_experiment import validate_docker_depth_config from db_backend import connect_postgres from owned_process import OwnedProcess from process_identity import exact_process_identity_state, open_process from postgres_runtime import PostgresState, canonical_database_url, controller_from_config from lifecycle_authority import ( DISCOVERY_PRODUCER_ROLE, DISCOVERY_PRODUCER_SOURCES, PHASE_ACTIVE, PHASE_ACTIVATING, PHASE_FAILED_HOLD, PHASE_STOPPING, LifecycleAuthorityError, build_code_manifest, code_manifest_sha256, dsn_sha256, strip_supervisor_credentials, supervised_child_environment, verify_code_manifest, ) from runtime_security import ( ClusterAuthorityLock, durable_unlink, harden_private_file, preflight_lifecycle_paths, private_directory_ready, private_file_ready, read_private_json, reject_reparse_components, require_private_directory, sha256_file, ) from supervisor_instance import ( CONTROL_SCHEMA, InstanceMetadataError, InstanceLockError, SupervisorInstanceLock, authenticate_request, build_instance_metadata, load_instance_metadata, load_shutdown_receipt, is_loopback_host, remove_instance_if_matches, remove_shutdown_receipt, update_instance_activation, verify_instance_process, write_shutdown_receipt, write_instance_metadata, ) SOURCE_ALIASES = {'docker': 'dockerhub'} KEYCHECK_SERVICE_NAMES = { 'anthropic', 'aws', 'azure', 'deepseek', 'dockerhub', 'gcp', 'gemini', 'github', 'gitlab', 'groq', 'huggingface', 'kimi', 'openai', 'openrouter', 'provider_resolver', 'qwen', 'replicate', 'xai', 'zai', } RECHECK_TYPE_FLAGS = { 'network': '--retry-network', 'net': '--retry-network', 'limited': '--retry-limited', 'limit': '--retry-limited', 'ratelimited': '--retry-limited', 'rate-limited': '--retry-limited', 'rate-limit': '--retry-limited', 'aliveratelimited': '--retry-limited', 'alive-rate-limited': '--retry-limited', 'alive-limited': '--retry-limited', 'validratelimited': '--retry-limited', 'valid-rate-limited': '--retry-limited', 'valid-limited': '--retry-limited', 'unknown': '--retry-unknown', 'restricted': '--retry-restricted', 'restriction': '--retry-restricted', 'nobalance': '--retry-no-balance', 'no-balance': '--retry-no-balance', 'noquota': '--retry-no-balance', 'no-quota': '--retry-no-balance', 'quota': '--retry-no-balance', 'valid': '--retry-valid', 'alive': '--retry-valid', 'legacy-vertex': '--import-legacy-gcp-vertex', 'legacyvertex': '--import-legacy-gcp-vertex', 'all': '--recheck-all', 'everything': '--recheck-all', } RECHECK_BOOL_OPTIONS = { 'no-resource-probe': '--no-resource-probe', 'no-summary': '--no-summary', 'summary-only': '--summary-only', } RECHECK_VALUE_OPTIONS = { 'max-keys': '--max-keys', 'max': '--max-keys', 'limit': '--max-keys', 'input': '--input', 'proxy-file': '--proxy-file', 'proxy': '--proxy-file', } DEFAULT_REFRESH_SEC = 5 DEFAULT_RESTART_DELAY = 30 SOURCE_INFRASTRUCTURE_HOLD_EXIT = 75 DEFAULT_MAX_RESTART_DELAY = 600 DEFAULT_RESTART_RESET_AFTER = 300 DEFAULT_TEMP_CLEANUP_INTERVAL = 900 ANSI_ALT_SCREEN = '\x1b[?1049h' ANSI_MAIN_SCREEN = '\x1b[?1049l' ANSI_HOME = '\x1b[H' ANSI_CLEAR_SCREEN = '\x1b[2J' DETACHED_PROCESS = 0x00000008 CREATE_NO_WINDOW = 0x08000000 MAX_CONTROL_REQUEST_BYTES = 64 * 1024 MAX_CONTROL_RESPONSE_BYTES = 2 * 1024 * 1024 MAX_CONTROL_PENDING_SOCKETS = 32 MAX_CONTROL_WORKERS = 4 CONTROL_READ_TIMEOUT_SEC = 1.0 MAX_LOG_TAIL_BYTES = 1024 * 1024 MAX_LOG_TAIL_LINES = 5000 RUNTIME_SNAPSHOT_SCHEMA = 2 MAX_MANAGED_SOURCE_DELAY_SECONDS = 365 * 24 * 60 * 60 MANAGED_SOURCE_LIFECYCLE_ACTIONS = ('start', 'stop', 'restart', 'pause', 'resume') MANAGED_SOURCE_SETTING_ACTIONS = ( 'once', 'set-mode', 'set-interval', 'set-restart', 'set-restart-delay', ) DISCOVERY_CYCLE_STATUSES = frozenset({ 'running', 'completed', 'completed_with_retries', 'query_invalid', 'failed', 'source_failed', 'backlog_only', 'paused', 'auth_failed', 'rate_limited', }) DISCOVERY_ERROR_CATEGORIES = frozenset({ 'auth_failed', 'failed', 'paused', 'query_invalid', 'rate_limited', 'source_failed', 'runtime_error', }) RUNTIME_BOOTSTRAP_ENV = 'TRUF_RUNTIME_BOOTSTRAP' RUNTIME_BOOTSTRAP_VALUE = '1' def child_bootstrap_command(kind, arguments, provider_entrypoint=None): command = [ sys.executable, '-I', '-S', '-B', os.path.join(os.path.dirname(os.path.abspath(__file__)), 'child_bootstrap.py'), str(kind), ] if provider_entrypoint: command.append(str(provider_entrypoint).replace('\\', '/')) command.append('--') command.extend(str(value) for value in arguments) return command def load_yaml(path, *, managed_postgres=None, final_cutover=None): try: import yaml except ImportError as e: raise SystemExit('PyYAML is required. Run: python -m pip install -r requirements.txt') from e with open(path, 'r', encoding='utf-8') as f: loaded = yaml.safe_load(f) if loaded is None: loaded = {} validated = validate_docker_depth_config( loaded, managed_postgres=managed_postgres, final_cutover=final_cutover, ) return apply_path_config(validated.config, path) def resolve_path(base_file, value): if not value: return value if os.path.isabs(str(value)): return str(value) return os.path.join(os.path.dirname(os.path.abspath(base_file)), str(value)) def load_postgres_env(config_path=None, layout=None, enforce_canonical=False): candidates = [] root_dir = (layout or {}).get('root_dir') if root_dir: candidates.append(os.path.join(root_dir, '.env.postgres')) if config_path: config_dir = os.path.dirname(os.path.abspath(config_path)) candidates.append(os.path.join(config_dir, '..', '.env.postgres')) candidates.append(os.path.join(config_dir, '.env.postgres')) seen = set() loaded = None for path in candidates: path = os.path.abspath(path) if path in seen: continue seen.add(path) if not os.path.exists(path): continue try: with open(path, 'r', encoding='utf-8') as f: for line in f: text = line.strip() if not text or text.startswith('#') or '=' not in text: continue key, value = text.split('=', 1) key = key.strip() value = value.strip().strip('"').strip("'") if key and value and not os.getenv(key): os.environ[key] = value except OSError: continue loaded = path if os.getenv('SCANNER_DB_URL') or os.getenv('DATABASE_URL'): break password = os.getenv('TRUF_POSTGRES_PASSWORD') if password: user = os.getenv('TRUF_POSTGRES_USER') or 'truf' database = os.getenv('TRUF_POSTGRES_DB') or 'truf' port = os.getenv('TRUF_POSTGRES_PORT') or '5432' os.environ['SCANNER_DB_URL'] = ( f'postgresql://{quote(user, safe="")}:{quote(password, safe="")}@127.0.0.1:{port}/{quote(database, safe="")}' ) break if enforce_canonical: try: url = canonical_database_url() except ValueError as exc: raise SystemExit(f'Managed PostgreSQL database URL failed closed: {exc}') from exc if url: for key in list(os.environ): if key.upper().startswith('PG'): os.environ.pop(key, None) os.environ['SCANNER_DB_URL'] = url os.environ['DATABASE_URL'] = url os.environ['SCANNER_DASHBOARD_DB_URL'] = url os.environ['TRUF_MANAGED_POSTGRES_DSN'] = url return loaded def normalize_source_name(source): source = str(source or '').strip().lower() return SOURCE_ALIASES.get(source, source) def parse_source_list(value): if not value: return None if isinstance(value, str): parts = value.split(',') else: parts = value return [normalize_source_name(item) for item in parts if str(item).strip()] def normalize_recheck_token(value): return str(value or '').strip().lower().lstrip('-').replace('_', '-') def append_unique(items, value): if value not in items: items.append(value) def parse_recheck_command(parts): if len(parts) < 2: return None, 'Usage: recheck [network|ratelimited|unknown|restricted|nobalance|valid|legacy-vertex|all] [--force] [--max-keys N] [--input PATH] [--proxy-file PATH]' first = normalize_source_name(parts[1]) first_norm = normalize_recheck_token(first) if first_norm in RECHECK_TYPE_FLAGS and first_norm != 'all': service = 'all' tokens = parts[1:] else: service = 'all' if first_norm == 'all' else first tokens = parts[2:] if service != 'all' and service not in KEYCHECK_SERVICE_NAMES: return None, f'Unknown keycheck service: {service}. Available: all, {", ".join(sorted(KEYCHECK_SERVICE_NAMES))}' runner_args = [] retry_flag_seen = False force = False index = 0 while index < len(tokens): token = tokens[index] norm = normalize_recheck_token(token) if not norm: index += 1 continue if norm in ('force', 'f'): force = True index += 1 continue if norm in RECHECK_TYPE_FLAGS: append_unique(runner_args, RECHECK_TYPE_FLAGS[norm]) retry_flag_seen = True index += 1 continue if norm in RECHECK_BOOL_OPTIONS: append_unique(runner_args, RECHECK_BOOL_OPTIONS[norm]) index += 1 continue matched_value_option = None matched_value = None for option_name, flag in RECHECK_VALUE_OPTIONS.items(): prefix = option_name + '=' if norm.startswith(prefix): matched_value_option = flag matched_value = token.split('=', 1)[1] break if matched_value_option: if matched_value == '': return None, f'{token} requires a value' runner_args.extend([matched_value_option, matched_value]) index += 1 continue if norm in RECHECK_VALUE_OPTIONS: if index + 1 >= len(tokens): return None, f'{token} requires a value' runner_args.extend([RECHECK_VALUE_OPTIONS[norm], tokens[index + 1]]) index += 2 continue return None, f'Unknown recheck option: {token}' if not retry_flag_seen: append_unique(runner_args, '--recheck-all') return {'service': service, 'runner_args': runner_args, 'force': force}, None def bool_value(value, default=False): if value is None: return default if isinstance(value, bool): return value return str(value).strip().lower() in ('1', 'true', 'yes', 'on') def list_value(value): if not value: return [] if isinstance(value, str): return [item for item in value.split() if item] return [str(item) for item in value] PROXY_ENV_KEYS = ( 'HTTP_PROXY', 'HTTPS_PROXY', 'ALL_PROXY', 'http_proxy', 'https_proxy', 'all_proxy', ) def safe_source_filename(source): return ''.join(ch if ch.isalnum() or ch in ('-', '_') else '_' for ch in source) def int_value(value, default): try: return int(value) except (TypeError, ValueError): return default def next_rotated_path(path): base, ext = os.path.splitext(path) seq = 1 while True: candidate = f'{base}.{seq:06d}{ext or ".log"}' if not os.path.exists(candidate): return candidate seq += 1 def prune_rotated_logs(path, keep): keep = max(0, int_value(keep, 0)) base, ext = os.path.splitext(path) parent = os.path.dirname(path) or '.' prefix = os.path.basename(base) + '.' suffix = ext or '.log' try: rotated = [] for name in os.listdir(parent): if name.startswith(prefix) and name.endswith(suffix): full = os.path.join(parent, name) if os.path.isfile(full): rotated.append(full) rotated.sort(key=lambda item: os.path.getmtime(item), reverse=True) for old in rotated[keep:]: try: os.remove(old) except OSError: pass except OSError: pass def rotate_log_if_needed(path, max_mb=64, keep=5): max_bytes = max(0, int_value(max_mb, 64)) * 1024 * 1024 if max_bytes <= 0 or not path or not os.path.exists(path): return try: if os.path.getsize(path) < max_bytes: return rotated_path = next_rotated_path(path) os.replace(path, rotated_path) prune_rotated_logs(path, keep) except OSError: pass def open_private_append(path): parent = os.path.dirname(os.path.abspath(path)) require_private_directory(parent, create=False) reject_reparse_components(path) flags = os.O_WRONLY | os.O_APPEND | os.O_CREAT if hasattr(os, 'O_BINARY'): flags |= os.O_BINARY if hasattr(os, 'O_NOFOLLOW'): flags |= os.O_NOFOLLOW created = False if os.path.lexists(path): if not private_file_ready(path): raise OSError(f'private log file ACL is not ready; run offline hardening: {path}') descriptor = os.open(path, flags, 0o600) else: try: descriptor = os.open(path, flags | os.O_EXCL, 0o600) created = True except FileExistsError: if not private_file_ready(path): raise OSError(f'private log file ACL is not ready; run offline hardening: {path}') descriptor = os.open(path, flags, 0o600) try: if created: harden_private_file(path) return os.fdopen(descriptor, 'a', encoding='utf-8', buffering=1) except BaseException: os.close(descriptor) raise def open_private_append_binary(path): parent = os.path.dirname(os.path.abspath(path)) require_private_directory(parent, create=False) reject_reparse_components(path) flags = os.O_WRONLY | os.O_APPEND | os.O_CREAT | getattr(os, 'O_BINARY', 0) | getattr(os, 'O_NOFOLLOW', 0) created = False if os.path.lexists(path): if not private_file_ready(path): raise OSError(f'private log file ACL is not ready; run offline hardening: {path}') descriptor = os.open(path, flags, 0o600) else: try: descriptor = os.open(path, flags | os.O_EXCL, 0o600) created = True except FileExistsError: if not private_file_ready(path): raise OSError(f'private log file ACL is not ready; run offline hardening: {path}') descriptor = os.open(path, flags, 0o600) try: if created: harden_private_file(path) return os.fdopen(descriptor, 'ab', buffering=0) except BaseException: os.close(descriptor) raise class BoundedRotatingLogPump: """Drain one child pipe into privately-owned bounded log segments.""" def __init__(self, path, max_bytes, keep): self.path = os.path.abspath(path) self.max_bytes = max(1, int(max_bytes)) self.keep = max(1, int(keep)) self._lock = threading.RLock() self.handle = open_private_append_binary(self.path) self.thread = None self.error = '' def _rotate(self): self.handle.flush() opened = os.fstat(self.handle.fileno()) current = os.stat(self.path, follow_symlinks=False) if (opened.st_dev, opened.st_ino) != (current.st_dev, current.st_ino): raise OSError('active source log path changed while its bounded writer was open') if opened.st_size > self.max_bytes: os.ftruncate(self.handle.fileno(), self.max_bytes) self.handle.flush() os.fsync(self.handle.fileno()) self.handle.close() self.handle = None rotated = next_rotated_path(self.path) os.replace(self.path, rotated) harden_private_file(rotated) prune_rotated_logs(self.path, self.keep) self.handle = open_private_append_binary(self.path) def write(self, data): data = data.encode('utf-8', errors='replace') if isinstance(data, str) else bytes(data or b'') with self._lock: view = memoryview(data) while view: size = os.fstat(self.handle.fileno()).st_size if size >= self.max_bytes: self._rotate() size = 0 portion = view[:max(1, self.max_bytes - size)] self.handle.write(portion) view = view[len(portion):] def flush(self): with self._lock: if self.handle is not None: self.handle.flush() def attach(self, stream): def drain(): try: while True: chunk = stream.read(64 * 1024) if not chunk: break if not self.error: try: self.write(chunk) except Exception as exc: self.error = f'{type(exc).__name__}: {exc}' finally: try: stream.close() except OSError: pass self.close() self.thread = threading.Thread(target=drain, name='source-log-pump', daemon=True) self.thread.start() def join(self, timeout=5): if self.thread is not None and self.thread.ident is not None: self.thread.join(timeout=max(0.0, float(timeout))) if self.thread.is_alive() and not self.error: self.error = 'log pipe did not close after child exit' else: self.close() def close(self): with self._lock: handle = self.handle self.handle = None if handle is not None: try: handle.flush() os.fsync(handle.fileno()) except OSError: pass handle.close() class BoundedRotatingTextWriter: encoding = 'utf-8' def __init__(self, path, max_bytes, keep): self.sink = BoundedRotatingLogPump(path, max_bytes, keep) self.closed = False def write(self, value): if self.closed: return 0 text = str(value or '') self.sink.write(text) return len(text) def flush(self): if not self.closed: self.sink.flush() def isatty(self): return False def close(self): if not self.closed: self.closed = True self.sink.close() def install_bounded_background_output(path, max_bytes, keep): writer = BoundedRotatingTextWriter(path, max_bytes, keep) previous = (sys.stdout, sys.stderr) sys.stdout = writer sys.stderr = writer for stream in dict.fromkeys(previous): try: stream.flush() stream.close() except Exception: pass return writer def append_bounded_log_record(path, max_bytes, keep, text): sink = BoundedRotatingLogPump(path, max_bytes, keep) try: sink.write(str(text).encode('utf-8', errors='replace')) finally: sink.close() def bounded_tail_lines(path, limit=40, max_bytes=MAX_LOG_TAIL_BYTES): if not path or not os.path.exists(path): return [] limit = min(MAX_LOG_TAIL_LINES, max(1, int(limit or 40))) max_bytes = max(1, int(max_bytes)) try: with open(path, 'rb') as handle: handle.seek(0, os.SEEK_END) size = handle.tell() start = max(0, size - max_bytes) handle.seek(start) data = handle.read(max_bytes) if start and data: newline = data.find(b'\n') data = data[newline + 1:] if newline >= 0 else b'' return [line.decode('utf-8', errors='replace') for line in data.splitlines()[-limit:]] except OSError: return [] def console_safe_text(value, stream=None): text = str(value) encoding = getattr(stream or sys.stdout, 'encoding', None) or 'utf-8' try: text.encode(encoding, errors='strict') return text except (LookupError, UnicodeEncodeError): try: return text.encode(encoding, errors='backslashreplace').decode(encoding, errors='replace') except LookupError: return text.encode('ascii', errors='backslashreplace').decode('ascii') def format_duration(seconds): if seconds is None: return '-' seconds = max(0, int(seconds)) hours, rem = divmod(seconds, 3600) minutes, sec = divmod(rem, 60) if hours: return f'{hours}h{minutes:02d}m' if minutes: return f'{minutes}m{sec:02d}s' return f'{sec}s' def format_exit_code(code): if code is None: return '-' if os.name == 'nt' and (code < 0 or code > 255): return f'0x{code & 0xFFFFFFFF:08X}' return str(code) def now_iso(): return datetime.now().isoformat(timespec='seconds') def safe_state_timestamp(value): if type(value) is not str or 'T' not in value or len(value) > 64: return None try: datetime.fromisoformat(value) except ValueError: return None return value def finalize_source_runs(database_url, source, status, reason): if not database_url: return 0, 0 connection = connect_postgres( database_url, connect_timeout_sec=5, statement_timeout_ms=10000, lock_timeout_ms=5000, idle_in_transaction_timeout_ms=10000, tcp_user_timeout_ms=5000, ) timestamp = datetime.now().astimezone().isoformat(timespec='seconds') try: cycles = connection.execute( '''UPDATE source_cycles SET ended_at = ?, status = ?, message = COALESCE(NULLIF(message, ''), ?) WHERE source = ? AND status = 'running' ''', (timestamp, status, reason, source), ) runs = connection.execute( '''UPDATE runs SET ended_at = ?, status = ?, error = COALESCE(NULLIF(error, ''), ?), updated_at = ? WHERE selected_source = ? AND status = 'running' ''', (timestamp, status, reason, timestamp, source), ) connection.commit() return int(cycles.rowcount or 0), int(runs.rowcount or 0) except Exception: connection.rollback() raise finally: connection.close() def terminal_width(default=120): try: return max(80, shutil.get_terminal_size((default, 24)).columns) except OSError: return default def truncate_text(value, width): value = str(value or '') if width <= 0: return '' if len(value) <= width: return value if width <= 1: return value[:width] return value[:width - 1] + '~' def enable_ansi_terminal(): if os.name != 'nt': return True try: import ctypes kernel32 = ctypes.windll.kernel32 handle = kernel32.GetStdHandle(-11) mode = ctypes.c_uint32() if not kernel32.GetConsoleMode(handle, ctypes.byref(mode)): return False return bool(kernel32.SetConsoleMode(handle, mode.value | 0x0004)) except Exception: return False def parse_time(value): if not value: return None try: parsed = datetime.fromisoformat(str(value).replace('Z', '+00:00')) return parsed except ValueError: return None def auth_item_kind(item): if not isinstance(item, dict): return 'ok' if item.get('disabled_until') == 'manual' or item.get('status') == 'dead' or item.get('disabled_reason') == 'auth_invalid': return 'dead' disabled_until = parse_time(item.get('disabled_until')) if disabled_until: now = datetime.now(disabled_until.tzinfo) if disabled_until.tzinfo else datetime.now() if disabled_until > now: return 'limited' return 'ok' def auth_counts_from_status(auth_status): counts = {'ok': 0, 'dead': 0, 'limited': 0, 'rate_limit_errors': 0, 'auth_invalid_errors': 0} for item in (auth_status or {}).values(): kind = auth_item_kind(item) counts[kind] += 1 counts['rate_limit_errors'] += int((item or {}).get('rate_limit_count', 0) or 0) counts['rate_limit_errors'] += int((item or {}).get('secondary_rate_limit_count', 0) or 0) counts['auth_invalid_errors'] += int((item or {}).get('auth_invalid_count', 0) or 0) if (item or {}).get('disabled_reason') == 'auth_invalid' and not (item or {}).get('auth_invalid_count'): counts['auth_invalid_errors'] += int((item or {}).get('failures', 1) or 1) return counts def source_options(source, supervisor_config, source_config): defaults = supervisor_config.get('defaults') or {} per_source = (supervisor_config.get('sources') or {}).get(source, {}) or {} options = dict(defaults) options.update(per_source) if 'enabled' not in options: options['enabled'] = source_config.get('enabled', False) return options def get_enabled_sources(config, selected_sources=None, supervisor_config=None): sources = config.get('sources') or {} selected = parse_source_list(selected_sources) if selected and 'all' in selected: selected = None if selected: selected_sources_only = [source for source in selected if source != 'keychecks'] missing = [source for source in selected_sources_only if source not in sources] if missing: raise SystemExit(f'Source(s) not present in config: {", ".join(missing)}') return selected_sources_only supervisor_sources = (supervisor_config or {}).get('sources') or {} enabled = [] for name, source in sources.items(): supervisor_source = supervisor_sources.get(name) or {} if source.get('enabled', False) or bool_value(supervisor_source.get('enabled'), False): enabled.append(name) return enabled class DependencyGate: def __init__(self, ready=True, database_url=''): self.ready = bool(ready) self.database_url = str(database_url or '') self._dependents = [] self.stop_failures = [] def register(self, dependent): if dependent not in self._dependents: self._dependents.append(dependent) def unregister(self, dependent): if dependent in self._dependents: self._dependents.remove(dependent) def set_ready(self, ready): ready = bool(ready) if ready == self.ready and not (not ready and self.stop_failures): return if ready and self.stop_failures: raise DependencyStopError(self.stop_failures) self.ready = ready failures = [] for dependent in list(self._dependents): try: result = dependent.dependency_available() if ready else dependent.dependency_unavailable() if not ready and result is False: failures.append(getattr(dependent, 'source', dependent.__class__.__name__)) except Exception as exc: failures.append(f'{getattr(dependent, "source", dependent.__class__.__name__)}: {exc}') self.stop_failures = failures if not ready else [] if failures: raise DependencyStopError(failures) return True def force_database_environment(self, env): if not self.database_url: return env for key in list(env): if key.upper().startswith('PG'): env.pop(key, None) env['SCANNER_DB_URL'] = self.database_url env['DATABASE_URL'] = self.database_url env['SCANNER_DASHBOARD_DB_URL'] = self.database_url env['KEYCHECK_DB_URL'] = self.database_url env['TRUF_MANAGED_POSTGRES_DSN'] = self.database_url return env class DependencyStopError(RuntimeError): def __init__(self, failures): self.failures = [str(item) for item in failures] super().__init__('database-dependent child stop failed: ' + ', '.join(self.failures)) class ManagedSource: def __init__( self, source, config_path, project_dir, results_dir, supervisor_config, source_config, global_force_once=False, dependency_gate=None, authority_check=None, start_gate=None, child_environment=None, ): self.source = source self.config_path = os.path.abspath(config_path) self.project_dir = project_dir self.results_dir = results_dir self.options = source_options(source, supervisor_config, source_config) self.once = bool_value(self.options.get('once'), False) or global_force_once self.repeat = bool_value(self.options.get('repeat'), self.once) self.restart = bool_value(self.options.get('restart'), True) self.enabled = bool_value(self.options.get('enabled'), True) self.interval = int(self.options.get('interval', self.options.get('cooldown', supervisor_config.get('interval', 0))) or 0) self.restart_delay = int(self.options.get('restart_delay', supervisor_config.get('restart_delay', DEFAULT_RESTART_DELAY)) or DEFAULT_RESTART_DELAY) self.max_restart_delay = int(self.options.get('max_restart_delay', supervisor_config.get('max_restart_delay', DEFAULT_MAX_RESTART_DELAY)) or DEFAULT_MAX_RESTART_DELAY) self.restart_reset_after = int(self.options.get('restart_reset_after', supervisor_config.get('restart_reset_after', DEFAULT_RESTART_RESET_AFTER)) or DEFAULT_RESTART_RESET_AFTER) self.extra_args = list_value(self.options.get('extra_args')) self.env_overrides = self.options.get('env') if isinstance(self.options.get('env'), dict) else {} self.use_per_source_state = bool_value(supervisor_config.get('per_source_state'), True) self.log_max_mb = int_value(self.options.get('log_max_mb', supervisor_config.get('log_max_mb', 64)), 64) self.log_keep = int_value(self.options.get('log_keep', supervisor_config.get('log_keep', 5)), 5) log_dir = resolve_path(config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) state_dir = resolve_path(config_path, supervisor_config.get('state_dir') or os.path.join(results_dir, 'state')) require_private_directory(log_dir, create=False) require_private_directory(state_dir, create=False) safe_name = safe_source_filename(source) self.log_path = os.path.join(log_dir, f'{safe_name}.log') self.state_path = os.path.join(state_dir, f'runner_state_{safe_name}.json') self.process = None self.payload_identity = None self.startup_cleanup_pending = False self.log_handle = None self.log_pump = None self.started_at = None self.last_exit_code = None self.last_exit_at = None self.restarts = 0 self.restart_streak = 0 self.next_start_at = 0 self.status = 'disabled' if not self.enabled else 'stopped' self.stop_requested = False self.manual_stop = True self.paused = False self.desired_state = 'stopped' self.runtime_blocked = False self._resume_immediately = False self.dependency_gate = dependency_gate self.authority_check = authority_check self.start_gate = start_gate self.child_environment = dict(child_environment or {}) self.manual_only = False self.last_action_error = '' if self.dependency_gate is not None: self.dependency_gate.register(self) def is_running(self): return self.process is not None and self.process.poll() is None def reconfigure(self, results_dir, supervisor_config, source_config, global_force_once=False): self.results_dir = results_dir self.options = source_options(self.source, supervisor_config, source_config) self.once = bool_value(self.options.get('once'), False) or global_force_once self.repeat = bool_value(self.options.get('repeat'), self.once) self.restart = bool_value(self.options.get('restart'), True) self.enabled = bool_value(self.options.get('enabled'), True) self.interval = int(self.options.get('interval', self.options.get('cooldown', supervisor_config.get('interval', 0))) or 0) self.restart_delay = int(self.options.get('restart_delay', supervisor_config.get('restart_delay', DEFAULT_RESTART_DELAY)) or DEFAULT_RESTART_DELAY) self.max_restart_delay = int(self.options.get('max_restart_delay', supervisor_config.get('max_restart_delay', DEFAULT_MAX_RESTART_DELAY)) or DEFAULT_MAX_RESTART_DELAY) self.restart_reset_after = int(self.options.get('restart_reset_after', supervisor_config.get('restart_reset_after', DEFAULT_RESTART_RESET_AFTER)) or DEFAULT_RESTART_RESET_AFTER) self.extra_args = list_value(self.options.get('extra_args')) self.env_overrides = self.options.get('env') if isinstance(self.options.get('env'), dict) else {} self.use_per_source_state = bool_value(supervisor_config.get('per_source_state'), True) self.log_max_mb = int_value(self.options.get('log_max_mb', supervisor_config.get('log_max_mb', 64)), 64) self.log_keep = int_value(self.options.get('log_keep', supervisor_config.get('log_keep', 5)), 5) log_dir = resolve_path(self.config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) state_dir = resolve_path(self.config_path, supervisor_config.get('state_dir') or os.path.join(results_dir, 'state')) require_private_directory(log_dir, create=False) require_private_directory(state_dir, create=False) safe_name = safe_source_filename(self.source) self.log_path = os.path.join(log_dir, f'{safe_name}.log') self.state_path = os.path.join(state_dir, f'runner_state_{safe_name}.json') if not self.enabled: self.status = 'disabled' def build_command(self): arguments = [ '--config', self.config_path, '--source', self.source, ] if self.once: arguments.append('--once') arguments.extend(self.extra_args) return child_bootstrap_command('scanner', arguments) def build_env(self): env = os.environ.copy() use_system_proxy = bool_value(self.options.get('use_system_proxy'), False) if not use_system_proxy: for key in PROXY_ENV_KEYS: env.pop(key, None) env['NO_PROXY'] = '*' env['no_proxy'] = '*' else: env.pop('NO_PROXY', None) env.pop('no_proxy', None) env['PYTHONUNBUFFERED'] = '1' env['PYTHONIOENCODING'] = 'utf-8' env['SCANNER_SOURCE'] = self.source env['SCANNER_SKIP_STARTUP_CLEANUP'] = '1' if self.use_per_source_state: env['RUNNER_STATE_FILE'] = self.state_path for key, value in self.env_overrides.items(): env[str(key)] = str(value) if self.source == 'janitor': for key in list(env): if key.upper().startswith('PG') or key.upper() in { 'TRUF_MANAGED_POSTGRES_DSN', 'SCANNER_DB_URL', 'DATABASE_URL', 'SCANNER_DASHBOARD_DB_URL', 'KEYCHECK_DB_URL', }: env.pop(key, None) elif self.dependency_gate is not None: self.dependency_gate.force_database_environment(env) env.update(self.child_environment) if self.source == 'janitor': env.pop('TRUF_MANAGED_POSTGRES_DSN', None) env.pop('SCANNER_DB_URL', None) env.pop('DATABASE_URL', None) return env def mode_label(self): if self.once and self.repeat: return 'once+repeat' if self.once: return 'once' return 'loop' def finalize_database_runs(self, status, reason): database_url = self.child_environment.get('SCANNER_DB_URL') if not database_url and self.dependency_gate is not None: database_url = self.dependency_gate.database_url return finalize_source_runs(database_url, self.source, status, reason) def schedule_start(self, delay=0): if self.startup_cleanup_pending: return if not self.enabled: return if self.start_gate is not None and not self.start_gate(): self.last_action_error = 'supervisor lifecycle start gate is closed' self.status = 'stopped' return self.desired_state = 'running' self.manual_stop = False self.paused = False self.stop_requested = False self.next_start_at = time.time() + max(0, int(delay or 0)) if self.dependency_gate is not None and not self.dependency_gate.ready: self.runtime_blocked = True self.status = 'blocked' else: self.status = 'waiting' def start(self, force=False, record_intent=True): if self.startup_cleanup_pending: self.status = 'failed' return False self.last_action_error = '' if self.start_gate is not None and not self.start_gate(): self.last_action_error = 'supervisor lifecycle start gate is closed' self.status = 'stopped' return False if self.authority_check is not None and not self.authority_check(): self.last_action_error = 'runtime script/config authority drifted' self.status = 'failed' return False if not self.enabled: return False if record_intent: self.desired_state = 'running' self.manual_stop = False self.paused = False self.stop_requested = False if self.desired_state != 'running': return False if self.dependency_gate is not None and not self.dependency_gate.ready: self.runtime_blocked = True self._resume_immediately = self._resume_immediately or force or self.status != 'waiting' self.status = 'blocked' return False if self.process and self.process.poll() is None: return True now = time.time() if not force and now < self.next_start_at: self.status = 'waiting' return False if force: self.restart_streak = 0 self.last_exit_code = None self.runtime_blocked = False self._resume_immediately = False self.manual_stop = False self.paused = False self.stop_requested = False self.next_start_at = 0 try: log_max_bytes = max(1, self.log_max_mb) * 1024 * 1024 self.log_pump = BoundedRotatingLogPump(self.log_path, log_max_bytes, self.log_keep) self.log_pump.write(f'\n=== supervisor start {now_iso()} source={self.source} once={self.once} ===\n'.encode('utf-8')) self.log_pump.write(('command: ' + ' '.join(command_for_log(self.build_command())) + '\n').encode('utf-8')) if self.use_per_source_state: self.log_pump.write(f'RUNNER_STATE_FILE={self.state_path}\n'.encode('utf-8')) creationflags = (subprocess.CREATE_NEW_PROCESS_GROUP | CREATE_NO_WINDOW) if os.name == 'nt' else 0 self.process = OwnedProcess( self.build_command(), cwd=self.project_dir, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, env=self.build_env(), creationflags=creationflags, ) try: with open_process(self.process.pid) as retained: self.payload_identity = retained.identity except OSError: self.payload_identity = None stream = getattr(self.process, 'stdout', None) if stream is not None: self.log_pump.attach(stream) else: self.log_pump.close() except BaseException as exc: self.startup_cleanup_pending = self.process is not None self.started_at = None self.status = 'failed' self.last_action_error = f'owned process launch failed: {type(exc).__name__}: {exc}' cleanup_error = None if self.process is not None: try: if self.process.poll() is None: self.process.terminate() self.process.wait(timeout=5) if self.process.poll() is not None: self.process = None self.payload_identity = None self.startup_cleanup_pending = False except BaseException as cleanup_exc: cleanup_error = cleanup_exc if self.startup_cleanup_pending: # Only the retained owner can confirm rollback; a failed wait # must not make this source (or its peers) eligible to start. self.desired_state = 'stopped' self.manual_stop = True self.stop_requested = True self.next_start_at = 0 self.last_action_error += '; owned child exit is unconfirmed; shutdown retry required' elif self.log_pump is not None: try: try: if self.log_pump.thread is None and self.log_pump.handle is not None: self.log_pump.write( f'=== supervisor launch failed {now_iso()} source={self.source}: {type(exc).__name__} ===\n'.encode('utf-8') ) finally: self.log_pump.join(timeout=1) except BaseException as cleanup_exc: cleanup_error = cleanup_exc else: self.log_pump = None if not isinstance(exc, Exception): raise if cleanup_error is not None and not isinstance(cleanup_error, Exception): raise cleanup_error return False self.started_at = now self.status = 'running' return True def poll(self): if self.startup_cleanup_pending: return if not self.enabled: return if self.dependency_gate is not None and not self.dependency_gate.ready: if self.desired_state == 'running' and not self.runtime_blocked: self.dependency_unavailable() return if not self.process: if self.status == 'waiting' and self.desired_state == 'running' and not self.stop_requested: self.start(record_intent=False) return code = self.process.poll() if code is None: if self.log_pump is not None and self.log_pump.error: self.last_action_error = f'bounded source log writer failed: {self.log_pump.error}' try: self.process.terminate() self.process.wait(timeout=5) except Exception: return code = self.process.poll() else: self.status = 'running' if self.started_at and time.time() - self.started_at >= self.restart_reset_after: self.restart_streak = 0 self.last_exit_code = None return self._handle_process_exit(code) def _handle_process_exit(self, code): run_duration = time.time() - self.started_at if self.started_at else 0 log_error = False if self.log_pump is not None: self.log_pump.join(timeout=5) if self.log_pump.error: log_error = True self.last_action_error = f'bounded source log writer failed: {self.log_pump.error}' if code == 0: code = -1 self.log_pump = None if not log_error: append_bounded_log_record( self.log_path, max(1, self.log_max_mb) * 1024 * 1024, self.log_keep, f'=== supervisor exit {now_iso()} source={self.source} code={code} ===\n', ) self.last_exit_code = code self.last_exit_at = time.time() self.process = None self.payload_identity = None self.started_at = None if self.stop_requested or self.manual_stop or self.paused: self.status = 'paused' if self.paused else 'stopped' return if code == SOURCE_INFRASTRUCTURE_HOLD_EXIT: self.status = 'failed' self.desired_state = 'stopped' self.manual_stop = True self.last_action_error = ( 'source entered infrastructure hold after unresolved durable handoff cleanup' ) return # A failed periodic one-shot keeps its normal cadence when crash restarts are disabled. should_repeat = self.once and self.repeat and (code == 0 or not self.restart) should_restart = self.restart and (code != 0 or not self.once) if should_repeat or should_restart: scheduled_repeat = should_repeat if scheduled_repeat or run_duration >= self.restart_reset_after: self.restart_streak = 0 delay = self.interval if scheduled_repeat else min(self.restart_delay * max(1, 2 ** min(self.restart_streak, 6)), self.max_restart_delay) self.restarts += 1 if not scheduled_repeat: self.restart_streak += 1 self.next_start_at = time.time() + max(0, delay) self.status = 'waiting' else: self.status = 'done' if code == 0 else 'failed' self.desired_state = 'stopped' self.manual_stop = True def stop(self, timeout=15, final=False, paused=False, preserve_desired=False, dependency_block=False): self.last_action_error = '' preserve_wait_until = ( self.next_start_at if dependency_block and self.status == 'waiting' and not self.is_running() else 0 ) if not preserve_desired: self.desired_state = 'paused' if paused else 'stopped' self.stop_requested = bool(final) self.manual_stop = self.desired_state != 'running' self.paused = self.desired_state == 'paused' if dependency_block and self.desired_state == 'running': self.runtime_blocked = True self._resume_immediately = self._resume_immediately or self.is_running() or self.status != 'waiting' self.next_start_at = preserve_wait_until if not self.process or self.process.poll() is not None: self.process = None self.payload_identity = None self.startup_cleanup_pending = False if self.log_pump is not None: self.log_pump.join(timeout=5) if self.log_pump.error: self.status = 'failed' self.last_action_error = f'bounded source log writer failed: {self.log_pump.error}' self.log_pump = None return False self.log_pump = None if dependency_block and self.desired_state == 'running': self.status = 'blocked' else: self.status = 'paused' if self.desired_state == 'paused' else 'stopped' return True self.status = 'stopping' try: self.process.terminate() try: self.process.wait(timeout=timeout) except subprocess.TimeoutExpired: self.process.kill() self.process.wait(timeout=5) except Exception as exc: self.status = 'failed' self.last_action_error = f'process stop failed: {exc}' return False if self.process.poll() is None: self.status = 'failed' self.last_action_error = 'process remained live after bounded stop' return False if self.startup_cleanup_pending: self.process = None self.payload_identity = None self.startup_cleanup_pending = False if self.log_pump is not None: self.log_pump.join(timeout=5) if self.log_pump.error: self.status = 'failed' self.last_action_error = f'bounded source log writer failed: {self.log_pump.error}' self.log_pump = None return False self.log_pump = None self.payload_identity = None if not dependency_block: try: self.finalize_database_runs('stopped', 'supervisor controlled source stop') except Exception as exc: self.status = 'failed' self.last_action_error = f'database run finalization failed: {exc}' return False append_bounded_log_record( self.log_path, max(1, self.log_max_mb) * 1024 * 1024, self.log_keep, f'=== supervisor stopped {now_iso()} source={self.source} ===\n', ) if dependency_block and self.desired_state == 'running': self.status = 'blocked' else: self.status = 'paused' if self.desired_state == 'paused' else 'stopped' return True def pause(self): return self.stop(paused=True) def resume(self): return self.start(force=True) def restart_now(self): if not self.stop(timeout=10): if not self.last_action_error: self.last_action_error = 'current process did not stop' return False return self.start(force=True) def dependency_unavailable(self): if self.process is not None: code = self.process.poll() if code is None: return self.stop(timeout=15, preserve_desired=True, dependency_block=True) self._handle_process_exit(code) if self.desired_state != 'running': return True return self.stop(timeout=15, preserve_desired=True, dependency_block=True) def dependency_available(self): if not self.runtime_blocked: return True self.runtime_blocked = False if self.desired_state != 'running': return True force = self._resume_immediately self.status = 'waiting' return self.start(force=force, record_intent=False) def last_log_line(self): lines = [line.strip() for line in bounded_tail_lines(self.log_path, 40, 64 * 1024) if line.strip()] return lines[-1] if lines else '' def tail_log_lines(self, limit=40): return bounded_tail_lines(self.log_path, limit, MAX_LOG_TAIL_BYTES) def state_data(self): if not self.use_per_source_state or not os.path.exists(self.state_path): return {} try: with open(self.state_path, 'r', encoding='utf-8') as f: return json.load(f) except (OSError, json.JSONDecodeError): return {} def auth_summary(self): state = self.state_data() source_state = (state.get('sources') or {}).get(self.source) or {} summary = source_state.get('auth_summary') or {} auth_status = source_state.get('auth_status') or {} if not summary and auth_status: counts = auth_counts_from_status(auth_status) summary = { 'pool': '', 'current': source_state.get('last_auth') or 'none', 'total': len(auth_status), **counts, } return summary def auth_label(self): summary = self.auth_summary() if not summary: return '-' current = str(summary.get('current') or 'none') if current == 'none' and not summary.get('total'): return '-' return ( f"{current} ok={int(summary.get('ok', 0) or 0)} " f"dead={int(summary.get('dead', 0) or 0)} " f"lim={int(summary.get('limited', 0) or 0)} " f"rl={int(summary.get('rate_limit_errors', 0) or 0)}" ) def structured_state(self): pid = self.process.pid if self.process is not None and self.process.poll() is None else None next_run = timestamp_iso(self.next_start_at) auth = self.auth_summary() safe_error = '' if self.startup_cleanup_pending: safe_error = 'cleanup_pending' elif self.last_action_error: safe_error = 'runtime_error' elif self.last_exit_code not in (None, 0): safe_error = 'child_exit' return { 'id': managed_source_id(self), 'source': self.source, 'role': managed_source_role(self), 'lifecycle_state': self.status, 'desired_state': self.desired_state, 'process_state': 'running' if pid is not None else 'stopped', 'pid': pid, 'enabled': bool(self.enabled), 'dependency_blocked': bool(self.runtime_blocked), 'startup_cleanup_pending': bool(self.startup_cleanup_pending), 'mode': self.mode_label(), 'interval_seconds': max(0, int(self.interval or 0)), 'restart_enabled': bool(self.restart), 'restart_delay_seconds': max(0, int(self.restart_delay or 0)), 'restart_count': max(0, int(self.restarts or 0)), 'restart_streak': max(0, int(self.restart_streak or 0)), 'last_exit_code': self.last_exit_code, 'last_exit_at': timestamp_iso(self.last_exit_at), 'next_scheduled_run_at': next_run, 'safe_error_category': safe_error, 'auth_summary': { key: max(0, int(auth.get(key, 0) or 0)) for key in ( 'total', 'ok', 'dead', 'limited', 'rate_limit_errors', 'auth_invalid_errors', ) } if auth else {}, 'allowed_actions': list(managed_source_allowed_actions(self)), } def row(self): pid = self.process.pid if self.process and self.process.poll() is None else '-' uptime = format_duration(time.time() - self.started_at) if self.started_at else '-' wait_for = format_duration(self.next_start_at - time.time()) if self.status == 'waiting' else '-' return [ self.source, self.status, str(pid), self.mode_label(), uptime, format_exit_code(self.last_exit_code), wait_for, f'{self.restart_streak}/{self.restarts}', self.auth_label(), self.log_path, self.last_log_line()[:90], self.desired_state, ] class ManagedDiscoveryProducer(ManagedSource): role = DISCOVERY_PRODUCER_ROLE def __init__( self, source, config_path, project_dir, results_dir, supervisor_config, source_config, global_force_once=False, dependency_gate=None, authority_check=None, start_gate=None, child_environment=None, ): if source not in DISCOVERY_PRODUCER_SOURCES: raise ValueError('discovery producer source is outside the canonical allowlist') super().__init__( source, config_path, project_dir, results_dir, supervisor_config, source_config, True, dependency_gate, authority_check, start_gate, child_environment, ) self.once = True self.repeat = True self._set_producer_log_path() @property def producer_id(self): return f'{DISCOVERY_PRODUCER_ROLE}:{self.source}' def _set_producer_log_path(self): self.log_path = os.path.join( os.path.dirname(self.log_path), f'{DISCOVERY_PRODUCER_ROLE}-{safe_source_filename(self.source)}.log', ) def reconfigure(self, results_dir, supervisor_config, source_config, global_force_once=False): super().reconfigure(results_dir, supervisor_config, source_config, True) self.once = True self.repeat = True self._set_producer_log_path() def build_command(self): return child_bootstrap_command(DISCOVERY_PRODUCER_ROLE, [ '--config', self.config_path, '--source', self.source, '--once', ]) def build_env(self): env = os.environ.copy() use_system_proxy = bool_value(self.options.get('use_system_proxy'), False) if not use_system_proxy: for key in PROXY_ENV_KEYS: env.pop(key, None) env['NO_PROXY'] = '*' env['no_proxy'] = '*' else: env.pop('NO_PROXY', None) env.pop('no_proxy', None) env['PYTHONUNBUFFERED'] = '1' env['PYTHONIOENCODING'] = 'utf-8' if self.use_per_source_state: env['RUNNER_STATE_FILE'] = self.state_path for key, value in self.env_overrides.items(): env[str(key)] = str(value) strip_supervisor_credentials(env) env.pop('SCANNER_SKIP_STARTUP_CLEANUP', None) database_url = self.dependency_gate.database_url if self.dependency_gate else '' if not database_url: raise RuntimeError('discovery producer requires managed PostgreSQL authority') env['SCANNER_DB_URL'] = database_url env['DATABASE_URL'] = database_url env['TRUF_MANAGED_POSTGRES_DSN'] = database_url env['TRUF_DISCOVERY_SOURCE'] = self.source env.update(self.child_environment) return env def structured_state(self): state = self.state_data() source_state = (state.get('sources') or {}).get(self.source) or {} raw_result = source_state.get('last_cycle_result') or {} raw_status = raw_result.get('status') or source_state.get('last_status') or '' status = raw_status if type(raw_status) is str and raw_status in DISCOVERY_CYCLE_STATUSES else 'unknown' result = { 'status': status, } for key in ('fetched_count', 'queued_new_count', 'queued_updated_count'): try: result[key] = max(0, int(raw_result.get(key, 0) or 0)) except (TypeError, ValueError): result[key] = 0 raw_error = source_state.get('last_error_category') or '' if type(raw_error) is str and raw_error in DISCOVERY_ERROR_CATEGORIES: safe_error = raw_error else: safe_error = 'runtime_error' if raw_error else '' if not safe_error and self.last_action_error: safe_error = 'runtime_error' elif not safe_error and self.last_exit_code not in (None, 0): safe_error = 'child_exit' structured = super().structured_state() structured.update({ 'last_cycle_result': result, 'last_successful_discovery_at': safe_state_timestamp( source_state.get('last_discovery_success_at') ), 'safe_error_category': safe_error, }) return structured class ManagedKeychecks(ManagedSource): def __init__( self, config_path, project_dir, results_dir, supervisor_config, keychecks_config, global_force_once=False, dependency_gate=None, authority_check=None, start_gate=None, child_environment=None, ): self.keychecks_config = dict(keychecks_config or {}) self.command_override = None self.recheck_restore = None self.recheck_label = '' super().__init__( 'keychecks', config_path, project_dir, results_dir, supervisor_config, {'enabled': self.keychecks_config.get('enabled', False)}, global_force_once, dependency_gate, authority_check, start_gate, child_environment, ) self.once = True self.repeat = bool_value(self.keychecks_config.get('repeat'), True) self.restart = bool_value(self.keychecks_config.get('restart'), False) self.enabled = bool_value(self.keychecks_config.get('enabled'), False) self.interval = int(self.keychecks_config.get('interval', self.keychecks_config.get('cooldown', 3600)) or 3600) self.status = 'disabled' if not self.enabled else 'stopped' self.summary_path = self.keychecks_config.get('summary_tsv') or os.path.join( (self.keychecks_config.get('keycheck_dir') or os.path.join(os.path.dirname(results_dir), 'keychecks')), 'summary.tsv', ) self.state_path = self.summary_path def reconfigure(self, results_dir, supervisor_config, source_config, global_force_once=False): self.keychecks_config = dict(source_config or {}) super().reconfigure(results_dir, supervisor_config, {'enabled': self.keychecks_config.get('enabled', False)}, global_force_once) self.once = True self.repeat = bool_value(self.keychecks_config.get('repeat'), True) self.restart = bool_value(self.keychecks_config.get('restart'), False) self.enabled = bool_value(self.keychecks_config.get('enabled'), False) self.interval = int(self.keychecks_config.get('interval', self.keychecks_config.get('cooldown', 3600)) or 3600) self.summary_path = self.keychecks_config.get('summary_tsv') or os.path.join( (self.keychecks_config.get('keycheck_dir') or os.path.join(os.path.dirname(results_dir), 'keychecks')), 'summary.tsv', ) self.state_path = self.summary_path if not self.enabled: self.status = 'disabled' def build_keycheck_command(self, services=None, extra_args=None): services = services if services is not None else self.keychecks_config.get('services', 'all') if isinstance(services, (list, tuple)): services = ','.join(str(item) for item in services) arguments = [ '--config', self.config_path, '--service', str(services or 'all'), ] for key, flag in ( ('input', '--input'), ('proxy_file', '--proxy-file'), ('summary_tsv', '--summary-tsv'), ('summary_json', '--summary-json'), ('alive_summary_tsv', '--alive-summary-tsv'), ): value = self.keychecks_config.get(key) if value: arguments.extend([flag, str(value)]) max_keys = int(self.keychecks_config.get('max_keys', 0) or 0) if max_keys: arguments.extend(['--max-keys', str(max_keys)]) for key, flag in ( ('retry_network', '--retry-network'), ('retry_limited', '--retry-limited'), ('retry_unknown', '--retry-unknown'), ('retry_restricted', '--retry-restricted'), ('retry_no_balance', '--retry-no-balance'), ('retry_valid', '--retry-valid'), ('recheck_all', '--recheck-all'), ('summary_only', '--summary-only'), ('no_summary', '--no-summary'), ): if bool_value(self.keychecks_config.get(key), False): arguments.append(flag) arguments.extend(list_value(self.keychecks_config.get('extra_args'))) arguments.extend(extra_args or []) return child_bootstrap_command('keycheck', arguments) def build_command(self): if self.command_override: return list(self.command_override) return self.build_keycheck_command() def build_env(self): env = os.environ.copy() env['PYTHONUNBUFFERED'] = '1' for key, value in (self.keychecks_config.get('env') or {}).items() if isinstance(self.keychecks_config.get('env'), dict) else []: env[str(key)] = str(value) if self.dependency_gate is not None: self.dependency_gate.force_database_environment(env) env.update(self.child_environment) return env def mode_label(self): if self.command_override: return 'recheck' return 'hourly' if self.repeat else 'once' def start_recheck(self, service, runner_args, force=False): if self.is_running(): if not force: return False, 'keychecks is already running; use `recheck ... --force` to stop it and start the requested recheck.' if not self.stop(timeout=10): return False, 'keychecks recheck refused because the current process did not stop.' self.recheck_restore = { 'repeat': self.repeat, 'restart': self.restart, 'interval': self.interval, } self.command_override = self.build_keycheck_command(service, runner_args) self.recheck_label = f"{service} {' '.join(runner_args)}".strip() self.once = True self.repeat = False self.restart = False self.restarts = 0 started = self.start(force=True) if started: return True, f'keychecks: recheck started: {self.recheck_label}' if self.runtime_blocked: return True, f'keychecks: recheck intent recorded; blocked until PostgreSQL is stably ready: {self.recheck_label}' detail = self.last_action_error or 'owned keycheck process did not start' self.restore_recheck_mode(schedule_next=False) return False, f'keychecks: recheck failed: {detail}' def restore_recheck_mode(self, schedule_next=False): if not self.recheck_restore: self.command_override = None self.recheck_label = '' return restore = self.recheck_restore self.command_override = None self.recheck_restore = None self.recheck_label = '' self.once = True self.repeat = bool_value(restore.get('repeat'), True) self.restart = bool_value(restore.get('restart'), False) self.interval = int(restore.get('interval', self.interval) or self.interval or 3600) if schedule_next and self.repeat and self.enabled and not self.manual_stop and not self.paused: self.next_start_at = time.time() + max(0, self.interval) self.status = 'waiting' def poll(self): recheck_was_active = bool(self.command_override) desired_before_poll = self.desired_state super().poll() if recheck_was_active and not self.is_running() and not self.runtime_blocked: if desired_before_poll == 'running': self.desired_state = 'running' self.manual_stop = False self.restore_recheck_mode(schedule_next=self.status in ('done', 'failed')) def stop(self, timeout=15, final=False, paused=False, preserve_desired=False, dependency_block=False): stopped = super().stop( timeout=timeout, final=final, paused=paused, preserve_desired=preserve_desired, dependency_block=dependency_block, ) if self.command_override and not preserve_desired: self.restore_recheck_mode(schedule_next=False) return stopped class ManagedPipelineWorker(ManagedSource): CHILD_KINDS = { 'result-ingester': 'result-ingester', 'jsonl-projector': 'jsonl-projector', 'janitor': 'janitor', 'worker-api': 'worker-api', } def __init__( self, worker_name, config_path, project_dir, results_dir, supervisor_config, worker_config, dependency_gate=None, authority_check=None, start_gate=None, child_environment=None, ): default_enabled = worker_name != 'worker-api' options = { 'enabled': bool_value((worker_config or {}).get('enabled'), default_enabled), 'restart': True, 'once': False, 'repeat': False, 'interval': 0, 'env': (worker_config or {}).get('env') or {}, } super().__init__( worker_name, config_path, project_dir, results_dir, supervisor_config, options, False, dependency_gate, authority_check, start_gate, child_environment, ) self.worker_config = dict(worker_config or {}) self.enabled = bool_value(self.worker_config.get('enabled'), default_enabled) self.status = 'disabled' if not self.enabled else 'stopped' self.once = False self.repeat = False self.restart = True self.use_per_source_state = False def build_command(self): return child_bootstrap_command( self.CHILD_KINDS[self.source], ['--config', self.config_path], ) def build_env(self): env = os.environ.copy() env['PYTHONUNBUFFERED'] = '1' env['PYTHONIOENCODING'] = 'utf-8' if self.source == 'janitor': for key in list(env): if key.upper().startswith('PG') or key.upper() in { 'TRUF_MANAGED_POSTGRES_DSN', 'SCANNER_DB_URL', 'DATABASE_URL', 'SCANNER_DASHBOARD_DB_URL', 'KEYCHECK_DB_URL', }: env.pop(key, None) elif self.dependency_gate is not None: self.dependency_gate.force_database_environment(env) env.update(self.child_environment) if self.source == 'janitor': env.pop('TRUF_MANAGED_POSTGRES_DSN', None) env.pop('SCANNER_DB_URL', None) env.pop('DATABASE_URL', None) return env def finalize_database_runs(self, status, reason): return 0, 0 def state_data(self): return {} def mode_label(self): return 'singleton' class ManagedDockerShadow(ManagedSource): def __init__( self, config_path, project_dir, results_dir, supervisor_config, shadow_config, dependency_gate=None, authority_check=None, start_gate=None, child_environment=None, ): options = {'enabled': bool_value((shadow_config or {}).get('enabled'), False)} super().__init__( 'docker-shadow', config_path, project_dir, results_dir, supervisor_config, options, False, dependency_gate, authority_check, start_gate, child_environment, ) self.enabled = bool_value((shadow_config or {}).get('enabled'), False) self.status = 'disabled' if not self.enabled else 'stopped' self.manual_only = True self.once = True self.repeat = False self.restart = False self.interval = 0 self.use_per_source_state = False def build_command(self): return child_bootstrap_command( 'docker-shadow', ['--config', self.config_path], ) def start(self, force=False, record_intent=True): if self.dependency_gate is not None and not self.dependency_gate.ready: self.last_action_error = 'managed PostgreSQL dependency is unavailable' self.desired_state = 'stopped' self.manual_stop = True self.runtime_blocked = False self.status = 'stopped' return False return super().start(force=force, record_intent=record_intent) def dependency_unavailable(self): self.desired_state = 'stopped' self.manual_stop = True self.runtime_blocked = False return self.stop(timeout=15) def dependency_available(self): self.runtime_blocked = False return True def finalize_database_runs(self, status, reason): return 0, 0 def state_data(self): return {} def mode_label(self): return 'manual-once' def build_table_lines(managed_sources, max_width=None): max_width = max_width or terminal_width() headers = ['source', 'status', 'desired', 'pid', 'mode', 'up', 'exit', 'next', 'rs', 'auth', 'last_log'] raw_rows = [source.row() for source in managed_sources] rows = [ [row[0], row[1], row[11], row[2], row[3], row[4], row[5], row[6], row[7], row[8], row[10]] for row in raw_rows ] fixed_widths = [] for index, header in enumerate(headers[:-1]): values = [row[index] for row in rows] cap = 36 if header == 'auth' else 20 if header == 'source' else 12 if header in ('mode', 'desired') else 10 if header == 'exit' else 9 fixed_widths.append(min(max([len(header)] + [len(str(value)) for value in values]) if values else len(header), cap)) separator_width = 3 * (len(headers) - 1) last_width = max(20, max_width - sum(fixed_widths) - separator_width) widths = fixed_widths + [last_width] def format_row(values): return ' | '.join(truncate_text(values[index], widths[index]).ljust(widths[index]) for index in range(len(headers))) lines = [ 'Supervisor status ' + now_iso(), format_row(headers), '-+-'.join('-' * width for width in widths), ] for row in rows: lines.append(format_row(row)) lines.append('') lines.append('Type `help` for commands. Use `command ` for full log/state paths.') return lines def scan_worker_snapshot(config): global_config = (config or {}).get('global') or {} base_limit = max(0, int_value(global_config.get('max_active_scans'), 0)) bonus_limit = max(0, min(1, int_value(global_config.get('opportunistic_scan_slots'), 0))) limit = base_limit + bonus_limit snapshot = { 'active': 0, 'limit': limit, 'base_active': 0, 'base_limit': base_limit, 'bonus_active': 0, 'bonus_limit': bonus_limit, 'trufflehog': 0, 'sources': {}, 'detail': '', } if limit <= 0: return snapshot db_path = str(global_config.get('scan_limiter_db') or '') if not db_path: state_dir = str(global_config.get('state_dir') or '') if state_dir: db_path = os.path.join(state_dir, 'scan_limiter.db') if not db_path or not os.path.isfile(db_path): return snapshot connection = None try: reject_reparse_components(db_path) normalized = os.path.abspath(db_path).replace('\\', '/') uri = f'file:{quote(normalized, safe="/:")}?mode=ro' connection = sqlite3.connect(uri, uri=True, timeout=0.2) connection.execute('PRAGMA query_only=ON') connection.execute('PRAGMA busy_timeout=200') columns = { row[1] for row in connection.execute('PRAGMA table_info(scan_slots)').fetchall() } slot_kind = "COALESCE(slot_kind, 'base')" if 'slot_kind' in columns else "'base'" rows = connection.execute( f"""SELECT COALESCE(NULLIF(owner_source, ''), 'unknown'), child_executable, {slot_kind} FROM scan_slots ORDER BY 1""" ).fetchall() sources = {} for source, child_executable, kind in rows: source = str(source) sources[source] = sources.get(source, 0) + 1 if kind == 'bonus': snapshot['bonus_active'] += 1 else: snapshot['base_active'] += 1 if os.path.basename(str(child_executable or '')).lower() == 'trufflehog.exe': snapshot['trufflehog'] += 1 snapshot['sources'] = sources snapshot['active'] = len(rows) except (OSError, ValueError, sqlite3.Error) as exc: snapshot['detail'] = str(exc)[:200] finally: if connection is not None: connection.close() return snapshot def current_scan_worker_snapshot(context, refresh_sec=2): context = context or {} now = time.monotonic() cached = context.get('_scan_worker_snapshot') if cached is not None and now < float(context.get('_scan_worker_snapshot_refresh_at') or 0): return cached snapshot = scan_worker_snapshot(context.get('config')) context['_scan_worker_snapshot'] = snapshot context['_scan_worker_snapshot_refresh_at'] = now + max(0.1, float(refresh_sec)) return snapshot def scan_worker_status_line(context): snapshot = current_scan_worker_snapshot(context) if snapshot['limit'] <= 0: return '' staging = max(0, snapshot['active'] - snapshot['trufflehog']) line = ( f"Scan workers: active={snapshot['active']}/{snapshot['limit']} " f"base={snapshot['base_active']}/{snapshot['base_limit']} " f"bonus={snapshot['bonus_active']}/{snapshot['bonus_limit']} " f"trufflehog={snapshot['trufflehog']} staging={staging}" ) if snapshot['sources']: sources = ', '.join(f'{source}:{count}' for source, count in snapshot['sources'].items()) line += ' sources=' + truncate_text(sources, 160) if snapshot['detail']: line += ' detail=' + snapshot['detail'] return line def pipeline_status_snapshot(context, refresh_sec=None, allow_refresh=True): context = context or {} if context.get('_defer_pipeline_status_until_next_tick'): cached = context.get('_pipeline_status_snapshot') if cached is not None: return cached return { 'ingester_ready': False, 'projector_ready': False, 'cutover_ready': False, 'ingester_state': 'starting', 'projector_state': 'starting', 'bundle_items': 0, 'bundle_bytes': 0, 'projection_items': 0, 'projection_bytes': 0, 'keycheck_items': 0, 'keycheck_bytes': 0, 'quarantine_items': 0, 'quarantine_bytes': 0, 'detail': 'Pipeline workers are starting', } if refresh_sec is None: refresh_sec = (context.get('supervisor_config') or {}).get( 'pipeline_status_refresh_sec', 2, ) or 2 now = time.monotonic() cached = context.get('_pipeline_status_snapshot') if cached is not None and now < float(context.get('_pipeline_status_refresh_at') or 0): return cached if not allow_refresh: if cached is not None: return cached return { 'ingester_ready': False, 'projector_ready': False, 'cutover_ready': False, 'ingester_state': 'starting', 'projector_state': 'starting', 'bundle_items': 0, 'bundle_bytes': 0, 'projection_items': 0, 'projection_bytes': 0, 'keycheck_items': 0, 'keycheck_bytes': 0, 'quarantine_items': 0, 'quarantine_bytes': 0, 'detail': 'Pipeline status refresh is pending', } snapshot = { 'ingester_ready': False, 'projector_ready': False, 'cutover_ready': False, 'ingester_state': 'missing', 'projector_state': 'missing', 'bundle_items': 0, 'bundle_bytes': 0, 'projection_items': 0, 'projection_bytes': 0, 'keycheck_items': 0, 'keycheck_bytes': 0, 'quarantine_items': 0, 'quarantine_bytes': 0, 'detail': '', } gate = context.get('dependency_gate') database_url = getattr(gate, 'database_url', '') if gate is not None else '' if not database_url or (gate is not None and not gate.ready): snapshot['detail'] = 'PostgreSQL is not stably ready' else: connection = None try: connection = connect_postgres( database_url, connect_timeout_sec=2, statement_timeout_ms=3000, lock_timeout_ms=1000, idle_in_transaction_timeout_ms=3000, tcp_user_timeout_ms=3000, ) leases = connection.execute( '''SELECT worker_name, state, lease_expires_at FROM pipeline_leases WHERE worker_name IN ('result_ingester','jsonl_projector')''' ).fetchall() capacity = connection.execute( 'SELECT * FROM pipeline_capacity WHERE id = 1' ).fetchone() cutover = connection.execute( 'SELECT marker, evidence_sha256 FROM runtime_final_cutover WHERE id = 1' ).fetchone() connection.commit() lease_map = {row['worker_name']: row for row in leases} current = datetime.now(timezone.utc).isoformat(timespec='seconds') ingester = lease_map.get('result_ingester') projector = lease_map.get('jsonl_projector') snapshot['ingester_state'] = str(ingester['state']) if ingester else 'missing' snapshot['projector_state'] = str(projector['state']) if projector else 'missing' snapshot['cutover_ready'] = bool( cutover and cutover['marker'] == 'postgres-normalized-v2-authority' and re.fullmatch(r'[a-f0-9]{64}', str(cutover['evidence_sha256'] or '')) ) snapshot['ingester_ready'] = bool( ingester and ingester['state'] == 'ready' and str(ingester['lease_expires_at'] or '') > current and snapshot['cutover_ready'] ) snapshot['projector_ready'] = bool( projector and projector['state'] == 'ready' and str(projector['lease_expires_at'] or '') > current and snapshot['cutover_ready'] ) if not snapshot['cutover_ready']: snapshot['detail'] = 'Final PostgreSQL v2 cutover marker is absent or invalid' if capacity: for key in ( 'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes', 'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes', ): snapshot[key] = int(capacity[key] or 0) except Exception as exc: snapshot['detail'] = f'{type(exc).__name__}: {exc}'[:200] if connection is not None: try: connection.rollback() except Exception: pass finally: if connection is not None: connection.close() context['_pipeline_status_snapshot'] = snapshot context['_pipeline_status_refresh_at'] = now + max(0.2, float(refresh_sec)) return snapshot def pipeline_status_line(context): snapshot = pipeline_status_snapshot(context) line = ( f"Pipeline: ingester={snapshot['ingester_state']} projector={snapshot['projector_state']} " f"bundles={snapshot['bundle_items']}/{snapshot['bundle_bytes']}B " f"projection={snapshot['projection_items']}/{snapshot['projection_bytes']}B " f"candidates={snapshot['keycheck_items']}/{snapshot['keycheck_bytes']}B " f"quarantine={snapshot['quarantine_items']}/{snapshot['quarantine_bytes']}B" ) if snapshot['detail']: line += ' detail=' + snapshot['detail'] return line def build_runtime_table_lines(managed_sources, context=None, max_width=None): lines = [] scan_line = scan_worker_status_line(context) if scan_line: lines.extend((scan_line, '')) lines.extend((pipeline_status_line(context), '')) lines.extend(build_table_lines(managed_sources, max_width=max_width)) return lines def table_signature(managed_sources): signature = [] for source in managed_sources: process = source.process pid = process.pid if process is not None and process.poll() is None else None signature.append(( source.source, source.status, source.desired_state, pid, source.last_exit_code, source.restarts, source.restart_streak, source.runtime_blocked, source.last_action_error, )) return tuple(signature) def print_table(managed_sources, clear=True, context=None): if clear: os.system('cls' if os.name == 'nt' else 'clear') for line in build_runtime_table_lines(managed_sources, context=context): print(line) HELP_TEXT = r''' Commands: help | h | ? Show this help. status | s Print a fresh status table once. auth Print auth pool health from runner state. Example: auth github watch | w Open live status view in an alternate screen. Press q to return to prompt. This keeps scrollback clean and avoids table spam. reload | r Live reload is disabled by config authority binding. Use coordinated shutdown and restart. recheck [types...] [options] Start a one-shot keycheck_runner recheck using configured keycheck probes/service_args. If no type is given, defaults to --recheck-all. Types: network, ratelimited/aliveratelimited, unknown, restricted, nobalance, valid, legacy-vertex, all. Options: --force, --max-keys N, --input PATH, --proxy-file PATH, --no-resource-probe. Examples: recheck all network recheck gemini --network --ratelimited recheck replicate --valid --max-keys 5 start Start a stopped source using its current mode from the table. Example: start pypi stop Stop the child process and keep it stopped. Supervisor will not auto-restart it. Example: stop dockerhub restart Stop then start immediately. Example: restart github pause Stop and mark as paused. Same process behavior as stop, but visually distinct. Example: pause npm resume Unpause and start immediately. Example: resume npm once Switch source to once mode and start it with --once. It will not repeat unless mode is changed. Example: once pypi mode loop|once|repeat loop = no --once; child console_runner loops internally using config cooldown. once = pass --once; one source cycle, then stop when child exits. repeat = pass --once; supervisor repeats one-shot cycles after interval. Example: mode pypi repeat set interval Set repeat interval for once+repeat mode. Example: set pypi interval 120 set restart on|off Enable/disable restart after crash or loop child exit. Example: set dockerhub restart off set restart_delay Set initial failure restart delay. Example: set github restart_delay 60 logs [lines] Print the last N log lines from the configured runtime logs directory. Default: 40. Example: logs pypi 80 command Print child command, log file, and state file. Example: command github dashboard status|stop|start|restart Manage the dashboard process without stopping supervisor or source processes. Example: dashboard stop shutdown Request coordinated shutdown of sources, keychecks, dashboard, and identity-verified PostgreSQL. Authenticated shutdown remains available after config hash drift; other mutations are rejected. quit | exit | q In foreground supervisor: stop all children and exit. In --attach: detach only; background supervisor keeps running. Statuses: stopped Not running. This is the default state. running Child console_runner.py is currently running. waiting Supervisor will start/restart after the `next` countdown. blocked Desired state is running, but stable PostgreSQL readiness is unavailable. paused Explicitly paused by command. done One-shot child exited successfully and will not repeat. failed Child exited with non-zero code and restart is disabled. The `rs` column is current failure streak / total automatic restarts. Modes: loop Child runs without --once and handles its own loop/cooldown. once Child runs with --once and stops after one source cycle. once+repeat Child runs with --once; supervisor restarts it after `interval` seconds. '''.strip() def timestamp_iso(value): if not value: return None try: return datetime.fromtimestamp(float(value), timezone.utc).isoformat(timespec='seconds') except (OSError, OverflowError, TypeError, ValueError): return None def managed_source_id(source): if isinstance(source, ManagedDiscoveryProducer): return source.producer_id return str(source.source) def managed_source_role(source): if isinstance(source, ManagedDiscoveryProducer): return DISCOVERY_PRODUCER_ROLE if isinstance(source, ManagedKeychecks): return 'keycheck' if isinstance(source, ManagedPipelineWorker): return source.CHILD_KINDS.get(source.source, 'pipeline-worker') if isinstance(source, ManagedDockerShadow): return 'docker-shadow' return 'scanner' def managed_source_allowed_actions(source): if getattr(source, 'manual_only', False): return ('start', 'stop') actions = list(MANAGED_SOURCE_LIFECYCLE_ACTIONS) if isinstance(source, ManagedDiscoveryProducer): actions.append('set-interval') elif isinstance(source, ManagedPipelineWorker): actions.extend(('set-restart', 'set-restart-delay')) else: actions.extend(MANAGED_SOURCE_SETTING_ACTIONS) return tuple(actions) def managed_source_registry(managed_sources): registry = {} for source in managed_sources: source_id = managed_source_id(source) if not source_id or source_id in registry: raise ValueError('managed source IDs are not unique') registry[source_id] = source return registry def source_map(managed_sources): return {source.source: source for source in managed_sources} def select_sources(managed_sources, selector): selector = normalize_source_name(selector) if selector == 'all': return managed_sources sources = source_map(managed_sources) item = sources.get(selector) if not item: print(f'Unknown source: {selector}. Available: {", ".join(sorted(sources))}') return [] return [item] def autostart_sources(managed_sources): return [source for source in managed_sources if not getattr(source, 'manual_only', False)] def set_mode(source, mode): mode = str(mode or '').lower() if mode == 'loop': source.once = False source.repeat = True elif mode == 'once': source.once = True source.repeat = False elif mode == 'repeat': source.once = True source.repeat = True else: print(f'Unknown mode: {mode}. Use loop, once, or repeat.') return False return True def print_source_command(source): print(f'[{source.source}]') print('command:', ' '.join(command_for_log(source.build_command()))) print('mode:', source.mode_label()) print('log:', source.log_path) print('state:', source.state_path if source.use_per_source_state else '(config default)') print('desired:', source.desired_state, 'dependency_blocked:', source.runtime_blocked) print('restart:', source.restart, 'interval:', source.interval, 'restart_delay:', source.restart_delay) def print_auth_status(source): state = source.state_data() source_state = (state.get('sources') or {}).get(source.source) or {} summary = source_state.get('auth_summary') or source.auth_summary() auth_status = source_state.get('auth_status') or {} if not summary and not auth_status: print(f'[{source.source}] auth: no auth pool state') print('state:', source.state_path if source.use_per_source_state else '(config default)') return print( f"[{source.source}] pool={summary.get('pool') or '-'} current={summary.get('current') or '-'} " f"total={int(summary.get('total', len(auth_status)) or 0)} " f"ok={int(summary.get('ok', 0) or 0)} dead={int(summary.get('dead', 0) or 0)} " f"lim={int(summary.get('limited', 0) or 0)} rl={int(summary.get('rate_limit_errors', 0) or 0)} " f"auth_invalid={int(summary.get('auth_invalid_errors', 0) or 0)}" ) print('state:', source.state_path if source.use_per_source_state else '(config default)') if not auth_status: return print('name\tstatus\tdisabled_until\treason\trl\tauth_invalid\tfailures\tlast_error') for name in sorted(auth_status): item = auth_status.get(name) or {} status = item.get('status') or auth_item_kind(item) rl = int(item.get('rate_limit_count', 0) or 0) + int(item.get('secondary_rate_limit_count', 0) or 0) auth_invalid = int(item.get('auth_invalid_count', 0) or 0) if item.get('disabled_reason') == 'auth_invalid' and not auth_invalid: auth_invalid = int(item.get('failures', 1) or 1) print('\t'.join([ str(name), str(status), str(item.get('disabled_until') or '-'), str(item.get('disabled_reason') or '-'), str(rl), str(auth_invalid), str(int(item.get('failures', 0) or 0)), str(item.get('last_error') or '').replace('\t', ' ')[:180], ])) def runtime_authority( config_path, supervisor_path=None, trufflehog_path=None, expected_config_sha256=None, expected_supervisor_sha256=None, expected_code_manifest_sha256=None, code_manifest=None, policy_paths=None, include_trufflehog=True, ): supervisor_path = os.path.abspath(supervisor_path or __file__) config_sha256 = sha256_file(config_path) if expected_config_sha256 and config_sha256 != expected_config_sha256: raise RuntimeError('supervisor config changed after it was loaded') supervisor_sha256 = sha256_file(supervisor_path) if expected_supervisor_sha256 and supervisor_sha256 != expected_supervisor_sha256: raise RuntimeError('supervisor script changed after parent authority capture') code_manifest = code_manifest or build_code_manifest( trufflehog_path=trufflehog_path, policy_paths=policy_paths, include_trufflehog=include_trufflehog, ) manifest_sha256 = code_manifest_sha256(code_manifest) if expected_code_manifest_sha256 and manifest_sha256 != expected_code_manifest_sha256: raise RuntimeError('supervisor code manifest changed after parent authority capture') verify_code_manifest(code_manifest, manifest_sha256) return { 'supervisor_path': supervisor_path, 'supervisor_sha256': supervisor_sha256, 'config_path': os.path.abspath(config_path), 'config_sha256': config_sha256, 'code_manifest': code_manifest, 'code_manifest_sha256': manifest_sha256, } def configured_policy_paths(config): global_config = (config or {}).get('global') or {} values = [global_config.get('trufflehog_config')] for source in ((config or {}).get('sources') or {}).values(): if isinstance(source, dict) and source.get('trufflehog_config'): values.append(resolve_optional_path(source['trufflehog_config'], global_config)) return list(dict.fromkeys(value for value in values if value)) def runtime_authority_error(context, require_private_acl=False): authority = (context or {}).get('authority') if not authority: return '' try: verify_code_manifest( authority['code_manifest'], authority['code_manifest_sha256'], require_private_acl=require_private_acl, ) if sha256_file(authority['supervisor_path']) != authority['supervisor_sha256']: return 'supervisor script authority hash drifted' if sha256_file(authority['config_path']) != authority['config_sha256']: return 'supervisor config authority hash drifted' except (OSError, ValueError, LifecycleAuthorityError) as exc: return f'runtime code authority validation failed: {exc}' return '' def inhibit_for_authority_drift(context, detail): context = context or {} context['runtime_failed'] = True if not context.get('authority_drift'): context['authority_drift'] = str(detail) context['start_gate_open'] = False controller = context.get('postgres_controller') if controller is not None: controller.inhibit_lifecycle(detail) gate = context.get('dependency_gate') if gate is not None: try: gate.set_ready(False) except DependencyStopError as exc: context['fatal_child_stop_error'] = str(exc) grace = max(0.0, float((context.get('supervisor_config') or {}).get('authority_drift_shutdown_grace_sec', 5) or 0)) context['authority_drift_shutdown_at'] = time.monotonic() + grace shutdown_event = context.get('shutdown_event') if ( shutdown_event is not None and time.monotonic() >= float(context.get('authority_drift_shutdown_at') or 0) ): shutdown_event.set() return False def check_runtime_authority(context, trigger_shutdown=True, require_private_acl=False): detail = runtime_authority_error(context, require_private_acl=require_private_acl) if not detail: return True if trigger_shutdown: return inhibit_for_authority_drift(context, detail) return False def lifecycle_phase(context): return str((context or {}).get('lifecycle_phase') or (context or {}).get('activation_state') or PHASE_ACTIVE).upper() def lifecycle_start_allowed(context): context = context or {} return ( not context.get('shutdown_requested', False) and lifecycle_phase(context) == PHASE_ACTIVE and bool(context.get('start_gate_open', True)) and not any( getattr(source, 'startup_cleanup_pending', False) is True for source in context.get('managed_sources', ()) ) ) def shutdown_checkpoint(context): """Consume shutdown flags and uncertain rollback under the control lock.""" context = context or {} if lifecycle_phase(context) not in (PHASE_STOPPING, PHASE_FAILED_HOLD): pending = next(( source for source in context.get('managed_sources', ()) if getattr(source, 'startup_cleanup_pending', False) is True ), None) if pending is not None: try: begin_stopping(context) finally: enter_failed_hold(context, pending.last_action_error) elif context.get('shutdown_requested'): begin_stopping(context) shutdown_event = context.get('shutdown_event') return ( bool(context.get('shutdown_requested')) or lifecycle_phase(context) in (PHASE_STOPPING, PHASE_FAILED_HOLD) or (shutdown_event is not None and shutdown_event.is_set()) ) def begin_stopping(context): """Close every lifecycle start gate before a shutdown acknowledgement.""" context = {} if context is None else context phase = lifecycle_phase(context) already_stopping = phase in (PHASE_STOPPING, PHASE_FAILED_HOLD) context['lifecycle_phase'] = phase if already_stopping else PHASE_STOPPING context['activation_state'] = context['lifecycle_phase'] context['start_gate_open'] = False if not already_stopping: context['authority_release_safe'] = False try: shutdown_event = context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() if already_stopping: retry_event = context.get('shutdown_retry_event') if retry_event is not None: retry_event.set() return controller = context.get('postgres_controller') if controller is not None and hasattr(controller, 'inhibit_lifecycle'): controller.inhibit_lifecycle('supervisor lifecycle is STOPPING') instance_file = context.get('instance_file') instance_id = context.get('instance_id') if instance_file and instance_id: update_instance_activation(instance_file, instance_id, PHASE_STOPPING) except BaseException: context['runtime_failed'] = True raise def enter_failed_hold(context, detail='coordinated shutdown remains unsafe'): """Keep process and OS authority live after any unconfirmed child stop.""" context = {} if context is None else context context['lifecycle_phase'] = PHASE_FAILED_HOLD context['activation_state'] = PHASE_FAILED_HOLD context['start_gate_open'] = False context['authority_release_safe'] = False context['runtime_failed'] = True context['shutdown_failure'] = str(detail or 'coordinated shutdown remains unsafe') instance_file = context.get('instance_file') instance_id = context.get('instance_id') if instance_file and instance_id: update_instance_activation(instance_file, instance_id, PHASE_FAILED_HOLD) def shutdown_authority_error(context, instance_id=None, token=None, control_address_value=None): """Use private credentials/process identity even when on-disk code drifted.""" context = context or {} authority = context.get('authority') instance_file = context.get('instance_file') if not instance_file: return '' try: metadata = load_instance_metadata(instance_file) if not secrets.compare_digest(metadata['instance_id'], str(instance_id or '')): return 'private shutdown metadata instance mismatch' if not secrets.compare_digest(metadata['token'], str(token or '')): return 'private shutdown metadata credential mismatch' expected_control = metadata['control'] actual_host, actual_port = control_address_value or ('', 0) if str(expected_control['host']) != str(actual_host) or int(expected_control['port']) != int(actual_port): return 'private shutdown metadata control endpoint mismatch' retained = verify_instance_process( metadata, authority.get('supervisor_path') if authority else os.path.abspath(__file__), authority.get('config_path') if authority else context.get('config_path'), allow_config_drift=True, allow_code_drift=True, ) retained.close() except (OSError, ValueError) as exc: return f'private shutdown authority verification failed: {exc}' return '' def command_failure(context, message): if context is not None: context['_command_failed'] = True print(message) def command_is_mutating(parts): if not parts: return False action = parts[0].lower() if action == 'dashboard': return len(parts) < 2 or parts[1].lower() != 'status' return action in { 'reload', 'r', 'recheck', 'start', 'stop', 'restart', 'pause', 'resume', 'once', 'mode', 'set', 'shutdown', } def handle_recheck_command(parts, managed_sources, context=None): parsed, error = parse_recheck_command(parts) if error: command_failure(context, error) return True keychecks = source_map(managed_sources).get('keychecks') if not keychecks: command_failure(context, 'keychecks source is not managed. Enable keychecks.enabled or run supervisor with --sources keychecks/all.') return True if not isinstance(keychecks, ManagedKeychecks): command_failure(context, 'keychecks source is not a ManagedKeychecks instance') return True ok, message = keychecks.start_recheck(parsed['service'], parsed['runner_args'], force=parsed['force']) print(message) if not ok: if context is not None: context['_command_failed'] = True if ok: print('command:', ' '.join(keychecks.build_command())) print('log:', keychecks.log_path) return True class ManagedDashboard: def __init__( self, config_path, project_dir, supervisor_config, results_dir, queue_dir, dependency_gate=None, force=False, process_factory=OwnedProcess, health_probe=None, executor=None, clock=None, authority_check=None, start_gate=None, child_environment=None, ): self.source = 'dashboard' enabled, self.host, self.port, self.url = dashboard_settings(supervisor_config) self.enabled = bool(enabled or force) self.config_path = os.path.abspath(config_path) self.project_dir = project_dir self.supervisor_config = supervisor_config self.results_dir = results_dir self.queue_dir = queue_dir self.dependency_gate = dependency_gate self.process_factory = process_factory self.health_probe = health_probe or self._http_health_probe self._executor = executor or ThreadPoolExecutor(max_workers=1, thread_name_prefix='dashboard-health') self._owns_executor = executor is None self._clock = clock or time.monotonic self.authority_check = authority_check self.start_gate = start_gate self.child_environment = dict(child_environment or {}) dashboard_config = supervisor_config.get('dashboard') if isinstance(supervisor_config.get('dashboard'), dict) else {} self.startup_grace_sec = max(0.1, float(dashboard_config.get('startup_grace_sec', 30) or 30)) self.health_interval_sec = max(0.1, float(dashboard_config.get('health_interval_sec', 2) or 2)) self.health_timeout_sec = max(0.1, float(dashboard_config.get('health_timeout_sec', 1) or 1)) self.restart_base_sec = max(0.1, float(dashboard_config.get('restart_base_sec', 2) or 2)) self.restart_max_sec = max(self.restart_base_sec, float(dashboard_config.get('restart_max_sec', 60) or 60)) stable_health = dashboard_config.get('stable_health_sec', 60) self.stable_health_sec = max(0.0, float(60 if stable_health is None else stable_health)) self.env_overrides = dashboard_config.get('env') if isinstance(dashboard_config.get('env'), dict) else {} log_dir = resolve_path(config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) self.log_path = resolve_path(config_path, supervisor_config.get('dashboard_log') or os.path.join(log_dir, 'dashboard.log')) self.process = None self.log_pump = None self.desired_state = 'running' if self.enabled else 'stopped' self.status = ( 'blocked' if self.enabled and dependency_gate and not dependency_gate.ready else 'pending' if self.enabled else 'disabled' ) self.detail = 'waiting for dependency readiness' if self.enabled and dependency_gate and not dependency_gate.ready else '' self.failures = 0 self.started_at = None self.healthy_since = None self.next_health_at = 0.0 self.next_start_at = 0.0 self._future = None self._future_generation = None self._generation = 0 self.fatal_stop_failure = False if dependency_gate is not None: dependency_gate.register(self) @property def healthy(self): return self.status == 'healthy' and self.process is not None and self.process.poll() is None def _http_health_probe(self, url, timeout): request = urllib.request.Request(url.rstrip('/') + '/_stcore/health', method='GET') opener = urllib.request.build_opener(urllib.request.ProxyHandler({})) with opener.open(request, timeout=max(0.1, float(timeout))) as response: return int(getattr(response, 'status', response.getcode())) == 200 def build_command(self): return child_bootstrap_command('dashboard', [ '--server.headless', 'true', '--server.address', self.host, '--server.port', str(self.port), '--', '--config', self.config_path, '--results-dir', self.results_dir, '--queue-dir', self.queue_dir, ]) def build_env(self): env = os.environ.copy() env['PYTHONUNBUFFERED'] = '1' env['PYTHONIOENCODING'] = 'utf-8' for key, value in self.env_overrides.items(): env[str(key)] = str(value) if self.dependency_gate is not None: self.dependency_gate.force_database_environment(env) env.update(self.child_environment) env['TRUF_DASHBOARD_CANONICAL_LAUNCH'] = '1' env['TRUF_DASHBOARD_HOST'] = self.host return env def _launch(self, now): if self.start_gate is not None and not self.start_gate(): self.detail = 'supervisor lifecycle start gate is closed' self.status = 'stopped' return False if self.authority_check is not None and not self.authority_check(): self.detail = 'runtime script/config authority drifted' self.status = 'failed' return False require_private_directory(os.path.dirname(self.log_path), create=False) self._generation += 1 if self._future is not None: self._future.cancel() self._future = None self._future_generation = None log_max_bytes = max(1, int(self.supervisor_config.get('log_max_mb', 64) or 64)) * 1024 * 1024 self.log_pump = BoundedRotatingLogPump( self.log_path, log_max_bytes, self.supervisor_config.get('log_keep', 5), ) try: command = self.build_command() self.log_pump.write(f'\n=== dashboard start {now_iso()} ===\n') self.log_pump.write(f'url: {self.url}\n') self.log_pump.write('command: ' + ' '.join(command) + '\n') creationflags = (subprocess.CREATE_NEW_PROCESS_GROUP | CREATE_NO_WINDOW) if os.name == 'nt' else 0 self.process = self.process_factory( command, cwd=self.project_dir, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, env=self.build_env(), creationflags=creationflags, ) stream = getattr(self.process, 'stdout', None) if stream is not None: self.log_pump.attach(stream) else: self.log_pump.close() except BaseException: if self.process is not None and self.process.poll() is None: self._stop_current() elif self.log_pump is not None: self.log_pump.close() self.log_pump = None raise self.started_at = now self.healthy_since = None self.next_health_at = now self.next_start_at = 0.0 self.status = 'pending' self.detail = 'process started; health has not been confirmed' return True def start(self, force=False, record_intent=True): if self.start_gate is not None and not self.start_gate(): self.detail = 'supervisor lifecycle start gate is closed' self.status = 'stopped' return False if record_intent: self.desired_state = 'running' if force: self.enabled = True if not self.enabled or self.desired_state != 'running': self.status = 'disabled' if not self.enabled else 'stopped' return False if self.dependency_gate is not None and not self.dependency_gate.ready: self.status = 'blocked' self.detail = 'waiting for stable PostgreSQL readiness' return False if self.process is not None and self.process.poll() is None: return True now = self._clock() if not force and now < self.next_start_at: self.status = 'backoff' return False if force: self.next_start_at = 0.0 try: return self._launch(now) except Exception as exc: if self.process is not None and self.process.poll() is None: self.status = 'failed' self.detail = 'dashboard process launch cleanup failed' self.fatal_stop_failure = True return False self.process = None self._schedule_failure(now, f'dashboard process launch failed: {type(exc).__name__}') return False def _stop_current(self, timeout=15): process = self.process if process is None or process.poll() is not None: self.process = None if self.log_pump is not None: self.log_pump.join(timeout=5) self.log_pump = None return True try: process.terminate() try: process.wait(timeout=timeout) except subprocess.TimeoutExpired: process.kill() process.wait(timeout=5) except Exception as exc: self.detail = f'dashboard process stop failed: {type(exc).__name__}' self.fatal_stop_failure = True return False if process.poll() is None: self.detail = 'dashboard process remained live after bounded stop' self.fatal_stop_failure = True return False if self.log_pump is not None: self.log_pump.join(timeout=5) if self.log_pump.error: self.detail = f'bounded dashboard log writer failed: {self.log_pump.error}' self.log_pump = None self.fatal_stop_failure = True return False self.log_pump = None self.process = None return True def _schedule_failure(self, now, detail): self.failures += 1 delay = min(self.restart_base_sec * (2 ** min(self.failures - 1, 20)), self.restart_max_sec) self.next_start_at = now + delay self.status = 'backoff' self.detail = detail self.started_at = None self.healthy_since = None def _consume_health(self, now): if self._future is None or not self._future.done(): return future = self._future generation = self._future_generation self._future = None self._future_generation = None if generation != self._generation: return try: healthy = bool(future.result()) except BaseException: healthy = False if self.process is None or self.process.poll() is not None: return if healthy: if self.healthy_since is None: self.healthy_since = now self.status = 'healthy' self.detail = 'Streamlit health endpoint is ready' if now - self.healthy_since >= self.stable_health_sec: self.failures = 0 self.next_health_at = now + self.health_interval_sec return self.healthy_since = None if self.started_at is not None and now - self.started_at < self.startup_grace_sec: self.status = 'pending' self.detail = 'waiting for Streamlit health during startup grace' self.next_health_at = now + self.health_interval_sec return stopped = self._stop_current() detail = 'dashboard health check failed after startup grace' if not stopped: self.status = 'failed' self.detail = detail + '; current owned process did not stop' raise RuntimeError(self.detail) self._schedule_failure(now, detail) def poll(self): now = self._clock() self._consume_health(now) if self.start_gate is not None and not self.start_gate(): return self.status if self.desired_state != 'running': return self.status if self.dependency_gate is not None and not self.dependency_gate.ready: if self.process is not None: self.dependency_unavailable() return self.status if self.process is None: self.start(record_intent=False) return self.status code = self.process.poll() if code is None and self.log_pump is not None and self.log_pump.error: self.detail = f'bounded dashboard log writer failed: {self.log_pump.error}' if not self._stop_current(): self.status = 'failed' raise RuntimeError(self.detail) self._schedule_failure(now, self.detail) return self.status if code is not None: if self.log_pump is not None: self.log_pump.join(timeout=5) if self.log_pump.error: self.detail = f'bounded dashboard log writer failed: {self.log_pump.error}' self.log_pump = None self.process = None self._generation += 1 if self._future is not None: self._future.cancel() self._future = None self._future_generation = None self._schedule_failure(now, f'dashboard process exited with code {format_exit_code(code)}') return self.status if self._future is None and now >= self.next_health_at: self._future = self._executor.submit(self.health_probe, self.url, self.health_timeout_sec) self._future_generation = self._generation return self.status def stop(self, timeout=15, final=False, preserve_desired=False): if not preserve_desired: self.desired_state = 'stopped' if self._future is not None: self._future.cancel() self._future = None self._future_generation = None self._generation += 1 stopped = self._stop_current(timeout=timeout) if stopped: self.fatal_stop_failure = False self.status = 'blocked' if preserve_desired and self.desired_state == 'running' else 'stopped' self.detail = 'waiting for stable PostgreSQL readiness' if self.status == 'blocked' else '' else: self.status = 'failed' return stopped def dependency_unavailable(self): if self.process is not None and self.process.poll() is None: return self.stop(preserve_desired=True) if self.desired_state != 'running': return True return self.stop(preserve_desired=True) def dependency_available(self): if self.desired_state == 'running': return self.start(force=True, record_intent=False) return True def snapshot(self): pid = self.process.pid if self.process is not None and self.process.poll() is None else None return { 'status': self.status, 'desired': self.desired_state, 'healthy': self.healthy, 'pid': pid, 'failures': self.failures, 'detail': self.detail, 'url': self.url, } def close(self): if self._owns_executor: self._executor.shutdown(wait=False, cancel_futures=True) def handle_dashboard_command(parts, context): manager = context.get('dashboard_manager') if context else None if manager is None: command_failure(context, 'Dashboard control is unavailable in this supervisor context') return True action = parts[1].lower() if len(parts) > 1 else 'status' if action == 'status': snapshot = manager.snapshot() pid = f" pid={snapshot['pid']}" if snapshot.get('pid') else '' print(f"dashboard: {snapshot['status']}{pid}: {snapshot['detail']}") return True if action == 'stop': if manager.stop(): print('dashboard: stopped') else: command_failure(context, 'dashboard: stop failed') return True if action == 'start': started = manager.start(force=True) snapshot = manager.snapshot() print(f"dashboard: {snapshot['status']}: {snapshot['detail']}") if not started and snapshot['status'] != 'blocked': if context is not None: context['_command_failed'] = True return True if action == 'restart': if not manager.stop(): command_failure(context, 'dashboard: restart refused because the current process did not stop') return True started = manager.start(force=True) snapshot = manager.snapshot() print(f"dashboard: {snapshot['status']}: {snapshot['detail']}") if not started and snapshot['status'] != 'blocked': if context is not None: context['_command_failed'] = True return True command_failure(context, 'Usage: dashboard status|stop|start|restart') return True def load_supervisor_runtime( config_path, selected_sources_arg=None, *, managed_postgres=None, final_cutover=None, ): config = load_yaml( config_path, managed_postgres=managed_postgres, final_cutover=final_cutover, ) return supervisor_runtime_from_config( config_path, config, selected_sources_arg, ) def supervisor_runtime_from_config(config_path, config, selected_sources_arg=None): global_config = config.get('global') or {} project_dir = global_config.get('project_dir') or os.path.dirname(os.path.abspath(config_path)) results_dir = global_config.get('results_dir') supervisor_config = dict(config.get('supervisor') or {}) supervisor_config.setdefault('interval', int(global_config.get('cooldown', 300) or 300)) selected_sources = get_enabled_sources( config, selected_sources_arg or supervisor_config.get('enabled_sources') or supervisor_config.get('sources_enabled'), supervisor_config, ) return config, project_dir, results_dir, supervisor_config, selected_sources def validate_managed_runtime_startup(config_path): from runtime_document_io import validate_managed_runtime_files return validate_managed_runtime_files(config_path) def keychecks_config_for(config_path, config): global_config = config.get('global') or {} keychecks_config = dict(config.get('keychecks') or {}) keycheck_dir = keychecks_config.get('keycheck_dir') or global_config.get('keycheck_dir') or os.path.join(os.path.dirname(global_config.get('results_dir') or ''), 'keychecks') keychecks_config['keycheck_dir'] = keycheck_dir keychecks_config.setdefault('summary_tsv', os.path.join(keycheck_dir, 'summary.tsv')) keychecks_config.setdefault('summary_json', os.path.join(keycheck_dir, 'summary.json')) keychecks_config.setdefault('alive_summary_tsv', os.path.join(keycheck_dir, 'alive_summary.tsv')) for key in ('input', 'proxy_file', 'summary_tsv', 'summary_json', 'alive_summary_tsv', 'keycheck_dir'): if keychecks_config.get(key): keychecks_config[key] = resolve_optional_path(keychecks_config[key], global_config) return keychecks_config def should_manage_keychecks(config, supervisor_config, selected_sources_arg=None): keychecks_config = config.get('keychecks') or {} if not bool_value(keychecks_config.get('enabled'), False): return False return True def reload_supervisor_config(managed_sources, context): config_path = context['config_path'] selected_sources_arg = context.get('selected_sources_arg') global_force_once = bool_value(context.get('global_force_once'), False) dependency_gate = context.get('dependency_gate') try: config, project_dir, results_dir, supervisor_config, selected_sources = load_supervisor_runtime( config_path, selected_sources_arg, managed_postgres=bool(context.get('with_postgres')), ) except Exception as e: print(f'Reload failed: {e}') return False current = source_map(managed_sources) include_keychecks = should_manage_keychecks(config, supervisor_config, selected_sources_arg) selected_set = set(selected_sources) if include_keychecks: selected_set.add('keychecks') next_sources = [] added = [] updated = [] kept_running = [] removed = [] for source_name in selected_sources: source_config = (config.get('sources') or {}).get(source_name, {}) options = source_options(source_name, supervisor_config, source_config) enabled = bool_value(options.get('enabled'), True) existing = current.get(source_name) if existing: if not enabled and existing.is_running(): kept_running.append(source_name) next_sources.append(existing) elif not enabled: if dependency_gate is not None: dependency_gate.unregister(existing) removed.append(source_name) else: was_running = existing.is_running() existing.reconfigure(results_dir, supervisor_config, source_config, global_force_once) updated.append(source_name + (' (restart to apply to running child)' if was_running else '')) next_sources.append(existing) elif enabled: item = ManagedSource( source_name, config_path, project_dir, results_dir, supervisor_config, source_config, global_force_once, dependency_gate, ) added.append(source_name) next_sources.append(item) if include_keychecks: keychecks_config = keychecks_config_for(config_path, config) existing = current.get('keychecks') if existing: was_running = existing.is_running() existing.reconfigure(results_dir, supervisor_config, keychecks_config, global_force_once) updated.append('keychecks' + (' (restart to apply to running child)' if was_running else '')) next_sources.append(existing) else: item = ManagedKeychecks( config_path, project_dir, results_dir, supervisor_config, keychecks_config, global_force_once, dependency_gate, ) added.append('keychecks') next_sources.append(item) for source in managed_sources: if source.source in selected_set: continue if source.is_running(): kept_running.append(source.source) next_sources.append(source) else: if dependency_gate is not None: dependency_gate.unregister(source) removed.append(source.source) managed_sources[:] = next_sources context['config'] = config context['project_dir'] = project_dir context['results_dir'] = results_dir context['supervisor_config'] = supervisor_config context['selected_sources'] = selected_sources context['keychecks_config'] = keychecks_config_for(config_path, config) print('Reloaded config.yaml') if added: print('Added:', ', '.join(added)) if updated: print('Updated:', ', '.join(updated)) if removed: print('Removed:', ', '.join(removed)) if kept_running: print('Kept running until stopped/restarted:', ', '.join(kept_running)) if not any((added, updated, removed, kept_running)): print('No source changes') return True def handle_command(command, managed_sources, context=None): if context is not None: context['_command_failed'] = False command = command.strip() if not command: return True try: parts = shlex.split(command) except ValueError as e: print(f'Invalid command: {e}') return True if not parts: return True action = parts[0].lower() if command_is_mutating(parts) and action != 'shutdown' and not lifecycle_start_allowed(context): command_failure(context, f'supervisor lifecycle is {lifecycle_phase(context)}; mutation is refused') return True if command_is_mutating(parts) and action != 'shutdown' and not check_runtime_authority(context): command_failure(context, (context or {}).get('authority_drift') or 'runtime authority drifted') return True if action in ('help', 'h', '?'): print(HELP_TEXT) return True if action in ('quit', 'exit', 'q'): return False if action == 'shutdown': shutdown_event = context.get('shutdown_event') if context else None if shutdown_event is None: command_failure(context, 'Shutdown control is unavailable in this supervisor context') else: begin_stopping(context) print('Coordinated shutdown requested') return True if action in ('status', 's'): poll_managed_sources(managed_sources, context) print_table(managed_sources, clear=False, context=context) return True if action == 'auth': if len(parts) < 2: print('Usage: auth ') return True targets = select_sources(managed_sources, parts[1]) for source in targets: print_auth_status(source) return True if action in ('watch', 'w'): print('watch is only available in foreground supervisor or --attach prompt.') return True if action in ('reload', 'r'): command_failure(context, 'Live reload is disabled by config authority binding; use authenticated coordinated shutdown and restart the supervisor.') return True if action == 'dashboard': return handle_dashboard_command(parts, context) if action == 'recheck': return handle_recheck_command(parts, managed_sources, context) if action in ('start', 'stop', 'restart', 'pause', 'resume', 'once', 'command'): if len(parts) < 2: print(f'Usage: {action} ') return True targets = select_sources(managed_sources, parts[1]) for source in targets: if isinstance(source, ManagedDiscoveryProducer) and action == 'once': command_failure(context, f'{source.source}: producer mode is fixed to once+repeat') continue if getattr(source, 'manual_only', False) and parts[1].lower() == 'all' and action == 'start': continue if getattr(source, 'manual_only', False) and action not in ('start', 'stop', 'command'): command_failure( context, f'{source.source}: only explicit start, stop, command, and logs are supported', ) continue if ( isinstance(source, ManagedPipelineWorker) and source.source in ('result-ingester', 'jsonl-projector') and action in ('stop', 'restart', 'pause', 'once') and any( item.is_running() for item in managed_sources if not isinstance(item, ManagedPipelineWorker) and item.source != 'keychecks' ) ): command_failure( context, f'{source.source}: manual lifecycle change is refused while scanner sources are running', ) continue if action == 'start': started = source.start(force=True) if started: print(f'{source.source}: started') elif source.runtime_blocked: print(f'{source.source}: running intent recorded; dependency blocked') else: command_failure(context, f'{source.source}: start failed: {source.last_action_error or "process did not start"}') elif action == 'stop': if source.stop(): print(f'{source.source}: stopped') else: command_failure(context, f'{source.source}: stop failed: {source.last_action_error or "process remained live"}') elif action == 'restart': started = source.restart_now() if started: print(f'{source.source}: restarted') elif source.last_action_error: command_failure(context, f'{source.source}: restart failed: {source.last_action_error}') else: print(f'{source.source}: restart intent recorded; dependency blocked') elif action == 'pause': if source.pause(): print(f'{source.source}: paused') else: command_failure(context, f'{source.source}: pause failed: {source.last_action_error or "process remained live"}') elif action == 'resume': started = source.resume() if started: print(f'{source.source}: resumed') elif source.runtime_blocked: print(f'{source.source}: resume intent recorded; dependency blocked') else: command_failure(context, f'{source.source}: resume failed: {source.last_action_error or "process did not start"}') elif action == 'once': set_mode(source, 'once') started = source.start(force=True) if started: print(f'{source.source}: once started') elif source.runtime_blocked: print(f'{source.source}: once intent recorded; dependency blocked') else: command_failure(context, f'{source.source}: once start failed: {source.last_action_error or "process did not start"}') elif action == 'command': print_source_command(source) return True if action == 'mode': if len(parts) < 3: print('Usage: mode loop|once|repeat') return True for source in select_sources(managed_sources, parts[1]): if isinstance(source, ManagedDiscoveryProducer): command_failure(context, f'{source.source}: producer mode changes are forbidden') continue if getattr(source, 'manual_only', False): command_failure(context, f'{source.source}: mode changes are forbidden') continue if set_mode(source, parts[2]): print(f'{source.source}: mode={source.mode_label()}') return True if action == 'set': if len(parts) < 4: print('Usage: set interval|restart|restart_delay ') return True key = parts[2].lower() value = parts[3] for source in select_sources(managed_sources, parts[1]): if isinstance(source, ManagedDiscoveryProducer) and key != 'interval': command_failure( context, f'{source.source}: only producer interval may be changed at runtime', ) continue if getattr(source, 'manual_only', False): command_failure(context, f'{source.source}: runtime option changes are forbidden') continue if key == 'interval': try: source.interval = max(0, int(value)) except ValueError: print('interval must be an integer number of seconds') continue print(f'{source.source}: interval={source.interval}') elif key == 'restart_delay': try: source.restart_delay = max(0, int(value)) except ValueError: print('restart_delay must be an integer number of seconds') continue print(f'{source.source}: restart_delay={source.restart_delay}') elif key == 'restart': source.restart = bool_value(value, source.restart) print(f'{source.source}: restart={source.restart}') else: print(f'Unknown set key: {key}') return True if action in ('logs', 'tail'): if len(parts) < 2: print('Usage: logs [lines]') return True targets = select_sources(managed_sources, parts[1]) if not targets: return True try: limit = int(parts[2]) if len(parts) > 2 else 40 except ValueError: print('lines must be an integer') return True source = targets[0] print(f'--- {source.log_path} (last {limit}) ---') for line in source.tail_log_lines(limit): print(console_safe_text(line)) print('--- end log ---') return True command_failure(context, f'Unknown command: {action}. Type `help`.') return True def read_single_key(): if os.name == 'nt': import msvcrt if not msvcrt.kbhit(): return None ch = msvcrt.getwch() if ch in ('\x00', '\xe0'): if msvcrt.kbhit(): msvcrt.getwch() return None return ch import select ready, _, _ = select.select([sys.stdin], [], [], 0) if ready: return sys.stdin.read(1) return None def poll_managed_sources(managed_sources, context): for source in managed_sources: if not lifecycle_start_allowed(context): return source.poll() if context is not None and not context.get('background_child') and source.status == 'failed': context['runtime_failed'] = True def watch_local(managed_sources, poll_sec=0.5, context=None): if not sys.stdin.isatty() or not sys.stdout.isatty() or not enable_ansi_terminal(): print('watch requires an interactive ANSI terminal') return sys.stdout.write(ANSI_ALT_SCREEN) sys.stdout.flush() try: while True: guard = (context or {}).get('control_lock') or nullcontext() pipeline_snapshot = pipeline_status_snapshot(context) with guard: if shutdown_checkpoint(context): return tick_supervisor_runtime(context, pipeline_snapshot=pipeline_snapshot) poll_managed_sources(managed_sources, context) sys.stdout.write(ANSI_HOME + ANSI_CLEAR_SCREEN) print('\n'.join(build_runtime_table_lines(managed_sources, context=context))) print('\nwatch mode: press q to return') sys.stdout.flush() deadline = time.time() + max(0.1, float(poll_sec or 0.5)) while time.time() < deadline: with guard: if shutdown_checkpoint(context): return key = read_single_key() if key and key.lower() == 'q': return time.sleep(0.05) finally: sys.stdout.write(ANSI_MAIN_SCREEN) sys.stdout.flush() def watch_remote(metadata, poll_sec=0.5): if not sys.stdin.isatty() or not sys.stdout.isatty() or not enable_ansi_terminal(): print('watch requires an interactive ANSI terminal') return sys.stdout.write(ANSI_ALT_SCREEN) sys.stdout.flush() try: while True: snapshot = get_control_snapshot(metadata) sys.stdout.write(ANSI_HOME + ANSI_CLEAR_SCREEN) print((snapshot.get('table') or '').rstrip()) print('\nwatch mode: press q to return') sys.stdout.flush() deadline = time.time() + max(0.1, float(poll_sec or 0.5)) while time.time() < deadline: key = read_single_key() if key and key.lower() == 'q': return time.sleep(0.05) finally: sys.stdout.write(ANSI_MAIN_SCREEN) sys.stdout.flush() def interactive_loop_blocking(managed_sources, autostart=False, clear=True, context=None, poll_sec=0.5): guard = (context or {}).get('control_lock') or nullcontext() with guard: if shutdown_checkpoint(context): return if autostart: for source in autostart_sources(managed_sources): if shutdown_checkpoint(context) or not lifecycle_start_allowed(context): break source.start(force=True) print_table(managed_sources, clear=clear, context=context) print('\n' + HELP_TEXT + '\n') commands = queue.Queue() allow_read = threading.Event() allow_read.set() def read_commands(): while True: allow_read.wait() allow_read.clear() try: commands.put(input('supervisor> ')) except (EOFError, KeyboardInterrupt): commands.put(None) return threading.Thread(target=read_commands, daemon=True).start() while True: pipeline_snapshot = pipeline_status_snapshot(context) with guard: if shutdown_checkpoint(context): break tick_supervisor_runtime(context, pipeline_snapshot=pipeline_snapshot) poll_managed_sources(managed_sources, context) if shutdown_checkpoint(context): break try: command = commands.get(timeout=max(0.05, float(poll_sec or 0.2))) except queue.Empty: continue if command is None: print() break if command.strip().lower() in ('watch', 'w'): watch_local(managed_sources, poll_sec, context) allow_read.set() continue with guard: if shutdown_checkpoint(context): break keep_running = handle_command(command, managed_sources, context) if not keep_running: break allow_read.set() print('Stopping child processes...') def interactive_loop(managed_sources, autostart=False, clear=True, context=None, poll_sec=0.2): interactive_loop_blocking(managed_sources, autostart, clear, context, poll_sec) def tick_supervisor_runtime(context, pipeline_snapshot=None): context = context or {} context.pop('_defer_pipeline_status_until_next_tick', None) if shutdown_checkpoint(context): return if not lifecycle_start_allowed(context): if context.get('authority_drift'): inhibit_for_authority_drift(context, context['authority_drift']) controller = context.get('postgres_controller') if controller is not None: controller.tick() return now = time.monotonic() if now >= float(context.get('next_authority_check_at') or 0): interval = max(0.2, float((context.get('supervisor_config') or {}).get('authority_check_interval_sec', 5) or 5)) context['next_authority_check_at'] = now + interval if not check_runtime_authority(context): return elif context.get('authority_drift'): return controller = context.get('postgres_controller') gate = context.get('dependency_gate') if controller is None: if gate is not None: gate.set_ready(True) else: previous = context.get('postgres_reported_state') controller.tick() snapshot = controller.snapshot() current = snapshot['state'] if current != previous: detail = f": {snapshot['detail']}" if snapshot.get('detail') else '' print(f'PostgreSQL controller: {current}{detail}') context['postgres_reported_state'] = current if gate is not None: try: gate_was_ready = gate.ready gate.set_ready(controller.ready) if not gate_was_ready and gate.ready: context['_defer_pipeline_status_until_next_tick'] = True context['_pipeline_status_refresh_at'] = 0 except DependencyStopError as exc: controller.inhibit_lifecycle(str(exc)) context['fatal_child_stop_error'] = str(exc) shutdown_event = context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() raise source_gate = context.get('source_dependency_gate') if ( source_gate is not None and not source_gate.ready and not context.get('_defer_pipeline_status_until_next_tick') ): snapshot = ( pipeline_snapshot if pipeline_snapshot is not None else pipeline_status_snapshot(context) ) if snapshot.get('ingester_ready'): source_gate.set_ready(True) context['pipeline_initial_ready'] = True print('Result ingester reported ready; source admission gate is open.') dashboard = context.get('dashboard_manager') if dashboard is not None: try: dashboard.poll() except Exception as exc: if controller is not None: controller.inhibit_lifecycle(f'dashboard stop failure: {exc}') context['fatal_child_stop_error'] = str(exc) shutdown_event = context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() raise def non_interactive_loop(managed_sources, refresh_sec=DEFAULT_REFRESH_SEC, status_file=None, poll_sec=1.0, context=None, lock=None, stay_alive=False, status_heartbeat_sec=60): last_status_signature = None last_status_write = 0 guard = lock or (context or {}).get('control_lock') or nullcontext() while True: pipeline_snapshot = pipeline_status_snapshot(context) with guard: if shutdown_checkpoint(context): return tick_supervisor_runtime(context, pipeline_snapshot=pipeline_snapshot) poll_managed_sources(managed_sources, context) if shutdown_checkpoint(context): return if context.get('_defer_pipeline_status_until_next_tick'): time.sleep(max(0.2, float(poll_sec or 1.0))) continue with guard: if shutdown_checkpoint(context): return controller = context.get('postgres_controller') if context else None postgres_signature = None if controller is not None: snapshot = controller.snapshot() postgres_signature = (snapshot['state'], snapshot['failures'], snapshot['detail']) dashboard = context.get('dashboard_manager') if context else None dashboard_signature = None if dashboard is not None: dashboard_snapshot = dashboard.snapshot() dashboard_signature = tuple(dashboard_snapshot.get(key) for key in ('status', 'desired', 'healthy', 'pid', 'failures', 'detail')) scan_snapshot = current_scan_worker_snapshot(context) scan_signature = ( scan_snapshot['active'], scan_snapshot['limit'], scan_snapshot['trufflehog'], tuple(scan_snapshot['sources'].items()), scan_snapshot['detail'], ) pipeline_signature = tuple( pipeline_snapshot.get(key) for key in ( 'ingester_state', 'projector_state', 'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes', 'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes', 'detail', ) ) current_signature = ( table_signature(managed_sources), postgres_signature, dashboard_signature, scan_signature, pipeline_signature, ) now = time.time() heartbeat_due = status_heartbeat_sec and now - last_status_write >= max(1, int(status_heartbeat_sec)) if current_signature != last_status_signature or heartbeat_due: write_status_file(status_file, managed_sources, context) last_status_signature = current_signature last_status_write = now with guard: all_done = all(source.status in ('done', 'failed', 'disabled', 'stopped') for source in managed_sources) dashboard = context.get('dashboard_manager') if context else None if dashboard is not None and dashboard.desired_state == 'running': all_done = False if all_done and not stay_alive: break time.sleep(max(0.2, float(poll_sec or 1.0))) def _coordinated_shutdown_locked(managed_sources, context): context = {} if context is None else context context['coordinated_shutdown_complete'] = False context['authority_release_safe'] = False children_complete = True drain_timeout = max( 0.0, float((context.get('supervisor_config') or {}).get('source_handoff_drain_timeout_sec', 30) or 0), ) deadline = time.monotonic() + drain_timeout while time.monotonic() < deadline: snapshot = scan_worker_snapshot(context.get('config')) if snapshot['active'] == 0: break time.sleep(0.1) ordered = sorted(managed_sources, key=lambda source: { 'worker-api': 1, 'keychecks': 2, 'jsonl-projector': 3, 'result-ingester': 4, 'janitor': 5, }.get(source.source, 0)) for source in ordered: try: if not source.stop(final=True): children_complete = False print(f'{source.source}: shutdown did not confirm process exit') except Exception as exc: children_complete = False print(f'{source.source}: shutdown failed: {exc}') dashboard = context.get('dashboard_manager') if dashboard is not None: try: if not dashboard.stop(final=True): children_complete = False print('dashboard: shutdown did not confirm process exit') except Exception as exc: children_complete = False print(f'dashboard: shutdown failed: {exc}') for source in ordered: try: if source.is_running() is not False: children_complete = False print(f'{source.source}: process exit is unconfirmed after bounded shutdown') except Exception: children_complete = False if dashboard is not None: process = dashboard.process try: if process is not None and process.poll() is None: children_complete = False print('dashboard: process is still live after bounded shutdown') except Exception: children_complete = False controller = context.get('postgres_controller') if controller is None: if dashboard is not None and children_complete: dashboard.close() context['coordinated_shutdown_complete'] = children_complete context['authority_release_safe'] = children_complete return children_complete if not children_complete: print('PostgreSQL controller close/stop was deferred because one or more managed children did not stop.') return False supervisor_config = context.get('supervisor_config') or {} timeout = max(5.0, float(supervisor_config.get('postgres_shutdown_timeout_sec', 120) or 120)) deadline = time.monotonic() + timeout if bool(getattr(controller, 'lifecycle_action_required', True)): controller.request_stop() while not controller.terminal and time.monotonic() < deadline: controller.tick() time.sleep(min(0.05, max(0.0, deadline - time.monotonic()))) if not controller.terminal: print(f'PostgreSQL controller stop did not complete within {timeout:g}s; no unverified process action was taken.') try: close_result = controller.close(wait=False, timeout_sec=max(0.0, deadline - time.monotonic())) except TypeError: close_result = controller.close(wait=False) if controller.detail: print(f'PostgreSQL controller: {controller.state.value}: {controller.detail}') close_safe = bool(close_result) and bool(getattr(controller, 'authority_release_safe', False)) if dashboard is not None and close_safe: dashboard.close() context['coordinated_shutdown_complete'] = ( children_complete and bool(getattr(controller, 'terminal', False)) and bool(getattr(controller, 'stop_succeeded', False)) and close_safe ) context['authority_release_safe'] = context['coordinated_shutdown_complete'] return context['coordinated_shutdown_complete'] def coordinated_shutdown(managed_sources, context): context = {} if context is None else context guard = context.get('control_lock') or nullcontext() with guard: context['authority_release_safe'] = False try: begin_stopping(context) complete = _coordinated_shutdown_locked(managed_sources, context) if not complete: enter_failed_hold(context, 'one or more owned processes did not confirm stopped disposition') return complete except BaseException: try: enter_failed_hold(context, 'coordinated shutdown was interrupted or failed') except BaseException: pass raise def retain_unsafe_authority(managed_sources, context, max_attempts=None): """Retry bounded stops while retaining locks, metadata, and authenticated control.""" context = {} if context is None else context if context.get('authority_release_safe', False): return True guard = context.get('control_lock') or nullcontext() try: with guard: enter_failed_hold(context, context.get('shutdown_failure')) except BaseException: pass retry_delay = 30.0 attempts = 0 retry_event = context.get('shutdown_retry_event') while True: try: attempts += 1 supervisor_config = context.get('supervisor_config') or {} retry_delay = max(1.0, float(supervisor_config.get('postgres_stop_failed_retry_sec', 30) or 30)) # FAILED_HOLD permits only read-only status and another shutdown # request, so retries need not monopolize authenticated control. if _coordinated_shutdown_locked(managed_sources, context): context['shutdown_failure'] = '' context['authority_release_safe'] = True return True with guard: enter_failed_hold(context, 'bounded shutdown retry did not confirm every owned process stopped') print('FAILED_HOLD: retaining all OS locks and authenticated control until every owned process is stopped.') write_status_file(context.get('status_file'), managed_sources, context) if max_attempts is not None and attempts >= max(1, int(max_attempts)): return False if retry_event is not None: retry_event.wait(retry_delay) retry_event.clear() else: time.sleep(retry_delay) except BaseException as exc: context['authority_release_safe'] = False try: with guard: enter_failed_hold(context, f'shutdown retry failed: {type(exc).__name__}: {exc}') print(f'Authority retention ignored shutdown interruption: {type(exc).__name__}: {exc}') except BaseException: pass if max_attempts is not None and attempts >= max(1, int(max_attempts)): return False try: if retry_event is not None: retry_event.wait(retry_delay) retry_event.clear() else: time.sleep(retry_delay) except BaseException: pass def retain_unsafe_postgres_authority(context): """Compatibility wrapper for pre-activation PostgreSQL compensation.""" if (context or {}).get('authority_release_safe', False): return True return retain_unsafe_authority([], context) def dashboard_settings(supervisor_config): dashboard_config = supervisor_config.get('dashboard') if not dashboard_config: return False, '127.0.0.1', 5000, None if isinstance(dashboard_config, dict): enabled = bool_value(dashboard_config.get('enabled'), False) port = int(dashboard_config.get('port', 5000)) host = str(dashboard_config.get('address') or dashboard_config.get('host') or '127.0.0.1') else: enabled = bool_value(dashboard_config, False) port = 5000 host = '127.0.0.1' if not is_loopback_host(host): raise ValueError(f'dashboard address must be loopback-only: {host}') if not 0 < port <= 65535: raise ValueError(f'invalid dashboard port: {port}') return enabled, host, port, f'http://{host}:{port}' def background_paths(config_path, results_dir, supervisor_config, explicit_instance_file=None, explicit_pid_file=None): log_dir = resolve_path(config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) require_private_directory(log_dir, create=False) control_dir = resolve_path( config_path, supervisor_config.get('control_dir') or os.path.join(os.path.dirname(log_dir), 'control'), ) if os.path.normcase(os.path.abspath(control_dir)) == os.path.normcase(os.path.abspath(log_dir)): raise ValueError('private supervisor control directory must be separate from the log directory') require_private_directory(control_dir, create=False) if explicit_instance_file: instance_file = resolve_path(config_path, explicit_instance_file) if os.path.normcase(os.path.dirname(os.path.abspath(instance_file))) != os.path.normcase(os.path.abspath(control_dir)): raise ValueError('custom supervisor instance metadata must remain in the private control directory') else: instance_file = resolve_path(config_path, supervisor_config.get('instance_file') or os.path.join(control_dir, 'supervisor.instance.json')) if os.path.normcase(os.path.dirname(os.path.abspath(instance_file))) != os.path.normcase(os.path.abspath(control_dir)): raise ValueError('supervisor instance metadata must reside directly in the private control directory') reject_reparse_components(os.path.dirname(os.path.abspath(instance_file))) if not private_directory_ready(os.path.dirname(os.path.abspath(instance_file))): raise ValueError('supervisor instance parent directory is not private') status_file = resolve_path(config_path, supervisor_config.get('status_file') or os.path.join(log_dir, 'supervisor.status.txt')) log_file = resolve_path(config_path, supervisor_config.get('supervisor_log') or os.path.join(log_dir, 'supervisor.log')) legacy_pid_file = resolve_path(config_path, explicit_pid_file) if explicit_pid_file else os.path.join(log_dir, 'supervisor.pid') background_lock_path(config_path, results_dir, supervisor_config) return instance_file, log_file, status_file, legacy_pid_file def background_lock_path(config_path, results_dir, supervisor_config): log_dir = resolve_path(config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) control_dir = resolve_path( config_path, supervisor_config.get('control_dir') or os.path.join(os.path.dirname(log_dir), 'control'), ) if os.path.normcase(os.path.abspath(control_dir)) == os.path.normcase(os.path.abspath(log_dir)): raise ValueError('private supervisor control directory must be separate from the log directory') require_private_directory(control_dir, create=False) lock_file = resolve_path(config_path, supervisor_config.get('lock_file') or os.path.join(control_dir, 'supervisor.lock')) if os.path.normcase(os.path.dirname(os.path.abspath(lock_file))) != os.path.normcase(os.path.abspath(control_dir)): raise ValueError('supervisor singleton lock must reside directly in the private control directory') reject_reparse_components(os.path.dirname(os.path.abspath(lock_file))) return lock_file def legacy_log_instance_path(config_path, results_dir, supervisor_config): log_dir = resolve_path(config_path, supervisor_config.get('log_dir') or os.path.join(results_dir, 'logs')) return os.path.join(log_dir, 'supervisor.instance.json') def read_pid_file(path): try: with open(path, 'r', encoding='utf-8') as f: return int(f.read().strip()) except (OSError, ValueError): return None def write_status_file(status_file, managed_sources, context=None): if not status_file: return parent = os.path.dirname(status_file) if parent: require_private_directory(parent, create=False) if os.path.lexists(status_file) and not private_file_ready(status_file): raise OSError(f'private status file ACL is not ready; run offline hardening: {status_file}') tmp_path = f'{status_file}.{os.getpid()}.{threading.get_ident()}.tmp' with open(tmp_path, 'x', encoding='utf-8') as f: controller = context.get('postgres_controller') if context else None if controller is not None: snapshot = controller.snapshot() f.write(f"PostgreSQL: {snapshot['state']} failures={snapshot['failures']} detail={snapshot['detail']}\n\n") dashboard = context.get('dashboard_manager') if context else None if dashboard is not None: snapshot = dashboard.snapshot() f.write( f"Dashboard: {snapshot['status']} desired={snapshot['desired']} " f"healthy={snapshot['healthy']} detail={snapshot['detail']}\n\n" ) f.write('\n'.join(build_runtime_table_lines(managed_sources, context=context))) f.write('\n') harden_private_file(tmp_path) for attempt in range(5): try: os.replace(tmp_path, status_file) if not private_file_ready(status_file): raise OSError(f'private status file ACL changed during update: {status_file}') return except PermissionError as e: if attempt == 4: print(f'Unable to update status file {status_file}: {e}') break time.sleep(0.1 * (attempt + 1)) try: if os.path.exists(tmp_path): os.remove(tmp_path) except OSError: pass def control_address(supervisor_config): host = str(supervisor_config.get('control_host') or '127.0.0.1') port = int(supervisor_config.get('control_port') or 8765) if not is_loopback_host(host): raise ValueError(f'control address must be loopback-only: {host}') if not 0 <= port <= 65535: raise ValueError(f'invalid control port: {port}') return host, port def send_control_request(metadata, action, command=None, timeout=60, **extra): control = metadata['control'] request = { 'schema': CONTROL_SCHEMA, 'instance_id': metadata['instance_id'], 'token': metadata['token'], 'action': str(action), } if command is not None: request['command'] = str(command) request.update(extra) payload = json.dumps(request, ensure_ascii=True, separators=(',', ':')).encode('utf-8') + b'\n' if len(payload) > MAX_CONTROL_REQUEST_BYTES: raise ValueError('control request is too large') timeout_seconds = max(0.1, float(timeout)) deadline = time.monotonic() + timeout_seconds def remaining_timeout(): remaining = deadline - time.monotonic() if remaining <= 0: raise TimeoutError('control request timed out') return remaining with socket.create_connection( (control['host'], control['port']), timeout=remaining_timeout(), ) as sock: sock.settimeout(remaining_timeout()) sock.sendall(payload) sock.shutdown(socket.SHUT_WR) chunks = [] total = 0 while True: sock.settimeout(remaining_timeout()) chunk = sock.recv(65536) if not chunk: break total += len(chunk) if total > MAX_CONTROL_RESPONSE_BYTES: raise ValueError('control response is too large') chunks.append(chunk) raw = b''.join(chunks) try: response = json.loads(raw.decode('utf-8')) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise ValueError('invalid control response') from exc if not isinstance(response, dict) or response.get('schema') != CONTROL_SCHEMA: raise ValueError('invalid control response schema') if response.get('instance_id') != metadata['instance_id']: raise ValueError('control response instance mismatch') if type(response.get('ok')) is not bool: raise ValueError('invalid control response status') expected_fields = ( {'schema', 'instance_id', 'ok', 'result'} if response['ok'] else {'schema', 'instance_id', 'ok', 'error'} ) if set(response) != expected_fields: raise ValueError('invalid control response fields') if not response['ok']: if type(response.get('error')) is not str: raise ValueError('invalid control response error') raise RuntimeError(str(response.get('error') or 'control request failed')) return response.get('result') def send_control_command(metadata, command, timeout=60): return send_control_request(metadata, 'command', command=command, timeout=timeout) def get_control_snapshot(metadata): result = send_control_request(metadata, 'snapshot') if not isinstance(result, dict): raise ValueError('invalid supervisor snapshot') return result _STRUCTURED_SOURCE_FIELDS = frozenset({ 'id', 'source', 'role', 'lifecycle_state', 'desired_state', 'process_state', 'pid', 'enabled', 'dependency_blocked', 'startup_cleanup_pending', 'mode', 'interval_seconds', 'restart_enabled', 'restart_delay_seconds', 'restart_count', 'restart_streak', 'last_exit_code', 'last_exit_at', 'next_scheduled_run_at', 'safe_error_category', 'auth_summary', 'allowed_actions', }) _STRUCTURED_PRODUCER_FIELDS = frozenset({ 'last_cycle_result', 'last_successful_discovery_at', }) _STRUCTURED_AUTH_FIELDS = frozenset({ 'total', 'ok', 'dead', 'limited', 'rate_limit_errors', 'auth_invalid_errors', }) _STRUCTURED_DASHBOARD_FIELDS = frozenset({ 'id', 'status', 'desired_state', 'process_state', 'healthy', 'pid', 'restart_count', 'safe_error_category', 'allowed_actions', }) def _nonnegative_control_integer(value): return type(value) is int and value >= 0 def _valid_structured_source_state(source, source_id=None): if not isinstance(source, dict): return False fields = frozenset(source) producer = fields == _STRUCTURED_SOURCE_FIELDS.union(_STRUCTURED_PRODUCER_FIELDS) if fields != _STRUCTURED_SOURCE_FIELDS and not producer: return False if source_id is not None and source.get('id') != source_id: return False if any( type(source.get(key)) is not str for key in ( 'id', 'source', 'role', 'lifecycle_state', 'desired_state', 'process_state', 'mode', 'safe_error_category', ) ): return False if source['process_state'] not in ('running', 'stopped'): return False if not ( source['pid'] is None or (type(source['pid']) is int and source['pid'] > 0) ): return False if any( type(source.get(key)) is not bool for key in ( 'enabled', 'dependency_blocked', 'startup_cleanup_pending', 'restart_enabled', ) ): return False if any( not _nonnegative_control_integer(source.get(key)) for key in ( 'interval_seconds', 'restart_delay_seconds', 'restart_count', 'restart_streak', ) ): return False if source['last_exit_code'] is not None and type(source['last_exit_code']) is not int: return False if any( source[key] is not None and type(source[key]) is not str for key in ('last_exit_at', 'next_scheduled_run_at') ): return False auth = source['auth_summary'] if not isinstance(auth, dict) or ( auth and ( frozenset(auth) != _STRUCTURED_AUTH_FIELDS or any(not _nonnegative_control_integer(value) for value in auth.values()) ) ): return False if type(source['allowed_actions']) is not list or any( type(action) is not str for action in source['allowed_actions'] ): return False if not producer: return True cycle = source['last_cycle_result'] return ( source['role'] == DISCOVERY_PRODUCER_ROLE and isinstance(cycle, dict) and set(cycle) == { 'status', 'fetched_count', 'queued_new_count', 'queued_updated_count', } and type(cycle.get('status')) is str and all( _nonnegative_control_integer(cycle.get(key)) for key in ('fetched_count', 'queued_new_count', 'queued_updated_count') ) and ( source['last_successful_discovery_at'] is None or type(source['last_successful_discovery_at']) is str ) ) def _valid_structured_dashboard_state(dashboard): return ( isinstance(dashboard, dict) and frozenset(dashboard) == _STRUCTURED_DASHBOARD_FIELDS and dashboard.get('id') == 'dashboard' and all( type(dashboard.get(key)) is str for key in ( 'id', 'status', 'desired_state', 'process_state', 'safe_error_category', ) ) and dashboard.get('process_state') in ('running', 'stopped') and type(dashboard.get('healthy')) is bool and ( dashboard.get('pid') is None or (type(dashboard.get('pid')) is int and dashboard['pid'] > 0) ) and _nonnegative_control_integer(dashboard.get('restart_count')) and type(dashboard.get('allowed_actions')) is list and all(type(action) is str for action in dashboard['allowed_actions']) ) def get_runtime_snapshot(metadata, timeout=60): result = send_control_request(metadata, 'runtime-snapshot', timeout=timeout) runtime_fields = { 'pid', 'phase', 'manages_postgres', 'start_gate_open', 'shutdown_requested', 'runtime_failed', } postgres_base_fields = {'state', 'ready', 'failures', 'safe_error_category'} postgres_live_fields = postgres_base_fields.union({ 'stop_succeeded', 'lifecycle_inert', 'automatic_inhibited', 'authority_release_safe', 'inflight_start', 'lifecycle_action_required', }) pipeline_fields = { 'ingester_ready', 'projector_ready', 'cutover_ready', 'ingester_state', 'projector_state', 'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes', 'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes', } scan_worker_fields = { 'active', 'limit', 'base_active', 'base_limit', 'bonus_active', 'bonus_limit', 'trufflehog', 'sources', } valid = isinstance(result, dict) and set(result) == { 'snapshot_schema', 'runtime', 'postgres', 'dashboard', 'sources', 'pipeline', 'scan_workers', } if valid: runtime = result['runtime'] postgres = result['postgres'] pipeline = result['pipeline'] scan_workers = result['scan_workers'] valid = ( result['snapshot_schema'] == RUNTIME_SNAPSHOT_SCHEMA and isinstance(runtime, dict) and set(runtime) == runtime_fields and type(runtime.get('pid')) is int and runtime['pid'] > 0 and type(runtime.get('phase')) is str and all( type(runtime.get(key)) is bool for key in ( 'manages_postgres', 'start_gate_open', 'shutdown_requested', 'runtime_failed', ) ) and isinstance(postgres, dict) and set(postgres) in (postgres_base_fields, postgres_live_fields) and type(postgres.get('state')) is str and type(postgres.get('ready')) is bool and _nonnegative_control_integer(postgres.get('failures')) and type(postgres.get('safe_error_category')) is str ) if valid and set(postgres) == postgres_live_fields: valid = all( type(postgres.get(key)) is bool for key in postgres_live_fields - postgres_base_fields ) valid = valid and _valid_structured_dashboard_state(result['dashboard']) valid = valid and type(result['sources']) is list and all( _valid_structured_source_state(source) for source in result['sources'] ) valid = valid and isinstance(pipeline, dict) and set(pipeline) == pipeline_fields if valid: valid = ( type(pipeline['ingester_ready']) is bool and type(pipeline['projector_ready']) is bool and type(pipeline['cutover_ready']) is bool and type(pipeline['ingester_state']) is str and type(pipeline['projector_state']) is str and all( _nonnegative_control_integer(pipeline[key]) for key in pipeline_fields - { 'ingester_ready', 'projector_ready', 'cutover_ready', 'ingester_state', 'projector_state', } ) and isinstance(scan_workers, dict) and set(scan_workers) == scan_worker_fields and all( _nonnegative_control_integer(scan_workers[key]) for key in scan_worker_fields - {'sources'} ) and isinstance(scan_workers['sources'], dict) and all( type(key) is str and _nonnegative_control_integer(value) for key, value in scan_workers['sources'].items() ) ) if not valid: raise ValueError('invalid structured supervisor snapshot') return result def send_managed_source_action(metadata, source_id, source_action, timeout=60, **parameters): result = send_control_request( metadata, 'managed-source-action', timeout=timeout, source_id=source_id, source_action=source_action, **parameters, ) source = result.get('source') if isinstance(result, dict) else None if ( not isinstance(result, dict) or set(result) != {'source_action', 'outcome', 'source'} or result.get('source_action') != source_action or result.get('outcome') not in ('completed', 'dependency-blocked') or not _valid_structured_source_state(source, source_id=source_id) ): raise ValueError('invalid managed source action response') return result def send_managed_source_log_tail(metadata, source_id, line_count, timeout=60): result = send_control_request( metadata, 'managed-source-log-tail', timeout=timeout, source_id=source_id, line_count=line_count, ) if ( not isinstance(result, dict) or set(result) != {'source_id', 'line_count', 'lines', 'response_truncated'} or result.get('source_id') != source_id or type(result.get('lines')) is not list or any(type(line) is not str for line in result['lines']) or type(result.get('line_count')) is not int or result['line_count'] != len(result['lines']) or result['line_count'] > line_count or type(result.get('response_truncated')) is not bool or len(json.dumps( result['lines'], ensure_ascii=True, separators=(',', ':'), ).encode('utf-8')) > MAX_LOG_TAIL_BYTES ): raise ValueError('invalid managed source log tail') return result def send_dashboard_action(metadata, dashboard_action, timeout=60): result = send_control_request( metadata, 'dashboard-action', timeout=timeout, dashboard_action=dashboard_action, ) dashboard = result.get('dashboard') if isinstance(result, dict) else None if ( not isinstance(result, dict) or set(result) != {'dashboard_action', 'outcome', 'dashboard'} or result.get('dashboard_action') != dashboard_action or result.get('outcome') not in ('completed', 'dependency-blocked') or not _valid_structured_dashboard_state(dashboard) ): raise ValueError('invalid dashboard action response') return result _CONTROL_BASE_FIELDS = frozenset({'schema', 'instance_id', 'token', 'action'}) def require_exact_control_fields(request, fields): if set(request) != _CONTROL_BASE_FIELDS.union(fields): raise ValueError('control request fields are invalid for this action') def structured_postgres_state(controller): if controller is None: return { 'state': PostgresState.DISABLED.value, 'ready': True, 'failures': 0, 'safe_error_category': '', } snapshot = controller.snapshot() state = str(snapshot.get('state') or PostgresState.DISABLED.value) return { 'state': state, 'ready': bool(snapshot.get('ready')), 'failures': max(0, int(snapshot.get('failures', 0) or 0)), 'stop_succeeded': bool(snapshot.get('stop_succeeded')), 'lifecycle_inert': bool(snapshot.get('lifecycle_inert')), 'automatic_inhibited': bool(snapshot.get('automatic_inhibited')), 'authority_release_safe': bool(snapshot.get('authority_release_safe')), 'inflight_start': bool(snapshot.get('inflight_start')), 'lifecycle_action_required': bool(snapshot.get('lifecycle_action_required')), 'safe_error_category': 'postgres_error' if 'FAILED' in state.upper() else '', } def structured_dashboard_state(manager): if manager is None: return { 'id': 'dashboard', 'status': 'disabled', 'desired_state': 'stopped', 'process_state': 'stopped', 'healthy': False, 'pid': None, 'restart_count': 0, 'safe_error_category': '', 'allowed_actions': ['start', 'stop', 'restart'], } snapshot = manager.snapshot() status = str(snapshot.get('status') or 'disabled') pid = snapshot.get('pid') if type(snapshot.get('pid')) is int else None safe_error = '' if status in ('failed', 'backoff'): safe_error = 'dashboard_error' return { 'id': 'dashboard', 'status': status, 'desired_state': str(snapshot.get('desired') or 'stopped'), 'process_state': 'running' if pid is not None else 'stopped', 'healthy': bool(snapshot.get('healthy')), 'pid': pid, 'restart_count': max(0, int(snapshot.get('failures', 0) or 0)), 'safe_error_category': safe_error, 'allowed_actions': ['start', 'stop', 'restart'], } def structured_runtime_snapshot(managed_sources, context): context = context or {} pipeline = pipeline_status_snapshot(context, allow_refresh=False) scan_workers = current_scan_worker_snapshot(context) return { 'snapshot_schema': RUNTIME_SNAPSHOT_SCHEMA, 'runtime': { 'pid': os.getpid(), 'phase': lifecycle_phase(context), 'manages_postgres': bool(context.get('with_postgres')), 'start_gate_open': bool(context.get('start_gate_open', True)), 'shutdown_requested': bool(context.get('shutdown_requested')), 'runtime_failed': bool(context.get('runtime_failed')), }, 'postgres': structured_postgres_state(context.get('postgres_controller')), 'dashboard': structured_dashboard_state(context.get('dashboard_manager')), 'sources': [source.structured_state() for source in managed_sources], 'pipeline': { key: pipeline[key] for key in ( 'ingester_ready', 'projector_ready', 'cutover_ready', 'ingester_state', 'projector_state', 'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes', 'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes', ) }, 'scan_workers': { key: scan_workers[key] for key in ( 'active', 'limit', 'base_active', 'base_limit', 'bonus_active', 'bonus_limit', 'trufflehog', 'sources', ) }, } def _pipeline_lifecycle_change_is_blocked(source, source_action, managed_sources): return ( isinstance(source, ManagedPipelineWorker) and source.source in ('result-ingester', 'jsonl-projector') and source_action in ('stop', 'restart', 'pause', 'once') and any( item.is_running() for item in managed_sources if not isinstance(item, ManagedPipelineWorker) and item.source != 'keychecks' ) ) def run_managed_source_action(request, managed_sources, context): source_action = request.get('source_action') expected_fields = {'source_id', 'source_action'} if source_action == 'set-mode': expected_fields.add('mode') elif source_action == 'set-interval': expected_fields.add('interval_seconds') elif source_action == 'set-restart': expected_fields.add('restart_enabled') elif source_action == 'set-restart-delay': expected_fields.add('restart_delay_seconds') require_exact_control_fields(request, expected_fields) if source_action not in MANAGED_SOURCE_LIFECYCLE_ACTIONS + MANAGED_SOURCE_SETTING_ACTIONS: raise ValueError('unsupported managed source action') source_id = request.get('source_id') if not isinstance(source_id, str) or not source_id or len(source_id) > 128: raise ValueError('managed source ID is invalid') source = managed_source_registry(managed_sources).get(source_id) if source is None: raise ValueError('managed source ID is unknown') if source_action not in managed_source_allowed_actions(source): raise ValueError('managed source action is unavailable for this source') if _pipeline_lifecycle_change_is_blocked(source, source_action, managed_sources): raise RuntimeError('managed pipeline lifecycle change is refused while scanner sources are running') outcome = 'completed' try: if source_action == 'start': success = source.start(force=True) elif source_action == 'stop': success = source.stop() elif source_action == 'restart': success = source.restart_now() elif source_action == 'pause': success = source.pause() elif source_action == 'resume': success = source.resume() elif source_action == 'once': if source.is_running() and not source.stop(timeout=10): success = False else: set_mode(source, 'once') success = source.start(force=True) elif source_action == 'set-mode': mode = request.get('mode') if mode not in ('loop', 'once', 'repeat'): raise ValueError('managed source mode is invalid') if isinstance(source, ManagedKeychecks) and mode == 'loop': raise ValueError('managed source mode is unavailable for this source') if source.is_running(): raise RuntimeError('managed source mode change requires a stopped source') set_mode(source, mode) success = True elif source_action == 'set-interval': value = request.get('interval_seconds') if ( type(value) is not int or value < 1 or value > MAX_MANAGED_SOURCE_DELAY_SECONDS ): raise ValueError('managed source interval is invalid') source.interval = value success = True elif source_action == 'set-restart': value = request.get('restart_enabled') if type(value) is not bool: raise ValueError('managed source restart setting is invalid') source.restart = value success = True else: value = request.get('restart_delay_seconds') if ( type(value) is not int or value < 1 or value > MAX_MANAGED_SOURCE_DELAY_SECONDS ): raise ValueError('managed source restart delay is invalid') source.restart_delay = value success = True except (ValueError, RuntimeError): raise except Exception as exc: raise RuntimeError('managed source action failed') from exc if not success: if source.runtime_blocked and source.desired_state == 'running': outcome = 'dependency-blocked' else: raise RuntimeError('managed source action failed') status_file = (context or {}).get('status_file') if status_file: try: write_status_file(status_file, managed_sources, context) except Exception: pass return { 'source_action': source_action, 'outcome': outcome, 'source': source.structured_state(), } def run_dashboard_action(request, context): require_exact_control_fields(request, {'dashboard_action'}) dashboard_action = request.get('dashboard_action') if dashboard_action not in ('start', 'stop', 'restart'): raise ValueError('unsupported dashboard action') manager = (context or {}).get('dashboard_manager') if manager is None: raise RuntimeError('dashboard control is unavailable') try: if dashboard_action == 'start': success = manager.start(force=True) elif dashboard_action == 'stop': success = manager.stop() else: success = manager.stop() and manager.start(force=True) except Exception as exc: raise RuntimeError('dashboard action failed') from exc snapshot = structured_dashboard_state(manager) if not success and snapshot['status'] != 'blocked': raise RuntimeError('dashboard action failed') return { 'dashboard_action': dashboard_action, 'outcome': 'dependency-blocked' if snapshot['status'] == 'blocked' else 'completed', 'dashboard': snapshot, } def run_managed_source_log_tail(request, managed_sources): require_exact_control_fields(request, {'source_id', 'line_count'}) source_id = request.get('source_id') line_count = request.get('line_count') if not isinstance(source_id, str) or not source_id or len(source_id) > 128: raise ValueError('managed source ID is invalid') if type(line_count) is not int or not 1 <= line_count <= MAX_LOG_TAIL_LINES: raise ValueError('managed source log line count is invalid') source = managed_source_registry(managed_sources).get(source_id) if source is None: raise ValueError('managed source ID is unknown') try: lines = source.tail_log_lines(line_count) except Exception as exc: raise RuntimeError('managed source log tail failed') from exc if type(lines) is not list or any(type(line) is not str for line in lines): raise RuntimeError('managed source log tail failed') lines = lines[-line_count:] selected = [] encoded_bytes = 2 for line in reversed(lines): item_bytes = len(json.dumps(line, ensure_ascii=True).encode('utf-8')) separator_bytes = 1 if selected else 0 if encoded_bytes + separator_bytes + item_bytes > MAX_LOG_TAIL_BYTES: break selected.append(line) encoded_bytes += separator_bytes + item_bytes selected.reverse() return { 'source_id': source_id, 'line_count': len(selected), 'lines': selected, 'response_truncated': len(selected) != len(lines), } def run_control_command(command, managed_sources, context, lock=None): command = (command or '').strip() if not command: return {'success': True, 'output': ''} if command in ('__status__', 'status-raw'): guard = lock or nullcontext() with guard: poll_managed_sources(managed_sources, context) return {'success': True, 'output': '\n'.join(build_runtime_table_lines(managed_sources, context=context)) + '\n'} if command.lower() in ('quit', 'exit', 'q'): return {'success': True, 'output': 'Detached from background supervisor. Use --stop-background to stop it.\n'} output = StringIO() guard = lock or nullcontext() with guard: with redirect_stdout(output): keep_running = handle_command(command, managed_sources, context) if not keep_running: print('Ignored quit/exit for background supervisor. Use --stop-background to stop it.') status_file = context.get('status_file') if context else None write_status_file(status_file, managed_sources, context) failed = bool((context or {}).pop('_command_failed', False)) return {'success': not failed, 'output': output.getvalue()} class nullcontext: def __enter__(self): return None def __exit__(self, exc_type, exc, tb): return False class SupervisorControlHandler(socketserver.StreamRequestHandler): def handle(self): self.connection.settimeout(5.0) try: raw = self.rfile.readline(MAX_CONTROL_REQUEST_BYTES + 1) except OSError: raw = b'' if len(raw) > MAX_CONTROL_REQUEST_BYTES or not raw.endswith(b'\n'): response = self.server.error_response('invalid or oversized control request') else: try: request = json.loads(raw.decode('utf-8')) response = self.server.run_request(request) except (UnicodeDecodeError, json.JSONDecodeError): response = self.server.error_response('invalid JSON control request') except Exception as exc: response = self.server.error_response( f'control request failed: {type(exc).__name__}: {exc}' ) signal_shutdown = bool(response.pop('_signal_shutdown', False)) if isinstance(response, dict) else False encoded = json.dumps(response, ensure_ascii=True, default=str, separators=(',', ':')).encode('utf-8') + b'\n' if len(encoded) > MAX_CONTROL_RESPONSE_BYTES: encoded = json.dumps(self.server.error_response('control response is too large'), separators=(',', ':')).encode('utf-8') + b'\n' try: self.wfile.write(encoded) self.wfile.flush() except (BrokenPipeError, ConnectionAbortedError, ConnectionResetError): pass finally: if signal_shutdown: shutdown_event = self.server.context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() class SupervisorControlServer(socketserver.TCPServer): allow_reuse_address = False request_queue_size = 16 def server_bind(self): if os.name == 'nt' and hasattr(socket, 'SO_EXCLUSIVEADDRUSE'): self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_EXCLUSIVEADDRUSE, 1) return super().server_bind() def __init__(self, server_address, managed_sources, context, lock, instance_id, token): super().__init__(server_address, SupervisorControlHandler) self.managed_sources = managed_sources self.context = context self.lock = lock self.instance_id = str(instance_id) self.token = str(token) self._intake_selector = selectors.DefaultSelector() self._pending_lock = threading.Lock() self._pending = {} self._intake_stopping = threading.Event() self._worker_slots = threading.BoundedSemaphore(MAX_CONTROL_WORKERS) self._worker_executor = ThreadPoolExecutor(max_workers=MAX_CONTROL_WORKERS, thread_name_prefix='supervisor-control') self._active_workers = 0 self._intake_thread = threading.Thread( target=self._control_intake_loop, name='supervisor-control-intake', daemon=True, ) self._intake_thread.start() @property def pending_control_connections(self): with self._pending_lock: return len(self._pending) @property def active_control_workers(self): with self._pending_lock: return self._active_workers @staticmethod def _close_control_socket(request): try: request.shutdown(socket.SHUT_RDWR) except OSError: pass try: request.close() except OSError: pass def _remove_pending_locked(self, request): self._pending.pop(request, None) try: self._intake_selector.unregister(request) except (KeyError, OSError, ValueError): pass def process_request(self, request, client_address): request.setblocking(False) evicted = None with self._pending_lock: if len(self._pending) >= MAX_CONTROL_PENDING_SOCKETS: evicted = next(iter(self._pending)) self._remove_pending_locked(evicted) self._pending[request] = { 'address': client_address, 'buffer': bytearray(), 'deadline': time.monotonic() + CONTROL_READ_TIMEOUT_SEC, } try: self._intake_selector.register(request, selectors.EVENT_READ) except BaseException: self._pending.pop(request, None) self._close_control_socket(request) raise if evicted is not None: self._close_control_socket(evicted) def _take_pending(self, request): with self._pending_lock: state = self._pending.get(request) if state is not None: self._remove_pending_locked(request) return state def _send_control_response(self, request, response): signal_shutdown = bool(response.pop('_signal_shutdown', False)) if isinstance(response, dict) else False encoded = json.dumps(response, ensure_ascii=True, default=str, separators=(',', ':')).encode('utf-8') + b'\n' if len(encoded) > MAX_CONTROL_RESPONSE_BYTES: encoded = json.dumps(self.error_response('control response is too large'), separators=(',', ':')).encode('utf-8') + b'\n' try: request.setblocking(True) request.settimeout(1.0) request.sendall(encoded) except OSError: pass finally: self._close_control_socket(request) if signal_shutdown: shutdown_event = self.context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() def _run_authenticated_request(self, request, value): try: try: response = self.run_request(value) except Exception as exc: if value.get('action') in ( 'runtime-snapshot', 'managed-source-action', 'dashboard-action', 'managed-source-log-tail', ): response = self.error_response('control request failed') else: response = self.error_response( f'control request failed: {type(exc).__name__}: {exc}' ) self._send_control_response(request, response) finally: with self._pending_lock: self._active_workers -= 1 self._worker_slots.release() def _dispatch_control_request(self, request, raw): try: value = json.loads(raw.decode('utf-8')) except (UnicodeDecodeError, json.JSONDecodeError): self._send_control_response(request, self.error_response('invalid JSON control request')) return if not authenticate_request(value, self.instance_id, self.token): self._send_control_response(request, self.error_response('authentication failed')) return if not self._worker_slots.acquire(blocking=False): self._send_control_response(request, self.error_response('control worker limit reached')) return with self._pending_lock: self._active_workers += 1 try: self._worker_executor.submit(self._run_authenticated_request, request, value) except BaseException: with self._pending_lock: self._active_workers -= 1 self._worker_slots.release() self._close_control_socket(request) raise def _read_pending_control(self, request): with self._pending_lock: state = self._pending.get(request) if state is None: return try: chunk = request.recv(min(65536, MAX_CONTROL_REQUEST_BYTES + 1 - len(state['buffer']))) except BlockingIOError: return except OSError: chunk = b'' if not chunk: self._take_pending(request) self._close_control_socket(request) return state['buffer'].extend(chunk) newline = state['buffer'].find(b'\n') if newline < 0 and len(state['buffer']) <= MAX_CONTROL_REQUEST_BYTES: return self._take_pending(request) if newline < 0 or newline + 1 > MAX_CONTROL_REQUEST_BYTES: self._send_control_response(request, self.error_response('invalid or oversized control request')) return self._dispatch_control_request(request, bytes(state['buffer'][:newline])) def _control_intake_loop(self): while not self._intake_stopping.is_set(): with self._pending_lock: has_pending = bool(self._pending) if not has_pending: self._intake_stopping.wait(0.05) continue try: events = self._intake_selector.select(0.05) except (OSError, ValueError): break for key, _ in events: self._read_pending_control(key.fileobj) now = time.monotonic() with self._pending_lock: expired = [request for request, state in self._pending.items() if state['deadline'] <= now] for request in expired: self._remove_pending_locked(request) for request in expired: self._close_control_socket(request) def server_close(self): self._intake_stopping.set() with self._pending_lock: pending = list(self._pending) for request in pending: self._remove_pending_locked(request) for request in pending: self._close_control_socket(request) if self._intake_thread.is_alive(): self._intake_thread.join(timeout=2) try: self._intake_selector.close() except OSError: pass self._worker_executor.shutdown(wait=True, cancel_futures=False) super().server_close() def response(self, ok, result=None, error=None): value = {'schema': CONTROL_SCHEMA, 'instance_id': self.instance_id, 'ok': bool(ok)} if ok: value['result'] = result else: value['error'] = str(error or 'request failed') return value def error_response(self, error): return self.response(False, error=error) def run_request(self, request): if not authenticate_request(request, self.instance_id, self.token): return self.error_response('authentication failed') action = str(request.get('action') or '') if action == 'handshake': guard = self.lock or nullcontext() with guard: dashboard = self.context.get('dashboard_manager') authority = self.context.get('authority') or {} return self.response(True, { 'instance_id': self.instance_id, 'pid': os.getpid(), 'manages_postgres': bool(self.context.get('with_postgres')), 'activation_state': lifecycle_phase(self.context), 'config_sha256': authority.get('config_sha256', ''), 'supervisor_sha256': authority.get('supervisor_sha256', ''), 'code_manifest_sha256': authority.get('code_manifest_sha256', ''), 'canonical_dsn_sha256': self.context.get('canonical_dsn_sha256', ''), 'dashboard': dashboard.snapshot() if dashboard else {'status': 'disabled', 'healthy': False}, }) if action == 'activation': guard = self.lock or nullcontext() with guard: return self.response(True, { 'activation_state': lifecycle_phase(self.context), }) if action == 'activate': guard = self.lock or nullcontext() with guard: if shutdown_checkpoint(self.context): return self.error_response('supervisor shutdown is already requested') if lifecycle_phase(self.context) == PHASE_ACTIVE: return self.response(True, {'activation_state': PHASE_ACTIVE}) if lifecycle_phase(self.context) != PHASE_ACTIVATING: return self.error_response(f'supervisor cannot activate from {lifecycle_phase(self.context)}') if not check_runtime_authority(self.context, trigger_shutdown=False): self.context['runtime_failed'] = True detail = runtime_authority_error(self.context) or 'runtime authority drifted before activation' shutdown_event = self.context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() return self.error_response(detail) callback = self.context.get('activation_callback') try: if callback is not None: callback() except BaseException as exc: self.context['runtime_failed'] = True self.context['start_gate_open'] = False shutdown_event = self.context.get('shutdown_event') if shutdown_event is not None: shutdown_event.set() return self.error_response(f'activation metadata update failed: {exc}') if shutdown_checkpoint(self.context): return self.error_response('supervisor shutdown is already requested') if lifecycle_phase(self.context) != PHASE_ACTIVE: self.context['lifecycle_phase'] = PHASE_ACTIVE self.context['activation_state'] = PHASE_ACTIVE self.context['start_gate_open'] = True activation_event = self.context.get('activation_event') if activation_event is not None: activation_event.set() return self.response(True, {'activation_state': PHASE_ACTIVE}) if action == 'shutdown': guard = self.lock or nullcontext() with guard: authority_error = shutdown_authority_error( self.context, self.instance_id, self.token, self.server_address[:2], ) if authority_error: return self.error_response(authority_error) if bool(request.get('with_postgres')) and not self.context.get('with_postgres'): return self.error_response('supervisor does not manage PostgreSQL') if self.context.get('shutdown_event') is None: return self.error_response('shutdown event is unavailable') try: begin_stopping(self.context) except Exception as exc: return self.error_response(f'unable to enter STOPPING phase: {exc}') return self.response(True, 'coordinated shutdown requested') command_shutdown = action == 'command' and str(request.get('command') or '').strip().lower() == 'shutdown' command_status = action == 'command' and str(request.get('command') or '').strip().lower() in ('status', 's', 'status-raw', '__status__') typed_mutation = action in ('managed-source-action', 'dashboard-action') hold_status = lifecycle_phase(self.context) in (PHASE_STOPPING, PHASE_FAILED_HOLD) and ( action in ('snapshot', 'runtime-snapshot', 'managed-source-log-tail', 'status') or command_status ) if lifecycle_phase(self.context) != PHASE_ACTIVE and not command_shutdown and not hold_status: return self.error_response(f'supervisor runtime is {lifecycle_phase(self.context)}; request is refused') if ( lifecycle_phase(self.context) == PHASE_ACTIVE and not command_shutdown and not typed_mutation and not check_runtime_authority(self.context) ): if action in ( 'runtime-snapshot', 'managed-source-action', 'dashboard-action', 'managed-source-log-tail', ): return self.error_response('runtime authority drifted') return self.error_response(self.context.get('authority_drift') or 'runtime authority drifted') if action == 'snapshot': guard = self.lock or nullcontext() with guard: poll_managed_sources(self.managed_sources, self.context) controller = self.context.get('postgres_controller') dashboard = self.context.get('dashboard_manager') return self.response(True, { 'activation_state': lifecycle_phase(self.context), 'signature': table_signature(self.managed_sources), 'discovery_producers': [ producer.structured_state() for producer in self.context.get('discovery_producers', ()) ], 'table': '\n'.join(build_runtime_table_lines(self.managed_sources, context=self.context)) + '\n', 'postgres': controller.snapshot() if controller else {'state': PostgresState.DISABLED.value, 'ready': True}, 'dashboard': dashboard.snapshot() if dashboard else {'status': 'disabled', 'healthy': False}, }) if action == 'runtime-snapshot': try: require_exact_control_fields(request, set()) except ValueError as exc: return self.error_response(str(exc)) guard = self.lock or nullcontext() with guard: poll_managed_sources(self.managed_sources, self.context) return self.response(True, structured_runtime_snapshot( self.managed_sources, self.context, )) if action == 'managed-source-action': guard = self.lock or nullcontext() with guard: if lifecycle_phase(self.context) != PHASE_ACTIVE: return self.error_response( f'supervisor runtime is {lifecycle_phase(self.context)}; request is refused' ) if not check_runtime_authority(self.context): return self.error_response('runtime authority drifted') try: result = run_managed_source_action( request, self.managed_sources, self.context, ) except (ValueError, RuntimeError) as exc: return self.error_response(str(exc)) return self.response(True, result) if action == 'managed-source-log-tail': guard = self.lock or nullcontext() with guard: try: result = run_managed_source_log_tail(request, self.managed_sources) except (ValueError, RuntimeError) as exc: return self.error_response(str(exc)) return self.response(True, result) if action == 'dashboard-action': guard = self.lock or nullcontext() with guard: if lifecycle_phase(self.context) != PHASE_ACTIVE: return self.error_response( f'supervisor runtime is {lifecycle_phase(self.context)}; request is refused' ) if not check_runtime_authority(self.context): return self.error_response('runtime authority drifted') try: result = run_dashboard_action(request, self.context) except (ValueError, RuntimeError) as exc: return self.error_response(str(exc)) return self.response(True, result) if action == 'status': outcome = run_control_command('__status__', self.managed_sources, self.context, self.lock) return self.response(True, outcome['output']) if action == 'command': command = str(request.get('command') or '').strip() if command.lower() == 'shutdown': guard = self.lock or nullcontext() with guard: authority_error = shutdown_authority_error( self.context, self.instance_id, self.token, self.server_address[:2], ) if authority_error: return self.error_response(authority_error) try: begin_stopping(self.context) except Exception as exc: return self.error_response(f'unable to enter STOPPING phase: {exc}') return self.response(True, 'coordinated shutdown requested\n') outcome = run_control_command(command, self.managed_sources, self.context, self.lock) if not outcome['success']: return self.error_response(outcome['output'].strip() or 'supervisor command failed') return self.response(True, outcome['output']) return self.error_response('unknown control action') def start_control_server(supervisor_config, managed_sources, context, lock, instance_id, token, start_thread=True): host, port = control_address(supervisor_config) server = SupervisorControlServer((host, port), managed_sources, context, lock, instance_id, token) if start_thread: threading.Thread(target=server.serve_forever, daemon=True).start() print(f'Control server listening on {server.server_address[0]}:{server.server_address[1]}') return server def process_running_status(pid): if not pid: return False if os.name == 'nt': try: process_query_limited_information = 0x1000 handle = _SUPERVISOR_OPEN_PROCESS(process_query_limited_information, False, int(pid)) if not handle: error = ctypes.get_last_error() return False if error in (87, 1168) else None try: exit_code = wintypes.DWORD() ok = _SUPERVISOR_GET_EXIT_CODE_PROCESS(handle, ctypes.byref(exit_code)) return (exit_code.value == 259) if ok else None finally: _SUPERVISOR_CLOSE_HANDLE(handle) except Exception: return None try: os.kill(pid, 0) return True except PermissionError: return None except OSError: return False def is_pid_running(pid): return process_running_status(pid) is True def background_child_command(args, config_path, launch_nonce, instance_file, expected_authority=None): app_dir = os.path.dirname(os.path.abspath(__file__)) command = [ sys.executable, '-I', '-S', '-B', os.path.join(app_dir, 'runtime_bootstrap.py'), 'supervisor', '--', '--runtime-bootstrap-entrypoint', os.path.abspath(__file__), '--config', config_path, '--non-interactive', '--background-child', '--no-clear', ] command.extend(['--instance-file', instance_file, f'--launch-nonce={launch_nonce}']) if expected_authority: command.extend([ '--expected-config-sha256', expected_authority['config_sha256'], '--expected-supervisor-sha256', expected_authority['supervisor_sha256'], '--expected-code-manifest-sha256', expected_authority['code_manifest_sha256'], ]) if args.sources: command.extend(['--sources', args.sources]) if args.once: command.append('--once') if args.autostart: command.append('--autostart') if args.status_interval: command.extend(['--status-interval', str(args.status_interval)]) if args.dashboard: command.append('--dashboard') if args.no_dashboard: command.append('--no-dashboard') if args.with_postgres: command.append('--with-postgres') return command def command_for_log(command): hidden_value_flags = { '--launch-nonce', '--expected-config-sha256', '--expected-supervisor-sha256', '--expected-code-manifest-sha256', '--token', '--docker-token', '--password', '--api-key', '--secret', '--db-url', '--database-url', '--scanner-db-url', } output = [] hide_next = False for value in command: text = str(value) if hide_next: output.append('') hide_next = False else: flag = text.split('=', 1)[0].lower() if flag in hidden_value_flags and '=' in text: output.append(text.split('=', 1)[0] + '=') else: output.append(text) hide_next = flag in hidden_value_flags return output def inspect_existing_instance(instance_file, config_path): if not os.path.exists(instance_file): return 'absent', None, None try: metadata = load_instance_metadata(instance_file) except (OSError, ValueError) as exc: try: raw = read_private_json(instance_file) if ( raw.get('schema') == 2 and raw.get('instance_id') and os.path.normcase(os.path.abspath(raw.get('instance_file') or '')) == os.path.normcase(os.path.abspath(instance_file)) and exact_process_identity_state( raw.get('pid'), raw.get('process_creation_time'), raw.get('executable'), ) in ('dead', 'reused') ): return 'stale', raw, f'legacy manifest is invalid and exact process identity is dead: {exc}' except (OSError, TypeError, ValueError): pass return 'invalid', None, str(exc) try: process = verify_instance_process(metadata, os.path.abspath(__file__), config_path) except (OSError, ValueError) as exc: try: drift_process = verify_instance_process( metadata, os.path.abspath(__file__), config_path, allow_config_drift=True, ) except (OSError, ValueError): drift_process = None if drift_process is not None: drift_process.close() return 'config_drift', metadata, 'live supervisor permits authenticated shutdown only' identity_state = exact_process_identity_state( metadata.get('pid'), metadata.get('process_creation_time'), metadata.get('executable'), ) if identity_state in ('dead', 'reused'): return 'stale', metadata, str(exc) return 'identity_mismatch', metadata, f'{identity_state}: {exc}' try: handshake = send_control_request(metadata, 'handshake', timeout=3) if not isinstance(handshake, dict) or handshake.get('instance_id') != metadata['instance_id']: raise ValueError('invalid authenticated handshake') return 'verified', metadata, process except (OSError, ValueError, RuntimeError) as exc: process.close() return 'unreachable', metadata, str(exc) def _terminate_unpublished_child(process): if process.poll() is not None: return True try: process.terminate() except OSError: pass try: process.wait(timeout=5) except subprocess.TimeoutExpired: process.kill() process.wait() return process.poll() is not None def _rollback_background_start( process, candidate, instance_file, launch_nonce, with_postgres, lock_path=None, activation_attempted=False, shutdown_timeout=180, ): exact_candidate = ( isinstance(candidate, dict) and candidate.get('launch_nonce') == launch_nonce and candidate.get('pid') == process.pid ) if process.poll() is None and exact_candidate: try: send_control_request(candidate, 'activation', timeout=2) except (OSError, ValueError, RuntimeError): pass if not exact_candidate: _terminate_unpublished_child(process) outcome = 'unpublished child was reaped by its retained parent handle' else: if process.poll() is None: try: send_control_request(candidate, 'shutdown', timeout=5, with_postgres=bool(with_postgres)) except (OSError, ValueError, RuntimeError): pass try: process.wait(timeout=max(5.0, float(shutdown_timeout))) except subprocess.TimeoutExpired: pass if process.poll() is None: return 'activation authority was uncertain; the exact child was left running because authenticated shutdown was not confirmed' outcome = 'exact candidate exited after authenticated coordinated shutdown' if ( process.poll() == 0 and exact_candidate ): remove_instance_if_matches(instance_file, candidate.get('instance_id'), lock_path=lock_path) remove_shutdown_receipt(instance_file, candidate.get('instance_id')) return outcome def start_background(args, config_path, results_dir, supervisor_config, expected_authority=None): expected_authority = expected_authority or runtime_authority(config_path) runtime_authority( config_path, expected_config_sha256=expected_authority['config_sha256'], expected_supervisor_sha256=expected_authority['supervisor_sha256'], expected_code_manifest_sha256=expected_authority['code_manifest_sha256'], code_manifest=expected_authority['code_manifest'], ) instance_file, log_file, status_file, legacy_pid_file = background_paths( config_path, results_dir, supervisor_config, args.instance_file, args.pid_file, ) singleton_lock_file = background_lock_path(config_path, results_dir, supervisor_config) legacy_instance_file = legacy_log_instance_path(config_path, results_dir, supervisor_config) if ( os.path.normcase(os.path.abspath(legacy_instance_file)) != os.path.normcase(os.path.abspath(instance_file)) and os.path.exists(legacy_instance_file) ): raise SystemExit( f'Refusing background launch while legacy log-directory instance metadata exists: {legacy_instance_file}' ) existing_state, existing_metadata, existing_detail = inspect_existing_instance(instance_file, config_path) if existing_state == 'verified': existing_process = existing_detail existing_process.close() raise SystemExit(f'Refusing duplicate verified background supervisor instance: {existing_metadata["instance_id"]}') if existing_state not in ('absent', 'stale'): raise SystemExit(f'Refusing to replace {existing_state} supervisor metadata at {instance_file}: {existing_detail}') legacy_pid = read_pid_file(legacy_pid_file) if legacy_pid: legacy_status = process_running_status(legacy_pid) label = 'running' if legacy_status is True else 'not running' if legacy_status is False else 'unknown' raise SystemExit( f'Refusing background launch while legacy status-only PID metadata exists ' f'({legacy_pid}, {label}): {legacy_pid_file}' ) rotate_log_if_needed(log_file, supervisor_config.get('log_max_mb', 64), supervisor_config.get('log_keep', 5)) log_handle = open_private_append(log_file) log_handle.write(f'\n=== background supervisor launch {now_iso()} ===\n') launch_nonce = secrets.token_urlsafe(32) command = background_child_command(args, config_path, launch_nonce, instance_file, expected_authority) log_handle.write('command: ' + ' '.join(command_for_log(command)) + '\n') log_handle.write(f'instance_file: {instance_file}\n') log_handle.write(f'status_file: {status_file}\n') dashboard_enabled, dashboard_host, dashboard_port, dashboard_url = dashboard_settings(supervisor_config) if dashboard_enabled: log_handle.write(f'dashboard_url: {dashboard_url}\n') log_handle.flush() creationflags = 0 if os.name == 'nt': creationflags = DETACHED_PROCESS | subprocess.CREATE_NEW_PROCESS_GROUP | CREATE_NO_WINDOW try: process = subprocess.Popen( command, cwd=os.path.dirname(config_path), stdin=subprocess.DEVNULL, stdout=log_handle, stderr=subprocess.STDOUT, env={ **os.environ, 'PYTHONUNBUFFERED': '1', 'TRUF_SUPERVISOR_LAUNCH_NONCE': launch_nonce, }, creationflags=creationflags, close_fds=True, ) finally: log_handle.close() timeout = max(1.0, float(supervisor_config.get('background_start_timeout_sec', 20) or 20)) deadline = time.monotonic() + timeout last_error = 'instance metadata was not written' metadata = None candidate = None handshake = None activation_attempted = False while time.monotonic() < deadline: if process.poll() is not None: last_error = f'child exited with code {process.returncode}' break try: candidate = load_instance_metadata(instance_file) if candidate['launch_nonce'] != launch_nonce or candidate['pid'] != process.pid: raise InstanceMetadataError('launch nonce or child PID mismatch') if candidate['manages_postgres'] != bool(args.with_postgres): raise InstanceMetadataError('PostgreSQL management mode mismatch') if candidate.get('instance_file') != os.path.normcase(os.path.realpath(os.path.abspath(instance_file))): raise InstanceMetadataError('child instance metadata path mismatch') for key in ('config_sha256', 'supervisor_sha256', 'code_manifest_sha256'): if candidate.get(key) != expected_authority[key]: raise InstanceMetadataError(f'parent/child {key} authority mismatch') retained = verify_instance_process(candidate, os.path.abspath(__file__), config_path) try: handshake = send_control_request(candidate, 'handshake', timeout=2) finally: retained.close() if not isinstance(handshake, dict) or handshake.get('instance_id') != candidate['instance_id']: raise InstanceMetadataError('authenticated startup handshake mismatch') if candidate.get('activation_state') != PHASE_ACTIVATING or handshake.get('activation_state') != PHASE_ACTIVATING: raise InstanceMetadataError('new background child was not published in ACTIVATING state') for key in ('config_sha256', 'supervisor_sha256', 'code_manifest_sha256'): if handshake.get(key) != expected_authority[key]: raise InstanceMetadataError(f'authenticated handshake {key} mismatch') if handshake.get('canonical_dsn_sha256', '') != candidate.get('canonical_dsn_sha256', ''): raise InstanceMetadataError('authenticated handshake DSN authority mismatch') metadata = candidate break except (OSError, ValueError, RuntimeError) as exc: last_error = str(exc) time.sleep(0.05) if metadata is not None: activation_attempted = True activation_timeout = max( 5.0, float(supervisor_config.get('background_activation_timeout_sec', timeout) or timeout), ) activation_deadline = time.monotonic() + activation_timeout try: result = send_control_request(metadata, 'activate', timeout=2) if not isinstance(result, dict) or result.get('activation_state') != PHASE_ACTIVE: last_error = 'authenticated activate response was invalid' except (OSError, ValueError, RuntimeError) as exc: last_error = f'activate response was unavailable: {exc}' confirmed = False while time.monotonic() < activation_deadline and process.poll() is None: try: handshake = send_control_request(metadata, 'handshake', timeout=2) if ( isinstance(handshake, dict) and handshake.get('instance_id') == metadata['instance_id'] and handshake.get('activation_state') == PHASE_ACTIVE ): metadata = load_instance_metadata(instance_file) if metadata.get('activation_state') != PHASE_ACTIVE: raise InstanceMetadataError('active handshake disagreed with private metadata') for key in ('config_sha256', 'supervisor_sha256', 'code_manifest_sha256'): if metadata.get(key) != expected_authority[key] or handshake.get(key) != expected_authority[key]: raise InstanceMetadataError(f'parent activation {key} recheck failed') if handshake.get('canonical_dsn_sha256', '') != metadata.get('canonical_dsn_sha256', ''): raise InstanceMetadataError('parent activation DSN authority recheck failed') retained = verify_instance_process(metadata, os.path.abspath(__file__), config_path) retained.close() confirmed = True break last_error = 'activation handshake remained ACTIVATING' except (OSError, ValueError, RuntimeError) as exc: last_error = f'activation confirmation failed: {exc}' time.sleep(0.05) if not confirmed: metadata = None if metadata is None: rollback = _rollback_background_start( process, candidate, instance_file, launch_nonce, args.with_postgres, lock_path=singleton_lock_file, activation_attempted=activation_attempted, shutdown_timeout=supervisor_config.get('background_shutdown_timeout_sec', 180), ) raise SystemExit( f'Background supervisor startup failed within {timeout:g}s: {last_error}; {rollback}. Check {log_file}' ) print(f'Started background supervisor: PID {metadata["pid"]}') print(f'Instance file: {instance_file}') print(f'Log file: {log_file}') print(f'Status file: {status_file}') if dashboard_enabled: dashboard = handshake.get('dashboard') if isinstance(handshake, dict) else {} print(f'Dashboard: {dashboard.get("status", "pending")} ({dashboard.get("detail", "health pending")})') print(f'Dashboard URL: {dashboard_url}') def stop_background(config_path, results_dir, supervisor_config, explicit_instance_file=None, explicit_pid_file=None, with_postgres=False): instance_file, log_file, status_file, legacy_pid_file = background_paths( config_path, results_dir, supervisor_config, explicit_instance_file, explicit_pid_file, ) singleton_lock_file = background_lock_path(config_path, results_dir, supervisor_config) try: metadata = load_instance_metadata(instance_file) process = verify_instance_process( metadata, os.path.abspath(__file__), config_path, allow_config_drift=True, allow_code_drift=True, ) except (OSError, ValueError) as exc: print(f'Refusing background stop without verified instance metadata: {exc}') legacy_instance_file = legacy_log_instance_path(config_path, results_dir, supervisor_config) if os.path.exists(legacy_instance_file) and os.path.abspath(legacy_instance_file) != os.path.abspath(instance_file): print(f'Legacy instance metadata at {legacy_instance_file} is detection-only and does not grant control authority.') legacy_pid = read_pid_file(legacy_pid_file) if legacy_pid: print(f'Legacy PID {legacy_pid} is status-only and will not be terminated.') return False try: if with_postgres and not metadata['manages_postgres']: print('Refusing --with-postgres stop because this verified supervisor does not manage PostgreSQL.') return False result = send_control_request(metadata, 'shutdown', timeout=5, with_postgres=bool(with_postgres)) print(str(result)) timeout = max(5.0, float(supervisor_config.get('background_shutdown_timeout_sec', 180) or 180)) if not process.wait(timeout): print(f'Coordinated shutdown did not finish within {timeout:g}s; no PID-based termination was attempted.') return False try: exit_code = process.exit_code() except (AttributeError, OSError, ValueError): exit_code = None try: receipt_code = load_shutdown_receipt(instance_file, metadata['instance_id']) except (OSError, ValueError): receipt_code = None if exit_code is not None and receipt_code is not None and exit_code != receipt_code: print('Background supervisor exit status disagreed with its authenticated shutdown receipt.') return False confirmed_code = exit_code if exit_code is not None else receipt_code if confirmed_code != 0: label = 'unavailable' if confirmed_code is None else str(confirmed_code) print(f'Background supervisor coordinated shutdown exit status was {label}; metadata was retained.') return False if not remove_instance_if_matches( instance_file, metadata['instance_id'], lock_path=singleton_lock_file, ): print('Background supervisor exited successfully, but matching metadata could not be removed safely.') return False remove_shutdown_receipt(instance_file, metadata['instance_id']) print('Background supervisor exited after coordinated shutdown.') return True except (OSError, ValueError, RuntimeError) as exc: print(f'Authenticated coordinated shutdown failed: {exc}') return False finally: process.close() def background_status(config_path, results_dir, supervisor_config, explicit_instance_file=None, explicit_pid_file=None): instance_file, log_file, status_file, legacy_pid_file = background_paths( config_path, results_dir, supervisor_config, explicit_instance_file, explicit_pid_file, ) dashboard_enabled, dashboard_host, dashboard_port, dashboard_url = dashboard_settings(supervisor_config) print(f'Instance file: {instance_file}') print(f'Log file: {log_file}') print(f'Status file: {status_file}') if dashboard_enabled: print(f'Dashboard URL: {dashboard_url}') verified = False try: metadata = load_instance_metadata(instance_file) process = verify_instance_process(metadata, os.path.abspath(__file__), config_path) try: handshake = send_control_request(metadata, 'handshake', timeout=3) print(f'PID: {metadata["pid"]}') print(f'Instance: {metadata["instance_id"]}') print('Status: verified and authenticated') dashboard = handshake.get('dashboard') if isinstance(handshake, dict) else None if isinstance(dashboard, dict): print( f"Dashboard status: {dashboard.get('status', 'unknown')} " f"healthy={dashboard.get('healthy', False)} detail={dashboard.get('detail', '')}" ) verified = True finally: process.close() except (OSError, ValueError, RuntimeError) as exc: print(f'Status: no verified authenticated instance ({exc})') legacy_instance_file = legacy_log_instance_path(config_path, results_dir, supervisor_config) if os.path.exists(legacy_instance_file) and os.path.abspath(legacy_instance_file) != os.path.abspath(instance_file): print(f'Legacy instance metadata (detection-only): {legacy_instance_file}') legacy_pid = read_pid_file(legacy_pid_file) if legacy_pid: legacy_status = process_running_status(legacy_pid) label = 'running' if legacy_status is True else 'not running' if legacy_status is False else 'unknown' print(f'Legacy PID (status-only): {legacy_pid} ({label})') if os.path.exists(status_file): print('\nLast supervisor table:') try: with open(status_file, 'r', encoding='utf-8') as f: print(f.read().rstrip()) except OSError as e: print(f'Unable to read status file: {e}') return verified def verified_instance_for_control(instance_file, config_path): metadata = load_instance_metadata(instance_file) process = verify_instance_process(metadata, os.path.abspath(__file__), config_path) try: send_control_request(metadata, 'handshake', timeout=3) return metadata, process except BaseException: process.close() raise def send_background_command(config_path, results_dir, supervisor_config, command, explicit_instance_file=None, explicit_pid_file=None): instance_file, log_file, status_file, legacy_pid_file = background_paths( config_path, results_dir, supervisor_config, explicit_instance_file, explicit_pid_file, ) try: metadata, process = verified_instance_for_control(instance_file, config_path) except (OSError, ValueError, RuntimeError) as exc: print(f'Refusing command without verified authenticated instance metadata: {exc}') return False try: if str(command).strip().lower() == 'shutdown': response = send_control_request(metadata, 'shutdown') else: response = send_control_command(metadata, command) print(str(response).rstrip()) return True except (OSError, ValueError, RuntimeError) as exc: print(f'Authenticated control command failed: {exc}') return False finally: process.close() def attach_background(config_path, results_dir, supervisor_config, explicit_instance_file=None, explicit_pid_file=None): instance_file, log_file, status_file, legacy_pid_file = background_paths( config_path, results_dir, supervisor_config, explicit_instance_file, explicit_pid_file, ) try: metadata, process = verified_instance_for_control(instance_file, config_path) initial_snapshot = get_control_snapshot(metadata) except (OSError, ValueError, RuntimeError) as exc: print(f'Unable to verify and authenticate background supervisor: {exc}') return False finally: if 'process' in locals(): process.close() print((initial_snapshot.get('table') or '').rstrip()) print('\nAttached to background supervisor. Type `quit` to detach, `help` for commands, `watch` for live view.') poll_sec = float(supervisor_config.get('attach_poll_sec', supervisor_config.get('poll_sec', 0.5)) or 0.5) while True: try: command = input('attach> ').strip() except (EOFError, KeyboardInterrupt): print() break if not command: continue if command.lower() in ('quit', 'exit', 'q'): break if command.lower() in ('watch', 'w'): try: watch_remote(metadata, poll_sec) except (OSError, ValueError, RuntimeError) as e: print(f'Connection failed: {e}') break continue try: if command.lower() == 'shutdown': response = send_control_request(metadata, 'shutdown') else: response = send_control_command(metadata, command) print(str(response).rstrip()) except (OSError, ValueError, RuntimeError) as e: print(f'Connection failed: {e}') break print('Detached. Background supervisor is still running.') return True def parse_args(): parser = argparse.ArgumentParser(description='Supervisor for running configured scanner sources in one console.') parser.add_argument('--config', default='config.yaml') parser.add_argument('--sources', help='Comma-separated source list. Defaults to enabled sources from config.yaml') parser.add_argument('--once', action='store_true', help='Force --once for every managed source') actions = parser.add_mutually_exclusive_group() actions.add_argument('--dry-run', action='store_true', help='Print child commands and exit') parser.add_argument('--status-interval', type=int, help='Seconds between status redraws') parser.add_argument('--no-clear', action='store_true', help='Do not clear console before status redraws') parser.add_argument('--autostart', action='store_true', help='Start selected sources immediately') parser.add_argument('--non-interactive', action='store_true', help='Run status loop without command prompt') parser.add_argument('--background-child', action='store_true', help=argparse.SUPPRESS) parser.add_argument('--launch-nonce', help=argparse.SUPPRESS) parser.add_argument('--expected-config-sha256', help=argparse.SUPPRESS) parser.add_argument('--expected-supervisor-sha256', help=argparse.SUPPRESS) parser.add_argument('--expected-code-manifest-sha256', help=argparse.SUPPRESS) parser.add_argument('--runtime-bootstrap-entrypoint', help=argparse.SUPPRESS) actions.add_argument('--background', action='store_true', help='Launch supervisor in the background and exit') actions.add_argument('--attach', action='store_true', help='Attach a foreground command prompt to running background supervisor') actions.add_argument('--stop-background', action='store_true', help='Request authenticated coordinated background shutdown') actions.add_argument('--background-status', action='store_true', help='Show verified background supervisor status') actions.add_argument('--cmd', help='Send one supervisor prompt command to a running background supervisor') parser.add_argument('--instance-file', help='Override instance metadata path; its parent must be the configured private control directory') parser.add_argument('--pid-file', help='Legacy status-only PID file path; never grants control authority') parser.add_argument('--dashboard', action='store_true', help='Launch dashboard regardless of config supervisor.dashboard') parser.add_argument('--no-dashboard', action='store_true', help='Disable dashboard launch') parser.add_argument('--with-postgres', action='store_true', help='Let the child controller manage verified bundled PostgreSQL') return parser.parse_args() def main(): args = parse_args() command_parts = [] if getattr(args, 'cmd', None): try: command_parts = shlex.split(args.cmd) except ValueError: command_parts = ['invalid'] read_only = ( getattr(args, 'dry_run', False) or getattr(args, 'background_status', False) or (getattr(args, 'cmd', None) and not command_is_mutating(command_parts)) ) control_requested = ( getattr(args, 'stop_background', False) or getattr(args, 'background_status', False) or getattr(args, 'attach', False) or bool(getattr(args, 'cmd', None)) ) runtime_launch = getattr(args, 'background', False) or not control_requested if not read_only and runtime_launch and not getattr(args, 'with_postgres', False): raise SystemExit('Supervisor unmanaged PostgreSQL mutation is retired; use --with-postgres.') if not read_only and os.getenv(RUNTIME_BOOTSTRAP_ENV) != RUNTIME_BOOTSTRAP_VALUE: raise SystemExit('Mutating supervisor runtime requires the canonical runtime bootstrap.') bootstrap_entrypoint = getattr(args, 'runtime_bootstrap_entrypoint', None) if getattr(args, 'background_child', False): expected_entrypoint = os.path.normcase(os.path.realpath(os.path.abspath(__file__))) actual_entrypoint = ( os.path.normcase(os.path.realpath(os.path.abspath(bootstrap_entrypoint))) if bootstrap_entrypoint and os.path.isabs(bootstrap_entrypoint) else '' ) if actual_entrypoint != expected_entrypoint: raise SystemExit('Background child requires its canonical bootstrap entrypoint binding.') config_path = os.path.abspath(args.config) config_hash_before = sha256_file(config_path) managed_runtime = ( not read_only and runtime_launch and getattr(args, 'with_postgres', False) ) if managed_runtime: documents = validate_managed_runtime_startup(config_path) validated_config = documents.config validated_config_sha256 = documents.config_sha256 documents = None if validated_config_sha256 != config_hash_before: raise SystemExit( 'Supervisor config changed while it was being validated; ' 'retry with a stable private config file.' ) validated_depth = validate_docker_depth_config( validated_config, managed_postgres=bool(getattr(args, 'with_postgres', False)), ) validated_config = apply_path_config(validated_depth.config, config_path) config, project_dir, results_dir, supervisor_config, selected_sources = ( supervisor_runtime_from_config( config_path, validated_config, getattr(args, 'sources', None), ) ) loaded_config_sha256 = sha256_file(config_path) if loaded_config_sha256 != validated_config_sha256: raise SystemExit( 'Supervisor config changed while it was being validated; ' 'retry with a stable private config file.' ) else: config, project_dir, results_dir, supervisor_config, selected_sources = load_supervisor_runtime( config_path, getattr(args, 'sources', None), managed_postgres=bool(getattr(args, 'with_postgres', False)), ) loaded_config_sha256 = sha256_file(config_path) if loaded_config_sha256 != config_hash_before: raise SystemExit( 'Supervisor config changed while it was being loaded; ' 'retry with a stable private config file.' ) try: preflight_lifecycle_paths( config_path, config, authority_profile='server' if managed_runtime else 'full', ) except (OSError, ValueError) as exc: raise SystemExit(str(exc)) from exc if getattr(args, 'dashboard', False): dashboard_config = supervisor_config.get('dashboard') if isinstance(supervisor_config.get('dashboard'), dict) else {} supervisor_config['dashboard'] = {**dashboard_config, 'enabled': True} if getattr(args, 'no_dashboard', False): dashboard_config = supervisor_config.get('dashboard') if isinstance(supervisor_config.get('dashboard'), dict) else {} supervisor_config['dashboard'] = {**dashboard_config, 'enabled': False} dashboard_settings(supervisor_config) authority = runtime_authority( config_path, trufflehog_path=(config.get('global') or {}).get('trufflehog_path'), policy_paths=configured_policy_paths(config), include_trufflehog=not managed_runtime, expected_config_sha256=loaded_config_sha256, ) if getattr(args, 'background', False): if not getattr(args, 'with_postgres', False): raise SystemExit('Supervisor unmanaged PostgreSQL mutation is retired; use --with-postgres.') start_background(args, config_path, results_dir, supervisor_config, authority) return 0 if getattr(args, 'stop_background', False): stopped = stop_background( config_path, results_dir, supervisor_config, getattr(args, 'instance_file', None), getattr(args, 'pid_file', None), with_postgres=getattr(args, 'with_postgres', False), ) return 0 if stopped else 1 if getattr(args, 'background_status', False): return 0 if background_status(config_path, results_dir, supervisor_config, getattr(args, 'instance_file', None), getattr(args, 'pid_file', None)) else 1 if getattr(args, 'attach', False): return 0 if attach_background(config_path, results_dir, supervisor_config, getattr(args, 'instance_file', None), getattr(args, 'pid_file', None)) else 1 if getattr(args, 'cmd', None): return 0 if send_background_command(config_path, results_dir, supervisor_config, args.cmd, getattr(args, 'instance_file', None), getattr(args, 'pid_file', None)) else 1 keychecks_config = keychecks_config_for(config_path, config) def create_sources( gate, authority_check=None, start_gate=None, child_environments=None, ): child_environments = child_environments or {} items = [] for source in selected_sources: producer = source in DISCOVERY_PRODUCER_SOURCES source_class = ManagedDiscoveryProducer if producer else ManagedSource items.append(source_class( source, config_path, project_dir, results_dir, supervisor_config, (config.get('sources') or {}).get(source, {}), getattr(args, 'once', False), gate, authority_check, start_gate, child_environments.get(DISCOVERY_PRODUCER_ROLE if producer else 'scanner'), )) if should_manage_keychecks(config, supervisor_config, getattr(args, 'sources', None)): keychecks = ManagedKeychecks( config_path, project_dir, results_dir, supervisor_config, keychecks_config, getattr(args, 'once', False), gate, authority_check, start_gate, child_environments.get('keycheck'), ) items.append(keychecks) return [source for source in items if source.enabled] if getattr(args, 'dry_run', False): dry_sources = create_sources(DependencyGate(ready=True)) if not dry_sources: raise SystemExit('No enabled supervisor sources selected') for source in dry_sources: print(f'[{source.source}]') print('command:', ' '.join(source.build_command())) print('log:', source.log_path) print('state:', source.state_path if source.use_per_source_state else '(config default)') print('once:', source.once, 'repeat:', source.repeat, 'restart:', source.restart, 'interval:', source.interval) print() return 0 if not getattr(args, 'with_postgres', False): raise SystemExit('Supervisor unmanaged PostgreSQL mutation is retired; use --with-postgres.') load_postgres_env( config_path, config.get('global') or {}, enforce_canonical=getattr(args, 'with_postgres', False), ) if getattr(args, 'with_postgres', False): managed_database_url = canonical_database_url() if not managed_database_url: raise RuntimeError('managed PostgreSQL lifecycle requires a canonical DSN') else: try: managed_database_url = canonical_database_url() except ValueError: managed_database_url = '' if not managed_database_url: raise RuntimeError('supervised mutation requires one caller-selected canonical PostgreSQL DSN') instance_file, supervisor_log_file, status_file, _ = background_paths( config_path, results_dir, supervisor_config, getattr(args, 'instance_file', None), getattr(args, 'pid_file', None), ) singleton_lock_file = background_lock_path(config_path, results_dir, supervisor_config) background_child = bool(getattr(args, 'background_child', False)) background_log_writer = None try: cluster_lock = ClusterAuthorityLock(config, endpoint_dsn=managed_database_url).acquire() except (BlockingIOError, OSError) as exc: raise SystemExit(f'Refusing duplicate lifecycle-owning supervisor or maintenance authority for this PostgreSQL cluster: {exc}') from exc try: instance_lock = SupervisorInstanceLock(instance_file, lock_path=singleton_lock_file).acquire() except InstanceLockError as exc: cluster_lock.release() raise SystemExit(f'Refusing duplicate lifecycle-owning supervisor: {exc}') from exc if background_child: try: background_log_writer = install_bounded_background_output( supervisor_log_file, max(1, int(supervisor_config.get('log_max_mb', 64) or 64)) * 1024 * 1024, supervisor_config.get('log_keep', 5), ) except BaseException: instance_lock.release() cluster_lock.release() raise refresh_sec = getattr(args, 'status_interval', None) or int(supervisor_config.get('refresh_sec', DEFAULT_REFRESH_SEC) or DEFAULT_REFRESH_SEC) heartbeat_sec = int(supervisor_config.get('heartbeat_sec', 60) or 0) poll_sec = float(supervisor_config.get('poll_sec', 0.2) or 0.2) autostart = getattr(args, 'autostart', False) or bool_value(supervisor_config.get('autostart'), False) interactive = not getattr(args, 'non_interactive', False) and bool_value(supervisor_config.get('interactive'), True) control_lock = threading.RLock() managed_sources = [] shutdown_event = threading.Event() shutdown_retry_event = threading.Event() activation_event = threading.Event() context = { 'config_path': config_path, 'selected_sources_arg': getattr(args, 'sources', None), 'global_force_once': getattr(args, 'once', False), 'config': config, 'project_dir': project_dir, 'results_dir': results_dir, 'supervisor_config': supervisor_config, 'selected_sources': selected_sources, 'managed_sources': managed_sources, 'discovery_producers': [], 'scanner_sources': [], 'keychecks_config': keychecks_config, 'queue_dir': (config.get('global') or {}).get('queue_dir') or os.path.join(os.path.dirname(results_dir), 'queues'), 'dashboard_manager': None, 'with_postgres': bool(getattr(args, 'with_postgres', False)), 'control_lock': control_lock, 'dependency_gate': None, 'source_dependency_gate': None, 'postgres_controller': None, 'background_child': background_child, 'shutdown_requested': False, 'runtime_failed': False, 'shutdown_event': shutdown_event, 'shutdown_retry_event': shutdown_retry_event, 'status_file': status_file, 'authority': authority, 'instance_file': instance_file, 'activation_state': PHASE_ACTIVATING, 'lifecycle_phase': PHASE_ACTIVATING, 'start_gate_open': False, 'activation_event': activation_event, 'activation_lock': threading.RLock(), } control_server = None control_server_started = False instance_metadata = None dashboard_manager = None postgres_controller = None runtime_activated = False exit_code = 0 sigterm_installed = False previous_sigterm_handler = None def request_shutdown(_signum, _frame): context['shutdown_requested'] = True try: if os.name == 'posix': previous_sigterm_handler = signal.signal(signal.SIGTERM, request_shutdown) sigterm_installed = True if os.path.exists(instance_file): try: existing = load_instance_metadata(instance_file) except (OSError, ValueError) as exc: existing = read_private_json(instance_file) identity_state = exact_process_identity_state( existing.get('pid'), existing.get('process_creation_time'), existing.get('executable'), ) if ( existing.get('schema') != 2 or not existing.get('instance_id') or os.path.normcase(os.path.abspath(existing.get('instance_file') or '')) != os.path.normcase(os.path.abspath(instance_file)) or identity_state not in ('dead', 'reused') ): raise InstanceLockError( f'refusing invalid supervisor metadata with {identity_state} identity: {exc}' ) from exc durable_unlink(instance_file) remove_shutdown_receipt(instance_file, existing['instance_id']) existing = None if existing is None: pass elif ( os.path.normcase(os.path.abspath(existing.get('instance_file') or '')) != os.path.normcase(os.path.abspath(instance_file)) ): raise InstanceLockError('refusing supervisor metadata bound to another instance path') elif exact_process_identity_state( existing.get('pid'), existing.get('process_creation_time'), existing.get('executable'), ) not in ('dead', 'reused'): raise InstanceLockError('refusing to replace supervisor metadata whose exact process identity is live or unknown') elif not remove_instance_if_matches(instance_file, existing['instance_id'], instance_lock=instance_lock): raise InstanceLockError('unable to remove matching stale supervisor metadata under the lifetime lock') elif existing is not None: remove_shutdown_receipt(instance_file, existing['instance_id']) if background_child: launch_nonce = getattr(args, 'launch_nonce', None) if not launch_nonce: raise SystemExit('Background child requires a launch nonce') expected = { 'config_sha256': getattr(args, 'expected_config_sha256', None), 'supervisor_sha256': getattr(args, 'expected_supervisor_sha256', None), 'code_manifest_sha256': getattr(args, 'expected_code_manifest_sha256', None), } if not all(expected.values()): raise SystemExit('Background child requires immutable parent authority hashes') authority = runtime_authority( config_path, trufflehog_path=(config.get('global') or {}).get('trufflehog_path'), policy_paths=configured_policy_paths(config), include_trufflehog=not managed_runtime, expected_config_sha256=expected['config_sha256'], expected_supervisor_sha256=expected['supervisor_sha256'], expected_code_manifest_sha256=expected['code_manifest_sha256'], ) context['authority'] = authority else: launch_nonce = secrets.token_urlsafe(32) # All fallible non-lifecycle initialization is complete before ACTIVE. context['canonical_dsn_sha256'] = dsn_sha256(managed_database_url) dependency_gate = DependencyGate( ready=not getattr(args, 'with_postgres', False), database_url=managed_database_url, ) context['dependency_gate'] = dependency_gate source_dependency_gate = DependencyGate( ready=False, database_url=managed_database_url, ) context['source_dependency_gate'] = source_dependency_gate authority_check = lambda: check_runtime_authority(context, require_private_acl=True) start_gate = lambda: lifecycle_start_allowed(context) postgres_controller = controller_from_config(config, supervisor_config) if getattr(args, 'with_postgres', False) else None if postgres_controller is not None: postgres_controller.authority_check = lambda: lifecycle_start_allowed(context) and authority_check() context['postgres_controller'] = postgres_controller instance_id = secrets.token_urlsafe(24) token = secrets.token_urlsafe(48) child_metadata = { 'instance_file': instance_file, 'instance_id': instance_id, 'token': token, 'config_sha256': authority['config_sha256'], 'supervisor_sha256': authority['supervisor_sha256'], 'code_manifest_sha256': authority['code_manifest_sha256'], 'canonical_dsn_sha256': context['canonical_dsn_sha256'], } child_environments = { 'scanner': supervised_child_environment(child_metadata, managed_database_url, 'scanner'), DISCOVERY_PRODUCER_ROLE: supervised_child_environment( child_metadata, managed_database_url, DISCOVERY_PRODUCER_ROLE, ), 'keycheck': supervised_child_environment(child_metadata, managed_database_url, 'keycheck'), 'dashboard': supervised_child_environment(child_metadata, managed_database_url, 'dashboard'), 'result-ingester': supervised_child_environment(child_metadata, managed_database_url, 'result-ingester'), 'jsonl-projector': supervised_child_environment(child_metadata, managed_database_url, 'jsonl-projector'), 'worker-api': supervised_child_environment(child_metadata, managed_database_url, 'worker-api'), 'janitor': supervised_child_environment(child_metadata, '', 'janitor'), 'docker-shadow': supervised_child_environment(child_metadata, managed_database_url, 'docker-shadow'), } pipeline_workers = [ ManagedPipelineWorker( 'result-ingester', config_path, project_dir, results_dir, supervisor_config, supervisor_config.get('result_ingester') or {}, dependency_gate, authority_check, start_gate, child_environments['result-ingester'], ), ManagedPipelineWorker( 'jsonl-projector', config_path, project_dir, results_dir, supervisor_config, supervisor_config.get('jsonl_projector') or {}, dependency_gate, authority_check, start_gate, child_environments['jsonl-projector'], ), ManagedPipelineWorker( 'janitor', config_path, project_dir, results_dir, supervisor_config, supervisor_config.get('janitor') or {}, None, authority_check, start_gate, child_environments['janitor'], ), ManagedPipelineWorker( 'worker-api', config_path, project_dir, results_dir, supervisor_config, supervisor_config.get('worker_api') or {}, dependency_gate, authority_check, start_gate, child_environments['worker-api'], ), ] managed_sources.extend(worker for worker in pipeline_workers if worker.enabled) docker_shadow = ManagedDockerShadow( config_path, project_dir, results_dir, supervisor_config, supervisor_config.get('docker_shadow') or {}, dependency_gate, authority_check, start_gate, child_environments['docker-shadow'], ) if docker_shadow.enabled: managed_sources.append(docker_shadow) configured_sources = create_sources( source_dependency_gate, authority_check, start_gate, child_environments, ) managed_sources.extend(configured_sources) context['discovery_producers'] = [ source for source in configured_sources if isinstance(source, ManagedDiscoveryProducer) ] context['scanner_sources'] = [ source for source in configured_sources if not isinstance(source, ManagedDiscoveryProducer) and not isinstance(source, ManagedKeychecks) ] if not managed_sources: raise SystemExit('No enabled supervisor sources selected') dashboard_manager = ManagedDashboard( config_path, project_dir, supervisor_config, results_dir, context['queue_dir'], dependency_gate=dependency_gate, authority_check=authority_check, start_gate=start_gate, child_environment=child_environments['dashboard'], ) context['dashboard_manager'] = dashboard_manager control_server = start_control_server( supervisor_config, managed_sources, context, control_lock, instance_id, token, start_thread=False, ) host, port = control_server.server_address[:2] instance_metadata = build_instance_metadata( launch_nonce, os.path.abspath(__file__), config_path, host, port, getattr(args, 'with_postgres', False), instance_id=instance_id, token=token, activation_state=PHASE_ACTIVATING, expected_config_sha256=authority['config_sha256'], expected_supervisor_sha256=authority['supervisor_sha256'], code_manifest=authority['code_manifest'], expected_code_manifest_sha256=authority['code_manifest_sha256'], canonical_dsn_sha256=context['canonical_dsn_sha256'], lifecycle_mode='background' if background_child else 'foreground', instance_file=instance_file, ) write_instance_metadata(instance_file, instance_metadata) context['instance_id'] = instance_id def activate_runtime(): nonlocal runtime_activated if shutdown_checkpoint(context): return current = load_instance_metadata(instance_file) if current['instance_id'] != instance_id or current['activation_state'] != PHASE_ACTIVATING: raise InstanceMetadataError('activation metadata no longer names the exact ACTIVATING candidate') for key in ('config_sha256', 'supervisor_sha256', 'code_manifest_sha256'): if current.get(key) != authority[key]: raise InstanceMetadataError(f'activation metadata {key} mismatch') detail = runtime_authority_error(context) if detail: raise InstanceMetadataError(detail) if shutdown_checkpoint(context): return # ACTIVE may reach disk even if publication or event delivery raises. runtime_activated = True update_instance_activation(instance_file, instance_id, PHASE_ACTIVE) if shutdown_checkpoint(context): return context['lifecycle_phase'] = PHASE_ACTIVE context['activation_state'] = PHASE_ACTIVE context['start_gate_open'] = True activation_event.set() context['activation_callback'] = activate_runtime threading.Thread(target=control_server.serve_forever, daemon=True).start() control_server_started = True print(f'Control server listening on {host}:{port}') if background_child: activation_timeout = max( 5.0, float(supervisor_config.get( 'background_activation_timeout_sec', max(30.0, float(supervisor_config.get('background_start_timeout_sec', 20) or 20) * 2), )), ) activation_deadline = time.monotonic() + activation_timeout while time.monotonic() < activation_deadline: with control_lock: if shutdown_checkpoint(context) or activation_event.is_set(): break time.sleep(0.02) with control_lock: if not shutdown_checkpoint(context) and not activation_event.is_set(): exit_code = 1 print(f'Background child activation timed out after {activation_timeout:g}s; no lifecycle action was taken.') else: with control_lock: activate_runtime() with control_lock: shutdown_checkpoint(context) if runtime_activated and lifecycle_start_allowed(context): keychecks_auto_started = False with control_lock: for source in managed_sources: if shutdown_checkpoint(context) or not lifecycle_start_allowed(context): break if isinstance(source, ManagedPipelineWorker): source.start(force=True) if not shutdown_checkpoint(context) and lifecycle_start_allowed(context): dashboard_manager.start(record_intent=False) if not autostart and bool_value(keychecks_config.get('autostart'), False): for source in managed_sources: if shutdown_checkpoint(context) or not lifecycle_start_allowed(context): break if source.source == 'keychecks': source.start(force=True) keychecks_auto_started = True if interactive: interactive_loop( managed_sources, autostart=autostart, clear=not getattr(args, 'no_clear', False), context=context, poll_sec=poll_sec, ) elif not autostart and not keychecks_auto_started and not background_child: with control_lock: if not shutdown_checkpoint(context): exit_code = 1 print('Non-interactive mode requires --autostart or supervisor.autostart=true') else: if autostart: with control_lock: for source in autostart_sources(managed_sources): if shutdown_checkpoint(context) or not lifecycle_start_allowed(context): break source.start(force=True) non_interactive_loop( managed_sources, refresh_sec=refresh_sec, status_file=status_file, poll_sec=poll_sec, context=context, lock=control_lock, stay_alive=background_child, status_heartbeat_sec=heartbeat_sec, ) except KeyboardInterrupt: try: print('\nStopping child processes...') except BaseException: exit_code = exit_code or 1 except SystemExit as exc: exit_code = exit_code or (0 if exc.code is None else int(exc.code) if isinstance(exc.code, int) else 1) if exc.code and not isinstance(exc.code, int): print(str(exc.code)) except BaseException as exc: exit_code = 1 print(f'Supervisor runtime failed closed: {type(exc).__name__}: {exc}') finally: context['shutdown_requested'] = True context['authority_release_safe'] = False try: with control_lock: if any( getattr(source, 'startup_cleanup_pending', False) is True or (not background_child and source.status == 'failed') for source in managed_sources ): context['runtime_failed'] = True if runtime_activated: if not coordinated_shutdown(managed_sources, context): exit_code = exit_code or 1 else: begin_stopping(context) if dashboard_manager is not None: dashboard_manager.close() context['authority_release_safe'] = ( postgres_controller is None or bool(postgres_controller.close(wait=False)) ) except BaseException as exc: exit_code = exit_code or 1 context['authority_release_safe'] = False try: with control_lock: enter_failed_hold(context, f'initial shutdown failed: {type(exc).__name__}: {exc}') print(f'Coordinated shutdown failed closed: {exc}') except BaseException: pass while not context.get('authority_release_safe', False): exit_code = exit_code or 1 retain_unsafe_authority(managed_sources, context) try: try: if control_server and control_server_started: control_server.shutdown() finally: if control_server: control_server.server_close() except BaseException: exit_code = exit_code or 1 try: print('Supervisor stopped.') sys.stdout.flush() sys.stderr.flush() except BaseException: exit_code = exit_code or 1 try: if background_log_writer is not None: background_log_writer.close() except BaseException: exit_code = exit_code or 1 exit_code = exit_code or int(bool(context.get('runtime_failed') or context.get('authority_drift'))) if instance_metadata: try: write_shutdown_receipt(instance_file, instance_metadata['instance_id'], exit_code) except BaseException: exit_code = exit_code or 1 # The waiting stopper (or locked stale reconciliation) removes both files. instance_lock.release() cluster_lock.release() if sigterm_installed: signal.signal(signal.SIGTERM, previous_sigterm_handler) return exit_code if __name__ == '__main__': raise SystemExit(main())