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

725 lines
35 KiB
Python

import io
import json
import os
import sys
import tempfile
import time
import unittest
from contextlib import redirect_stderr, redirect_stdout
from types import SimpleNamespace
from unittest import mock
APP_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', 'app'))
if APP_DIR not in sys.path:
sys.path.insert(0, APP_DIR)
import worker_cli
from worker_cli import (
EXIT_INVALID_INVOCATION,
EXIT_NOT_RUNNING,
EXIT_STALE_OR_UNVERIFIABLE,
EXIT_STARTUP_FAILED,
EXIT_STOP_INCOMPLETE,
InvocationError,
command_attach,
command_doctor,
command_history,
command_install,
command_logs,
command_start,
command_status,
command_stop,
parse_args,
)
from worker_local_state import WorkerLocalState
from worker_contracts import WorkerPhase, decode_worker_event
from runtime_security import (
atomic_write_private_json,
ensure_private_directory,
harden_private_file,
)
def fixture_paths(root):
return {
'package_manifest': os.path.join(root, 'worker-package.json'),
'state_dir': os.path.join(root, 'state'),
'bundle_dir': os.path.join(root, 'data', 'bundles'),
'work_dir': os.path.join(root, 'data', 'work'),
}
def status_fixture(state='running'):
return {
'schema': 1, 'command': 'status', 'state': state, 'detail': 'fixture',
'instance': None, 'package': None, 'runtime': None, 'protocol': None,
'worker': None, 'slots': [], 'retention': None,
}
class WorkerCLITests(unittest.TestCase):
def test_parser_exposes_all_commands_and_preserves_bare_run_alias(self):
commands = {
parse_args([name] + (
['--server', 'https://worker.example', '--token', 'x' * 32]
if name == 'install' else []
)).command
for name in (
'install', 'run', 'start', 'stop', 'status', 'attach',
'watch', 'logs', 'history', 'doctor',
)
}
self.assertEqual(commands, {
'install', 'run', 'start', 'stop', 'status', 'attach',
'watch', 'logs', 'history', 'doctor',
})
alias = parse_args([
'--server', 'https://worker.example', '--token', 'x' * 32,
'--parallelism', '3',
])
self.assertEqual((alias.command, alias.parallelism), ('run', 3))
def test_parser_rejects_operational_path_overrides(self):
with self.assertRaises(InvocationError):
parse_args([
'run', '--server', 'https://worker.example', '--token', 'x' * 32,
'--state-dir', 'foreign',
])
def test_assignment_runner_is_hidden_and_accepts_only_its_local_root(self):
root = os.path.abspath(os.path.join('worker-work', 'worker-assignment-0-1-' + ('a' * 32)))
args = parse_args(['_assignment_runner', '--root', root])
self.assertEqual(args.command, '_assignment_runner')
self.assertEqual(args.root, root)
with self.assertRaises(InvocationError):
parse_args(['_assignment_runner', '--root', root, '--token', 'x' * 32])
def test_install_verifies_complete_package_before_persisting_configuration(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
args = SimpleNamespace(
server='https://worker.example', token='x' * 32, parallelism=2,
retention_days=14, retention_bytes=32 * 1024 * 1024,
log_bytes=128 * 1024, log_files=3,
)
with mock.patch.object(
worker_cli, 'verify_worker_package',
side_effect=ValueError('package invalid'),
), redirect_stderr(io.StringIO()):
self.assertEqual(command_install(args, paths), EXIT_STARTUP_FAILED)
self.assertFalse(os.path.exists(paths['state_dir']))
with mock.patch.object(
worker_cli, 'verify_worker_package', return_value={'manifest': {}},
), redirect_stdout(io.StringIO()):
self.assertEqual(command_install(args, paths), 0)
config = worker_cli._load_config(paths)
self.assertEqual(config['schema'], 2)
self.assertEqual(config['retention_days'], 14)
self.assertEqual(config['retention_bytes'], 32 * 1024 * 1024)
self.assertEqual(config['log_bytes'], 128 * 1024)
self.assertEqual(config['log_files'], 3)
def test_install_accepts_strict_private_yaml_without_credential_arguments(self):
with tempfile.TemporaryDirectory() as root:
path = os.path.join(root, 'worker.yaml')
with open(path, 'w', encoding='utf-8') as handle:
handle.write(
'server: https://worker.example\n'
f"token: {'x' * 32}\n"
'parallelism: 2\n'
)
harden_private_file(path)
args = parse_args(['install', '--config', path])
paths = fixture_paths(root)
with mock.patch.object(
worker_cli, 'verify_worker_package', return_value={'manifest': {}},
), redirect_stdout(io.StringIO()):
self.assertEqual(command_install(args, paths), 0)
config = worker_cli._load_config(paths)
self.assertEqual(config['server'], 'https://worker.example')
self.assertEqual(config['token'], 'x' * 32)
self.assertEqual(config['parallelism'], 2)
def test_install_yaml_rejects_ambiguous_duplicate_and_nonprivate_input(self):
with tempfile.TemporaryDirectory() as root:
path = os.path.join(root, 'worker.yaml')
with open(path, 'w', encoding='utf-8') as handle:
handle.write(
'server: https://worker.example\n'
f"token: {'x' * 32}\n"
'parallelism: 1\n'
)
if os.name != 'nt':
os.chmod(path, 0o644)
with self.assertRaisesRegex(InvocationError, 'private regular file'):
worker_cli._load_install_config(path)
harden_private_file(path)
args = parse_args([
'install', '--config', path, '--server', 'https://worker.example',
'--token', 'x' * 32,
])
with self.assertRaisesRegex(InvocationError, 'cannot be combined'):
worker_cli._configuration(args, fixture_paths(root), require_explicit=True)
with open(path, 'w', encoding='utf-8') as handle:
handle.write(
'server: https://worker.example\n'
'server: https://duplicate.example\n'
f"token: {'x' * 32}\n"
'parallelism: 1\n'
)
harden_private_file(path)
with self.assertRaisesRegex(InvocationError, 'invalid YAML'):
worker_cli._load_install_config(path)
def test_install_yaml_accepts_bounded_standard_input_and_rejects_extra_fields(self):
payload = (
'server: https://worker.example\n'
f"token: {'x' * 32}\n"
'parallelism: 3\n'
).encode('utf-8')
stdin = SimpleNamespace(buffer=io.BytesIO(payload))
with mock.patch.object(worker_cli.sys, 'stdin', stdin):
self.assertEqual(worker_cli._load_install_config('-')['parallelism'], 3)
stdin = SimpleNamespace(buffer=io.BytesIO(payload + b'extra: rejected\n'))
with mock.patch.object(worker_cli.sys, 'stdin', stdin), self.assertRaisesRegex(
InvocationError, 'must contain only',
):
worker_cli._load_install_config('-')
def test_legacy_schema_one_config_is_normalized_without_rewriting(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
ensure_private_directory(paths['state_dir'], reject_reparse=True)
ensure_private_directory(os.path.join(paths['state_dir'], 'control'), reject_reparse=True)
legacy = {
'schema': 1,
'server': 'https://worker.example',
'token': 'x' * 32,
'parallelism': 3,
'installed_at': '2026-09-23T12:00:00Z',
}
atomic_write_private_json(worker_cli._config_path(paths), legacy)
normalized = worker_cli._load_config(paths)
self.assertEqual(normalized['schema'], 2)
self.assertEqual(normalized['retention_days'], worker_cli.DEFAULT_RETENTION_DAYS)
self.assertEqual(normalized['log_files'], worker_cli.DEFAULT_LOG_FILES)
self.assertEqual(
json.load(open(worker_cli._config_path(paths), encoding='utf-8'))['schema'],
1,
)
second_pass = {
**legacy,
'retention_days': 21,
'retention_bytes': 64 * 1024 * 1024,
'log_bytes': 256 * 1024,
'log_files': 7,
}
atomic_write_private_json(worker_cli._config_path(paths), second_pass)
normalized = worker_cli._load_config(paths)
self.assertEqual(normalized['schema'], 2)
self.assertEqual(normalized['retention_days'], 21)
self.assertEqual(normalized['log_files'], 7)
self.assertEqual(
json.load(open(worker_cli._config_path(paths), encoding='utf-8'))['schema'],
1,
)
def test_status_human_and_closed_json_modes(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
for json_mode in (False, True):
output = io.StringIO()
with mock.patch.object(
worker_cli, 'status_document', return_value=status_fixture('stopped'),
), redirect_stdout(output):
code = command_status(SimpleNamespace(json=json_mode), paths)
self.assertEqual(code, EXIT_NOT_RUNNING)
if json_mode:
value = json.loads(output.getvalue())
self.assertEqual(set(value), {
'schema', 'command', 'state', 'detail', 'instance',
'package', 'runtime', 'protocol', 'worker', 'slots',
'retention',
})
else:
self.assertIn('Worker: stopped', output.getvalue())
def test_running_status_exposes_identity_cap_deadlines_progress_and_retention(self):
instance = {
'instance_id': 'fixture', 'pid': 42,
'package': {'manifest_sha256': 'a' * 64},
'runtime': {'mode': 'detached'},
'protocol': {'worker_protocol': 2},
}
classification = {
'state': 'running', 'instance': instance, 'record': {'fixture': True},
'detail': 'verified',
}
snapshot = {
'worker': {
'state': 'running', 'parallelism': 2, 'slot_cap': 2,
'configured_slots': 2, 'recovery_slots': 0,
'started_at': '2026-09-23T11:00:00Z',
'drain_deadline_at': None,
'aggregate': {'slot_count': 1, 'phases': {'backoff': 1}},
},
'projection': {'slots': [{
'slot_id': 0, 'sequence': 4, 'phase': 'backoff',
'phase_started_at': '2026-09-23T12:00:00Z',
'timestamp': '2026-09-23T12:00:01Z',
'reservation_id': None, 'source': None,
'scan_deadline_at': '2999-01-01T00:00:00Z',
'assignment_deadline_at': '2999-01-02T00:00:00Z',
'progress': {
'attempt': 1, 'reason': 'server_backoff',
'next_claim_at': '2026-09-23T12:00:10Z',
'child_state': 'sleeping',
'counters': {'objects_scanned': 12},
},
}]},
'retention': {'total_bytes': 10, 'total_files': 2},
}
with mock.patch.object(worker_cli, 'send_control_request', return_value=snapshot):
value = worker_cli.status_document(fixture_paths('unused'), classification)
self.assertIs(value['package'], instance['package'])
self.assertEqual(value['worker']['slot_cap'], 2)
self.assertEqual(value['slots'][0]['idle_reason'], 'server_backoff')
self.assertEqual(value['slots'][0]['next_claim_at'], '2026-09-23T12:00:10Z')
self.assertEqual(value['retention']['total_bytes'], 10)
output = io.StringIO()
with redirect_stdout(output):
worker_cli._human_status(value)
self.assertIn('scan_deadline=2999-01-01T00:00:00Z', output.getvalue())
self.assertIn('assignment_deadline=2999-01-02T00:00:00Z', output.getvalue())
self.assertIn('next_claim_at=2026-09-23T12:00:10Z', output.getvalue())
self.assertIn('child_state=sleeping', output.getvalue())
self.assertIn('objects_scanned', output.getvalue())
def test_attach_machine_mode_emits_coherent_snapshot_then_events(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
snapshot = status_fixture()
responses = iter((
{'schema': 1, 'events': [{
'schema': 1, 'sequence': 1, 'slot_id': 0,
'phase': 'claiming', 'reservation_id': None,
}], 'last_sequence': 1},
{'schema': 1, 'events': [], 'last_sequence': 1},
))
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=snapshot), \
mock.patch.object(worker_cli, 'send_control_request', side_effect=lambda *_a, **_k: next(responses)), \
mock.patch.object(worker_cli.time, 'monotonic', side_effect=[0.0, 0.0, 0.0, 0.1, 1.1]), \
mock.patch.object(worker_cli.time, 'sleep'), redirect_stdout(output):
code = command_attach(
SimpleNamespace(json=False, ndjson=True, follow_seconds=1.0),
fixture_paths('unused'),
)
self.assertEqual(code, 0)
values = [json.loads(line) for line in output.getvalue().splitlines()]
self.assertEqual(values[0]['type'], 'snapshot')
self.assertEqual(values[1]['sequence'], 1)
def test_attach_json_is_exactly_one_document_and_does_not_follow(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture()), \
mock.patch.object(worker_cli, 'send_control_request', return_value={
'events': [], 'last_sequence': 0,
}) as request, \
redirect_stdout(output):
self.assertEqual(command_attach(
SimpleNamespace(json=True, ndjson=False, follow_seconds=60),
fixture_paths('unused'),
), 0)
values = output.getvalue().splitlines()
self.assertEqual(len(values), 1)
self.assertEqual(json.loads(values[0])['command'], 'attach')
request.assert_not_called()
def test_attach_human_mode_detaches_without_stop(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture()), \
redirect_stdout(output):
code = command_attach(
SimpleNamespace(json=False, ndjson=False, follow_seconds=0),
fixture_paths('unused'),
)
self.assertEqual(code, 0)
self.assertIn('detaches without stopping', output.getvalue())
def test_watch_alias_uses_live_attach_without_stopping(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
output = io.StringIO()
args = parse_args(['watch', '--follow-seconds', '0'])
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture()), \
redirect_stdout(output):
self.assertEqual(command_attach(args, fixture_paths('unused')), 0)
self.assertIn('Watching; ', output.getvalue())
self.assertIn('without stopping', output.getvalue())
def test_attach_human_mode_refreshes_coherent_status_after_transition(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
responses = iter((
{'schema': 1, 'events': [{
'schema': 1, 'sequence': 1, 'slot_id': 0,
'phase': 'claiming', 'reservation_id': None,
}], 'last_sequence': 1},
{'schema': 1, 'events': [], 'last_sequence': 1},
))
slow_stdin = mock.Mock()
slow_stdin.read.side_effect = lambda _size: (time.sleep(0.5), '')[1]
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture()), \
mock.patch.object(worker_cli, 'send_control_request', side_effect=lambda *_a, **_k: next(responses)), \
mock.patch.object(worker_cli.time, 'monotonic', side_effect=[0.0, 0.0, 0.0, 0.1, 1.1]), \
mock.patch.object(worker_cli.sys, 'stdin', slow_stdin), \
redirect_stdout(output):
code = command_attach(
SimpleNamespace(json=False, ndjson=False, follow_seconds=1.0),
fixture_paths('unused'),
)
self.assertEqual(code, 0)
self.assertGreaterEqual(output.getvalue().count('Worker: running'), 2)
def test_logs_and_history_support_finite_json_and_ndjson(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
local = WorkerLocalState(paths['state_dir'])
local.log('fixture log')
local.append_history({
'history_id': 'receipt-1', 'instance_id': 'instance-1',
'slot_id': 0, 'reservation_id': 9, 'source': 'gitlab',
'outcome': 'bundle_accepted', 'receipt': {}, 'started_at': None,
'completed_at': '2026-09-23T12:00:00Z', 'duration_seconds': None,
'first_sequence': None, 'last_sequence': None, 'diagnostics': [],
'timeline': [], 'phase_durations': {},
})
output = io.StringIO()
with redirect_stdout(output):
command_logs(SimpleNamespace(
tail=10, follow=False, json=True, ndjson=False,
follow_seconds=0,
), paths)
self.assertIn('fixture log', json.loads(output.getvalue())['lines'][0])
output = io.StringIO()
with redirect_stdout(output):
command_logs(SimpleNamespace(
tail=10, follow=False, json=False, ndjson=True,
follow_seconds=0,
), paths)
self.assertEqual(json.loads(output.getvalue())['type'], 'log')
output = io.StringIO()
with redirect_stdout(output):
command_logs(SimpleNamespace(
tail=10, follow=False, json=False, ndjson=False,
follow_seconds=0,
), paths)
self.assertIn('fixture log', output.getvalue())
output = io.StringIO()
with redirect_stdout(output):
command_history(SimpleNamespace(
limit=10, reservation=None, json=False, ndjson=True,
), paths)
self.assertEqual(json.loads(output.getvalue())['reservation_id'], 9)
output = io.StringIO()
with redirect_stdout(output):
command_history(SimpleNamespace(
limit=10, reservation=None, json=True, ndjson=False,
), paths)
self.assertEqual(json.loads(output.getvalue())['assignments'][0]['reservation_id'], 9)
output = io.StringIO()
with redirect_stdout(output):
command_history(SimpleNamespace(
limit=10, reservation=None, json=False, ndjson=False,
), paths)
self.assertIn('reservation 9', output.getvalue())
def test_follow_json_alias_emits_ndjson_and_handles_layout_without_journal(self):
parsed = parse_args(['logs', '--follow', '--json'])
self.assertTrue(parsed.follow)
self.assertTrue(parsed.json)
with tempfile.TemporaryDirectory() as root:
output = io.StringIO()
errors = io.StringIO()
with redirect_stdout(output), redirect_stderr(errors):
code = command_logs(SimpleNamespace(
tail=10, follow=True, json=True, ndjson=False,
follow_seconds=0,
), fixture_paths(root))
self.assertEqual(code, EXIT_NOT_RUNNING)
self.assertEqual(output.getvalue(), '')
self.assertIn('not initialized', errors.getvalue())
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
local = WorkerLocalState(paths['state_dir'])
local.log('human prelude must not be emitted')
local.emit_phase(
'instance-1', 0, WorkerPhase.IDLE,
timestamp='2026-09-23T12:00:00Z',
)
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value={
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}), redirect_stdout(output):
command_logs(SimpleNamespace(
tail=10, follow=True, json=True, ndjson=False,
follow_seconds=0,
), paths)
values = [json.loads(line) for line in output.getvalue().splitlines()]
self.assertEqual(len(values), 1)
self.assertEqual(values[0]['type'], 'slot.phase')
self.assertNotIn('human prelude', output.getvalue())
for line in output.getvalue().splitlines():
decode_worker_event(line.encode('ascii'))
ndjson_output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value={
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}), redirect_stdout(ndjson_output):
command_logs(SimpleNamespace(
tail=10, follow=True, json=False, ndjson=True,
follow_seconds=0,
), paths)
self.assertTrue(ndjson_output.getvalue().splitlines())
for line in ndjson_output.getvalue().splitlines():
decode_worker_event(line.encode('ascii'))
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
ensure_private_directory(paths['state_dir'], reject_reparse=True)
ensure_private_directory(os.path.join(paths['state_dir'], 'control'), reject_reparse=True)
output = io.StringIO()
errors = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value={
'state': 'stopped', 'instance': None, 'detail': 'not running',
}), redirect_stdout(output), redirect_stderr(errors):
code = command_logs(SimpleNamespace(
tail=10, follow=True, json=True, ndjson=False,
follow_seconds=0,
), paths)
self.assertEqual(code, EXIT_NOT_RUNNING)
self.assertEqual(output.getvalue(), '')
self.assertIn('follow ended', errors.getvalue())
def test_startup_failure_has_distinct_exit_code(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
args = SimpleNamespace(
server='https://worker.example', token='x' * 32,
parallelism=1, startup_timeout=0.2,
)
process = mock.Mock()
process.pid = 123
process.poll.return_value = None
with mock.patch.object(
worker_cli, 'classify_instance', return_value={
'state': 'stopped', 'instance': None, 'detail': 'none',
},
), mock.patch.object(worker_cli, 'spawn_detached', return_value=process), \
mock.patch.object(
worker_cli, 'capture_spawned_process_identity',
return_value={
'pid': 123, 'creation_time': 'fixture',
'executable': 'python',
},
), mock.patch.object(worker_cli, 'terminate_spawned_process') as terminate, \
mock.patch.object(
worker_cli, 'wait_for_startup',
side_effect=worker_cli.WorkerSupervisorError('fixture failure'),
), redirect_stderr(io.StringIO()):
self.assertEqual(command_start(args, paths), EXIT_STARTUP_FAILED)
terminate.assert_called_once()
def test_concurrent_start_does_not_accept_starting_instance(self):
classification = {
'state': 'starting', 'instance': {'instance_id': 'fixture'},
'detail': 'verified starting control handshake',
}
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture('starting')), \
mock.patch.object(worker_cli, 'spawn_detached') as spawn, \
redirect_stdout(io.StringIO()):
code = command_start(SimpleNamespace(
server=None, token=None, parallelism=None, startup_timeout=1,
), fixture_paths('unused'))
self.assertEqual(code, EXIT_STALE_OR_UNVERIFIABLE)
spawn.assert_not_called()
def test_stop_requires_clean_matching_shutdown_receipt(self):
classification = {
'state': 'running', 'record': {'instance_id': 'fixture'},
'instance': {'instance_id': 'fixture'}, 'detail': 'verified',
}
args = SimpleNamespace(timeout=1.0, json=True)
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'send_control_request', return_value={'accepted': True}), \
mock.patch.object(worker_cli, 'load_shutdown_receipt', return_value={
'instance_id': 'fixture', 'drained': False, 'exit_code': 2,
}), redirect_stdout(io.StringIO()):
self.assertEqual(command_stop(args, fixture_paths('unused')), EXIT_STOP_INCOMPLETE)
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'send_control_request', return_value={'accepted': True}), \
mock.patch.object(worker_cli, 'load_shutdown_receipt', return_value={
'instance_id': 'fixture', 'drained': True, 'exit_code': 0,
}), redirect_stdout(io.StringIO()):
self.assertEqual(command_stop(args, fixture_paths('unused')), 0)
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'send_control_request', return_value={'accepted': True}), \
mock.patch.object(worker_cli.time, 'monotonic', side_effect=[0.0, 7.0]), \
redirect_stdout(io.StringIO()):
self.assertEqual(command_stop(args, fixture_paths('unused')), EXIT_STOP_INCOMPLETE)
def test_control_disconnects_render_reclassified_final_state(self):
running = {
'state': 'running', 'record': {'instance_id': 'fixture'},
'instance': {'instance_id': 'fixture'}, 'detail': 'verified',
}
stopped = {'state': 'stopped', 'instance': None, 'detail': 'exited'}
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
output = io.StringIO()
with mock.patch.object(
worker_cli, 'classify_instance', side_effect=[running, stopped],
), mock.patch.object(
worker_cli, 'send_control_request', side_effect=OSError('disconnect'),
), redirect_stdout(output):
code = command_stop(SimpleNamespace(timeout=1.0, json=True), paths)
self.assertEqual(code, EXIT_NOT_RUNNING)
self.assertEqual(json.loads(output.getvalue())['state'], 'control_disconnected')
output = io.StringIO()
with mock.patch.object(
worker_cli, 'classify_instance', side_effect=[running, stopped],
), mock.patch.object(
worker_cli, 'status_document', side_effect=[status_fixture(), status_fixture('stopped')],
), mock.patch.object(
worker_cli, 'send_control_request', side_effect=OSError('disconnect'),
), redirect_stdout(output):
code = command_attach(
SimpleNamespace(json=False, ndjson=True, follow_seconds=1), paths,
)
self.assertEqual(code, EXIT_NOT_RUNNING)
self.assertEqual(json.loads(output.getvalue().splitlines()[-1])['status']['state'], 'stopped')
local = WorkerLocalState(paths['state_dir'])
local.emit_phase('instance-1', 0, WorkerPhase.IDLE)
output = io.StringIO()
errors = io.StringIO()
with mock.patch.object(
worker_cli, 'classify_instance', side_effect=[running, stopped],
), mock.patch.object(
worker_cli, 'status_document', return_value=status_fixture('stopped'),
), mock.patch.object(
worker_cli, 'send_control_request', side_effect=OSError('disconnect'),
), redirect_stdout(output), redirect_stderr(errors):
code = command_logs(SimpleNamespace(
tail=10, follow=True, json=True, ndjson=False,
follow_seconds=1,
), paths)
self.assertEqual(code, EXIT_NOT_RUNNING)
for line in output.getvalue().splitlines():
decode_worker_event(line.encode('ascii'))
self.assertIn('control disconnected: stopped', errors.getvalue())
def test_human_attach_detaches_on_non_tty_eof(self):
classification = {
'state': 'running', 'record': {'fixture': True},
'instance': None, 'detail': 'verified',
}
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value=classification), \
mock.patch.object(worker_cli, 'status_document', return_value=status_fixture()), \
mock.patch.object(worker_cli.sys, 'stdin', io.StringIO('')), \
mock.patch.object(worker_cli, 'send_control_request') as request, \
redirect_stdout(output):
code = command_attach(
SimpleNamespace(json=False, ndjson=False, follow_seconds=None),
fixture_paths('unused'),
)
self.assertEqual(code, 0)
self.assertLessEqual(request.call_count, 1)
def test_doctor_reports_applicability_in_human_and_json_modes(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
for json_mode in (False, True):
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value={
'state': 'stopped', 'instance': None, 'detail': 'no instance record',
}), redirect_stdout(output):
code = command_doctor(SimpleNamespace(json=json_mode), paths)
self.assertEqual(code, EXIT_STARTUP_FAILED)
if json_mode:
value = json.loads(output.getvalue())
server = next(item for item in value['checks'] if item['name'] == 'server_reachability')
self.assertFalse(server['applicable'])
self.assertEqual(server['status'], 'not_applicable')
else:
self.assertIn('[not_applicable] server_reachability', output.getvalue())
self.assertFalse(os.path.exists(paths['state_dir']))
def test_doctor_writability_probe_uses_existing_paths_and_removes_probe(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
for path in (paths['state_dir'], paths['bundle_dir'], paths['work_dir']):
os.makedirs(path)
output = io.StringIO()
with mock.patch.object(worker_cli, 'classify_instance', return_value={
'state': 'stopped', 'instance': None, 'detail': 'no instance record',
}), redirect_stdout(output):
command_doctor(SimpleNamespace(json=True), paths)
value = json.loads(output.getvalue())
path_check = next(item for item in value['checks'] if item['name'] == 'paths')
self.assertTrue(path_check['applicable'])
self.assertEqual(path_check['status'], 'ok')
for path in (paths['state_dir'], paths['bundle_dir'], paths['work_dir']):
self.assertFalse(any(name.startswith('.doctor-') for name in os.listdir(path)))
def test_malformed_instance_status_is_structured_and_read_only(self):
with tempfile.TemporaryDirectory() as root:
paths = fixture_paths(root)
control = ensure_private_directory(
os.path.join(paths['state_dir'], 'control'), reject_reparse=True,
)
record = os.path.join(control, 'worker.instance.json')
atomic_write_private_json(record, {'schema': 999})
output = io.StringIO()
with redirect_stdout(output):
code = command_status(SimpleNamespace(json=True), paths)
value = json.loads(output.getvalue())
self.assertEqual(code, EXIT_STALE_OR_UNVERIFIABLE)
self.assertEqual(value['state'], 'unverifiable')
self.assertTrue(os.path.exists(record))
def test_invalid_invocation_exit_is_distinct(self):
errors = io.StringIO()
with redirect_stderr(errors):
self.assertEqual(worker_cli.main(['run', '--parallelism', '0']), EXIT_INVALID_INVOCATION)
self.assertIn('parallelism', errors.getvalue())
if __name__ == '__main__':
unittest.main()