1453 lines
63 KiB
Python
1453 lines
63 KiB
Python
"""Container entrypoint policy with mocked Linux metadata and temp-only data.
|
|
|
|
No operational CLI, PostgreSQL process, socket, chown, or host mount is used.
|
|
PurePosixPath models Linux policy; real filesystem fixtures use native Path.
|
|
"""
|
|
|
|
import ast
|
|
import builtins
|
|
from contextlib import contextmanager
|
|
import io
|
|
import json
|
|
import os
|
|
from pathlib import Path, PurePosixPath
|
|
import stat
|
|
import subprocess
|
|
import sys
|
|
import traceback
|
|
from types import ModuleType, SimpleNamespace
|
|
from unittest import mock
|
|
|
|
import pytest
|
|
import yaml
|
|
|
|
|
|
SOURCE = Path(__file__).resolve().parents[1] / 'app' / 'container_runtime.py'
|
|
PASSWORD = 'generated-fixture-password-' + 'x' * 38
|
|
PROVIDER_KEY = 'not-a-real-provider-credential'
|
|
MOUNTS = (
|
|
b'21 1 8:1 / /data rw - ext4 /dev/fixture rw\n'
|
|
b'22 1 0:2 / /run/truf rw - tmpfs tmpfs rw\n'
|
|
)
|
|
|
|
|
|
@pytest.fixture
|
|
def runtime():
|
|
# Loading source directly avoids a bytecode write in a read-only app tree.
|
|
module = ModuleType('container_runtime_under_test')
|
|
module.__file__ = str(SOURCE)
|
|
module.__dict__['__builtins__'] = dict(vars(builtins))
|
|
exec(compile(SOURCE.read_text(encoding='utf-8'), str(SOURCE), 'exec'), module.__dict__)
|
|
blocked = mock.Mock(side_effect=AssertionError('unexpected real runtime I/O'))
|
|
module.os = SimpleNamespace(**vars(os))
|
|
module.os.environ = {}
|
|
module.os.getenv = module.os.environ.get
|
|
for name in ('open', 'mkdir', 'chown', 'fchown', 'statvfs', 'access', 'execv'):
|
|
setattr(module.os, name, blocked)
|
|
module.os.umask = mock.Mock()
|
|
module.sys = SimpleNamespace(**vars(sys))
|
|
module.sys.path = list(sys.path)
|
|
module.subprocess = SimpleNamespace(run=blocked)
|
|
module.runpy = SimpleNamespace(run_path=blocked)
|
|
module.open = blocked
|
|
module.secrets = SimpleNamespace(
|
|
token_urlsafe=mock.Mock(return_value=PASSWORD),
|
|
token_hex=mock.Mock(return_value='a' * 32),
|
|
)
|
|
return module
|
|
|
|
|
|
@pytest.fixture
|
|
def linux(runtime):
|
|
entries = {}
|
|
for name in ('/', '/opt', '/opt/truf', '/opt/truf/app', '/data', '/run', '/run/truf'):
|
|
entries[name] = SimpleNamespace(st_mode=stat.S_IFDIR | 0o700, st_uid=10001, st_gid=10001)
|
|
for name in ('/.dockerenv', '/opt/truf/app/container_runtime.py', '/data/file'):
|
|
entries[name] = SimpleNamespace(st_mode=stat.S_IFREG | 0o600, st_uid=10001, st_gid=10001)
|
|
|
|
def inspect(path):
|
|
if path not in entries:
|
|
raise FileNotFoundError(path)
|
|
return entries[path]
|
|
|
|
lookup = mock.Mock(side_effect=inspect)
|
|
|
|
class FixturePath(PurePosixPath):
|
|
def lstat(self):
|
|
return lookup(str(self))
|
|
|
|
def is_file(self):
|
|
return str(self) in entries and stat.S_ISREG(self.lstat().st_mode)
|
|
|
|
runtime.Path = FixturePath
|
|
runtime.APP = FixturePath('/opt/truf/app')
|
|
runtime.DATA = FixturePath('/data')
|
|
runtime.RUN = FixturePath('/run/truf')
|
|
runtime.__file__ = '/opt/truf/app/container_runtime.py'
|
|
runtime.sys.platform = 'linux'
|
|
runtime.sys.flags = SimpleNamespace(isolated=1, no_site=1, dont_write_bytecode=1)
|
|
for name in ('getuid', 'geteuid', 'getgid'):
|
|
setattr(runtime.os, name, mock.Mock(return_value=10001))
|
|
runtime.os.ST_RDONLY = 1
|
|
runtime.os.statvfs = mock.Mock(return_value=SimpleNamespace(f_flag=1 | 4))
|
|
runtime.open = mock.mock_open(read_data=MOUNTS)
|
|
return SimpleNamespace(runtime=runtime, entries=entries, lookup=lookup)
|
|
|
|
|
|
@pytest.fixture
|
|
def layout(runtime, tmp_path, monkeypatch):
|
|
for name in ('APP', 'DATA', 'RUN'):
|
|
path = tmp_path / name.lower()
|
|
path.mkdir(mode=0o700)
|
|
setattr(runtime, name, path)
|
|
runtime.DEFAULT_CONFIG = runtime.APP / 'config.linux.yaml'
|
|
runtime.PROVISIONED = runtime.DATA / '.provisioned.json'
|
|
runtime.INITIALIZED = runtime.DATA / 'initialized.json'
|
|
runtime.INITIALIZE_LOCK = runtime.DATA / 'initialize.lock'
|
|
runtime.PASSWORD = runtime.DATA / 'postgres-password'
|
|
runtime.PROVIDER_SECRETS = runtime.DATA / 'config/secrets.yaml'
|
|
|
|
def private(path, *, directory=False):
|
|
path = Path(path)
|
|
assert path.is_relative_to(tmp_path), 'fixture escaped its temporary directory'
|
|
if not path.exists():
|
|
raise FileNotFoundError(path)
|
|
assert path.is_dir() if directory else path.is_file()
|
|
return path
|
|
|
|
def write_new(path, payload, **kwargs):
|
|
private(path.parent, directory=True)
|
|
with path.open('xb') as handle:
|
|
handle.write(payload)
|
|
|
|
def open_lock(path, flags, mode=0o600):
|
|
assert Path(path) == runtime.DATA / '.provision.lock'
|
|
return os.open(path, flags, mode)
|
|
|
|
def mkdir(path, mode):
|
|
assert Path(path).is_relative_to(runtime.DATA)
|
|
os.mkdir(path, mode)
|
|
|
|
runtime.private_path = mock.Mock(side_effect=private)
|
|
runtime._write_new = mock.Mock(side_effect=write_new)
|
|
runtime.os.open = mock.Mock(side_effect=open_lock)
|
|
runtime.os.O_NOFOLLOW = getattr(os, 'O_NOFOLLOW', 0)
|
|
runtime.os.mkdir = mock.Mock(side_effect=mkdir)
|
|
runtime.os.chown = mock.Mock()
|
|
runtime.os.fchown = mock.Mock()
|
|
lock = SimpleNamespace(LOCK_EX=2, LOCK_NB=4, flock=mock.Mock())
|
|
monkeypatch.setitem(sys.modules, 'fcntl', lock)
|
|
return SimpleNamespace(runtime=runtime, flock=lock.flock)
|
|
|
|
|
|
@pytest.fixture
|
|
def prepared(layout, monkeypatch):
|
|
r = layout.runtime
|
|
for name in r.DIRECTORIES:
|
|
(r.DATA / name).mkdir(parents=True, exist_ok=True, mode=0o700)
|
|
r.PROVISIONED.write_text(json.dumps({'format': r.FORMAT, 'uid': r.UID, 'gid': r.GID}), encoding='ascii')
|
|
r.DEFAULT_CONFIG.write_text('{}\n', encoding='ascii')
|
|
r.PASSWORD.write_text(PASSWORD + '\n', encoding='ascii')
|
|
r.PROVIDER_SECRETS.write_bytes(b'{}\n')
|
|
h = SimpleNamespace(runtime=r, trace=[], locked=False, cluster_locked=False)
|
|
h.dsn = 'postgresql://truf:' + PASSWORD + '@127.0.0.1:5432/truf'
|
|
r.os.environ['SCANNER_DB_URL'] = h.dsn
|
|
h.identity = {'system_identifier': '123456789', 'pg_major': 16}
|
|
h.paths = {
|
|
'identity_path': str(r.DATA / 'runtime-linux/postgres/cluster_identity.json'),
|
|
'data_dir': str(r.DATA / 'postgres-linux'),
|
|
}
|
|
# Fixed POSIX strings are compared as policy only, never used for host I/O.
|
|
h.config = {'global': {
|
|
'root_dir': '/opt/truf', 'project_dir': str(r.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',
|
|
'results_dir': '/data/runtime-linux/results', 'control_dir': str(r.RUN / 'control'),
|
|
'secrets_file': str(r.PROVIDER_SECRETS),
|
|
}, 'supervisor': {
|
|
'control_dir': str(r.RUN / 'control'),
|
|
'instance_file': str(r.RUN / 'control/supervisor.instance.json'),
|
|
'janitor': {'enabled': True},
|
|
}}
|
|
|
|
@contextmanager
|
|
def lock(path):
|
|
assert path == str(r.INITIALIZE_LOCK)
|
|
assert not h.locked
|
|
h.locked = True
|
|
h.trace.append('lock')
|
|
try:
|
|
yield
|
|
finally:
|
|
h.locked = False
|
|
h.trace.append('unlock')
|
|
|
|
@contextmanager
|
|
def cluster_lock(config, create_parent=False, endpoint_dsn=None):
|
|
assert config is h.config and create_parent is False
|
|
assert endpoint_dsn == h.dsn and h.locked
|
|
if h.cluster_locked:
|
|
raise BlockingIOError('another runtime owns cluster authority')
|
|
h.cluster_locked = True
|
|
h.trace.append('cluster.lock')
|
|
try:
|
|
yield
|
|
finally:
|
|
assert h.locked
|
|
h.cluster_locked = False
|
|
h.trace.append('cluster.unlock')
|
|
|
|
def replace(source, destination):
|
|
assert h.locked and h.cluster_locked
|
|
assert Path(source).parent == r.PROVIDER_SECRETS.parent
|
|
assert destination == str(r.PROVIDER_SECRETS)
|
|
h.trace.append('replace')
|
|
os.replace(source, destination)
|
|
|
|
def fsync_directory(path):
|
|
assert h.locked and h.cluster_locked
|
|
assert path == str(r.PROVIDER_SECRETS.parent)
|
|
h.trace.append('fsync')
|
|
|
|
def publish(path, value):
|
|
assert h.locked and path == str(r.INITIALIZED)
|
|
h.trace.append('publish')
|
|
with Path(path).open('x', encoding='ascii') as handle:
|
|
json.dump(value, handle)
|
|
|
|
def run(command, *, check):
|
|
assert check is True and h.locked
|
|
action = command[8] if command[6] == 'postgres-runtime' else command[6]
|
|
if action == 'initialize-empty':
|
|
assert not r.INITIALIZED.exists(), 'initialization marker was published early'
|
|
h.trace.append(action)
|
|
if action == 'initialize-empty':
|
|
Path(h.paths['identity_path']).write_text(json.dumps(h.identity), encoding='ascii')
|
|
(Path(h.paths['data_dir']) / 'PG_VERSION').write_text('16\n', encoding='ascii')
|
|
return SimpleNamespace(returncode=0)
|
|
|
|
h.pg = SimpleNamespace(
|
|
postgres_runtime_paths=mock.Mock(return_value=h.paths),
|
|
_load_config=mock.Mock(return_value=h.config),
|
|
)
|
|
h.security = SimpleNamespace(
|
|
PrivateFileLock=mock.Mock(side_effect=lock),
|
|
ClusterAuthorityLock=mock.Mock(side_effect=cluster_lock),
|
|
read_private_json=mock.Mock(side_effect=lambda path: json.loads(Path(path).read_text(encoding='ascii'))),
|
|
write_private_json_exclusive=mock.Mock(side_effect=publish),
|
|
preflight_lifecycle_paths=mock.Mock(),
|
|
durable_replace=mock.Mock(side_effect=replace),
|
|
fsync_directory=mock.Mock(side_effect=fsync_directory),
|
|
)
|
|
|
|
def validate_documents(config_path, *, secrets_bytes=None):
|
|
assert Path(config_path) == r.DEFAULT_CONFIG or Path(config_path).parent == r.DATA / 'config'
|
|
if secrets_bytes is None:
|
|
return SimpleNamespace(config=h.config, secrets={}, config_sha256='d' * 64)
|
|
try:
|
|
value = yaml.safe_load(secrets_bytes)
|
|
except yaml.YAMLError:
|
|
raise RuntimeError('invalid provider credential YAML') from None
|
|
if not isinstance(value, dict):
|
|
raise RuntimeError('provider credential YAML must be a mapping')
|
|
return SimpleNamespace(config=h.config, secrets=value, config_sha256='d' * 64)
|
|
|
|
h.document_io = SimpleNamespace(
|
|
load_managed_runtime_config=mock.Mock(return_value=SimpleNamespace(
|
|
config=h.config, secrets=None, config_sha256='d' * 64,
|
|
)),
|
|
validate_managed_runtime_files=mock.Mock(side_effect=validate_documents),
|
|
)
|
|
h.paths_module = SimpleNamespace(
|
|
apply_path_config=mock.Mock(side_effect=lambda config, _path: config),
|
|
)
|
|
monkeypatch.setitem(sys.modules, 'postgres_runtime', h.pg)
|
|
monkeypatch.setitem(sys.modules, 'paths', h.paths_module)
|
|
monkeypatch.setitem(sys.modules, 'runtime_security', h.security)
|
|
monkeypatch.setitem(sys.modules, 'runtime_document_io', h.document_io)
|
|
h.enable_dependencies = mock.Mock()
|
|
r.runpy = SimpleNamespace(run_path=mock.Mock(return_value={'_enable_dependency_paths': h.enable_dependencies}))
|
|
r.subprocess = SimpleNamespace(run=mock.Mock(side_effect=run))
|
|
return h
|
|
|
|
|
|
@pytest.mark.parametrize('directory', [False, True])
|
|
def test_private_path_accepts_exact_owner_only_mode(linux, directory):
|
|
path = '/data' if directory else '/data/file'
|
|
assert linux.runtime.private_path(path, directory=directory) == PurePosixPath(path)
|
|
assert mock.call('/') in linux.lookup.call_args_list
|
|
|
|
|
|
@pytest.mark.parametrize('path', ['relative', '/data/../file'])
|
|
def test_private_path_rejects_nonabsolute_or_parent_traversal_before_stat(linux, path):
|
|
with pytest.raises(RuntimeError, match='absolute and normalized'):
|
|
linux.runtime.private_path(path)
|
|
linux.lookup.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('path', ['/', '/data', '/data/file'])
|
|
def test_private_path_rejects_symlink_at_every_component(linux, path):
|
|
linux.entries[path].st_mode = stat.S_IFLNK | 0o600
|
|
with pytest.raises(RuntimeError, match='symlinks'):
|
|
linux.runtime.private_path('/data/file')
|
|
|
|
|
|
@pytest.mark.parametrize('error', [FileNotFoundError('missing path'), PermissionError('uninspectable path')])
|
|
def test_private_path_fails_closed_on_missing_or_uninspectable_components(linux, error):
|
|
linux.lookup.side_effect = error
|
|
with pytest.raises(type(error)):
|
|
linux.runtime.private_path('/data/file')
|
|
linux.runtime.os.chown.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('directory,kind,mode,uid,gid', [
|
|
(False, stat.S_IFREG, 0o600, 0, 10001),
|
|
(False, stat.S_IFREG, 0o600, 10001, 0),
|
|
(False, stat.S_IFREG, 0o640, 10001, 10001),
|
|
(False, stat.S_IFREG, 0o604, 10001, 10001),
|
|
(False, stat.S_IFREG, 0o400, 10001, 10001),
|
|
(False, stat.S_IFREG, 0o4600, 10001, 10001),
|
|
(False, stat.S_IFDIR, 0o600, 10001, 10001),
|
|
(False, stat.S_IFIFO, 0o600, 10001, 10001),
|
|
(False, stat.S_IFSOCK, 0o600, 10001, 10001),
|
|
(True, stat.S_IFREG, 0o700, 10001, 10001),
|
|
(True, stat.S_IFDIR, 0o750, 10001, 10001),
|
|
(True, stat.S_IFDIR, 0o1700, 10001, 10001),
|
|
])
|
|
def test_private_path_rejects_wrong_owner_group_type_or_exact_mode(linux, directory, kind, mode, uid, gid):
|
|
linux.entries['/data/file'] = SimpleNamespace(st_mode=kind | mode, st_uid=uid, st_gid=gid)
|
|
with pytest.raises(RuntimeError, match='ownership, type, or mode'):
|
|
linux.runtime.private_path('/data/file', directory=directory)
|
|
|
|
|
|
@pytest.mark.parametrize('filesystem', [b'ext4', b'xfs', b'btrfs', b'zfs'])
|
|
def test_container_requires_private_image_native_data_and_dedicated_tmpfs(linux, filesystem):
|
|
r = linux.runtime
|
|
r.open = mock.mock_open(read_data=MOUNTS.replace(b'ext4', filesystem))
|
|
r.require_container()
|
|
r.os.statvfs.assert_called_once_with(r.APP)
|
|
r.open.assert_called_once_with('/proc/self/mountinfo', 'rb')
|
|
r.open.return_value.read.assert_called_once_with(1024 * 1024 + 1)
|
|
r.os.umask.assert_called_once_with(0o077)
|
|
|
|
|
|
@pytest.mark.parametrize('gate', [
|
|
'platform', 'entrypoint', 'dockerenv', 'isolated', 'no_site', 'dont_write_bytecode',
|
|
'getuid', 'geteuid', 'getgid', 'writable_image',
|
|
])
|
|
def test_container_refuses_before_mount_inventory_when_identity_gate_fails(linux, gate):
|
|
r = linux.runtime
|
|
if gate == 'platform':
|
|
r.sys.platform = 'win32'
|
|
elif gate == 'entrypoint':
|
|
r.__file__ = '/tmp/container_runtime.py'
|
|
elif gate == 'dockerenv':
|
|
del linux.entries['/.dockerenv']
|
|
elif gate in ('isolated', 'no_site', 'dont_write_bytecode'):
|
|
setattr(r.sys.flags, gate, 0)
|
|
elif gate == 'writable_image':
|
|
r.os.statvfs.return_value.f_flag = 0
|
|
else:
|
|
getattr(r.os, gate).return_value = 0
|
|
with pytest.raises(RuntimeError):
|
|
r.require_container()
|
|
r.open.assert_not_called()
|
|
r.os.umask.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('uid,provisioning,allowed', [
|
|
(0, True, True), (0, False, False), (10001, True, False), (10002, False, False),
|
|
])
|
|
def test_root_is_permitted_only_for_explicit_provisioning(linux, uid, provisioning, allowed):
|
|
r = linux.runtime
|
|
r.os.getuid.return_value = r.os.geteuid.return_value = uid
|
|
if allowed:
|
|
r.require_container(provisioning=provisioning)
|
|
else:
|
|
with pytest.raises(RuntimeError, match='UID'):
|
|
r.require_container(provisioning=provisioning)
|
|
r.open.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('path', ['/opt/truf', '/opt/truf/app', '/opt/truf/app/container_runtime.py', '/data', '/run/truf'])
|
|
def test_container_rejects_nonprivate_required_paths(linux, path):
|
|
linux.entries[path].st_mode |= 0o040
|
|
with pytest.raises(RuntimeError, match='ownership, type, or mode'):
|
|
linux.runtime.require_container()
|
|
linux.runtime.os.umask.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('mounts', [
|
|
b'', MOUNTS.replace(b'/data', b'/'), MOUNTS.replace(b'/run/truf', b'/run'),
|
|
MOUNTS.replace(b'ext4', b'9p'), MOUNTS.replace(b'ext4', b'drvfs'),
|
|
MOUNTS.replace(b'ext4', b'overlay'), MOUNTS.replace(b'ext4', b'cifs'),
|
|
MOUNTS.replace(b'tmpfs', b'ext4'), b'x' * (1024 * 1024 + 1),
|
|
], ids=['missing', 'shared-data-root', 'shared-run-root', '9p', 'drvfs', 'overlay', 'cifs', 'not-tmpfs', 'oversized'])
|
|
def test_container_rejects_missing_shared_or_host_style_storage(linux, mounts):
|
|
linux.runtime.open = mock.mock_open(read_data=mounts)
|
|
with pytest.raises(RuntimeError):
|
|
linux.runtime.require_container()
|
|
linux.runtime.os.umask.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('provisioning', [False, True])
|
|
def test_new_private_file_is_exclusive_nofollow_and_durable(runtime, provisioning):
|
|
path = PurePosixPath('/data/config/fixture')
|
|
runtime.private_path = mock.Mock(side_effect=lambda path, **kwargs: path)
|
|
runtime.os.O_NOFOLLOW, runtime.os.O_DIRECTORY = 0x20000, 0x10000
|
|
runtime.os.open = mock.Mock(side_effect=[11, 12])
|
|
runtime.os.fdopen = mock.MagicMock()
|
|
handle = runtime.os.fdopen.return_value.__enter__.return_value
|
|
handle.fileno.return_value = 11
|
|
runtime.os.fchown = mock.Mock()
|
|
runtime.os.fsync = mock.Mock()
|
|
runtime.os.close = mock.Mock()
|
|
runtime._write_new(path, b'private fixture', provisioning=provisioning)
|
|
assert runtime.os.open.call_args_list == [
|
|
mock.call(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | runtime.os.O_NOFOLLOW, 0o600),
|
|
mock.call(path.parent, os.O_RDONLY | runtime.os.O_DIRECTORY),
|
|
]
|
|
handle.write.assert_called_once_with(b'private fixture')
|
|
handle.flush.assert_called_once()
|
|
assert runtime.os.fsync.call_args_list == [mock.call(11), mock.call(12)]
|
|
runtime.os.close.assert_called_once_with(12)
|
|
assert runtime.os.fchown.call_args_list == ([mock.call(11, 10001, 10001)] if provisioning else [])
|
|
assert runtime.private_path.call_args_list == [mock.call(path.parent, directory=True), mock.call(path)]
|
|
|
|
|
|
def test_new_private_file_never_overwrites_existing_file(runtime):
|
|
runtime.private_path = mock.Mock()
|
|
runtime.os.O_NOFOLLOW = 0x20000
|
|
runtime.os.open = mock.Mock(side_effect=FileExistsError('already exists'))
|
|
runtime.os.fdopen = mock.Mock()
|
|
with pytest.raises(FileExistsError):
|
|
runtime._write_new(PurePosixPath('/data/config/fixture'), b'new')
|
|
runtime.os.fdopen.assert_not_called()
|
|
|
|
|
|
def test_fresh_provision_generates_private_password_once_without_initializing_pg(layout, capsys):
|
|
r = layout.runtime
|
|
(r.DATA / 'home').mkdir(mode=0o700)
|
|
r.provision()
|
|
assert all((r.DATA / name).is_dir() for name in r.DIRECTORIES)
|
|
assert r.PASSWORD.read_text(encoding='ascii') == PASSWORD + '\n'
|
|
assert r.PROVIDER_SECRETS.read_bytes() == b'{}\n'
|
|
assert json.loads(r.PROVISIONED.read_text(encoding='ascii')) == {'format': r.FORMAT, 'uid': 10001, 'gid': 10001}
|
|
assert r._write_new.call_args_list[-1].args[0] == r.PROVISIONED
|
|
assert all(call.kwargs == {'provisioning': True} for call in r._write_new.call_args_list)
|
|
assert not r.INITIALIZED.exists()
|
|
assert not list((r.DATA / 'postgres-linux').iterdir())
|
|
before = {path: path.read_bytes() for path in r.DATA.rglob('*') if path.is_file()}
|
|
r.provision()
|
|
assert before == {path: path.read_bytes() for path in before}
|
|
r.secrets.token_urlsafe.assert_called_once_with(48)
|
|
assert r._write_new.call_count == 4
|
|
assert layout.flock.call_count == 2
|
|
r.subprocess.run.assert_not_called()
|
|
output = capsys.readouterr()
|
|
assert PASSWORD not in output.out + output.err
|
|
|
|
|
|
@pytest.mark.parametrize('relative', ['postgres-linux/PG_VERSION', 'config/partial', 'home/existing', 'unrelated'])
|
|
def test_provision_refuses_nonempty_or_partial_volume_without_repair(layout, relative):
|
|
r = layout.runtime
|
|
existing = r.DATA / relative
|
|
existing.parent.mkdir(parents=True, exist_ok=True)
|
|
existing.write_bytes(b'untouched')
|
|
with pytest.raises(RuntimeError, match='nonempty|partially initialized'):
|
|
r.provision()
|
|
assert existing.read_bytes() == b'untouched'
|
|
r._write_new.assert_not_called()
|
|
r.secrets.token_urlsafe.assert_not_called()
|
|
r.subprocess.run.assert_not_called()
|
|
|
|
|
|
def test_provision_lock_conflict_never_creates_layout_or_credentials(layout):
|
|
layout.flock.side_effect = BlockingIOError('another provisioner owns the volume')
|
|
with pytest.raises(BlockingIOError):
|
|
layout.runtime.provision()
|
|
layout.runtime.os.mkdir.assert_not_called()
|
|
layout.runtime._write_new.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('missing', ['PASSWORD', 'PROVIDER_SECRETS'])
|
|
def test_provision_does_not_regenerate_incomplete_marked_volume(layout, missing):
|
|
r = layout.runtime
|
|
r.provision()
|
|
getattr(r, missing).unlink()
|
|
r._write_new.reset_mock()
|
|
r.secrets.token_urlsafe.reset_mock()
|
|
with pytest.raises(FileNotFoundError):
|
|
r.provision()
|
|
r._write_new.assert_not_called()
|
|
r.secrets.token_urlsafe.assert_not_called()
|
|
assert r.PROVISIONED.is_file()
|
|
|
|
|
|
def test_provision_adds_only_managed_files_to_legacy_marked_volume(layout):
|
|
r = layout.runtime
|
|
r.provision()
|
|
managed_files = r.DATA / 'managed-files'
|
|
managed_files.rmdir()
|
|
r.os.mkdir.reset_mock()
|
|
r.os.chown.reset_mock()
|
|
r._write_new.reset_mock()
|
|
r.secrets.token_urlsafe.reset_mock()
|
|
|
|
r.provision()
|
|
|
|
assert managed_files.is_dir()
|
|
r.os.mkdir.assert_called_once_with(managed_files, 0o700)
|
|
r.os.chown.assert_called_once_with(managed_files, 10001, 10001)
|
|
r._write_new.assert_not_called()
|
|
r.secrets.token_urlsafe.assert_not_called()
|
|
|
|
|
|
def test_provision_refuses_wrong_existing_managed_files_path(layout):
|
|
r = layout.runtime
|
|
r.provision()
|
|
managed_files = r.DATA / 'managed-files'
|
|
managed_files.rmdir()
|
|
managed_files.write_bytes(b'not a directory')
|
|
r._write_new.reset_mock()
|
|
|
|
with pytest.raises(AssertionError):
|
|
r.provision()
|
|
|
|
assert managed_files.read_bytes() == b'not a directory'
|
|
r._write_new.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('payload', [b'{}', b'[]', b'{"format":"other"}', b'{', b' ' * 4097],
|
|
ids=['missing-format', 'nonmapping', 'wrong-format', 'partial-json', 'oversized'])
|
|
def test_bad_provision_marker_prevents_environment_or_application_imports(prepared, payload):
|
|
r = prepared.runtime
|
|
r.PROVISIONED.write_bytes(payload)
|
|
with pytest.raises((RuntimeError, ValueError)):
|
|
r.prepare_environment(r.DEFAULT_CONFIG)
|
|
assert not list(r.RUN.iterdir())
|
|
r.runpy.run_path.assert_not_called()
|
|
prepared.pg._load_config.assert_not_called()
|
|
|
|
|
|
def test_prepare_environment_scrubs_dsn_and_runtime_overrides_before_imports(prepared, capsys):
|
|
h, r = prepared, prepared.runtime
|
|
overrides = ('PGHOSTADDR', 'PgServiceFile', 'TRUF_POSTGRES_PASSWORD', 'TRUF_CONTAINER_CONFIG',
|
|
'SCANNER_DB_URL', 'SCAN_SLOT_LIMIT', 'TRUFFLEHOG_PATH', 'KEYCHECK_SECRET',
|
|
'DATABASE_URL', 'PYTHONPATH', 'PYTHONHOME')
|
|
r.os.environ.update({name: 'untrusted-override' for name in overrides})
|
|
r.os.environ['LANG'] = 'C.UTF-8'
|
|
|
|
def bootstrap(path):
|
|
assert 'untrusted-override' not in r.os.environ.values()
|
|
return {'_enable_dependency_paths': h.enable_dependencies}
|
|
|
|
r.runpy.run_path.side_effect = bootstrap
|
|
assert r.prepare_environment(r.DEFAULT_CONFIG) is h.config
|
|
for name in ('authority', 'control', 'tmp'):
|
|
assert (r.RUN / name).is_dir()
|
|
assert mock.call(r.RUN / name, directory=True) in r.private_path.call_args_list
|
|
url = 'postgresql://truf:' + PASSWORD + '@127.0.0.1:5432/truf'
|
|
for name in ('SCANNER_DB_URL', 'DATABASE_URL', 'TRUF_MANAGED_POSTGRES_DSN'):
|
|
assert r.os.environ[name] == url
|
|
assert r.os.environ['TRUF_POSTGRES_PASSWORD'] == PASSWORD
|
|
assert r.os.environ['PATH'] == '/usr/local/bin:/usr/bin:/bin:/usr/lib/postgresql/16/bin'
|
|
assert r.os.environ['HOME'] == '/data/home'
|
|
assert {r.os.environ[name] for name in ('TMP', 'TEMP', 'TMPDIR')} == {str(r.RUN / 'tmp')}
|
|
assert r.os.environ['LANG'] == 'C.UTF-8'
|
|
h.enable_dependencies.assert_called_once_with('supervisor')
|
|
r.runpy.run_path.assert_called_once_with(str(r.APP / 'child_bootstrap.py'))
|
|
h.security.preflight_lifecycle_paths.assert_called_once_with(
|
|
str(r.DEFAULT_CONFIG), h.config, authority_profile='server',
|
|
)
|
|
h.document_io.validate_managed_runtime_files.assert_called_once_with(str(r.DEFAULT_CONFIG))
|
|
h.pg._load_config.assert_not_called()
|
|
assert PASSWORD not in capsys.readouterr().out
|
|
|
|
|
|
@pytest.mark.parametrize('password', ['', 'x' * 31, 'x' * 129, 'x' * 31 + ':', 'x' * 32 + '\nembedded'],
|
|
ids=['empty', 'short', 'long', 'unsafe-character', 'multiline'])
|
|
def test_prepare_rejects_invalid_generated_password_before_imports(prepared, password):
|
|
r = prepared.runtime
|
|
r.PASSWORD.write_text(password, encoding='ascii')
|
|
with pytest.raises(RuntimeError, match='password is invalid'):
|
|
r.prepare_environment(r.DEFAULT_CONFIG)
|
|
r.runpy.run_path.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('field', [
|
|
'root_dir', 'project_dir', 'runtime_dir', 'postgres_data_dir', 'postgres_bin_dir',
|
|
'result_bundle_dir', 'work_dir', 'control_dir', 'secrets_file', 'supervisor.control_dir',
|
|
])
|
|
def test_prepare_refuses_config_escaping_fixed_storage_contract(prepared, field):
|
|
section, key = ('supervisor', 'control_dir') if field.startswith('supervisor.') else ('global', field)
|
|
prepared.config[section][key] = '/foreign/storage'
|
|
with pytest.raises(RuntimeError, match='storage contract|ephemeral storage'):
|
|
prepared.runtime.prepare_environment(prepared.runtime.DEFAULT_CONFIG)
|
|
prepared.security.preflight_lifecycle_paths.assert_not_called()
|
|
prepared.document_io.validate_managed_runtime_files.assert_called_once_with(
|
|
str(prepared.runtime.DEFAULT_CONFIG),
|
|
)
|
|
|
|
|
|
@pytest.mark.parametrize('allowed', [False, True])
|
|
def test_only_default_image_config_or_private_data_config_is_accepted(prepared, allowed):
|
|
r = prepared.runtime
|
|
config = (r.DATA / 'config' if allowed else r.APP) / 'alternate.yaml'
|
|
config.write_text('{}\n', encoding='ascii')
|
|
if allowed:
|
|
r.prepare_environment(config)
|
|
prepared.document_io.validate_managed_runtime_files.assert_called_once_with(
|
|
str(config),
|
|
)
|
|
else:
|
|
with pytest.raises(RuntimeError, match='configuration must be'):
|
|
r.prepare_environment(config)
|
|
r.runpy.run_path.assert_not_called()
|
|
|
|
|
|
def test_bootstrap_target_is_the_real_runtime_safety_migration_module():
|
|
tree = ast.parse(SOURCE.with_name('runtime_bootstrap.py').read_text(encoding='utf-8'))
|
|
targets = next(ast.literal_eval(node.value) for node in tree.body
|
|
if isinstance(node, ast.Assign)
|
|
and any(isinstance(target, ast.Name) and target.id == 'TARGETS' for target in node.targets))
|
|
assert targets['migrate-runtime-safety'] == 'migrate_runtime_safety.py'
|
|
|
|
|
|
def test_initialize_publishes_marker_only_after_migration_and_confirmed_stop(prepared, capsys):
|
|
h, r = prepared, prepared.runtime
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert h.trace == ['lock', 'initialize-empty', 'maintenance-start', 'migrate-runtime-safety',
|
|
'maintenance-stop', 'publish', 'unlock']
|
|
prefix = [sys.executable, '-u', '-I', '-S', '-B', str(r.APP / 'runtime_bootstrap.py')]
|
|
assert r.subprocess.run.call_args_list == [
|
|
mock.call(prefix + ['postgres-runtime', '--', 'initialize-empty', '--config', str(r.DEFAULT_CONFIG)], check=True),
|
|
mock.call(prefix + ['postgres-runtime', '--', 'maintenance-start', '--config', str(r.DEFAULT_CONFIG)], check=True),
|
|
mock.call(prefix + ['migrate-runtime-safety', '--', '--config', str(r.DEFAULT_CONFIG),
|
|
'--initialize-base', '--apply', '--sources-stopped'], check=True),
|
|
mock.call(prefix + ['postgres-runtime', '--', 'maintenance-stop', '--config', str(r.DEFAULT_CONFIG)], check=True),
|
|
]
|
|
assert json.loads(r.INITIALIZED.read_text(encoding='ascii')) == {'format': r.FORMAT, **h.identity}
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert h.trace[-5:] == [
|
|
'lock', 'maintenance-start', 'migrate-runtime-safety', 'maintenance-stop', 'unlock',
|
|
]
|
|
assert r.subprocess.run.call_count == 7
|
|
h.security.write_private_json_exclusive.assert_called_once()
|
|
assert PASSWORD not in capsys.readouterr().out
|
|
|
|
|
|
@pytest.mark.parametrize('stage', ['maintenance-start', 'migrate-runtime-safety', 'maintenance-stop'])
|
|
def test_existing_cluster_upgrade_failure_preserves_marker_and_stops_maintenance(
|
|
prepared, stage):
|
|
h, r = prepared, prepared.runtime
|
|
Path(h.paths['identity_path']).write_text(json.dumps(h.identity), encoding='ascii')
|
|
marker = {'format': r.FORMAT, **h.identity}
|
|
r.INITIALIZED.write_text(json.dumps(marker), encoding='ascii')
|
|
run = r.subprocess.run.side_effect
|
|
|
|
def fail(command, **kwargs):
|
|
result = run(command, **kwargs)
|
|
if h.trace[-1] == stage:
|
|
raise subprocess.CalledProcessError(1, command)
|
|
return result
|
|
|
|
r.subprocess.run.side_effect = fail
|
|
with pytest.raises(subprocess.CalledProcessError):
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert ('maintenance-stop' in h.trace) is True
|
|
assert json.loads(r.INITIALIZED.read_text(encoding='ascii')) == marker
|
|
h.security.write_private_json_exclusive.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('stage', ['initialize-empty', 'maintenance-start', 'migrate-runtime-safety', 'maintenance-stop'])
|
|
def test_failed_initialization_always_stops_maintenance_and_never_marks_success(prepared, stage, capsys):
|
|
h, r = prepared, prepared.runtime
|
|
run = r.subprocess.run.side_effect
|
|
|
|
def fail(command, **kwargs):
|
|
result = run(command, **kwargs)
|
|
if h.trace[-1] == stage:
|
|
raise subprocess.CalledProcessError(1, command)
|
|
return result
|
|
|
|
r.subprocess.run.side_effect = fail
|
|
with pytest.raises(subprocess.CalledProcessError):
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert ('maintenance-stop' in h.trace) is (stage != 'initialize-empty')
|
|
assert not r.INITIALIZED.exists()
|
|
h.security.write_private_json_exclusive.assert_not_called()
|
|
assert h.trace[-1] == 'unlock'
|
|
assert 'confirmed stopped' not in capsys.readouterr().out
|
|
|
|
|
|
@pytest.mark.parametrize('stage', ['before', 'initialize-empty', 'maintenance-start'])
|
|
def test_shutdown_during_initialization_never_publishes_early_marker(prepared, stage):
|
|
h, r = prepared, prepared.runtime
|
|
run = r.subprocess.run.side_effect
|
|
|
|
def request_shutdown(command, **kwargs):
|
|
result = run(command, **kwargs)
|
|
if h.trace[-1] == stage:
|
|
r._shutdown_requested = True
|
|
return result
|
|
|
|
r._shutdown_requested = stage == 'before'
|
|
r.subprocess.run.side_effect = request_shutdown
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert ('maintenance-stop' in h.trace) is (stage == 'maintenance-start')
|
|
assert 'migrate-runtime-safety' not in h.trace
|
|
assert not r.INITIALIZED.exists()
|
|
h.security.write_private_json_exclusive.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('identity_present', [False, True])
|
|
def test_partial_initialization_is_not_adopted_or_repaired(prepared, identity_present):
|
|
h, r = prepared, prepared.runtime
|
|
existing = Path(h.paths['identity_path']) if identity_present else Path(h.paths['data_dir']) / 'PG_VERSION'
|
|
existing.write_bytes(b'untouched')
|
|
with pytest.raises(RuntimeError, match='partial initialization'):
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
assert existing.read_bytes() == b'untouched'
|
|
r.subprocess.run.assert_not_called()
|
|
h.security.write_private_json_exclusive.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('difference', [{'system_identifier': 'other'}, {'system_identifier': None}, {'pg_major': 15}],
|
|
ids=['system-identifier', 'missing-system-identifier', 'postgres-major'])
|
|
def test_initialized_marker_must_match_bound_cluster_identity(prepared, difference):
|
|
h, r = prepared, prepared.runtime
|
|
Path(h.paths['identity_path']).write_text(json.dumps(h.identity), encoding='ascii')
|
|
r.INITIALIZED.write_text(json.dumps({'format': r.FORMAT, **h.identity, **difference}), encoding='ascii')
|
|
with pytest.raises(RuntimeError, match='marker|cluster'):
|
|
r.initialize(r.DEFAULT_CONFIG, h.config)
|
|
r.subprocess.run.assert_not_called()
|
|
h.security.write_private_json_exclusive.assert_not_called()
|
|
|
|
|
|
@pytest.fixture
|
|
def healthy(prepared, monkeypatch):
|
|
h, r = prepared, prepared.runtime
|
|
r.INITIALIZED.write_text(json.dumps({'format': r.FORMAT, **h.identity}), encoding='ascii')
|
|
r.os.environ['SCANNER_DB_URL'] = h.dsn
|
|
h.metadata = {'instance_id': 'fixture-instance', 'activation_state': 'ACTIVE'}
|
|
h.snapshot = {'activation_state': 'ACTIVE', 'postgres': {'state': 'READY', 'ready': True},
|
|
'signature': [(name, 'running', 'running', 123, None, 0, 0, '', '')
|
|
for name in ('result-ingester', 'jsonl-projector', 'janitor')]}
|
|
h.db = mock.Mock(enabled=True)
|
|
h.db.conn.is_postgres = True
|
|
h.db.pipeline_worker_health.return_value = {'healthy': True}
|
|
h.db.conn.execute.return_value.fetchone.return_value = {'data_directory': '/data/postgres-linux', 'version': '160015'}
|
|
h.database = mock.Mock(return_value=h.db)
|
|
h.load_metadata = mock.Mock(return_value=h.metadata)
|
|
h.control_snapshot = mock.Mock(return_value=h.snapshot)
|
|
monkeypatch.setitem(sys.modules, 'scanner_db', SimpleNamespace(ScannerDB=h.database))
|
|
monkeypatch.setitem(sys.modules, 'lifecycle_authority', SimpleNamespace(
|
|
DISCOVERY_PRODUCER_SOURCES=('gitlab', 'dockerhub', 'huggingface'),
|
|
))
|
|
monkeypatch.setitem(sys.modules, 'supervisor', SimpleNamespace(get_control_snapshot=h.control_snapshot))
|
|
monkeypatch.setitem(sys.modules, 'supervisor_instance', SimpleNamespace(load_instance_metadata=h.load_metadata))
|
|
r.os.access = mock.Mock(return_value=True)
|
|
r.os.statvfs = mock.Mock(return_value=SimpleNamespace(f_bavail=1))
|
|
return h
|
|
|
|
|
|
@pytest.mark.parametrize('payload', [None, b'{', b'{}', b'[]', b'{"format":"other"}', b' ' * 4097],
|
|
ids=['absent', 'partial-json', 'missing-format', 'nonmapping', 'wrong-format', 'oversized'])
|
|
def test_health_does_not_import_database_or_control_before_valid_marker(healthy, payload):
|
|
h, r = healthy, healthy.runtime
|
|
if payload is None:
|
|
r.INITIALIZED.unlink()
|
|
else:
|
|
r.INITIALIZED.write_bytes(payload)
|
|
imports = mock.Mock(side_effect=AssertionError('application import before initialization marker'))
|
|
r.__dict__['__builtins__']['__import__'] = imports
|
|
with pytest.raises((RuntimeError, ValueError, FileNotFoundError)):
|
|
r.health(h.config)
|
|
imports.assert_not_called()
|
|
h.database.assert_not_called()
|
|
h.load_metadata.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('state,postgres', [
|
|
('ACTIVATING', {'state': 'READY', 'ready': True}),
|
|
('STOPPING', {'state': 'READY', 'ready': True}),
|
|
('ACTIVE', {'state': 'STARTING', 'ready': True}),
|
|
('ACTIVE', {'state': 'READY', 'ready': False}),
|
|
('ACTIVE', {'state': 'READY', 'ready': 1}), ('ACTIVE', {}),
|
|
])
|
|
def test_health_requires_active_supervisor_and_explicit_postgres_ready(healthy, state, postgres):
|
|
healthy.snapshot.update(activation_state=state, postgres=postgres)
|
|
with pytest.raises(RuntimeError, match='not ready'):
|
|
healthy.runtime.health(healthy.config)
|
|
healthy.database.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('index,value', [(1, 'stopped'), (2, 'stopped'), (3, None), (7, 'blocked'), (8, 'uncertain owner')])
|
|
def test_health_rejects_missing_or_unready_pipeline_processes_before_db(healthy, index, value):
|
|
row = list(healthy.snapshot['signature'][0])
|
|
row[index] = value
|
|
healthy.snapshot['signature'][0] = tuple(row)
|
|
with pytest.raises(RuntimeError, match='worker|ownership'):
|
|
healthy.runtime.health(healthy.config)
|
|
healthy.database.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('name', ['result-ingester', 'jsonl-projector', 'janitor'])
|
|
def test_health_requires_each_enabled_pipeline_worker(healthy, name):
|
|
healthy.snapshot['signature'] = [row for row in healthy.snapshot['signature'] if row[0] != name]
|
|
with pytest.raises(RuntimeError, match='worker'):
|
|
healthy.runtime.health(healthy.config)
|
|
healthy.database.assert_not_called()
|
|
|
|
|
|
def test_health_allows_explicitly_disabled_janitor(healthy):
|
|
healthy.config['supervisor']['janitor']['enabled'] = False
|
|
healthy.snapshot['signature'].pop()
|
|
assert healthy.runtime.health(healthy.config)['healthy'] is True
|
|
|
|
|
|
def test_health_requires_worker_api_only_when_explicitly_enabled(healthy):
|
|
healthy.config['supervisor']['worker_api'] = {'enabled': True}
|
|
with pytest.raises(RuntimeError, match='worker'):
|
|
healthy.runtime.health(healthy.config)
|
|
healthy.database.assert_not_called()
|
|
|
|
healthy.snapshot['signature'].append(
|
|
('worker-api', 'running', 'running', 124, None, 0, 0, '', ''),
|
|
)
|
|
result = healthy.runtime.health(healthy.config)
|
|
assert 'worker-api' in result['workers']
|
|
|
|
|
|
def test_strict_health_requires_enabled_worker_api_before_database(healthy):
|
|
h, r = healthy, healthy.runtime
|
|
r._probe_worker_api = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='required but disabled'):
|
|
r.health(h.config, require_worker_api=True)
|
|
r._probe_worker_api.assert_not_called()
|
|
h.database.assert_not_called()
|
|
|
|
|
|
def test_strict_health_probes_exact_unauthorized_worker_endpoint(healthy, monkeypatch):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['supervisor']['worker_api'] = {
|
|
'enabled': True, 'address': '127.0.0.1', 'port': 8766,
|
|
}
|
|
h.snapshot['signature'].append(
|
|
('worker-api', 'running', 'running', 124, None, 0, 0, '', ''),
|
|
)
|
|
response = mock.Mock(status=401)
|
|
response.getheader.return_value = 'Bearer'
|
|
response.read.return_value = json.dumps({
|
|
'error': {
|
|
'code': 'unauthorized',
|
|
'message': 'worker credentials are invalid',
|
|
},
|
|
}).encode('ascii')
|
|
connection = mock.Mock()
|
|
connection.getresponse.return_value = response
|
|
factory = mock.Mock(return_value=connection)
|
|
monkeypatch.setattr(r.http.client, 'HTTPConnection', factory)
|
|
|
|
result = r.health(h.config, require_worker_api=True)
|
|
|
|
assert result['healthy'] is True
|
|
factory.assert_called_once_with('127.0.0.1', 8766, timeout=2)
|
|
connection.request.assert_called_once_with(
|
|
'POST', '/api/v1/worker/claim', body=b'',
|
|
headers={
|
|
'Authorization': 'Bearer ' + r.WORKER_HEALTH_TOKEN,
|
|
'Content-Length': '0',
|
|
},
|
|
)
|
|
response.read.assert_called_once_with(4097)
|
|
connection.close.assert_called_once_with()
|
|
|
|
|
|
@pytest.mark.parametrize('failure', [
|
|
'connect', 'status', 'challenge', 'body', 'oversized',
|
|
])
|
|
def test_strict_worker_probe_fails_closed_before_database(healthy, monkeypatch, failure):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['supervisor']['worker_api'] = {
|
|
'enabled': True, 'address': '127.0.0.1', 'port': 8766,
|
|
}
|
|
h.snapshot['signature'].append(
|
|
('worker-api', 'running', 'running', 124, None, 0, 0, '', ''),
|
|
)
|
|
response = mock.Mock(status=401)
|
|
response.getheader.return_value = 'Bearer'
|
|
response.read.return_value = json.dumps({
|
|
'error': {
|
|
'code': 'unauthorized',
|
|
'message': 'worker credentials are invalid',
|
|
},
|
|
}).encode('ascii')
|
|
connection = mock.Mock()
|
|
connection.getresponse.return_value = response
|
|
factory = mock.Mock(return_value=connection)
|
|
if failure == 'connect':
|
|
factory.side_effect = OSError('private network detail')
|
|
elif failure == 'status':
|
|
response.status = 200
|
|
elif failure == 'challenge':
|
|
response.getheader.return_value = 'Basic'
|
|
elif failure == 'body':
|
|
response.read.return_value = b'{"error":{"code":"foreign"}}'
|
|
else:
|
|
response.read.return_value = b'x' * 4097
|
|
monkeypatch.setattr(r.http.client, 'HTTPConnection', factory)
|
|
|
|
with pytest.raises(RuntimeError, match='worker API is not ready') as raised:
|
|
r.health(h.config, require_worker_api=True)
|
|
|
|
assert 'private network detail' not in str(raised.value)
|
|
h.database.assert_not_called()
|
|
if failure != 'connect':
|
|
connection.close.assert_called_once_with()
|
|
|
|
|
|
@pytest.mark.parametrize('status,error', [('failed', ''), ('stopped', 'uncertain owner')])
|
|
def test_health_rejects_failed_or_uncertain_managed_source(healthy, status, error):
|
|
healthy.snapshot['signature'].append(('fixture-source', status, 'stopped', None, None, 0, 0, '', error))
|
|
with pytest.raises(RuntimeError, match='failed|ownership'):
|
|
healthy.runtime.health(healthy.config)
|
|
healthy.database.assert_not_called()
|
|
|
|
|
|
def test_strict_health_accepts_running_and_successfully_waiting_discovery(healthy):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['sources'] = {
|
|
'gitlab': {'enabled': True},
|
|
'dockerhub': {'enabled': True},
|
|
'huggingface': {'enabled': False},
|
|
}
|
|
h.snapshot['signature'].extend((
|
|
('gitlab', 'running', 'running', 201, None, 0, 0, '', ''),
|
|
('dockerhub', 'waiting', 'running', None, 0, 0, 0, '', ''),
|
|
))
|
|
|
|
result = r.health(h.config, require_discovery_producers=True)
|
|
|
|
assert {'gitlab', 'dockerhub'}.issubset(result['workers'])
|
|
|
|
|
|
@pytest.mark.parametrize('row', [
|
|
None,
|
|
('gitlab', 'blocked', 'running', None, None, 0, 0, '', ''),
|
|
('gitlab', 'waiting', 'running', None, None, 0, 0, '', ''),
|
|
('gitlab', 'waiting', 'running', None, 1, 0, 1, '', ''),
|
|
('gitlab', 'running', 'paused', 201, None, 0, 0, '', ''),
|
|
('gitlab', 'running', 'running', 201, None, 0, 0, 'blocked', ''),
|
|
], ids=['absent', 'blocked', 'never-ran', 'failed-exit', 'not-desired', 'runtime-blocked'])
|
|
def test_strict_health_rejects_unready_discovery_before_database(healthy, row):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['sources'] = {'gitlab': {'enabled': True}}
|
|
if row is not None:
|
|
h.snapshot['signature'].append(row)
|
|
|
|
with pytest.raises(RuntimeError, match='discovery producer'):
|
|
r.health(h.config, require_discovery_producers=True)
|
|
|
|
h.database.assert_not_called()
|
|
|
|
|
|
def test_strict_health_honors_supervisor_enabled_discovery(healthy):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['sources'] = {'gitlab': {'enabled': False}}
|
|
h.config['supervisor']['sources'] = {'gitlab': {'enabled': True}}
|
|
with pytest.raises(RuntimeError, match='discovery producer'):
|
|
r.health(h.config, require_discovery_producers=True)
|
|
h.database.assert_not_called()
|
|
|
|
|
|
def test_strict_health_honors_supervisor_disabled_discovery(healthy):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['sources'] = {'gitlab': {'enabled': True}}
|
|
h.config['supervisor']['sources'] = {'gitlab': {'enabled': False}}
|
|
|
|
result = r.health(h.config, require_discovery_producers=True)
|
|
|
|
assert result['healthy'] is True
|
|
|
|
|
|
def test_health_checks_readonly_schema_cutover_leases_identity_then_storage(healthy, capsys):
|
|
h, r = healthy, healthy.runtime
|
|
result = r.health(h.config)
|
|
assert result == {'healthy': True, 'activation_state': 'ACTIVE', 'postgres': 'READY',
|
|
'workers': ['janitor', 'jsonl-projector', 'result-ingester']}
|
|
h.load_metadata.assert_called_once_with(h.config['supervisor']['instance_file'])
|
|
h.control_snapshot.assert_called_once_with(h.metadata)
|
|
h.database.assert_called_once_with(db_url=r.os.environ['SCANNER_DB_URL'], initialize=False)
|
|
assert h.db.mock_calls[:7] == [
|
|
mock.call.set_application_name('truf-container-health'),
|
|
mock.call.conn.execute('SET default_transaction_read_only = on'), mock.call.conn.commit(),
|
|
mock.call.require_runtime_safety_schema(), mock.call.require_final_cutover(),
|
|
mock.call.pipeline_worker_health('result_ingester', 'fixture-instance'),
|
|
mock.call.pipeline_worker_health('jsonl_projector', 'fixture-instance'),
|
|
]
|
|
h.db.close.assert_called_once()
|
|
assert r.os.access.call_args_list == [mock.call(h.config['global'][name], os.W_OK | os.X_OK)
|
|
for name in ('work_dir', 'result_bundle_dir', 'results_dir')]
|
|
assert PASSWORD not in json.dumps(result) + capsys.readouterr().out
|
|
|
|
|
|
@pytest.mark.parametrize('failure', ['disabled', 'not_postgres', 'schema', 'cutover', 'ingester_lease', 'projector_lease', 'directory', 'version'])
|
|
def test_health_fails_closed_and_closes_database_on_each_database_gate(healthy, failure):
|
|
h = healthy
|
|
if failure == 'disabled':
|
|
h.db.enabled = False
|
|
elif failure == 'not_postgres':
|
|
h.db.conn.is_postgres = False
|
|
elif failure in ('schema', 'cutover'):
|
|
method = h.db.require_runtime_safety_schema if failure == 'schema' else h.db.require_final_cutover
|
|
method.side_effect = RuntimeError('missing durable schema or cutover')
|
|
elif failure.endswith('_lease'):
|
|
h.db.pipeline_worker_health.side_effect = lambda name, instance: {
|
|
'healthy': name != ('result_ingester' if failure == 'ingester_lease' else 'jsonl_projector')}
|
|
else:
|
|
row = h.db.conn.execute.return_value.fetchone.return_value
|
|
row['data_directory' if failure == 'directory' else 'version'] = '/foreign/data' if failure == 'directory' else '170001'
|
|
with pytest.raises(RuntimeError):
|
|
h.runtime.health(h.config)
|
|
h.db.close.assert_called_once()
|
|
h.runtime.os.access.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('name', ['work_dir', 'result_bundle_dir', 'results_dir'])
|
|
@pytest.mark.parametrize('failure', ['unwritable', 'full'])
|
|
def test_health_requires_writable_nonfull_persistent_storage(healthy, name, failure):
|
|
h, r = healthy, healthy.runtime
|
|
target = h.config['global'][name]
|
|
|
|
def access(path, mode):
|
|
h.db.close.assert_called_once()
|
|
return path != target or failure != 'unwritable'
|
|
|
|
r.os.access.side_effect = access
|
|
r.os.statvfs.side_effect = lambda path: SimpleNamespace(f_bavail=int(path != target or failure != 'full'))
|
|
with pytest.raises(RuntimeError, match='storage'):
|
|
r.health(h.config)
|
|
|
|
|
|
@pytest.mark.parametrize('active', ['supervisor', 'postgres'])
|
|
def test_secret_import_refuses_active_runtime_before_reading_stdin(prepared, active):
|
|
h, r = prepared, prepared.runtime
|
|
path = Path(h.config['supervisor']['instance_file']) if active == 'supervisor' else Path(h.paths['data_dir']) / 'postmaster.pid'
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
path.write_bytes(b'active or uncertain owner')
|
|
r.sys.stdin = SimpleNamespace(buffer=mock.Mock())
|
|
with pytest.raises(RuntimeError, match='stop the runtime'):
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
h.security.ClusterAuthorityLock.assert_called_once_with(h.config, create_parent=False, endpoint_dsn=h.dsn)
|
|
assert h.trace == ['lock', 'cluster.lock', 'cluster.unlock', 'unlock']
|
|
assert not h.locked and not h.cluster_locked
|
|
r.sys.stdin.buffer.read.assert_not_called()
|
|
r._write_new.assert_not_called()
|
|
h.security.durable_replace.assert_not_called()
|
|
h.document_io.validate_managed_runtime_files.assert_not_called()
|
|
|
|
|
|
def test_secret_import_cluster_lock_conflict_precedes_stopped_check_and_stdin(prepared):
|
|
h, r = prepared, prepared.runtime
|
|
h.cluster_locked = True
|
|
r.sys.stdin = SimpleNamespace(buffer=mock.Mock())
|
|
with pytest.raises(BlockingIOError, match='cluster authority'):
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
h.security.ClusterAuthorityLock.assert_called_once_with(h.config, create_parent=False, endpoint_dsn=h.dsn)
|
|
assert h.trace == ['lock', 'unlock']
|
|
assert not h.locked and h.cluster_locked
|
|
h.pg.postgres_runtime_paths.assert_not_called()
|
|
r.sys.stdin.buffer.read.assert_not_called()
|
|
r._write_new.assert_not_called()
|
|
r.secrets.token_hex.assert_not_called()
|
|
h.security.durable_replace.assert_not_called()
|
|
h.security.fsync_directory.assert_not_called()
|
|
h.document_io.validate_managed_runtime_files.assert_not_called()
|
|
assert r.PROVIDER_SECRETS.read_bytes() == b'{}\n'
|
|
|
|
|
|
def test_secret_import_is_locked_atomic_private_and_never_prints_credentials(prepared, capsys):
|
|
h, r = prepared, prepared.runtime
|
|
payload = ('provider:\n token: ' + PROVIDER_KEY + '\n').encode('ascii')
|
|
stream = io.BytesIO(payload)
|
|
|
|
def read(size):
|
|
assert h.locked and h.cluster_locked
|
|
h.trace.append('stdin')
|
|
return stream.read(size)
|
|
|
|
r.sys.stdin = SimpleNamespace(buffer=mock.Mock())
|
|
r.sys.stdin.buffer.read.side_effect = read
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
h.security.ClusterAuthorityLock.assert_called_once_with(h.config, create_parent=False, endpoint_dsn=h.dsn)
|
|
assert h.trace == ['lock', 'cluster.lock', 'stdin', 'replace', 'fsync', 'cluster.unlock', 'unlock']
|
|
assert not h.locked and not h.cluster_locked
|
|
r.sys.stdin.buffer.read.assert_called_once_with(1024 * 1024 + 1)
|
|
h.document_io.validate_managed_runtime_files.assert_called_once_with(
|
|
str(r.DEFAULT_CONFIG), secrets_bytes=payload,
|
|
)
|
|
assert yaml.safe_load(r.PROVIDER_SECRETS.read_bytes()) == {'provider': {'token': PROVIDER_KEY}}
|
|
temporary = r._write_new.call_args.args[0]
|
|
assert temporary.parent == r.PROVIDER_SECRETS.parent and not temporary.exists()
|
|
h.security.durable_replace.assert_called_once_with(str(temporary), str(r.PROVIDER_SECRETS))
|
|
h.security.fsync_directory.assert_called_once_with(str(r.PROVIDER_SECRETS.parent))
|
|
output = capsys.readouterr()
|
|
assert PROVIDER_KEY not in output.out + output.err
|
|
assert 'credentials replaced' in output.out
|
|
|
|
|
|
@pytest.mark.parametrize('payload', [b'', b'x' * (1024 * 1024 + 1), b'[]', b'null', PROVIDER_KEY.encode('ascii'), ('token: [' + PROVIDER_KEY).encode('ascii')],
|
|
ids=['empty', 'oversized', 'sequence', 'null', 'scalar', 'invalid-yaml'])
|
|
def test_invalid_secret_input_is_bounded_and_does_not_leak_or_replace(prepared, payload, capsys):
|
|
h, r = prepared, prepared.runtime
|
|
r.sys.stdin = SimpleNamespace(buffer=mock.Mock(wraps=io.BytesIO(payload)))
|
|
with pytest.raises(RuntimeError) as failure:
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
assert h.trace == ['lock', 'cluster.lock', 'cluster.unlock', 'unlock']
|
|
assert not h.locked and not h.cluster_locked
|
|
assert PROVIDER_KEY not in str(failure.value)
|
|
r.sys.stdin.buffer.read.assert_called_once_with(1024 * 1024 + 1)
|
|
r._write_new.assert_not_called()
|
|
h.security.durable_replace.assert_not_called()
|
|
assert r.PROVIDER_SECRETS.read_bytes() == b'{}\n'
|
|
output = capsys.readouterr()
|
|
assert PROVIDER_KEY not in output.out + output.err
|
|
|
|
|
|
def test_secret_import_validation_error_clears_candidate_traceback_locals(prepared):
|
|
h, r = prepared, prepared.runtime
|
|
payload = ('auth_pools: [' + PROVIDER_KEY).encode('ascii')
|
|
r.sys.stdin = SimpleNamespace(buffer=io.BytesIO(payload))
|
|
h.document_io.validate_managed_runtime_files.side_effect = RuntimeError(
|
|
'candidate rejected',
|
|
)
|
|
|
|
with pytest.raises(RuntimeError) as raised:
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
|
|
formatted = ''.join(traceback.format_exception(
|
|
type(raised.value), raised.value, raised.value.__traceback__,
|
|
))
|
|
assert PROVIDER_KEY not in formatted
|
|
current = raised.value.__traceback__
|
|
while current is not None:
|
|
if current.tb_frame.f_code.co_filename == str(SOURCE):
|
|
assert PROVIDER_KEY not in repr(current.tb_frame.f_locals)
|
|
current = current.tb_next
|
|
r._write_new.assert_not_called()
|
|
h.security.durable_replace.assert_not_called()
|
|
|
|
|
|
def test_secret_import_rejects_config_hash_drift_before_temporary_write(prepared):
|
|
h, r = prepared, prepared.runtime
|
|
payload = b'auth_pools: {}\n'
|
|
r.sys.stdin = SimpleNamespace(buffer=io.BytesIO(payload))
|
|
h.document_io.validate_managed_runtime_files.side_effect = None
|
|
h.document_io.validate_managed_runtime_files.return_value = SimpleNamespace(
|
|
config=h.config, secrets={'auth_pools': {}}, config_sha256='e' * 64,
|
|
)
|
|
|
|
with pytest.raises(RuntimeError, match='configuration changed'):
|
|
r.import_secrets(
|
|
r.DEFAULT_CONFIG,
|
|
h.config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
|
|
r._write_new.assert_not_called()
|
|
r.secrets.token_hex.assert_not_called()
|
|
h.security.durable_replace.assert_not_called()
|
|
|
|
|
|
def test_secret_import_rechecks_config_hash_before_replacement(prepared):
|
|
h, r = prepared, prepared.runtime
|
|
payload = b'auth_pools: {}\n'
|
|
r.sys.stdin = SimpleNamespace(buffer=io.BytesIO(payload))
|
|
h.document_io.load_managed_runtime_config.return_value = SimpleNamespace(
|
|
config=h.config, secrets=None, config_sha256='e' * 64,
|
|
)
|
|
|
|
with pytest.raises(RuntimeError, match='configuration changed'):
|
|
r.import_secrets(
|
|
r.DEFAULT_CONFIG,
|
|
h.config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
|
|
assert not r._write_new.call_args.args[0].exists()
|
|
h.security.durable_replace.assert_not_called()
|
|
h.security.fsync_directory.assert_not_called()
|
|
|
|
|
|
def test_write_new_clears_secret_payload_from_traceback_locals(runtime, tmp_path):
|
|
secret = PROVIDER_KEY.encode('ascii')
|
|
target = tmp_path / 'secret.yaml'
|
|
runtime.private_path = mock.Mock(side_effect=lambda path, **_kwargs: Path(path))
|
|
runtime.os.open = os.open
|
|
runtime.os.fdopen = os.fdopen
|
|
runtime.os.close = os.close
|
|
runtime.os.fsync = mock.Mock(side_effect=OSError('fsync refused'))
|
|
runtime.os.O_NOFOLLOW = getattr(os, 'O_NOFOLLOW', 0)
|
|
runtime.os.O_DIRECTORY = getattr(os, 'O_DIRECTORY', 0)
|
|
|
|
with pytest.raises(OSError) as raised:
|
|
runtime._write_new(target, secret)
|
|
|
|
current = raised.value.__traceback__
|
|
while current is not None:
|
|
if current.tb_frame.f_code.co_filename == str(SOURCE):
|
|
assert PROVIDER_KEY not in repr(current.tb_frame.f_locals)
|
|
current = current.tb_next
|
|
|
|
|
|
def test_initialize_rejects_config_hash_drift_before_mutation(prepared):
|
|
h, r = prepared, prepared.runtime
|
|
h.document_io.load_managed_runtime_config.return_value = SimpleNamespace(
|
|
config=h.config, secrets=None, config_sha256='e' * 64,
|
|
)
|
|
|
|
with pytest.raises(RuntimeError, match='configuration changed'):
|
|
r.initialize(
|
|
r.DEFAULT_CONFIG,
|
|
h.config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
|
|
h.pg.postgres_runtime_paths.assert_not_called()
|
|
r.subprocess.run.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('failure', ['private_destination', 'replace'])
|
|
def test_secret_import_error_removes_temporary_file_and_preserves_old_credentials(prepared, failure):
|
|
h, r = prepared, prepared.runtime
|
|
r.sys.stdin = SimpleNamespace(buffer=io.BytesIO(('token: ' + PROVIDER_KEY).encode('ascii')))
|
|
if failure == 'private_destination':
|
|
private = r.private_path.side_effect
|
|
|
|
def reject_destination(path, **kwargs):
|
|
if path == r.PROVIDER_SECRETS:
|
|
raise OSError('unsafe destination refused')
|
|
return private(path, **kwargs)
|
|
|
|
r.private_path.side_effect = reject_destination
|
|
else:
|
|
h.security.durable_replace.side_effect = OSError('replace refused')
|
|
with pytest.raises(OSError, match='refused'):
|
|
r.import_secrets(r.DEFAULT_CONFIG, h.config)
|
|
assert h.trace == ['lock', 'cluster.lock', 'cluster.unlock', 'unlock']
|
|
assert not h.locked and not h.cluster_locked
|
|
assert r.PROVIDER_SECRETS.read_bytes() == b'{}\n'
|
|
assert not r._write_new.call_args.args[0].exists()
|
|
h.security.fsync_directory.assert_not_called()
|
|
if failure == 'private_destination':
|
|
h.security.durable_replace.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('action', ['provision', 'initialize', 'run', 'health', 'status', 'import-secrets', 'import-snapshot'])
|
|
def test_main_container_gate_precedes_every_action(runtime, action):
|
|
runtime.require_container = mock.Mock(side_effect=RuntimeError('unprepared container'))
|
|
runtime.prepare_environment = mock.Mock()
|
|
runtime.provision = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='unprepared container'):
|
|
runtime.main([action])
|
|
runtime.require_container.assert_called_once_with(provisioning=action == 'provision')
|
|
runtime.prepare_environment.assert_not_called()
|
|
runtime.provision.assert_not_called()
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_main_provision_does_not_enter_nonroot_environment_or_database(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.provision = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock(side_effect=AssertionError('provision entered runtime preparation'))
|
|
runtime.initialize = mock.Mock()
|
|
assert runtime.main(['provision']) == 0
|
|
runtime.require_container.assert_called_once_with(provisioning=True)
|
|
runtime.provision.assert_called_once_with()
|
|
runtime.prepare_environment.assert_not_called()
|
|
runtime.initialize.assert_not_called()
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_main_import_secrets_uses_preliminary_environment_validation(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
config = {'fixture': True}
|
|
runtime.prepare_environment = mock.Mock(return_value=(config, 'd' * 64))
|
|
runtime.import_secrets = mock.Mock()
|
|
assert runtime.main(['import-secrets']) == 0
|
|
runtime.prepare_environment.assert_called_once_with(
|
|
str(runtime.DEFAULT_CONFIG),
|
|
validate_documents=False,
|
|
return_config_sha256=True,
|
|
)
|
|
runtime.import_secrets.assert_called_once_with(
|
|
str(runtime.DEFAULT_CONFIG),
|
|
config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_main_initialize_binds_validated_config_hash(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
config = {'fixture': True}
|
|
runtime.prepare_environment = mock.Mock(return_value=(config, 'd' * 64))
|
|
runtime.initialize = mock.Mock()
|
|
previous = object()
|
|
runtime.signal = SimpleNamespace(SIGTERM=15, signal=mock.Mock(return_value=previous))
|
|
|
|
assert runtime.main(['initialize']) == 0
|
|
|
|
runtime.prepare_environment.assert_called_once_with(
|
|
str(runtime.DEFAULT_CONFIG),
|
|
validate_documents=True,
|
|
return_config_sha256=True,
|
|
)
|
|
runtime.initialize.assert_called_once_with(
|
|
runtime.DEFAULT_CONFIG,
|
|
config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_main_run_rechecks_config_hash_before_exec(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
config = {'fixture': True}
|
|
runtime.prepare_environment = mock.Mock(return_value=(config, 'd' * 64))
|
|
runtime.initialize = mock.Mock()
|
|
previous = object()
|
|
runtime.signal = SimpleNamespace(SIGTERM=15, signal=mock.Mock(return_value=previous))
|
|
document_io = SimpleNamespace(load_managed_runtime_config=mock.Mock(
|
|
return_value=SimpleNamespace(config=config, config_sha256='e' * 64),
|
|
))
|
|
|
|
with mock.patch.dict(sys.modules, {'runtime_document_io': document_io}):
|
|
with pytest.raises(RuntimeError, match='configuration changed'):
|
|
runtime.main(['run'])
|
|
|
|
runtime.initialize.assert_called_once_with(
|
|
runtime.DEFAULT_CONFIG,
|
|
config,
|
|
expected_config_sha256='d' * 64,
|
|
)
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('arguments', [
|
|
[], ['--manifest-sha256', 'invalid'],
|
|
['--manifest-sha256', 'a' * 64, '--config', '/data/config/other.yaml'],
|
|
])
|
|
def test_main_snapshot_requires_approved_digest_and_default_config(runtime, arguments):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock()
|
|
runtime.initialize = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='snapshot import'):
|
|
runtime.main(['import-snapshot', *arguments])
|
|
runtime.require_container.assert_called_once_with(provisioning=False)
|
|
runtime.prepare_environment.assert_not_called()
|
|
runtime.initialize.assert_not_called()
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_main_snapshot_dispatch_preserves_module_signal_state_without_initialize(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock(return_value={})
|
|
runtime.initialize = mock.Mock()
|
|
previous = object()
|
|
runtime.signal = SimpleNamespace(SIGTERM=15, signal=mock.Mock(return_value=previous))
|
|
runtime.sys.modules = {runtime.__name__: runtime}
|
|
|
|
def restore(module, digest):
|
|
assert module is runtime and digest == 'a' * 64
|
|
handler = runtime.signal.signal.call_args.args[1]
|
|
handler(15, None)
|
|
assert module._shutdown_requested is True
|
|
return 130
|
|
|
|
importer = SimpleNamespace(import_snapshot=mock.Mock(side_effect=restore))
|
|
with mock.patch.dict(sys.modules, {'container_import': importer}):
|
|
assert runtime.main(['import-snapshot', '--manifest-sha256', 'a' * 64]) == 130
|
|
importer.import_snapshot.assert_called_once_with(runtime, 'a' * 64)
|
|
runtime.signal.signal.assert_called_with(15, previous)
|
|
runtime.initialize.assert_not_called()
|
|
runtime.os.execv.assert_not_called()
|
|
|
|
|
|
def test_manifest_digest_is_rejected_for_other_actions(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock()
|
|
runtime.provision = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='restricted to snapshot import'):
|
|
runtime.main(['provision', '--manifest-sha256', 'a' * 64])
|
|
runtime.prepare_environment.assert_not_called()
|
|
runtime.provision.assert_not_called()
|
|
|
|
|
|
def test_worker_api_requirement_flag_is_restricted_to_health(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='restricted to health and status'):
|
|
runtime.main(['run', '--require-worker-api'])
|
|
runtime.require_container.assert_called_once_with(provisioning=False)
|
|
runtime.prepare_environment.assert_not_called()
|
|
|
|
|
|
def test_discovery_requirement_flag_is_restricted_to_health(runtime):
|
|
runtime.require_container = mock.Mock()
|
|
runtime.prepare_environment = mock.Mock()
|
|
with pytest.raises(RuntimeError, match='restricted to health and status'):
|
|
runtime.main(['run', '--require-discovery-producers'])
|
|
runtime.require_container.assert_called_once_with(provisioning=False)
|
|
runtime.prepare_environment.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize('action', ['health', 'status'])
|
|
def test_strict_status_dispatches_worker_api_requirement(healthy, action, capsys):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['supervisor']['worker_api'] = {
|
|
'enabled': True, 'address': '127.0.0.1', 'port': 8766,
|
|
}
|
|
h.snapshot['signature'].append(
|
|
('worker-api', 'running', 'running', 124, None, 0, 0, '', ''),
|
|
)
|
|
r.require_container = mock.Mock()
|
|
r.prepare_environment = mock.Mock(return_value=h.config)
|
|
r._probe_worker_api = mock.Mock()
|
|
|
|
assert r.main([action, '--require-worker-api']) == 0
|
|
|
|
assert json.loads(capsys.readouterr().out)['healthy'] is True
|
|
r._probe_worker_api.assert_called_once_with(h.config)
|
|
|
|
|
|
@pytest.mark.parametrize('action', ['health', 'status'])
|
|
def test_strict_status_dispatches_discovery_requirement(healthy, action, capsys):
|
|
h, r = healthy, healthy.runtime
|
|
h.config['sources'] = {'gitlab': {'enabled': True}}
|
|
h.snapshot['signature'].append(
|
|
('gitlab', 'running', 'running', 201, None, 0, 0, '', ''),
|
|
)
|
|
r.require_container = mock.Mock()
|
|
r.prepare_environment = mock.Mock(return_value=h.config)
|
|
|
|
assert r.main([action, '--require-discovery-producers']) == 0
|
|
|
|
assert json.loads(capsys.readouterr().out)['healthy'] is True
|
|
|
|
|
|
@pytest.mark.parametrize('action', ['health', 'status'])
|
|
def test_status_stdout_contains_readiness_but_no_credentials(healthy, action, capsys):
|
|
h, r = healthy, healthy.runtime
|
|
r.require_container = mock.Mock()
|
|
r.prepare_environment = mock.Mock(return_value=h.config)
|
|
assert r.main([action]) == 0
|
|
output = capsys.readouterr()
|
|
assert json.loads(output.out)['healthy'] is True
|
|
assert PASSWORD not in output.out + output.err
|
|
assert 'postgresql://' not in output.out + output.err
|
|
r.os.execv.assert_not_called()
|