import json import os from pathlib import Path import sys import tempfile import time import unittest import uuid from types import SimpleNamespace from unittest import mock ROOT = Path(__file__).resolve().parents[1] APP_DIR = ROOT / 'app' sys.path.insert(0, str(APP_DIR)) import host_agent_lifecycle from host_agent_protocol import HostAgentAction from host_agent_state import HostStateError class HostAgentDeploymentProfileTests(unittest.TestCase): def test_profile_is_exact_root_owned_policy_with_standalone_default(self): with tempfile.TemporaryDirectory() as temporary: missing = Path(temporary) / 'missing' self.assertIs( host_agent_lifecycle._deployment_profile(str(missing)), host_agent_lifecycle.STANDALONE_PROFILE, ) selected = Path(temporary) / 'profile' selected.write_text('shared-host-edge-v1\n', encoding='ascii') selected.chmod(0o444) self.assertIs( host_agent_lifecycle._deployment_profile(str(selected)), host_agent_lifecycle.SHARED_HOST_PROFILE, ) selected.chmod(0o666) with self.assertRaises(host_agent_lifecycle.HostLifecycleError): host_agent_lifecycle._deployment_profile(str(selected)) selected.write_text('unsupported\n', encoding='ascii') selected.chmod(0o444) with self.assertRaises(host_agent_lifecycle.HostLifecycleError): host_agent_lifecycle._deployment_profile(str(selected)) class _Clock: def __init__(self): self.value = 0.0 def __call__(self): return self.value def sleep(self, seconds): self.value += float(seconds) class _Runner: OLD_RUNTIME = '1' * 64 OLD_EDGE = '2' * 64 NEW_RUNTIME = '3' * 64 NEW_EDGE = '4' * 64 ROLLBACK_RUNTIME = '5' * 64 ROLLBACK_EDGE = '6' * 64 RUNTIME_IMAGE = 'sha256:' + 'a' * 64 EDGE_IMAGE = 'sha256:' + 'b' * 64 def __init__(self, events): self.events = events self.commands = [] self.current = {'runtime': self.OLD_RUNTIME, 'edge': self.OLD_EDGE} self.states = { self.OLD_RUNTIME: self._state('runtime', self.OLD_RUNTIME, self.OLD_RUNTIME), self.OLD_EDGE: self._state('edge', self.OLD_EDGE, self.OLD_RUNTIME), } self.images = { 'truf-local:runtime': self.RUNTIME_IMAGE, 'truf-local:edge': self.EDGE_IMAGE, } self.health = { 'healthy': True, 'activation_state': 'ACTIVE', 'postgres': 'READY', 'workers': ['janitor', 'jsonl-projector', 'result-ingester', 'worker-api'], } self.edge_stop_exit_code = 0 self.fail_caddy = False self.reuse_runtime_id = False self.drift_runtime_after_stop = False self.replace_runtime_after_stop = False self.new_state_mutations = {'runtime': {}, 'edge': {}} self.runtime_up_count = 0 self.edge_up_count = 0 self.fail_forward_health = False self.fail_forward_health_command = False self.forward_health_failures_remaining = 0 self.fail_rollback_health = False def _state(self, service, container_id, runtime_id): return { 'id': container_id, 'image': self.RUNTIME_IMAGE if service == 'runtime' else self.EDGE_IMAGE, 'status': 'running', 'running': True, 'paused': False, 'restarting': False, 'dead': False, 'pid': 100 if service == 'runtime' else 200, 'exit_code': 0, 'oom_killed': False, 'restarts': 0, 'user': '10001:10001', 'entrypoint': ( [ '/usr/bin/tini', '--', '/usr/local/bin/python3', '-u', '-I', '-S', '-B', '/opt/truf/app/container_runtime.py', ] if service == 'runtime' else ['/usr/local/bin/truf-edge-entrypoint'] ), 'command': ( ['run', '--config', '/data/config/config.yaml'] if service == 'runtime' else None ), 'stop_timeout': 600 if service == 'runtime' else 30, 'stop_signal': 'SIGTERM' if service == 'runtime' else '', 'mounts': ( 'volume|truf-docker_data|/var/lib/docker/volumes/truf-docker_data/_data|/data|true;' 'bind||/etc/truf/runtime|/data/config|false;' 'bind||/etc/truf/worker-packages|/data/worker-packages|false;' 'bind||/var/lib/truf/runtime-document-candidates|/data/runtime-document-candidates|true;' 'bind||/run/truf/host-agent.sock|/run/truf/host-agent.sock|false;' 'bind||/var/lib/truf/host-agent/results|/data/host-agent-results|false;' 'bind||/run/truf-postgres|/run/truf-postgres|true;' if service == 'runtime' else ( 'volume|truf-docker_edge_data|/var/lib/docker/volumes/truf-docker_edge_data/_data|/data|true;' 'volume|truf-docker_edge_config|/var/lib/docker/volumes/truf-docker_edge_config/_data|/config|true;' 'bind||/var/log/truf-edge|/var/log/caddy|true;' 'bind||/etc/truf-edge/denylist|/etc/caddy/denylist|false;' ) ), 'readonly': True, 'privileged': False, 'network': ( 'truf-docker_default' if service == 'runtime' else f'container:{runtime_id}' ), 'pid_mode': '', 'ipc_mode': 'private', 'userns_mode': '', 'cgroupns_mode': 'private', 'uts_mode': '', 'group_add': None, 'oci_runtime': 'runc', 'devices': 0, 'device_requests': 0, 'device_cgroup_rules': 0, 'ports': ( {'443/tcp': [{'HostIp': '', 'HostPort': '443'}]} if service == 'runtime' else {} ), 'tmpfs': ( dict(host_agent_lifecycle._RUNTIME_TMPFS) if service == 'runtime' else dict(host_agent_lifecycle._EDGE_TMPFS) ), 'cpus': 2_000_000_000 if service == 'runtime' else 1_000_000_000, 'memory': 6 * 1024 ** 3 if service == 'runtime' else 256 * 1024 ** 2, 'pids_limit': 512 if service == 'runtime' else 128, 'shm_size': 256 * 1024 ** 2 if service == 'runtime' else 64 * 1024 ** 2, 'log_config': { 'Type': 'json-file', 'Config': { 'max-size': '16m' if service == 'runtime' else '8m', 'max-file': '4', }, }, 'cap_drop': ['ALL'], 'cap_add': [] if service == 'runtime' else ['NET_BIND_SERVICE'], 'security_opt': ['no-new-privileges:true'], 'restart_policy': ( {'Name': 'on-failure', 'MaximumRetryCount': 3} if service == 'runtime' else {'Name': 'unless-stopped', 'MaximumRetryCount': 0} ), 'project': 'truf-docker', 'service': service, 'oneoff': 'False', 'config_hash': ( 'c' * 64 if service == 'runtime' else { self.OLD_EDGE: 'd' * 64, self.NEW_EDGE: 'e' * 64, self.ROLLBACK_EDGE: 'f' * 64, }[container_id] ), 'config_files': '/opt/truf/compose.yaml,/opt/truf/compose.edge.yaml', 'working_dir': '/opt/truf', 'health_test': ( list(host_agent_lifecycle._RUNTIME_HEALTH_TEST) if service == 'runtime' else None ), 'health': 'healthy' if service == 'runtime' else None, } def _stop(self, service): container_id = self.current[service] state = self.states[container_id] state.update( status='exited', running=False, pid=0, exit_code=self.edge_stop_exit_code if service == 'edge' else 0, ) def __call__(self, command, timeout): command = tuple(command) self.commands.append((command, timeout)) if command[:len(host_agent_lifecycle._COMPOSE)] == host_agent_lifecycle._COMPOSE: arguments = command[len(host_agent_lifecycle._COMPOSE):] if arguments == ('config', '--quiet'): self.events.append('config') return b'' if arguments[:3] == ('ps', '--all', '--quiet'): current = self.current[arguments[3]] return ((current + '\n') if current is not None else '').encode('ascii') if arguments[:3] == ('stop', '--timeout', '30'): self.events.append('stop-edge') self._stop('edge') return b'' if arguments[:3] == ('stop', '--timeout', '600'): self.events.append('stop-runtime') self._stop('runtime') if self.drift_runtime_after_stop: self.images['truf-local:runtime'] = 'sha256:' + 'c' * 64 if self.replace_runtime_after_stop: self.current['runtime'] = self.NEW_RUNTIME self.states[self.NEW_RUNTIME] = self._state( 'runtime', self.NEW_RUNTIME, self.NEW_RUNTIME, ) return b'' if arguments[:2] == ('rm', '--force'): service = arguments[2] self.events.append('rm-' + service) self.current[service] = None return b'' if arguments[0] == 'up': service = arguments[-1] self.events.append('up-' + service) if service == 'runtime': self.runtime_up_count += 1 container_id = ( self.OLD_RUNTIME if self.reuse_runtime_id else ( self.NEW_RUNTIME if self.runtime_up_count == 1 else self.ROLLBACK_RUNTIME ) ) self.states[container_id] = self._state( service, container_id, container_id, ) else: self.edge_up_count += 1 container_id = ( self.NEW_EDGE if self.edge_up_count == 1 else self.ROLLBACK_EDGE ) self.states[container_id] = self._state( service, container_id, self.current['runtime'], ) self.states[container_id].update(self.new_state_mutations[service]) self.current[service] = container_id return b'' if arguments[:4] == ('exec', '-T', '--user', '10001:10001'): service = arguments[4] if service == 'runtime': self.events.append('strict-health') if ( self.fail_forward_health_command and self.runtime_up_count == 1 ): raise host_agent_lifecycle.HostLifecycleError('command') if ( self.fail_rollback_health or self.fail_forward_health and self.runtime_up_count == 1 ): return b'{"healthy":false}' if ( self.runtime_up_count == 1 and self.forward_health_failures_remaining > 0 ): self.forward_health_failures_remaining -= 1 return b'{"healthy":false}' return json.dumps(self.health, sort_keys=True).encode('ascii') self.events.append('caddy-validate') if self.fail_caddy: raise RuntimeError('sensitive caddy output') return b'' raise AssertionError(f'unexpected compose arguments: {arguments!r}') if command[:3] == ('/usr/bin/docker', 'image', 'inspect'): return json.dumps(self.images[command[-1]]).encode('ascii') if command[:3] == ('/usr/bin/docker', 'container', 'inspect'): return json.dumps(self.states[command[-1]], sort_keys=True).encode('ascii') raise AssertionError(f'unexpected command: {command!r}') class _Session: def __init__(self, events): self._entered = True self.events = events self.request = SimpleNamespace( operation_id=str(uuid.uuid4()), action=HostAgentAction.APPLY_BOTH, active_config_sha256='a' * 64, active_secrets_sha256='b' * 64, candidate_config_sha256='d' * 64, candidate_secrets_sha256='e' * 64, ) self.result = { 'active_config_sha256': 'd' * 64, 'active_secrets_sha256': 'e' * 64, } self.revalidate_error = None self.publication_state = 'original' def backup(self): self.events.append('backup') def revalidate_for_stop(self): self.events.append('revalidate') if self.revalidate_error is not None: raise self.revalidate_error def replace(self, proof): self.proof = proof self.events.append('replace') self.publication_state = 'candidate' return dict(self.result) def restore_backups(self, proof): self.events.append('restore') self.publication_state = 'original' return self.original_identity() def original_identity(self): return { 'active_config_sha256': 'a' * 64, 'active_secrets_sha256': 'b' * 64, } class _State: def __init__(self, events): self.events = events self.phase = 'prepared' self.result = None self.forward_category = None self.safe_detail = None self.fail_advance_to = None self.commit_before_advance_failure = False self.advance_failure_cancellation = None self.fail_result = None self.hold_evidence = None def terminal_result(self): return self.result def initialize(self, publication_state): self.events.append('state-prepared') return { 'phase': self.phase, 'publication_state': publication_state, 'forward_category': self.forward_category, 'safe_detail': self.safe_detail, } def advance( self, expected_phase, next_phase, publication_state, **evidence, ): if self.phase != expected_phase: raise AssertionError((self.phase, expected_phase, next_phase)) if self.fail_advance_to == next_phase: if self.commit_before_advance_failure: self.phase = next_phase self.events.append('state-' + next_phase) raise HostStateError( 'uncertain', cancellation=self.advance_failure_cancellation, ) raise HostStateError('state') self.phase = next_phase self.forward_category = evidence.get( 'forward_category', self.forward_category, ) self.safe_detail = evidence.get('safe_detail', self.safe_detail) self.events.append('state-' + next_phase) return {'phase': next_phase, **evidence} def publish_result( self, result, *, safe_category, safe_detail, resulting_identity, ): if self.fail_result == result: raise HostStateError('filesystem') self.events.append('result-' + result) self.result = { 'schema': 1, 'result': result, 'safe_category': safe_category, 'safe_detail': safe_detail, 'resulting_identity': resulting_identity, } return self.result def publish_failed_hold(self, **evidence): self.events.append('hold') self.hold_evidence = dict(evidence) return evidence class HostAgentLifecycleTests(unittest.TestCase): def lifecycle(self): events = [] runner = _Runner(events) clock = _Clock() lifecycle = host_agent_lifecycle.FixedDeploymentLifecycle( _runner=runner, _clock=clock, _sleep=clock.sleep, ) return events, runner, lifecycle def test_forward_order_and_fixed_command_surface(self): events, runner, lifecycle = self.lifecycle() session = _Session(events) result = host_agent_lifecycle.execute_fixed_forward(session, lifecycle) self.assertEqual(result, session.result) required_order = [ 'backup', 'config', 'revalidate', 'stop-edge', 'stop-runtime', 'replace', 'rm-edge', 'rm-runtime', 'up-runtime', 'strict-health', 'up-edge', 'caddy-validate', ] positions = [events.index(name) for name in required_order] self.assertEqual(positions, sorted(positions)) rendered = [list(command) for command, _ in runner.commands] self.assertFalse(any( forbidden in command for command in rendered for forbidden in ('down', 'kill', '--volumes', '--remove-orphans', 'provision') )) up = [command for command in rendered if 'up' in command] self.assertEqual(len(up), 2) self.assertTrue(all('--no-build' in command for command in up)) self.assertTrue(all('--pull' in command and 'never' in command for command in up)) def test_fixed_operation_persists_success_without_rollback(self): events, runner, lifecycle = self.lifecycle() session = _Session(events) state = _State(events) result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'succeeded') self.assertEqual(state.phase, 'succeeded') self.assertNotIn('restore', events) self.assertNotIn('hold', events) self.assertEqual(runner.runtime_up_count, 1) self.assertEqual(runner.edge_up_count, 1) def test_terminal_phase_precedes_result_and_missing_result_replays(self): events, _runner, lifecycle = self.lifecycle() session = _Session(events) state = _State(events) state.fail_advance_to = 'succeeded' result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'rolled_back') self.assertNotIn('result-succeeded', events) self.assertEqual(state.result['result'], 'rolled_back') self.assertEqual(state.phase, 'rolled_back') replay_events = [] replay = _State(replay_events) replay.phase = 'succeeded' result = host_agent_lifecycle.execute_fixed_operation( _Session(replay_events), lifecycle, replay, ) self.assertEqual(result['result'], 'succeeded') self.assertIn('result-succeeded', replay_events) self.assertNotIn('backup', replay_events) self.assertNotIn('stop-edge', replay_events) def test_uncertain_terminal_phase_write_never_triggers_rollback(self): events, _runner, lifecycle = self.lifecycle() state = _State(events) state.fail_advance_to = 'succeeded' state.commit_before_advance_failure = True with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(raised.exception.category, 'state') self.assertEqual(state.phase, 'succeeded') self.assertNotIn('restore', events) self.assertNotIn('hold', events) self.assertNotIn('result-succeeded', events) def test_terminal_phase_cancellation_propagates_without_containment(self): events, _runner, lifecycle = self.lifecycle() state = _State(events) state.fail_advance_to = 'succeeded' state.commit_before_advance_failure = True state.advance_failure_cancellation = SystemExit() with self.assertRaises(SystemExit): host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(state.phase, 'succeeded') self.assertNotIn('restore', events) self.assertNotIn('hold', events) def test_uncertain_rollback_start_enters_failed_hold(self): events, runner, lifecycle = self.lifecycle() runner.fail_forward_health = True state = _State(events) state.fail_advance_to = 'rollback_started' state.commit_before_advance_failure = True result = host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(result['result'], 'failed_hold') self.assertEqual(state.phase, 'failed_hold') self.assertIn('hold', events) self.assertNotIn('restore', events) def test_forward_health_failure_rolls_back_exactly_once(self): events, runner, lifecycle = self.lifecycle() runner.fail_forward_health = True session = _Session(events) state = _State(events) result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'rolled_back') self.assertEqual(result['safe_category'], 'health_check_failed') self.assertEqual(state.phase, 'rolled_back') self.assertEqual(events.count('restore'), 1) self.assertEqual(runner.runtime_up_count, 2) self.assertEqual(runner.edge_up_count, 1) self.assertNotIn('hold', events) def test_forward_health_command_failure_is_classified_as_health(self): events, runner, lifecycle = self.lifecycle() runner.fail_forward_health_command = True result = host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, _State(events), ) self.assertEqual(result['result'], 'rolled_back') self.assertEqual(result['safe_category'], 'health_check_failed') self.assertEqual(events.count('restore'), 1) self.assertNotIn('hold', events) def test_transient_strict_health_failure_is_retried_before_success(self): events, runner, lifecycle = self.lifecycle() runner.forward_health_failures_remaining = 1 result = host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, _State(events), ) self.assertEqual(result['result'], 'succeeded') self.assertEqual(events.count('strict-health'), 2) self.assertNotIn('restore', events) self.assertNotIn('hold', events) def test_restart_recreation_failure_rolls_back_once_as_restart_failed(self): events, runner, lifecycle = self.lifecycle() session = _Session(events) session.request.action = HostAgentAction.RESTART session.request.candidate_config_sha256 = None session.request.candidate_secrets_sha256 = None session.result = session.original_identity() def restart_without_publication(proof): session.proof = proof events.append('replace') return dict(session.result) session.replace = restart_without_publication lifecycle.recreate_and_verify = mock.Mock( side_effect=host_agent_lifecycle.HostLifecycleError('command'), ) state = _State(events) result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'rolled_back') self.assertEqual(result['safe_category'], 'restart_failed') self.assertEqual(result['safe_detail'], 'restart_failed') self.assertEqual(state.phase, 'rolled_back') self.assertEqual(events.count('restore'), 1) self.assertEqual(runner.runtime_up_count, 1) self.assertEqual(runner.edge_up_count, 1) self.assertNotIn('hold', events) def test_rollback_health_failure_enters_hold_without_third_attempt(self): events, runner, lifecycle = self.lifecycle() runner.fail_forward_health = True runner.fail_rollback_health = True state = _State(events) result = host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(result['result'], 'failed_hold') self.assertEqual(state.phase, 'failed_hold') self.assertEqual(events.count('hold'), 1) self.assertEqual(events.count('restore'), 1) self.assertEqual(runner.runtime_up_count, 2) self.assertEqual(runner.edge_up_count, 0) def test_failed_hold_marker_precedes_containment_and_records_partial(self): events, runner, lifecycle = self.lifecycle() runner.fail_forward_health = True session = _Session(events) state = _State(events) def fail_restore(_proof): session.publication_state = 'partial' raise RuntimeError('rollback detail') session.restore_backups = fail_restore contain = lifecycle.contain_for_failed_hold def assert_fenced(*args): self.assertIn('hold', events) return contain(*args) with mock.patch.object( lifecycle, 'contain_for_failed_hold', side_effect=assert_fenced, ): result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'failed_hold') self.assertEqual(state.hold_evidence['publication_state'], 'partial') self.assertLess(events.index('hold'), events.index('state-failed_hold')) def test_system_cancellation_is_not_converted_to_operation_result(self): for cancellation in (KeyboardInterrupt, SystemExit): with self.subTest(cancellation=cancellation.__name__): events, _runner, lifecycle = self.lifecycle() session = _Session(events) session.revalidate_error = cancellation() state = _State(events) with self.assertRaises(cancellation): host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertIsNone(state.result) self.assertEqual(state.phase, 'prepared') self.assertNotIn('stop-edge', events) def test_post_mutation_cancellation_rolls_back_once_then_propagates(self): events, _runner, lifecycle = self.lifecycle() state = _State(events) lifecycle.recreate_and_verify = mock.Mock( side_effect=KeyboardInterrupt(), ) with self.assertRaises(KeyboardInterrupt): host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(state.phase, 'rolled_back') self.assertEqual(state.result['result'], 'rolled_back') self.assertEqual(events.count('restore'), 1) self.assertNotIn('hold', events) def test_terminal_write_failure_does_not_mask_saved_cancellation(self): events, _runner, lifecycle = self.lifecycle() state = _State(events) state.fail_result = 'rolled_back' lifecycle.recreate_and_verify = mock.Mock( side_effect=KeyboardInterrupt(), ) with self.assertRaises(KeyboardInterrupt): host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(state.phase, 'rolled_back') self.assertIsNone(state.result) self.assertEqual(events.count('restore'), 1) def test_replay_preflight_failure_terminalizes_failed_hold(self): events, _runner, lifecycle = self.lifecycle() state = _State(events) state.phase = 'forward_started' state.forward_category = 'apply_failed' lifecycle.preflight = mock.Mock(side_effect=RuntimeError('detail')) result = host_agent_lifecycle.execute_fixed_operation( _Session(events), lifecycle, state, ) self.assertEqual(result['result'], 'failed_hold') self.assertEqual(state.phase, 'failed_hold') self.assertIn('state-rollback_started', events) self.assertIn('hold', events) def test_same_operation_failed_hold_replay_only_finishes_result(self): events, _runner, lifecycle = self.lifecycle() session = _Session(events) session._failed_hold_replay = True state = _State(events) state.phase = 'rollback_started' state.forward_category = 'health_check_failed' result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'failed_hold') self.assertEqual(state.phase, 'failed_hold') self.assertNotIn('backup', events) self.assertNotIn('stop-edge', events) def test_pre_stop_failure_records_failed_without_lifecycle_mutation(self): events, runner, lifecycle = self.lifecycle() session = _Session(events) session.revalidate_error = RuntimeError('candidate detail') state = _State(events) result = host_agent_lifecycle.execute_fixed_operation( session, lifecycle, state, ) self.assertEqual(result['result'], 'failed') self.assertEqual(state.phase, 'failed') self.assertNotIn('stop-edge', events) self.assertNotIn('restore', events) self.assertEqual(runner.runtime_up_count, 0) def test_stop_failure_never_publishes_removes_or_recreates(self): events, runner, lifecycle = self.lifecycle() runner.edge_stop_exit_code = 1 with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'stop') self.assertNotIn('stop-runtime', events) self.assertNotIn('replace', events) self.assertFalse(any(name.startswith(('rm-', 'up-')) for name in events)) def test_document_drift_before_stop_never_stops_or_publishes(self): events, _runner, lifecycle = self.lifecycle() session = _Session(events) session.revalidate_error = RuntimeError('candidate drift detail') with self.assertRaises(RuntimeError): host_agent_lifecycle.execute_fixed_forward(session, lifecycle) self.assertIn('revalidate', events) self.assertNotIn('stop-edge', events) self.assertNotIn('replace', events) def test_image_drift_after_stop_refuses_removal_and_recreation(self): events, runner, lifecycle = self.lifecycle() runner.drift_runtime_after_stop = True with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertNotIn('replace', events) self.assertFalse(any(name.startswith(('rm-', 'up-')) for name in events)) def test_service_identity_drift_after_stop_refuses_publication(self): events, runner, lifecycle = self.lifecycle() runner.replace_runtime_after_stop = True with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertNotIn('replace', events) self.assertFalse(any(name.startswith(('rm-', 'up-')) for name in events)) def test_preflight_attests_fixed_container_contract(self): cases = { 'config_hash': 'not-a-hash', 'entrypoint': ['/bin/sh'], 'command': ['shell'], 'health_test': ['NONE'], 'mounts': 'bind||/host|/data|true;', 'stop_timeout': 1, 'stop_signal': 'SIGKILL', 'pid_mode': 'host', 'ipc_mode': 'host', 'userns_mode': 'host', 'cgroupns_mode': 'host', 'uts_mode': 'host', 'group_add': ['0'], 'oci_runtime': 'alternate', 'devices': 1, 'device_requests': 1, 'device_cgroup_rules': 1, 'ports': {}, 'tmpfs': {}, 'cpus': 1, 'memory': 1, 'pids_limit': 1, 'shm_size': 1, 'log_config': {'Type': 'none', 'Config': {}}, } for field, value in cases.items(): with self.subTest(field=field): events, runner, lifecycle = self.lifecycle() runner.states[runner.OLD_RUNTIME][field] = value with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward( _Session(events), lifecycle, ) self.assertEqual(raised.exception.category, 'identity') self.assertNotIn('stop-edge', events) def test_shared_host_profile_attests_host_network_without_public_ports(self): runner = _Runner([]) runtime_payload = runner._state( 'runtime', runner.OLD_RUNTIME, runner.OLD_RUNTIME, ) runtime_payload['network'] = 'host' runtime_payload['ports'] = {} runtime_payload['cpus'] = 900_000_000 runtime_payload['memory'] = 720 * 1024 ** 2 runtime_payload['mounts'] = runtime_payload['mounts'].replace( 'volume|truf-docker_data|', 'volume|truf-remote-server-data|', 1, ) runtime_payload['config_files'] = ( '/opt/truf/compose.yaml,/opt/truf/compose.shared-host.yaml' ) edge_payload = runner._state('edge', runner.OLD_EDGE, runner.OLD_RUNTIME) edge_payload['cap_add'] = [] edge_payload['config_files'] = runtime_payload['config_files'] runtime = host_agent_lifecycle._container_state( json.dumps(runtime_payload).encode('ascii'), ) edge = host_agent_lifecycle._container_state( json.dumps(edge_payload).encode('ascii'), ) lifecycle = host_agent_lifecycle.FixedDeploymentLifecycle( _runner=runner, _profile=host_agent_lifecycle.SHARED_HOST_PROFILE, ) lifecycle._require_running( runtime, 'runtime', runner.RUNTIME_IMAGE, runner.OLD_RUNTIME, fresh=False, ) lifecycle._require_running( edge, 'edge', runner.EDGE_IMAGE, runner.OLD_RUNTIME, fresh=False, ) self.assertIn( '/opt/truf/compose.shared-host.yaml', lifecycle._compose_command, ) self.assertEqual( lifecycle._profile.edge_caddyfile, '/etc/caddy/Caddyfile.shared-host', ) runtime_payload['ports'] = { '443/tcp': [{'HostIp': '', 'HostPort': '443'}], } with self.assertRaises(host_agent_lifecycle.HostLifecycleError): lifecycle._require_running( host_agent_lifecycle._container_state( json.dumps(runtime_payload).encode('ascii'), ), 'runtime', runner.RUNTIME_IMAGE, runner.OLD_RUNTIME, fresh=False, ) runtime_payload['ports'] = {} runtime_payload['cpus'] = 2_000_000_000 with self.assertRaises(host_agent_lifecycle.HostLifecycleError): lifecycle._require_running( host_agent_lifecycle._container_state( json.dumps(runtime_payload).encode('ascii'), ), 'runtime', runner.RUNTIME_IMAGE, runner.OLD_RUNTIME, fresh=False, ) def test_preflight_rejects_runtime_bind_source_substitution(self): events, runner, lifecycle = self.lifecycle() runner.states[runner.OLD_RUNTIME]['mounts'] = runner.states[ runner.OLD_RUNTIME ]['mounts'].replace('/etc/truf/runtime', '/tmp/attacker', 1) with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertNotIn('stop-edge', events) def test_recreated_container_must_retain_compose_config_hash(self): events, runner, lifecycle = self.lifecycle() runner.new_state_mutations['runtime']['config_hash'] = 'e' * 64 with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertIn('replace', events) self.assertNotIn('strict-health', events) self.assertNotIn('up-edge', events) def test_recreated_runtime_must_have_new_identity(self): events, runner, lifecycle = self.lifecycle() runner.reuse_runtime_id = True with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertNotIn('strict-health', events) self.assertNotIn('up-edge', events) def test_strict_health_requires_exact_core_workers(self): events, runner, lifecycle = self.lifecycle() runner.health['workers'].remove('worker-api') with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'health') self.assertNotIn('up-edge', events) def test_caddy_failure_is_generic_and_is_not_retried(self): events, runner, lifecycle = self.lifecycle() runner.fail_caddy = True with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(str(raised.exception), 'host runtime lifecycle failed') self.assertNotIn('sensitive', str(raised.exception)) self.assertEqual(events.count('caddy-validate'), 1) def test_permanent_edge_identity_failure_is_not_polled(self): events, runner, lifecycle = self.lifecycle() runner.new_state_mutations['edge']['user'] = '0:0' with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle.execute_fixed_forward(_Session(events), lifecycle) self.assertEqual(raised.exception.category, 'identity') self.assertEqual(events.count('up-edge'), 1) self.assertEqual(events.count('caddy-validate'), 0) def test_strict_health_command_is_bounded_by_remaining_deadline(self): _events, runner, lifecycle = self.lifecycle() lifecycle._strict_runtime_health(0.125) command, timeout = runner.commands[-1] self.assertIn('--require-worker-api', command) self.assertIn('--require-discovery-producers', command) self.assertEqual(timeout, 0.125) def test_forged_stopped_proof_is_rejected_before_commands(self): _events, runner, lifecycle = self.lifecycle() with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: lifecycle.recreate_and_verify(object()) self.assertEqual(raised.exception.category, 'state') self.assertEqual(runner.commands, []) def test_forward_requires_entered_apply_session(self): events, _runner, lifecycle = self.lifecycle() session = _Session(events) session._entered = False with self.assertRaises(host_agent_lifecycle.HostLifecycleError): host_agent_lifecycle.execute_fixed_forward(session, lifecycle) self.assertEqual(events, []) @unittest.skipUnless(os.name == 'posix', 'real process control requires POSIX') def test_subprocess_runner_enforces_absolute_timeout(self): started = time.monotonic() with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle._subprocess_runner( (sys.executable, '-c', 'import time; time.sleep(30)'), 0.05, ) self.assertEqual(raised.exception.category, 'timeout') self.assertLess(time.monotonic() - started, 8.0) @unittest.skipUnless(os.name == 'posix', 'real process control requires POSIX') def test_subprocess_runner_does_not_signal_successful_process_group(self): with mock.patch.object(host_agent_lifecycle.os, 'killpg') as kill_group: result = host_agent_lifecycle._subprocess_runner( (sys.executable, '-c', 'print("ok")'), 10, ) self.assertEqual(result, b'ok\n') kill_group.assert_not_called() @unittest.skipUnless(os.name == 'posix', 'real process control requires POSIX') def test_subprocess_runner_terminates_oversized_output(self): started = time.monotonic() with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle._subprocess_runner( ( sys.executable, '-c', 'import os,time; os.write(1,b"x"*32768); time.sleep(30)', ), 10, ) self.assertEqual(raised.exception.category, 'command') self.assertLess(time.monotonic() - started, 8.0) @unittest.skipUnless(os.name == 'posix', 'real process control requires POSIX') def test_subprocess_runner_kills_descendant_holding_stdout(self): started = time.monotonic() code = ( 'import subprocess,sys; ' 'subprocess.Popen([sys.executable,"-c",' '"import time; time.sleep(30)"])' ) with self.assertRaises(host_agent_lifecycle.HostLifecycleError) as raised: host_agent_lifecycle._subprocess_runner( (sys.executable, '-c', code), 10, ) self.assertEqual(raised.exception.category, 'command') self.assertLess(time.monotonic() - started, 8.0) if __name__ == '__main__': unittest.main()