1037 lines
42 KiB
Python
1037 lines
42 KiB
Python
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()
|