Files
2026-09-30 20:30:56 +03:00

6563 lines
274 KiB
Python

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 <service|all> [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 <source>` 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 <source|all>
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 <service|all> [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 <source|all>
Start a stopped source using its current mode from the table.
Example: start pypi
stop <source|all>
Stop the child process and keep it stopped. Supervisor will not auto-restart it.
Example: stop dockerhub
restart <source|all>
Stop then start immediately.
Example: restart github
pause <source|all>
Stop and mark as paused. Same process behavior as stop, but visually distinct.
Example: pause npm
resume <source|all>
Unpause and start immediately.
Example: resume npm
once <source|all>
Switch source to once mode and start it with --once. It will not repeat unless mode is changed.
Example: once pypi
mode <source|all> 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 <source|all> interval <seconds>
Set repeat interval for once+repeat mode.
Example: set pypi interval 120
set <source|all> restart on|off
Enable/disable restart after crash or loop child exit.
Example: set dockerhub restart off
set <source|all> restart_delay <seconds>
Set initial failure restart delay.
Example: set github restart_delay 60
logs <source> [lines]
Print the last N log lines from the configured runtime logs directory. Default: 40.
Example: logs pypi 80
command <source|all>
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 <source|all>')
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} <source|all>')
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 <source|all> 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 <source|all> interval|restart|restart_delay <value>')
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 <source> [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('<authority-value>')
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] + '=<authority-value>')
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())