#!/usr/bin/env python3 """Run the standalone edge E2E with only random, labelled Docker resources.""" import argparse import configparser import hashlib import ipaddress import json import os from pathlib import Path import re import secrets import shutil import ssl import stat import subprocess import sys import time TEST_IMAGE = 'truf-edge-e2e-runtime:test' EDGE_IMAGE = 'truf-edge-e2e:test' FAIL2BAN_IMAGE = 'truf-fail2ban-edge-e2e:test' RUN_LABEL = 'com.truf.edge-e2e.run' KIND_LABEL = 'com.truf.edge-e2e.kind' MAX_OUTPUT = 2 * 1024 * 1024 MAX_DOCKER_RESOURCES = 1024 MAX_SNAPSHOT_BYTES = 8 * 1024 * 1024 MAX_SAFE_EVIDENCE = 64 * 1024 HEX_64 = re.compile(r'(?:sha256:)?[a-f0-9]{64}') RESOURCE_NAME = re.compile(r'truf-edge-e2e-[a-f0-9]{16}-[a-z][a-z0-9-]{0,48}') CONTAINER_METADATA_FORMAT = ( '{"id":{{json .Id}},"name":{{json .Name}},"image":{{json .Image}},' '"status":{{json .State.Status}},"running":{{json .State.Running}},' '"paused":{{json .State.Paused}},"restarting":{{json .State.Restarting}},' '"dead":{{json .State.Dead}},"mounts":{{json .Mounts}}}' ) VOLUME_METADATA_FORMAT = ( '{"name":{{json .Name}},"driver":{{json .Driver}},"scope":{{json .Scope}},' '"created":{{json .CreatedAt}},"mountpoint":{{json .Mountpoint}},' '"labels":{{json .Labels}},"options":{{json .Options}}}' ) def isolated_test_subnet(run_id, scope='edge'): value = int.from_bytes( hashlib.sha256(f'{run_id}\0{scope}'.encode('ascii')).digest()[:2], 'big', ) & 0x1fff return f'198.{18 + (value >> 12)}.{(value >> 4) & 0xff}.{(value & 0xf) * 16}/28' class Failure(RuntimeError): pass def require(condition, label): if not condition: raise Failure(label) def json_bytes(value, label): require(len(value) <= MAX_OUTPUT, label + '_output_bound') try: result = json.loads(value.decode('utf-8', errors='strict')) except (UnicodeError, ValueError): raise Failure(label + '_json') from None return result def write_bytes(path, content): path = Path(path) path.parent.mkdir(parents=True, exist_ok=True) with open(path, 'xb') as handle: handle.write(content) class CommandRunner: def __init__(self, root, run_root, distro, deadline): self.root = Path(root) self.run_root = Path(run_root) self.deadline = deadline self.wsl = shutil.which('wsl.exe') or shutil.which('wsl') require(self.wsl, 'wsl_unavailable') self.distro = distro self.docker_prefix = [self.wsl, '-d', distro, '--', 'sudo', '-n', 'docker'] allowed = ( 'PATH', 'PATHEXT', 'SystemRoot', 'SYSTEMROOT', 'WINDIR', 'COMSPEC', 'TEMP', 'TMP', 'USERPROFILE', 'WSLENV', ) self.host_env = { name: os.environ[name] for name in allowed if name in os.environ } def execute(self, label, command, *, timeout=60, check=True): remaining = self.deadline - time.monotonic() require(remaining > 0, 'aggregate_timeout') timeout = max(0.1, min(float(timeout), remaining)) try: completed = subprocess.run( command, cwd=os.fspath(self.root), env=self.host_env, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=timeout, check=False, creationflags=getattr(subprocess, 'CREATE_NO_WINDOW', 0), ) except subprocess.TimeoutExpired: raise Failure(label + '_timeout') from None except OSError: raise Failure(label + '_unavailable') from None require( len(completed.stdout) <= MAX_OUTPUT and len(completed.stderr) <= MAX_OUTPUT, label + '_output_bound', ) if check: if completed.returncode != 0: error = Failure(label + '_failed') error.exit_code = completed.returncode raise error return completed.returncode, completed.stdout, completed.stderr def docker(self, label, arguments, *, timeout=60, check=True): return self.execute( label, [*self.docker_prefix, *arguments], timeout=timeout, check=check, ) def wsl_path(self, path): _, stdout, _ = self.execute( 'wsl_path', [self.wsl, '-d', self.distro, '--', 'wslpath', '-a', '-u', os.fspath(Path(path).resolve(strict=True)).replace('\\', '/')], timeout=30, ) value = stdout.decode('utf-8', errors='strict').strip() require(value.startswith('/') and '\x00' not in value, 'wsl_path_invalid') return value class Resources: def __init__(self, runner, run_id): self.runner = runner self.run_id = run_id self.created = {'containers': {}, 'volumes': {}, 'networks': {}} self.foreign_baseline = None def labels(self, kind): return [ '--label', f'{RUN_LABEL}={self.run_id}', '--label', f'{KIND_LABEL}={kind}', ] @staticmethod def singular(kind): return {'containers': 'container', 'volumes': 'volume', 'networks': 'network'}[kind] def list_names(self, kind, *, name=None, owned=False): singular = self.singular(kind) arguments = [singular, 'ls'] if kind == 'containers': arguments.append('--all') if name is not None: pattern = '^/' + name + '$' if kind == 'containers' else '^' + name + '$' arguments.extend(['--filter', 'name=' + pattern]) if owned: arguments.extend(['--filter', f'label={RUN_LABEL}={self.run_id}']) field = 'Names' if kind == 'containers' else 'Name' arguments.append('--format={{.' + field + '}}') _, stdout, _ = self.runner.docker('list_' + kind, arguments) names = {line for line in stdout.decode('ascii', errors='strict').splitlines() if line} require(len(names) <= 16, 'resource_inventory_bound') return names def guard(self, kind, name): require(RESOURCE_NAME.fullmatch(name), 'resource_name_guard') require(not self.list_names(kind, name=name), 'resource_name_collision') def inspect_owned(self, kind, name, expected=None): singular = self.singular(kind) if kind == 'containers': template = ( '{"id":{{json .Id}},"name":{{json .Name}},' '"run":{{json (index .Config.Labels "' + RUN_LABEL + '")}},' '"kind":{{json (index .Config.Labels "' + KIND_LABEL + '")}},' '"running":{{json .State.Running}},"paused":{{json .State.Paused}},' '"restarting":{{json .State.Restarting}}}' ) else: template = ( '{"id":{{json ' + ('.Id' if kind == 'networks' else '""') + '}},' '"name":{{json .Name}},"run":{{json (index .Labels "' + RUN_LABEL + '")}},' '"kind":{{json (index .Labels "' + KIND_LABEL + '")}}}' ) _, stdout, _ = self.runner.docker( 'inspect_owned_' + singular, [singular, 'inspect', '--format', template, name], ) value = json_bytes(stdout, 'inspect_owned_' + singular) require( value.get('name') == ('/' + name if kind == 'containers' else name) and value.get('run') == self.run_id and isinstance(value.get('kind'), str) and (kind == 'volumes' or HEX_64.fullmatch(str(value.get('id') or ''))), 'owned_' + singular + '_identity_guard', ) if expected is not None: require( (expected.get('id') is None or value.get('id') == expected['id']) and value['kind'] == expected['kind'], 'owned_' + singular + '_identity_changed', ) return value def metadata_names(self, kind): arguments = ([kind, 'ls', '--all', '--no-trunc', '--format={{.ID}}'] if kind == 'container' else [kind, 'ls', '--format={{.Name}}']) _, stdout, _ = self.runner.docker('foreign_' + kind + '_list', arguments) try: values = {line for line in stdout.decode('ascii').splitlines() if line} except UnicodeError: raise Failure('foreign_' + kind + '_inventory_encoding') from None require(len(values) <= MAX_DOCKER_RESOURCES, 'foreign_' + kind + '_inventory_bound') return values def container_metadata(self, identifier): require(re.fullmatch(r'[a-f0-9]{64}', identifier), 'foreign_container_id_guard') _, stdout, _ = self.runner.docker( 'foreign_container_inspect', ['container', 'inspect', '--format', CONTAINER_METADATA_FORMAT, identifier], ) value = json_bytes(stdout, 'foreign_container_inspect') require( value.get('id') == identifier and isinstance(value.get('name'), str) and HEX_64.fullmatch(str(value.get('image') or '')) and value.get('status') in ( 'created', 'running', 'paused', 'restarting', 'removing', 'exited', 'dead', ) and all(type(value.get(key)) is bool for key in ( 'running', 'paused', 'restarting', 'dead', )) and isinstance(value.get('mounts'), list) and len(value['mounts']) <= 128, 'foreign_container_metadata_guard', ) mounts = [{ key: mount.get(key) for key in ( 'Type', 'Name', 'Source', 'Destination', 'Driver', 'Mode', 'RW', 'Propagation', ) } for mount in value['mounts']] return { 'id': value['id'], 'name': value['name'], 'image': value['image'], 'status': value['status'], 'running': value['running'], 'paused': value['paused'], 'restarting': value['restarting'], 'dead': value['dead'], 'mounts_sha256': hashlib.sha256(json.dumps( mounts, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii')).hexdigest(), } def volume_metadata(self, name): _, stdout, _ = self.runner.docker( 'foreign_volume_inspect', ['volume', 'inspect', '--format', VOLUME_METADATA_FORMAT, name], ) value = json_bytes(stdout, 'foreign_volume_inspect') require( value.get('name') == name and isinstance(value.get('driver'), str) and isinstance(value.get('scope'), str) and isinstance(value.get('mountpoint'), str) and (value.get('created') is None or isinstance(value['created'], str)) and (value.get('labels') is None or isinstance(value['labels'], dict)) and (value.get('options') is None or isinstance(value['options'], dict)), 'foreign_volume_metadata_guard', ) digest = lambda item: hashlib.sha256(item).hexdigest() return { 'name': name, 'driver': value['driver'], 'scope': value['scope'], 'created': value['created'], 'mountpoint_sha256': digest(value['mountpoint'].encode('utf-8')), 'labels_sha256': digest(json.dumps( value['labels'], ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii')), 'options_sha256': digest(json.dumps( value['options'], ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii')), } def metadata_snapshot(self, *, exclude_owned=False): containers = {} owned_ids = { value['id']: (name, value) for name, value in self.created['containers'].items() if value.get('id') } for identifier in sorted(self.metadata_names('container')): if exclude_owned and identifier in owned_ids: name, expected = owned_ids[identifier] self.inspect_owned('containers', name, expected) continue containers[identifier] = self.container_metadata(identifier) volumes = {} for name in sorted(self.metadata_names('volume')): value = self.volume_metadata(name) expected = self.created['volumes'].get(name) if exclude_owned and expected is not None and expected.get('metadata') == value: continue volumes[name] = value snapshot = {'containers': containers, 'volumes': volumes} require(len(json.dumps(snapshot, ensure_ascii=True, sort_keys=True)) <= MAX_SNAPSHOT_BYTES, 'foreign_metadata_snapshot_bound') return snapshot def snapshot_foreign(self): require(self.foreign_baseline is None, 'foreign_snapshot_already_taken') self.foreign_baseline = self.metadata_snapshot() def assert_foreign_unchanged(self): require(self.foreign_baseline is not None, 'foreign_snapshot_missing') require(self.metadata_snapshot(exclude_owned=True) == self.foreign_baseline, 'foreign_docker_state_changed') def create_volume(self, name, role): self.guard('volumes', name) self.created['volumes'][name] = {'id': None, 'kind': role} _, stdout, _ = self.runner.docker( 'create_' + role, ['volume', 'create', *self.labels(role), name], ) require(stdout.decode('ascii', errors='strict').strip() == name, role + '_create') self.inspect_owned('volumes', name, self.created['volumes'][name]) self.created['volumes'][name]['metadata'] = self.volume_metadata(name) def create_network(self, name): self.guard('networks', name) self.created['networks'][name] = {'id': None, 'kind': 'internal-network'} _, stdout, _ = self.runner.docker( 'create_network', ['network', 'create', '--driver', 'bridge', '--internal', '--subnet', isolated_test_subnet(self.run_id), *self.labels('internal-network'), name], ) identifier = stdout.decode('ascii', errors='strict').strip() require(re.fullmatch(r'[a-f0-9]{64}', identifier), 'network_create') self.created['networks'][name]['id'] = identifier self.inspect_owned('networks', name, self.created['networks'][name]) def create_container(self, name, role, arguments): self.guard('containers', name) self.created['containers'][name] = {'id': None, 'kind': role} _, stdout, _ = self.runner.docker( 'create_' + role, ['container', 'create', '--pull=never', '--name', name, *self.labels(role), *arguments], timeout=120, ) identifier = stdout.decode('ascii', errors='strict').strip() require(re.fullmatch(r'[a-f0-9]{64}', identifier), role + '_container_id') self.created['containers'][name]['id'] = identifier self.inspect_owned('containers', name, self.created['containers'][name]) def inventory(self): result = {} for kind in self.created: try: result[kind] = sorted(self.list_names(kind, owned=True)) except Exception: result[kind] = sorted(self.created[kind]) return result def guarded_stop(self): stopped = 0 for name, expected in reversed(self.created['containers'].items()): try: value = self.inspect_owned('containers', name, expected) if expected['id'] is None: expected['id'] = value['id'] if value['running'] or value['paused'] or value['restarting']: value = self.inspect_owned('containers', name, expected) code, _, _ = self.runner.docker( 'guarded_stop_owned_container', ['container', 'stop', '--time', '30', value['id']], timeout=45, check=False, ) stopped += int(code == 0) except (Exception, KeyboardInterrupt): continue return stopped def cleanup(self): inventory = { kind: self.list_names(kind, owned=True) for kind in self.created } for kind, names in self.created.items(): require(inventory[kind] == set(names), 'cleanup_ownership_guard') for name, expected in reversed(self.created['containers'].items()): value = self.inspect_owned('containers', name, expected) if value['running'] or value['paused'] or value['restarting']: value = self.inspect_owned('containers', name, expected) self.runner.docker( 'stop_owned_container', ['container', 'stop', '--time', '30', value['id']], timeout=45, ) value = self.inspect_owned('containers', name, expected) self.runner.docker( 'remove_owned_container', ['container', 'rm', value['id']], timeout=45, ) for name, expected in reversed(self.created['volumes'].items()): self.inspect_owned('volumes', name, expected) self.runner.docker('remove_owned_volume', ['volume', 'rm', name]) for name, expected in reversed(self.created['networks'].items()): value = self.inspect_owned('networks', name, expected) self.runner.docker('remove_owned_network', ['network', 'rm', value['id']]) require( all(not self.list_names(kind, owned=True) for kind in self.created), 'cleanup_incomplete', ) self.assert_foreign_unchanged() class Verifier: def __init__(self, args): self.args = args self.script = Path(__file__).resolve(strict=True) self.root = self.script.parent.parent.resolve(strict=True) require(self.script == self.root / 'docker' / 'verify_edge_e2e.py', 'script_path') require('build/' in (self.root / '.gitignore').read_text(encoding='utf-8').splitlines(), 'build_not_gitignored') self.run_id = secrets.token_hex(8) self.prefix = 'truf-edge-e2e-' + self.run_id build = self.root / 'build' build.mkdir(exist_ok=True) self.run_root = build / ('edge-e2e-' + self.run_id) self.run_root.mkdir() self.deadline = time.monotonic() + args.timeout_seconds self.runner = CommandRunner(self.root, self.run_root, args.wsl_distro, self.deadline) self.resources = Resources(self.runner, self.run_id) self.names = { 'network': self.prefix + '-network', 'backend_data': self.prefix + '-backend-data', 'edge_data': self.prefix + '-edge-data', 'edge_config': self.prefix + '-edge-config', 'auth_logs': self.prefix + '-auth-logs', 'denylist': self.prefix + '-denylist', 'denylist_state': self.prefix + '-denylist-state', 'fail2ban_data': self.prefix + '-fail2ban-data', 'hash': self.prefix + '-hash', 'seed': self.prefix + '-seed', 'gate_seed': self.prefix + '-gate-seed', 'backend': self.prefix + '-backend', 'edge': self.prefix + '-edge', 'fail2ban': self.prefix + '-fail2ban', 'client': self.prefix + '-client', } self.prefix_secret = secrets.token_hex(32) self.edge_marker = secrets.token_hex(32) self.admin_user = 'edge-e2e-admin' self.admin_password = 'edge-e2e-admin-' + secrets.token_hex(12) self.tokens = { name: secrets.token_urlsafe(36) for name in ('good', 'wrong', 'revoked') } self.image_ids = {} self._prepare_host_files() def _prepare_host_files(self): self.tls_dir = self.run_root / 'tls' self.tls_dir.mkdir() self.edge_env = self.run_root / 'edge.env' source_cert = self.root / 'tests' / 'fixtures' / 'worker_tls_cert.pem' source_key = self.root / 'tests' / 'fixtures' / 'worker_tls_key.pem' decoded = ssl._ssl._test_decode_cert(os.fspath(source_cert)) require(('DNS', 'localhost') in decoded.get('subjectAltName', ()), 'tls_localhost_san') require(ssl.cert_time_to_seconds(decoded['notAfter']) > time.time(), 'tls_expired') shutil.copyfile(source_cert, self.tls_dir / source_cert.name) shutil.copyfile(source_key, self.tls_dir / source_key.name) write_bytes( self.tls_dir / 'static-tls.caddy', b'tls /etc/caddy/tls/worker_tls_cert.pem /etc/caddy/tls/worker_tls_key.pem\n', ) def write_edge_environment(self, password_hash): values = { 'TRUF_EDGE_HOST': 'localhost', 'TRUF_EDGE_TLS_INCLUDE': '/etc/caddy/tls/static-tls.caddy', 'TRUF_ADMIN_PREFIX': self.prefix_secret, 'TRUF_ADMIN_USER': self.admin_user, 'TRUF_ADMIN_PASSWORD_HASH': password_hash, 'TRUF_ADMIN_EDGE_MARKER': self.edge_marker, } content = ''.join(f'{name}={value}\n' for name, value in values.items()).encode('ascii') write_bytes(self.edge_env, content) def image(self, reference, label): _, stdout, _ = self.runner.docker('inspect_' + label, ['image', 'inspect', reference]) value = json_bytes(stdout, label + '_image') require(isinstance(value, list) and len(value) == 1, label + '_image_shape') details = value[0] require( details.get('Os') == 'linux' and details.get('Architecture') == 'amd64' and HEX_64.fullmatch(str(details.get('Id') or '')), label + '_image_platform', ) self.image_ids[label] = details['Id'] def preflight(self): require(os.name == 'nt', 'windows_host_required') self.resources.snapshot_foreign() self.image(TEST_IMAGE, 'test') self.image(EDGE_IMAGE, 'edge') self.image(FAIL2BAN_IMAGE, 'fail2ban') def hash_password(self): name = self.names['hash'] self.resources.create_container( name, 'password-hash', ['--network', 'none', '--read-only', '--cap-drop', 'ALL', '--cap-add', 'NET_BIND_SERVICE', '--entrypoint', '/usr/bin/caddy', EDGE_IMAGE, 'hash-password', '--plaintext', self.admin_password], ) _, stdout, _ = self.runner.docker( 'hash_password', ['container', 'start', '--attach', name], timeout=60, ) value = stdout.decode('ascii', errors='strict').strip() require( re.fullmatch(r'\$2[aby]\$(?:0[4-9]|[12][0-9]|3[01])\$[./A-Za-z0-9]{53}', value), 'password_hash_shape', ) return value def create_resources(self): self.resources.create_network(self.names['network']) for key in ( 'backend_data', 'edge_data', 'edge_config', 'auth_logs', 'denylist', 'denylist_state', 'fail2ban_data', ): self.resources.create_volume(self.names[key], key.replace('_', '-')) password_hash = self.hash_password() self.write_edge_environment(password_hash) self.resources.create_container( self.names['seed'], 'backend-volume-seed', ['--network', 'none', '--user', '0:0', '--mount', f'type=volume,source={self.names["backend_data"]},target=/data', '--entrypoint', '/usr/bin/install', TEST_IMAGE, '-d', '-o', '10001', '-g', '10001', '-m', '0700', '/data'], ) self.runner.docker( 'seed_backend_volume', ['container', 'start', '--attach', self.names['seed']], timeout=60, ) self.resources.create_container( self.names['gate_seed'], 'gate-volume-seed', ['--network', 'none', '--user', '0:0', '--read-only', '--cap-drop', 'ALL', '--cap-add', 'CHOWN', '--mount', f'type=volume,source={self.names["auth_logs"]},target=/logs', '--mount', f'type=volume,source={self.names["denylist"]},target=/denylist', '--mount', f'type=volume,source={self.names["denylist_state"]},target=/state', '--mount', f'type=volume,source={self.names["fail2ban_data"]},target=/fail2ban', '--entrypoint', '/bin/sh', EDGE_IMAGE, '-c', 'set -eu; : > /logs/admin-auth-failures.json; ' ': > /state/.edge-e2e-owned; : > /fail2ban/.edge-e2e-owned; ' "printf '%s\\n' '# Managed by truf-caddy-admin-denylist. Admin-route import only.' " '> /denylist/admin-denylist.caddy; ' 'chmod 0770 /logs; chmod 0660 /logs/admin-auth-failures.json; ' 'chmod 0750 /denylist; chmod 0640 /denylist/admin-denylist.caddy; ' 'chmod 0700 /state /fail2ban; ' 'chown 10001:10001 /logs/admin-auth-failures.json ' '/denylist/admin-denylist.caddy /logs /denylist /state /fail2ban'], ) self.runner.docker( 'seed_gate_volumes', ['container', 'start', '--attach', self.names['gate_seed']], timeout=60, ) backend = [ '--network', self.names['network'], '--user', '10001:10001', '--read-only', '--cap-drop', 'ALL', '--security-opt', 'no-new-privileges', '--pids-limit', '512', '--stop-timeout', '30', '--tmpfs', '/tmp:rw,nosuid,nodev,noexec,size=128m,mode=1777', '--mount', f'type=volume,source={self.names["backend_data"]},target=/data', '--env', 'TRUF_ADMIN_EDGE_MARKER=' + self.edge_marker, ] for name, token in self.tokens.items(): backend.extend(['--env', f'TRUF_EDGE_E2E_{name.upper()}_TOKEN={token}']) backend.extend([ '--entrypoint', '/usr/bin/tini', TEST_IMAGE, '--', '/usr/local/bin/python3', '-u', '-I', '-S', '-B', '/opt/truf/tests/edge_e2e_backend.py', ]) self.resources.create_container(self.names['backend'], 'backend', backend) tls = self.runner.wsl_path(self.tls_dir) edge = [ '--network', 'container:' + self.names['backend'], '--user', '10001:10001', '--read-only', '--cap-drop', 'ALL', '--cap-add', 'NET_BIND_SERVICE', '--security-opt', 'no-new-privileges', '--pids-limit', '128', '--stop-timeout', '30', '--tmpfs', '/tmp:rw,nosuid,nodev,noexec,size=16m,mode=1777', '--tmpfs', '/run:rw,nosuid,nodev,noexec,size=4m,mode=0700,uid=10001,gid=10001', '--mount', f'type=volume,source={self.names["edge_data"]},target=/data', '--mount', f'type=volume,source={self.names["edge_config"]},target=/config', '--mount', f'type=volume,source={self.names["auth_logs"]},target=/var/log/caddy', '--mount', f'type=volume,source={self.names["denylist"]},target=/etc/caddy/denylist,readonly', '--mount', f'type=bind,source={tls},target=/etc/caddy/tls,readonly', '--env', 'TRUF_EDGE_HOST=localhost', '--env', 'TRUF_EDGE_TLS_INCLUDE=/etc/caddy/tls/static-tls.caddy', '--env', 'TRUF_ADMIN_PREFIX=' + self.prefix_secret, '--env', 'TRUF_ADMIN_USER=' + self.admin_user, '--env', 'TRUF_ADMIN_PASSWORD_HASH=' + password_hash.replace('$', r'\$'), '--env', 'TRUF_ADMIN_EDGE_MARKER=' + self.edge_marker, EDGE_IMAGE, ] self.resources.create_container(self.names['edge'], 'edge', edge) client = [ '--network', self.names['network'], '--user', '10001:10001', '--read-only', '--cap-drop', 'ALL', '--security-opt', 'no-new-privileges', '--pids-limit', '128', '--tmpfs', '/tmp:rw,nosuid,nodev,noexec,size=16m,mode=1777', '--env', 'TRUF_ADMIN_PREFIX=' + self.prefix_secret, '--env', 'TRUF_ADMIN_USER=' + self.admin_user, '--env', 'TRUF_EDGE_E2E_ADMIN_PASSWORD=' + self.admin_password, '--env', 'TRUF_EDGE_E2E_CONNECT_HOST=' + self.names['backend'], ] for name, token in self.tokens.items(): client.extend(['--env', f'TRUF_EDGE_E2E_{name.upper()}_TOKEN={token}']) client.extend(['--entrypoint', '/bin/sleep', TEST_IMAGE, '1800']) self.resources.create_container(self.names['client'], 'same-ip-client', client) def create_fail2ban(self): tls = self.runner.wsl_path(self.tls_dir) edge_env = self.runner.wsl_path(self.edge_env) fail2ban = [ '--network', 'container:' + self.names['backend'], '--pid', 'container:' + self.names['edge'], '--user', '10001:10001', '--read-only', '--cap-drop', 'ALL', '--security-opt', 'no-new-privileges', '--pids-limit', '128', '--stop-timeout', '30', '--tmpfs', '/tmp:rw,nosuid,nodev,noexec,size=16m,mode=1777', '--tmpfs', '/run/fail2ban:rw,nosuid,nodev,noexec,size=4m,mode=0700,uid=10001,gid=10001', '--mount', f'type=volume,source={self.names["auth_logs"]},target=/var/log/truf-edge,readonly', '--mount', f'type=volume,source={self.names["auth_logs"]},target=/var/log/caddy', '--mount', f'type=volume,source={self.names["denylist"]},target=/etc/truf-edge/denylist', '--mount', f'type=volume,source={self.names["denylist"]},target=/etc/caddy/denylist,readonly', '--mount', f'type=volume,source={self.names["denylist_state"]},target=/var/lib/truf-edge', '--mount', f'type=volume,source={self.names["fail2ban_data"]},target=/var/lib/fail2ban', '--mount', f'type=bind,source={tls},target=/etc/caddy/tls,readonly', '--mount', f'type=bind,source={edge_env},target=/etc/truf-edge/edge.env,readonly', FAIL2BAN_IMAGE, ] self.resources.create_container(self.names['fail2ban'], 'fail2ban-daemon', fail2ban) def container_running(self, name): code, stdout, _ = self.runner.docker( 'container_running', ['container', 'inspect', '--format={{.State.Running}}', name], check=False, ) return code == 0 and stdout.strip() == b'true' def wait_until(self, label, function, seconds, *, alive=None): end = min(self.deadline, time.monotonic() + seconds) while time.monotonic() < end: value = function() if value is not None and value is not False: return value if alive is not None: require(alive(), label + '_container_exited') time.sleep(0.25) raise Failure(label + '_timeout') def wait_backend(self): expected = { 'schema': 1, 'backend': 'private-network', 'postgres': 'fresh', 'sources': ['gitlab', 'dockerhub', 'huggingface'], } def ready(): code, stdout, _ = self.runner.docker( 'backend_ready', ['container', 'exec', self.names['backend'], '/bin/cat', '/data/control/ready.json'], check=False, ) if code: return None value = json_bytes(stdout, 'backend_ready') require(value == expected, 'backend_ready_content') return value return self.wait_until( 'backend_ready', ready, 300, alive=lambda: self.container_running(self.names['backend']), ) def client_mode(self, mode, *, check=True): code, stdout, stderr = self.runner.docker( 'client_' + mode.replace('-', '_'), ['container', 'exec', self.names['client'], '/usr/local/bin/python3', '-u', '-I', '-S', '-B', '/opt/truf/tests/edge_e2e_client.py', mode], timeout=120, check=False, ) if code: if check: try: failure = json_bytes(stdout, 'client_failure') stage = failure.get('failure_stage') except Failure: stage = None if not isinstance(stage, str) or not re.fullmatch( r'[a-z][a-z0-9_]{0,39}', stage, ): stage = 'failed' raise Failure('client_' + mode.replace('-', '_') + '_' + stage) return None require(not stderr, 'client_' + mode.replace('-', '_') + '_stderr') value = json_bytes(stdout, 'client_' + mode.replace('-', '_')) require(value.get('mode') == mode, 'client_' + mode.replace('-', '_') + '_mode') require(value.get('tls') == 'validated-localhost-certificate', 'tls_validation') return value def wait_edge(self, *, require_security_headers=True): last_reason = ['client'] def ready(): value = self.client_mode('probe', check=False) if value is None: return None if value.get('status') == 404 and ( not require_security_headers or value.get('ready') is True ): return value status = value.get('status') missing = value.get('missing') if type(status) is int and status != 404: last_reason[0] = 'status' elif isinstance(missing, list) and missing: name = str(missing[0]).replace('-', '_') last_reason[0] = name if re.fullmatch(r'[a-z_]{1,40}', name) else 'headers' else: last_reason[0] = 'response' return None try: return self.wait_until( 'edge_ready', ready, 45, alive=lambda: self.container_running(self.names['edge']), ) except Failure as exc: if str(exc) == 'edge_ready_timeout': raise Failure('edge_ready_' + last_reason[0]) from None raise def fail2ban_client(self, arguments, label, *, check=True): code, stdout, stderr = self.runner.docker( label, ['container', 'exec', self.names['fail2ban'], '/usr/bin/fail2ban-client', *arguments], timeout=45, check=False, ) if check: require(code == 0, label + '_failed') if code: return None require(not stderr, label + '_stderr') return stdout def fail2ban_bans(self): stdout = self.fail2ban_client( ['get', 'truf-admin-auth', 'banip'], 'fail2ban_get_bans', check=False, ) if stdout is None: return None try: words = stdout.decode('ascii', errors='strict').split() addresses = {ipaddress.ip_address(word).compressed.lower() for word in words} except (UnicodeError, ValueError): raise Failure('fail2ban_ban_inventory') from None require(len(addresses) <= 16, 'fail2ban_ban_inventory_bound') return addresses def wait_fail2ban(self, expected=None, seconds=45, label='fail2ban_ready'): def ready(): bans = self.fail2ban_bans() if bans is None or (expected is not None and bans != set(expected)): return None return bans if bans else True value = self.wait_until( label, ready, seconds, alive=lambda: self.container_running(self.names['fail2ban']), ) return set() if value is True else value def fail2ban_set(self, arguments, label): self.fail2ban_client(['set', 'truf-admin-auth', *arguments], label) def gate_file(self, path, label): code, stdout, stderr = self.runner.docker( label, ['container', 'exec', self.names['fail2ban'], '/bin/cat', path], check=False, ) require(code == 0 and not stderr, label + '_read') require(len(stdout) <= MAX_OUTPUT, label + '_bound') return stdout def auth_records(self): content = self.gate_file( '/var/log/truf-edge/admin-auth-failures.json', 'read_auth_log', ) records = [] for line in content.splitlines(): if line.strip(): value = json_bytes(line, 'auth_log_line') require(isinstance(value, dict), 'auth_log_object') records.append(value) return records def wait_auth_records(self, count): return self.wait_until( 'auth_log_records', lambda: (records if len(records := self.auth_records()) == count else None), 15, ) def validate_auth_records(self, records, expected_count): require(len(records) == expected_count, 'auth_log_exact_count') parser = configparser.ConfigParser(interpolation=None) parser.read( self.root / 'deploy' / 'fail2ban' / 'filter.d-truf-admin-auth.conf', encoding='ascii', ) failregex = parser['Definition']['failregex'] pattern = re.compile(failregex.replace('', r'(?P[0-9A-Fa-f:.]+)')) timestamps = [] direct_ip = str(records[0].get('remote_ip') or '') try: direct_address = ipaddress.ip_address(direct_ip) except ValueError: raise Failure('auth_log_direct_ip') from None require(not direct_address.is_loopback, 'auth_log_direct_ip') for record in records: encoded = json.dumps( record, ensure_ascii=True, sort_keys=False, separators=(',', ':'), ) match = pattern.match(encoded) require(match is not None and match.group('host') == direct_ip, 'failregex_direct_ip') require( record.get('status') == 401 and record.get('event') == 'admin_auth_failure' and record.get('remote_ip') == direct_ip, 'auth_log_fields', ) require( not re.search( r'"(?:request|uri|headers|authorization|password|token|prefix)"\s*:', encoded, flags=re.IGNORECASE, ) and self.prefix_secret not in encoded, 'auth_log_redaction', ) timestamps.append(float(record['ts'])) require(max(timestamps) - min(timestamps) <= 600, 'fail2ban_findtime') jail = (self.root / 'deploy' / 'fail2ban' / 'jail.d-truf-admin-auth.local').read_text( encoding='ascii', ) require( 'maxretry = 2' in jail and 'findtime = 10m' in jail and 'bantime = 24h' in jail, 'fail2ban_jail_semantics', ) return direct_ip def docker_logs(self, container): _, stdout, stderr = self.runner.docker( 'docker_logs', ['container', 'logs', container], check=False, ) return stdout + (b'\n' if stdout and stderr else b'') + stderr def assert_logs_safe(self): forbidden = [self.admin_password, self.edge_marker, self.prefix_secret, *self.tokens.values()] logs = [self.docker_logs(name) for name in self.names.values() if name in self.resources.created['containers']] logs.append(self.gate_file( '/var/log/truf-edge/admin-auth-failures.json', 'safe_auth_log', )) for content in logs: require(len(content) <= MAX_OUTPUT, 'log_output_bound') for secret in forbidden: require(secret.encode('ascii') not in content, 'test_secret_in_log') def clear_run_artifacts(self): require( self.run_root.parent == self.root / 'build' and re.fullmatch(r'edge-e2e-[a-f0-9]{16}', self.run_root.name) and self.run_root.is_dir() and not self.run_root.is_symlink(), 'run_root_guard', ) def retry_writable(function, name, _error): details = os.lstat(name) if stat.S_ISLNK(details.st_mode): function(name) return mode = stat.S_IRUSR | stat.S_IWUSR if stat.S_ISDIR(details.st_mode): mode |= stat.S_IXUSR os.chmod(name, mode) function(name) for path in self.run_root.iterdir(): if path.is_symlink() or path.is_file(): try: path.unlink() except PermissionError: require(not path.is_symlink(), 'run_artifact_symlink_permission') os.chmod(path, stat.S_IRUSR | stat.S_IWUSR) path.unlink() elif path.is_dir(): shutil.rmtree(path, onerror=retry_writable) else: raise Failure('run_artifact_type') require(not any(self.run_root.iterdir()), 'run_artifact_cleanup') def retain_failure_evidence(self, value): content = json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii') + b'\n' require(len(content) <= MAX_SAFE_EVIDENCE, 'safe_evidence_bound') forbidden = [ self.admin_password, self.edge_marker, self.prefix_secret, *self.tokens.values(), ] for secret in forbidden: require(secret.encode('ascii') not in content, 'secret_present_in_evidence') self.clear_run_artifacts() write_bytes(self.run_root / 'failure.json', content) def run(self): self.preflight() self.create_resources() self.runner.docker('start_backend', ['container', 'start', self.names['backend']]) self.wait_backend() self.runner.docker('start_client', ['container', 'start', self.names['client']]) require(self.container_running(self.names['client']), 'client_not_running') self.runner.docker('start_edge', ['container', 'start', self.names['edge']]) self.wait_edge() self.create_fail2ban() self.runner.docker( 'start_fail2ban', ['container', 'start', self.names['fail2ban']], ) self.wait_fail2ban(expected=set()) self.wait_edge(require_security_headers=False) baseline = self.client_mode('baseline') admin_identity = json_bytes( self.runner.docker( 'read_admin_identity_evidence', ['container', 'exec', self.names['backend'], '/bin/cat', '/data/control/last-admin-request.json'], )[1], 'admin_identity_evidence', ) require( admin_identity == { 'schema': 1, 'path': '/admin-internal/operations/' + baseline['operation_id'], 'marker_authorized': True, 'operators': [self.admin_user], }, 'trusted_admin_operator_identity', ) worker_identity = json_bytes( self.runner.docker( 'read_worker_identity_evidence', ['container', 'exec', self.names['backend'], '/bin/cat', '/data/control/last-worker-request.json'], )[1], 'worker_identity_evidence', ) require( worker_identity == {'schema': 1, 'admin_headers_absent': True}, 'worker_admin_headers_absent', ) time.sleep(0.5) require(self.auth_records() == [], 'missing_credentials_were_counted') bad = self.client_mode('bad-auth') records = self.wait_auth_records(2) direct_ip = self.validate_auth_records(records, 2) self.wait_fail2ban( expected={direct_ip}, seconds=30, label='automatic_second_failure_ban', ) state = json_bytes( self.gate_file('/var/lib/truf-edge/admin-denylist.json', 'read_ban_state'), 'ban_state', ) expires_at = state.get('bans', {}).get(direct_ip) now = int(time.time()) require( isinstance(expires_at, int) and now + 86300 <= expires_at <= now + 86500, 'production_updater_ban_state', ) self.wait_edge() self.client_mode('assert-ban') fail2ban = self.resources.inspect_owned( 'containers', self.names['fail2ban'], self.resources.created['containers'][self.names['fail2ban']], ) self.runner.docker( 'kill_fail2ban_for_restart', ['container', 'kill', '--signal', 'SIGKILL', fail2ban['id']], timeout=45, ) edge = self.resources.inspect_owned( 'containers', self.names['edge'], self.resources.created['containers'][self.names['edge']], ) self.runner.docker( 'restart_edge', ['container', 'restart', '--time', '30', edge['id']], timeout=60, ) self.wait_edge() self.client_mode('assert-ban') fail2ban = self.resources.inspect_owned( 'containers', self.names['fail2ban'], self.resources.created['containers'][self.names['fail2ban']], ) require(not fail2ban['running'], 'fail2ban_restart_stop') self.runner.docker( 'restart_fail2ban', ['container', 'start', fail2ban['id']], timeout=60, ) self.wait_fail2ban( expected={direct_ip}, seconds=45, label='persisted_fail2ban_ban', ) self.client_mode('assert-ban') self.fail2ban_set(['unbanip', direct_ip], 'operator_unban') self.wait_fail2ban(expected=set(), seconds=30, label='operator_unban_applied') state = json_bytes( self.gate_file('/var/lib/truf-edge/admin-denylist.json', 'read_unban_state'), 'operator_unban_state', ) require(state.get('bans') == {}, 'operator_unban_state') self.wait_edge() self.client_mode('assert-unban') self.fail2ban_set(['bantime', '3'], 'set_short_test_bantime') repeated_bad = self.client_mode('bad-auth') records = self.wait_auth_records(4) require(self.validate_auth_records(records, 4) == direct_ip, 'repeated_direct_ip') self.wait_fail2ban( expected={direct_ip}, seconds=30, label='short_lived_automatic_ban', ) self.wait_edge() self.client_mode('assert-ban') self.wait_fail2ban( expected=set(), seconds=20, label='automatic_fail2ban_expiry', ) state = json_bytes( self.gate_file('/var/lib/truf-edge/admin-denylist.json', 'read_expired_state'), 'expired_state', ) require(state['bans'] == {}, 'expired_state_persistence') require( self.gate_file( '/etc/truf-edge/denylist/admin-denylist.caddy', 'read_expired_snippet', ) == b'# Managed by truf-caddy-admin-denylist. Admin-route import only.\n', 'expired_snippet_persistence', ) self.wait_edge() self.client_mode('assert-unban') require(len(self.auth_records()) == 4, 'auth_log_count_changed_after_ban_flow') audit = self.gate_file( '/var/lib/truf-edge/edge-e2e-reload.audit', 'read_reload_audit', ).decode('ascii', errors='strict').splitlines() require( len(audit) >= 10 and len(audit) <= 32 and set(audit) == {'validate:0', 'reload:0'} and audit.count('validate:0') == audit.count('reload:0'), 'constrained_reload_audit', ) self.assert_logs_safe() summary = { 'schema': 1, 'status': 'passed', 'run_id': self.run_id, 'image_ids': self.image_ids, 'baseline_responses': baseline['checked_responses'], 'bad_basic_attempts': bad['attempts'] + repeated_bad['attempts'], 'direct_ip_log_records': len(records), 'direct_ip': direct_ip, 'fail2ban_daemon': 'production-filter-jail-action-applied', 'caddy_restart': 'admin-ban-persisted', 'fail2ban_restart': 'jail-ban-restored', 'operator_unban': 'fail2ban-actionunban-restored', 'expiry': 'fail2ban-actionunban-restored-and-persisted', 'worker_same_ip': 'usable-throughout', 'admin_actor': 'authenticated-basic-user', 'constrained_reload_pairs': audit.count('reload:0'), 'foreign_docker_state': 'unchanged', 'cleanup': 'complete', } self.resources.cleanup() shutil.rmtree(self.run_root) return summary def parse_args(argv=None): parser = argparse.ArgumentParser( description='Verify the real standalone Caddy edge against a private test backend.', ) parser.add_argument( '--wsl-distro', default='Ubuntu-24.04', help='WSL distribution that owns the Docker socket (default: %(default)s)', ) parser.add_argument( '--timeout-seconds', type=int, default=1200, help='aggregate timeout (default: %(default)s)', ) args = parser.parse_args(argv) if not re.fullmatch(r'[A-Za-z0-9][A-Za-z0-9._+-]{0,127}', args.wsl_distro): parser.error('--wsl-distro contains unsupported characters') if not 300 <= args.timeout_seconds <= 3600: parser.error('--timeout-seconds must be between 300 and 3600') return args def main(argv=None): verifier = None try: verifier = Verifier(parse_args(argv)) summary = verifier.run() except KeyboardInterrupt as exc: label = 'interrupted' failure_class = type(exc).__name__ exit_code = None except Failure as exc: label = str(exc) failure_class = type(exc).__name__ exit_code = getattr(exc, 'exit_code', None) except Exception as exc: label = 'unexpected_exception' failure_class = type(exc).__name__ exit_code = None else: print(json.dumps(summary, ensure_ascii=True, sort_keys=True)) return 0 if verifier is not None: verifier.runner.deadline = time.monotonic() + 600 verifier.resources.guarded_stop() if not re.fullmatch(r'[a-z0-9_]{1,160}', label): label = 'verifier_failure' if not re.fullmatch(r'[A-Za-z][A-Za-z0-9_]{0,79}', failure_class): failure_class = 'Exception' if type(exit_code) is not int or not -(2 ** 31) <= exit_code < 2 ** 31: exit_code = None failure = { 'stage': label, 'class': failure_class, 'exit_code': exit_code, } if verifier is not None: try: verifier.retain_failure_evidence(failure) except Exception: pass print(json.dumps(failure, ensure_ascii=True, sort_keys=True), file=sys.stderr) return 1 if __name__ == '__main__': raise SystemExit(main())