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

595 lines
25 KiB
Python

"""Isolated entrypoint for the private, read-only Docker deployment."""
import argparse
import http.client
import json
import os
from pathlib import Path
import re
import runpy
import secrets
import signal
import stat
import subprocess
import sys
from urllib.parse import quote
sys.dont_write_bytecode = True
if not sys.dont_write_bytecode:
raise RuntimeError('container runtime could not disable bytecode writes')
APP = Path('/opt/truf/app')
DATA = Path('/data')
RUN = Path('/run/truf')
UID = GID = 10001
DEFAULT_CONFIG = APP / 'config.linux.yaml'
PROVISIONED = DATA / '.provisioned.json'
INITIALIZED = DATA / 'initialized.json'
INITIALIZE_LOCK = DATA / 'initialize.lock'
PASSWORD = DATA / 'postgres-password'
PROVIDER_SECRETS = DATA / 'config/secrets.yaml'
FORMAT = 'truf-container-data-v1'
WORKER_HEALTH_TOKEN = '0' * 64
DIRECTORIES = (
'home', 'config', 'managed-files', 'runtime-linux', 'runtime-linux/results',
'runtime-linux/queues', 'runtime-linux/state',
'runtime-linux/state/gharchive_cache', 'runtime-linux/logs',
'runtime-linux/keychecks', 'runtime-linux/postman_cache',
'runtime-linux/result_spool', 'runtime-linux/postgres',
'runtime-linux/postgres/logs', 'postgres-linux', 'scanner-work',
'scanner-result-bundles', 'scanner-result-bundles/tmp',
'scanner-result-bundles/ready', 'scanner-result-bundles/quarantine',
)
_shutdown_requested = False
def private_path(path, *, directory=False):
path = Path(path)
if not path.is_absolute() or '..' in path.parts:
raise RuntimeError('private path must be absolute and normalized')
for component in (*reversed(path.parents), path):
if stat.S_ISLNK(component.lstat().st_mode):
raise RuntimeError('private paths must not contain symlinks')
details = path.lstat()
expected_type = stat.S_ISDIR if directory else stat.S_ISREG
if (not expected_type(details.st_mode) or details.st_uid != UID
or details.st_gid != GID
or stat.S_IMODE(details.st_mode) != (0o700 if directory else 0o600)):
raise RuntimeError('private path ownership, type, or mode is invalid: ' + str(path))
return path
def require_container(*, provisioning=False):
if (sys.platform != 'linux' or Path(__file__) != APP / 'container_runtime.py'
or not Path('/.dockerenv').is_file()):
raise RuntimeError('runtime commands are restricted to the prepared Docker image')
if not (sys.flags.isolated and sys.flags.no_site and sys.flags.dont_write_bytecode):
raise RuntimeError('container entrypoint requires python -I -S -B')
if os.getuid() != os.geteuid() or os.geteuid() != (0 if provisioning else UID):
raise RuntimeError('unexpected container runtime UID')
if not provisioning and os.getgid() != GID:
raise RuntimeError('unexpected container runtime GID')
if not os.statvfs(APP).f_flag & os.ST_RDONLY:
raise RuntimeError('the application image must be mounted read-only')
private_path(APP.parent, directory=True)
private_path(APP, directory=True)
private_path(APP / 'container_runtime.py')
with open('/proc/self/mountinfo', 'rb') as handle:
payload = handle.read(1024 * 1024 + 1)
if len(payload) > 1024 * 1024:
raise RuntimeError('mount inventory exceeds its bound')
mounts = {}
for line in payload.splitlines():
fields = line.split()
if len(fields) > 6 and b'-' in fields:
mounts[fields[4]] = fields[fields.index(b'-') + 1]
if mounts.get(b'/data') not in (b'ext4', b'xfs', b'btrfs', b'zfs'):
raise RuntimeError('/data must be an independent native Linux data volume')
if mounts.get(b'/run/truf') != b'tmpfs':
raise RuntimeError('/run/truf must be an independent private tmpfs')
private_path(DATA, directory=True)
private_path(RUN, directory=True)
os.umask(0o077)
def _read_json(path):
path = private_path(path)
if path.stat().st_size > 4096:
raise RuntimeError('container marker exceeds its bound')
value = json.loads(path.read_text(encoding='utf-8'))
if not isinstance(value, dict) or value.get('format') != FORMAT:
raise RuntimeError('unrecognized container data marker')
return value
def _write_new(path, payload, *, provisioning=False):
try:
private_path(path.parent, directory=True)
descriptor = os.open(
path,
os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW,
0o600,
)
with os.fdopen(descriptor, 'wb') as handle:
if provisioning:
os.fchown(handle.fileno(), UID, GID)
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
descriptor = os.open(path.parent, os.O_RDONLY | os.O_DIRECTORY)
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
private_path(path)
finally:
payload = None
def provision():
import fcntl
lock = DATA / '.provision.lock'
try:
descriptor = os.open(lock, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600)
os.fchown(descriptor, UID, GID)
except FileExistsError:
private_path(lock)
descriptor = os.open(lock, os.O_WRONLY | os.O_NOFOLLOW)
with os.fdopen(descriptor, 'wb') as handle:
fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
if PROVISIONED.exists():
_read_json(PROVISIONED)
managed_files = DATA / 'managed-files'
try:
os.mkdir(managed_files, 0o700)
os.chown(managed_files, UID, GID)
except FileExistsError:
pass
for name in DIRECTORIES:
private_path(DATA / name, directory=True)
private_path(PASSWORD)
private_path(PROVIDER_SECRETS)
print('Container data layout is already provisioned; nothing was changed.', flush=True)
return
if set(os.listdir(DATA)) - {'home', '.provision.lock'}:
raise RuntimeError('refusing to provision nonempty or partially initialized data')
home = DATA / 'home'
if home.exists() and any(home.iterdir()):
raise RuntimeError('refusing to provision a nonempty home directory')
for name in DIRECTORIES:
path = DATA / name
if not path.exists():
os.mkdir(path, 0o700)
os.chown(path, UID, GID)
private_path(path, directory=True)
_write_new(PASSWORD, (secrets.token_urlsafe(48) + '\n').encode('ascii'), provisioning=True)
_write_new(PROVIDER_SECRETS, b'{}\n', provisioning=True)
_write_new(DATA / 'runtime-linux/proxy.txt', b'', provisioning=True)
_write_new(PROVISIONED, (json.dumps({'format': FORMAT, 'uid': UID, 'gid': GID}) + '\n').encode('ascii'), provisioning=True)
print('Fresh private container data layout provisioned; PostgreSQL is not initialized yet.', flush=True)
def prepare_environment(
config_path, *, validate_documents=True, return_config_sha256=False,
):
_read_json(PROVISIONED)
for name in ('authority', 'control', 'tmp'):
path = RUN / name
try:
path.mkdir(mode=0o700)
except FileExistsError:
pass
private_path(path, directory=True)
config_path = Path(config_path)
if config_path != DEFAULT_CONFIG and config_path.parent != DATA / 'config':
raise RuntimeError('configuration must be image-owned or in /data/config')
private_path(config_path)
password = private_path(PASSWORD).read_text(encoding='ascii').rstrip('\n')
if re.fullmatch(r'[A-Za-z0-9_-]{32,128}', password) is None:
raise RuntimeError('the generated PostgreSQL password is invalid')
for name in tuple(os.environ):
upper = name.upper()
if (upper.startswith(('PG', 'TRUF_', 'SCANNER_', 'SCAN_', 'TRUFFLEHOG_', 'KEYCHECK_'))
or upper in ('DATABASE_URL', 'PYTHONPATH', 'PYTHONHOME')):
del os.environ[name]
url = 'postgresql://truf:' + quote(password, safe='') + '@127.0.0.1:5432/truf'
os.environ.update({
'PATH': '/usr/local/bin:/usr/bin:/bin:/usr/lib/postgresql/16/bin',
'HOME': '/data/home', 'TMPDIR': str(RUN / 'tmp'),
'TMP': str(RUN / 'tmp'), 'TEMP': str(RUN / 'tmp'),
'TRUF_CONTAINER_CONFIG': str(config_path),
'TRUF_POSTGRES_DB': 'truf', 'TRUF_POSTGRES_USER': 'truf',
'TRUF_POSTGRES_PORT': '5432', 'TRUF_POSTGRES_PASSWORD': password,
'SCANNER_DB_URL': url, 'DATABASE_URL': url, 'TRUF_MANAGED_POSTGRES_DSN': url,
'TRUF_DB_CONNECT_TIMEOUT_SEC': '3', 'TRUF_DB_STATEMENT_TIMEOUT_MS': '5000',
'TRUF_DB_LOCK_TIMEOUT_MS': '2000', 'TRUF_DB_IDLE_TRANSACTION_TIMEOUT_MS': '10000',
})
bootstrap = runpy.run_path(str(APP / 'child_bootstrap.py'))
bootstrap['_enable_dependency_paths']('supervisor')
sys.path.insert(0, str(APP))
from paths import apply_path_config
from runtime_document_io import (
load_managed_runtime_config,
validate_managed_runtime_files,
)
from runtime_security import preflight_lifecycle_paths
def check_resolved(candidate):
expected = {
'root_dir': '/opt/truf', 'project_dir': str(APP),
'runtime_dir': '/data/runtime-linux', 'postgres_data_dir': '/data/postgres-linux',
'postgres_bin_dir': '/usr/lib/postgresql/16/bin',
'result_bundle_dir': '/data/scanner-result-bundles', 'work_dir': '/data/scanner-work',
'control_dir': str(RUN / 'control'), 'secrets_file': str(PROVIDER_SECRETS),
}
if any(candidate['global'].get(name) != value for name, value in expected.items()):
raise RuntimeError('configuration escapes the fixed container storage contract')
if candidate['supervisor'].get('control_dir') != str(RUN / 'control'):
raise RuntimeError('supervisor control must remain on private ephemeral storage')
preflight_lifecycle_paths(
str(config_path), candidate, authority_profile='server',
)
documents = (
validate_managed_runtime_files(str(config_path))
if validate_documents
else load_managed_runtime_config(str(config_path))
)
config = documents.config
config_sha256 = documents.config_sha256
documents = None
config = apply_path_config(config, str(config_path))
check_resolved(config)
if return_config_sha256:
return config, config_sha256
return config
def _bootstrap_command(target, *arguments):
return [sys.executable, '-u', '-I', '-S', '-B', str(APP / 'runtime_bootstrap.py'), target, '--', *arguments]
def initialize(config_path, config, *, expected_config_sha256=None):
from postgres_runtime import postgres_runtime_paths
from runtime_document_io import load_managed_runtime_config
from runtime_security import PrivateFileLock, read_private_json, write_private_json_exclusive
def require_stable_config():
if expected_config_sha256 is None:
return
current = load_managed_runtime_config(str(config_path))
current_sha256 = current.config_sha256
current = None
if current_sha256 != expected_config_sha256:
raise RuntimeError('configuration changed after managed runtime validation')
require_stable_config()
paths = postgres_runtime_paths(config)
def migrate(*, initialize_base):
if _shutdown_requested:
return False
try:
# Even an uncertain maintenance start must enter the verified stop path.
require_stable_config()
subprocess.run(_bootstrap_command(
'postgres-runtime', 'maintenance-start', '--config', str(config_path),
), check=True)
if _shutdown_requested:
return False
require_stable_config()
arguments = [
'migrate-runtime-safety', '--config', str(config_path),
]
if initialize_base:
arguments.append('--initialize-base')
arguments.extend(('--apply', '--sources-stopped'))
subprocess.run(_bootstrap_command(*arguments), check=True)
return True
finally:
subprocess.run(_bootstrap_command(
'postgres-runtime', 'maintenance-stop', '--config', str(config_path),
), check=True)
with PrivateFileLock(str(INITIALIZE_LOCK)):
if INITIALIZED.exists():
marker = _read_json(INITIALIZED)
identity = read_private_json(paths['identity_path'])
if (not marker.get('system_identifier')
or marker.get('system_identifier') != identity.get('system_identifier')
or marker.get('pg_major') != 16 or identity.get('pg_major') != 16):
raise RuntimeError('initialization marker does not match the bound cluster')
migrated = migrate(initialize_base=False)
if migrated:
print('Existing PostgreSQL migrated and confirmed stopped.', flush=True)
return
if Path(paths['identity_path']).exists() or any(Path(paths['data_dir']).iterdir()):
raise RuntimeError('partial initialization requires offline inspection; no automatic repair is allowed')
if _shutdown_requested:
return
require_stable_config()
subprocess.run(_bootstrap_command('postgres-runtime', 'initialize-empty', '--config', str(config_path)), check=True)
if _shutdown_requested:
return
if not migrate(initialize_base=True):
return
identity = read_private_json(paths['identity_path'])
require_stable_config()
write_private_json_exclusive(str(INITIALIZED), {
'format': FORMAT, 'system_identifier': identity['system_identifier'],
'pg_major': identity['pg_major'],
})
print('Independent PostgreSQL initialized, migrated, cut over, and confirmed stopped.', flush=True)
def _probe_worker_api(config):
worker = config['supervisor']['worker_api']
connection = None
try:
connection = http.client.HTTPConnection(
worker['address'], int(worker['port']), timeout=2,
)
connection.request(
'POST', '/api/v1/worker/claim', body=b'',
headers={
'Authorization': 'Bearer ' + WORKER_HEALTH_TOKEN,
'Content-Length': '0',
},
)
response = connection.getresponse()
payload = response.read(4097)
if (
len(payload) > 4096
or response.status != 401
or response.getheader('WWW-Authenticate') != 'Bearer'
or json.loads(payload) != {
'error': {
'code': 'unauthorized',
'message': 'worker credentials are invalid',
},
}
):
raise RuntimeError('worker API is not ready')
except Exception:
raise RuntimeError('worker API is not ready') from None
finally:
if connection is not None:
connection.close()
def health(config, *, require_worker_api=False, require_discovery_producers=False):
# Never create a competing SQL session while first-install migration is exclusive.
_read_json(INITIALIZED)
from scanner_db import ScannerDB
from lifecycle_authority import DISCOVERY_PRODUCER_SOURCES
from supervisor import get_control_snapshot
from supervisor_instance import load_instance_metadata
metadata = load_instance_metadata(config['supervisor']['instance_file'])
snapshot = get_control_snapshot(metadata)
if (snapshot.get('activation_state') != 'ACTIVE'
or snapshot.get('postgres', {}).get('state') != 'READY'
or snapshot.get('postgres', {}).get('ready') is not True):
raise RuntimeError('supervisor and PostgreSQL are not ready')
signatures = {row[0]: row for row in snapshot.get('signature', ())}
required = ['result-ingester', 'jsonl-projector']
if config['supervisor'].get('janitor', {}).get('enabled', True):
required.append('janitor')
worker_api_enabled = config['supervisor'].get('worker_api', {}).get(
'enabled', False,
)
if require_worker_api and not worker_api_enabled:
raise RuntimeError('worker API is required but disabled')
if worker_api_enabled:
required.append('worker-api')
for name in required:
row = signatures.get(name)
if not row or tuple(row[1:3]) != ('running', 'running') or not row[3] or row[7]:
raise RuntimeError('a required pipeline worker is not running')
if any(row[1] == 'failed' or row[8] for row in signatures.values()):
raise RuntimeError('a managed source is failed or retains uncertain ownership')
if require_discovery_producers:
source_config = config.get('sources') or {}
supervisor_sources = config['supervisor'].get('sources') or {}
enabled_producers = [
name for name in DISCOVERY_PRODUCER_SOURCES
if (
(supervisor_sources.get(name) or {}).get('enabled')
if 'enabled' in (supervisor_sources.get(name) or {})
else (source_config.get(name) or {}).get('enabled', False)
)
]
for name in enabled_producers:
row = signatures.get(name)
ready = bool(
row
and row[2] == 'running'
and not row[7]
and not row[8]
and (
(row[1] == 'running' and row[3])
or (row[1] == 'waiting' and row[4] == 0)
)
)
if not ready:
raise RuntimeError('a required discovery producer is not ready')
if require_worker_api:
_probe_worker_api(config)
db = ScannerDB(db_url=os.environ['SCANNER_DB_URL'], initialize=False)
try:
if not db.enabled or not db.conn.is_postgres:
raise RuntimeError('PostgreSQL application connection is unavailable')
db.set_application_name('truf-container-health')
db.conn.execute('SET default_transaction_read_only = on')
db.conn.commit()
db.require_runtime_safety_schema()
db.require_final_cutover()
for name in ('result_ingester', 'jsonl_projector'):
if not db.pipeline_worker_health(name, metadata['instance_id'])['healthy']:
raise RuntimeError('a required durable worker lease is not ready')
row = db.conn.execute("SELECT current_setting('data_directory') AS data_directory, current_setting('server_version_num') AS version").fetchone()
db.conn.commit()
if row['data_directory'] != '/data/postgres-linux' or int(row['version']) // 10000 != 16:
raise RuntimeError('PostgreSQL identity does not match the container')
finally:
db.close()
for name in ('work_dir', 'result_bundle_dir', 'results_dir'):
path = config['global'][name]
if not os.access(path, os.W_OK | os.X_OK) or os.statvfs(path).f_bavail == 0:
raise RuntimeError('required persistent storage is not writable or is full')
return {'healthy': True, 'activation_state': 'ACTIVE', 'postgres': 'READY', 'workers': sorted(signatures)}
def import_secrets(config_path, config, *, expected_config_sha256=None):
from postgres_runtime import postgres_runtime_paths
from runtime_document_io import (
load_managed_runtime_config,
validate_managed_runtime_files,
)
from runtime_security import ClusterAuthorityLock, PrivateFileLock, durable_replace, fsync_directory
import yaml
with PrivateFileLock(str(INITIALIZE_LOCK)), ClusterAuthorityLock(
config, create_parent=False, endpoint_dsn=os.environ['SCANNER_DB_URL'],
):
if (Path(config['supervisor']['instance_file']).exists()
or (Path(postgres_runtime_paths(config)['data_dir']) / 'postmaster.pid').exists()):
raise RuntimeError('stop the runtime before replacing provider credentials')
payload = sys.stdin.buffer.read(1024 * 1024 + 1)
try:
if not payload or len(payload) > 1024 * 1024:
raise RuntimeError('credential input must be a nonempty YAML document below 1 MiB')
validated = validate_managed_runtime_files(
str(config_path), secrets_bytes=payload,
)
except BaseException:
payload = None
raise
payload = None
validated_config_sha256 = validated.config_sha256
if (
expected_config_sha256 is not None
and validated_config_sha256 != expected_config_sha256
):
validated = None
raise RuntimeError('configuration changed while provider credentials were being validated')
try:
serialized = yaml.safe_dump(
validated.secrets, allow_unicode=True,
).encode('utf-8')
except BaseException:
validated = None
raise
validated = None
temporary = DATA / ('config/secrets-import-' + secrets.token_hex(16))
try:
try:
_write_new(temporary, serialized)
finally:
serialized = None
private_path(PROVIDER_SECRETS)
current = load_managed_runtime_config(str(config_path))
current_sha256 = current.config_sha256
current = None
if current_sha256 != validated_config_sha256:
raise RuntimeError(
'configuration changed while provider credentials were being validated'
)
durable_replace(str(temporary), str(PROVIDER_SECRETS))
fsync_directory(str(PROVIDER_SECRETS.parent))
finally:
temporary.unlink(missing_ok=True)
print('Private provider credentials replaced; no credentials were printed.', flush=True)
def main(argv=None):
global _shutdown_requested
parser = argparse.ArgumentParser(description=__doc__, allow_abbrev=False)
parser.add_argument('action', choices=('provision', 'initialize', 'run', 'health', 'status', 'import-secrets', 'import-snapshot'))
parser.add_argument('--config', default=os.environ.get('TRUF_CONTAINER_CONFIG', str(DEFAULT_CONFIG)))
parser.add_argument('--manifest-sha256')
parser.add_argument('--require-worker-api', action='store_true')
parser.add_argument('--require-discovery-producers', action='store_true')
args = parser.parse_args(argv)
require_container(provisioning=args.action == 'provision')
if (
(args.require_worker_api or args.require_discovery_producers)
and args.action not in ('health', 'status')
):
raise RuntimeError('strict health requirements are restricted to health and status')
if args.action == 'import-snapshot':
if re.fullmatch(r'[0-9a-f]{64}', args.manifest_sha256 or '') is None:
raise RuntimeError('snapshot import requires --manifest-sha256 with the approved manifest digest')
if Path(args.config) != DEFAULT_CONFIG:
raise RuntimeError('snapshot import must begin with the image-owned default configuration')
elif args.manifest_sha256 is not None:
raise RuntimeError('--manifest-sha256 is restricted to snapshot import')
if args.action == 'provision':
provision()
return 0
bind_config = args.action in ('initialize', 'run', 'import-secrets')
prepared = prepare_environment(
args.config,
validate_documents=args.action != 'import-secrets',
return_config_sha256=bind_config,
)
if bind_config:
config, config_sha256 = prepared
else:
config = prepared
if args.action in ('health', 'status'):
print(json.dumps(health(
config, require_worker_api=args.require_worker_api,
require_discovery_producers=args.require_discovery_producers,
), sort_keys=True), flush=True)
return 0
if args.action == 'import-secrets':
import_secrets(
args.config, config,
expected_config_sha256=config_sha256,
)
return 0
def request_shutdown(_signum, _frame):
global _shutdown_requested
_shutdown_requested = True
previous = signal.signal(signal.SIGTERM, request_shutdown)
try:
if args.action == 'import-snapshot':
from container_import import import_snapshot
return import_snapshot(sys.modules[__name__], args.manifest_sha256)
initialize(
Path(args.config), config,
expected_config_sha256=config_sha256,
)
if _shutdown_requested or args.action == 'initialize':
return 0
from runtime_document_io import load_managed_runtime_config
current = load_managed_runtime_config(str(args.config))
current_sha256 = current.config_sha256
current = None
if current_sha256 != config_sha256:
raise RuntimeError('configuration changed after managed runtime validation')
command = _bootstrap_command(
'supervisor', '--runtime-bootstrap-entrypoint', str(APP / 'supervisor.py'),
'--config', args.config, '--with-postgres', '--non-interactive',
'--autostart', '--no-dashboard',
)
os.execv(sys.executable, command)
finally:
signal.signal(signal.SIGTERM, previous)
if __name__ == '__main__':
try:
raise SystemExit(main())
except Exception as exc:
print('Container runtime rejected: ' + str(exc), file=sys.stderr, flush=True)
raise SystemExit(1) from None