6563 lines
274 KiB
Python
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())
|