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()