import json import os import sys import tempfile import time import unittest 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) from worker_contracts import ( AssignmentOutcome, DiagnosticCategory, DiagnosticHTTPContext, DiagnosticKind, DiagnosticProcessContext, ScanOutcome, WorkerPhase, build_diagnostic_envelope, make_body_material, make_log_material, ) import worker_local_state from runtime_security import atomic_write_private_json from worker_local_state import WorkerLocalState, WorkerLocalStateError UTC = '2026-09-23T12:00:00Z' class WorkerLocalStateTests(unittest.TestCase): def test_journal_recovery_ignores_truncated_tail_and_rebuilds_projection(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) first = state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) second = state.emit_phase( 'instance-1', 0, WorkerPhase.CLAIMING, timestamp='2026-09-23T12:00:01Z', ) self.assertEqual((first['sequence'], second['sequence']), (1, 2)) with open(state.event_path, 'ab') as handle: handle.write(b'{"incomplete":') with open(state.status_path, 'wb') as handle: handle.write(b'not-json') recovered = WorkerLocalState(root) self.assertEqual(recovered.next_sequence, 3) self.assertEqual(recovered.snapshot()['slots'][0]['phase'], 'claiming') self.assertTrue(open(recovered.event_path, 'rb').read().endswith(b'\n')) projection = json.load(open(recovered.status_path, encoding='utf-8')) self.assertEqual(projection['sequence'], 2) def test_generated_parent_timestamp_never_predates_explicit_runner_event(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) runner_timestamp = '2026-09-23T12:00:20Z' runner_event = state.emit_phase( 'instance-1', 0, WorkerPhase.BUNDLING, reservation_id=7, timestamp=runner_timestamp, phase_started_at=runner_timestamp, ) with mock.patch.object( worker_local_state, 'utc_now', return_value='2026-09-23T12:00:17Z', ): uploading = state.emit_phase( 'instance-1', 0, WorkerPhase.UPLOADING, reservation_id=7, ) awaiting = state.emit_phase( 'instance-1', 0, WorkerPhase.AWAITING_RECEIPT, reservation_id=7, ) self.assertEqual(runner_event['timestamp'], runner_timestamp) self.assertEqual(uploading['timestamp'], runner_timestamp) self.assertEqual(uploading['phase_started_at'], runner_timestamp) self.assertEqual(awaiting['timestamp'], runner_timestamp) self.assertEqual(awaiting['phase_started_at'], runner_timestamp) def test_projection_does_not_modify_authoritative_slot_state(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) slot = os.path.join(root, 'slot-0.json') payload = b'{"assignment":{"reservation":{"reservation_id":7}},"phase":"assigned"}\n' with open(slot, 'wb') as handle: handle.write(payload) state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) state.recover() self.assertEqual(open(slot, 'rb').read(), payload) def test_read_only_inspection_never_rewrites_projection_or_journal(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) before_projection = open(state.status_path, 'rb').read() before_journal = open(state.event_path, 'rb').read() reader = WorkerLocalState(root, read_only=True) self.assertEqual(reader.snapshot()['sequence'], 1) self.assertEqual(open(state.status_path, 'rb').read(), before_projection) self.assertEqual(open(state.event_path, 'rb').read(), before_journal) with self.assertRaisesRegex(WorkerLocalStateError, 'cannot mutate'): reader.emit_phase('instance-1', 0, WorkerPhase.CLAIMING) def test_history_is_append_only_and_idempotent(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) record = { 'history_id': 'receipt-1', 'instance_id': 'instance-1', 'slot_id': 0, 'reservation_id': 7, 'source': 'dockerhub', 'outcome': 'bundle_accepted', 'receipt': {'receipt_id': 'receipt-1'}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': 1, 'last_sequence': 4, 'diagnostics': [], 'timeline': [], 'phase_durations': {}, } self.assertTrue(state.append_history(record)) self.assertFalse(state.append_history(record)) self.assertEqual(len(state.history()), 1) self.assertEqual(state.history(reservation_id=7)[0]['history_id'], 'receipt-1') with open(state.history_path, 'ab') as handle: handle.write(b'{"truncated":') recovered = WorkerLocalState(root) second = dict(record) second['history_id'] = 'receipt-2' self.assertTrue(recovered.append_history(second)) self.assertEqual( [item['history_id'] for item in recovered.history()], ['receipt-1', 'receipt-2'], ) def test_history_uses_one_closed_schema_validator_for_append_recovery_and_read(self): record = { 'history_id': 'receipt-1', 'instance_id': 'instance-1', 'slot_id': 0, 'reservation_id': 7, 'source': 'dockerhub', 'outcome': 'bundle_accepted', 'receipt': {}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': None, 'last_sequence': None, 'diagnostics': [], 'timeline': [], 'phase_durations': {}, } with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) with self.assertRaisesRegex(WorkerLocalStateError, 'shape'): state.append_history({**record, 'unexpected': True}) state.append_history(record) invalid = { 'schema': 1, **record, 'slot_id': -1, 'history_id': 'receipt-2', } with open(state.history_path, 'ab') as handle: handle.write(json.dumps( invalid, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii') + b'\n') with self.assertRaisesRegex(WorkerLocalStateError, 'slot identity'): WorkerLocalState(root, read_only=True) with self.assertRaisesRegex(WorkerLocalStateError, 'slot identity'): state.history() def test_legacy_schema_one_history_is_explicitly_normalized_to_schema_two(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) legacy = { 'schema': 1, 'history_id': 'legacy-1', 'slot_id': 0, 'reservation_id': 7, 'source': 'gitlab', 'outcome': 'bundle_accepted', 'receipt': {}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': 1, 'last_sequence': 1, 'diagnostics': [], 'timeline': [{ 'sequence': 1, 'timestamp': UTC, 'phase': 'assigned', }], 'phase_durations': {'assigned': 0.0}, } atomic_write_private_json(state.history_path, legacy) reader = WorkerLocalState(root, read_only=True) migrated = reader.history()[0] self.assertEqual(migrated['schema'], 2) self.assertEqual(migrated['instance_id'], 'legacy-instance-unavailable') self.assertEqual(migrated['timeline'][0]['progress'], {}) self.assertEqual( migrated['timeline'][0]['instance_id'], 'legacy-instance-unavailable', ) second_pass = { **legacy, 'instance_id': 'second-pass-instance', 'timeline': [{ 'sequence': 1, 'timestamp': UTC, 'instance_id': 'second-pass-instance', 'phase': 'assigned', 'progress': {'recovered': True}, }], } atomic_write_private_json(state.history_path, second_pass) migrated = WorkerLocalState(root, read_only=True).history()[0] self.assertEqual(migrated['schema'], 2) self.assertEqual(migrated['instance_id'], 'second-pass-instance') self.assertTrue(migrated['timeline'][0]['progress']['recovered']) hybrid = { **second_pass, 'timeline': [{ 'sequence': 1, 'timestamp': UTC, 'instance_id': 'second-pass-instance', 'phase': 'assigned', }], } atomic_write_private_json(state.history_path, hybrid) with self.assertRaisesRegex(WorkerLocalStateError, 'timeline event'): WorkerLocalState(root, read_only=True) def test_event_cursor_reads_only_requested_tail_without_writer_lock_scan(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) for index in range(2, 102): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=f'2026-09-23T12:00:{min(index, 59):02d}Z', progress={'sample': index}, ) original = worker_local_state.decode_worker_event with mock.patch.object( worker_local_state, 'decode_worker_event', wraps=original, ) as decode: values = state.events_after(100, 1) self.assertEqual(values[0]['sequence'], 101) self.assertLessEqual(decode.call_count, worker_local_state.EVENT_CURSOR_STRIDE) def test_assignment_timeline_crosses_instance_recovery_boundaries(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) state.emit_phase( 'instance-1', 0, WorkerPhase.CLAIMING, reservation_id=7, source='gitlab', timestamp='2026-09-23T12:00:01Z', ) state.emit_phase( 'instance-1', 0, WorkerPhase.ASSIGNED, reservation_id=7, source='gitlab', timestamp='2026-09-23T12:00:02Z', ) state.emit_phase( 'instance-2', 0, WorkerPhase.CLAIMING, reservation_id=7, source='gitlab', timestamp='2026-09-23T12:00:03Z', progress={'recovered': True}, ) state.emit_phase( 'instance-2', 0, WorkerPhase.ASSIGNED, reservation_id=7, source='gitlab', timestamp='2026-09-23T12:00:04Z', progress={'recovered': True}, ) timeline = state.assignment_timeline(0, 7) self.assertEqual( [item['instance_id'] for item in timeline], ['instance-1', 'instance-1', 'instance-2', 'instance-2'], ) self.assertTrue(timeline[-1]['progress']['recovered']) def test_diagnostic_archive_preserves_body_stdout_and_stderr_material(self): envelope = build_diagnostic_envelope( occurrence_id='occurrence-1', reservation_id=7, scan_event_id='a' * 32, slot_id=0, source='dockerhub', phase=WorkerPhase.SCANNING, kind=DiagnosticKind.PROVIDER_HTTP, category=DiagnosticCategory.AUTHORIZATION, code='docker.manifest_http_403', summary='request denied', retryable=False, attempt=1, assignment_outcome=AssignmentOutcome.ACCEPTED, scan_outcome=ScanOutcome.ERROR, occurred_at=UTC, captured_at=UTC, http=DiagnosticHTTPContext( operation='manifest.get', status_code=403, content_type='application/json', request_id='request-1', body=make_body_material(b'{"error":"denied"}'), ), process=DiagnosticProcessContext( name='trufflehog', exit_code=1, signal=None, timed_out=False, stdout=make_log_material(b'stdout\n'), stderr=make_log_material(b'stderr\n'), ), ) with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) archived = state.archive_diagnostic(envelope) self.assertEqual(state.diagnostic_references(7), [archived]) self.assertEqual( open(os.path.join(root, *archived['artifacts']['body']['path'].split('/')), 'rb').read(), b'{"error":"denied"}', ) self.assertEqual( open(os.path.join(root, *archived['artifacts']['stdout']['path'].split('/')), 'rb').read(), b'stdout\n', ) self.assertTrue(archived['artifacts']['stdout']['path'].endswith('.stdout.log')) self.assertEqual( open(os.path.join(root, *archived['artifacts']['stderr']['path'].split('/')), 'rb').read(), b'stderr\n', ) self.assertTrue(archived['artifacts']['stderr']['path'].endswith('.stderr.log')) record = json.load(open(os.path.join(root, *archived['record'].split('/')), encoding='utf-8')) self.assertFalse(record['artifacts']['body']['truncated']) state.append_history({ 'history_id': 'diagnostic-history', 'instance_id': 'instance-1', 'slot_id': 0, 'reservation_id': 7, 'source': 'dockerhub', 'outcome': 'bundle_accepted', 'receipt': {}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': None, 'last_sequence': None, 'diagnostics': [archived], 'timeline': [], 'phase_durations': {}, }) self.assertTrue(state.history()[0]['diagnostics'][0]['available']) os.remove(os.path.join(root, *archived['record'].split('/'))) for artifact in archived['artifacts'].values(): if artifact: os.remove(os.path.join(root, *artifact['path'].split('/'))) retained = state.history()[0]['diagnostics'][0] self.assertFalse(retained['available']) self.assertFalse(any(retained['artifact_availability'].values())) def test_full_artifact_encoding_uses_complete_bytes_not_bounded_excerpt(self): full_body = (b'a' * (16 * 1024)) + b'\xff' envelope = build_diagnostic_envelope( occurrence_id='encoding-1', reservation_id=7, scan_event_id='a' * 32, slot_id=0, source='dockerhub', phase=WorkerPhase.RESOLVING, kind=DiagnosticKind.PROVIDER_HTTP, category=DiagnosticCategory.PROVIDER, code='fixture', summary='fixture', retryable=False, attempt=1, assignment_outcome=AssignmentOutcome.ACCEPTED, scan_outcome=ScanOutcome.ERROR, occurred_at=UTC, captured_at=UTC, http=DiagnosticHTTPContext( operation='fixture', status_code=500, content_type=None, request_id=None, body=make_body_material(full_body), ), ) self.assertEqual(envelope.http.body.encoding.value, 'text') self.assertTrue(envelope.http.body.truncated) with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root) archived = state.archive_diagnostic(envelope, {'body': full_body}) artifact = archived['artifacts']['body'] self.assertEqual(artifact['encoding'], 'base64') self.assertFalse(artifact['truncated']) self.assertEqual( open(os.path.join(root, *artifact['path'].split('/')), 'rb').read(), full_body, ) def test_logs_rotate_and_retention_accounts_and_removes_expired_artifacts(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, log_bytes=64 * 1024, log_files=2, retention_days=1, retention_bytes=1024 * 1024, ) state.log('a' * 40000) state.log('b' * 40000) self.assertTrue(os.path.exists(state.log_path + '.1')) expired = os.path.join(state.diagnostics_dir, 'expired.bin') with open(expired, 'wb') as handle: handle.write(b'x' * 100) os.utime(expired, (time.time() - 172800, time.time() - 172800)) before = state.retention_usage() cleanup = state.cleanup_retention() after = state.retention_usage() self.assertGreaterEqual(before['total_bytes'], after['total_bytes']) self.assertGreaterEqual(cleanup['removed_files'], 1) self.assertFalse(os.path.exists(expired)) def test_retention_never_evicts_active_assignment_diagnostics_or_data_roots(self): envelope = build_diagnostic_envelope( occurrence_id='active-1', reservation_id=7, scan_event_id='a' * 32, slot_id=0, source='dockerhub', phase=WorkerPhase.SCANNING, kind=DiagnosticKind.PROVIDER_HTTP, category=DiagnosticCategory.AUTHORIZATION, code='fixture', summary='fixture', retryable=False, attempt=1, assignment_outcome=AssignmentOutcome.UNFINISHED, scan_outcome=ScanOutcome.ERROR, occurred_at=UTC, captured_at=UTC, http=DiagnosticHTTPContext( operation='fixture', status_code=500, content_type=None, request_id=None, body=make_body_material(b'active body'), ), ) with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root, retention_days=1, retention_bytes=1024 * 1024) archived = state.archive_diagnostic(envelope) slot_path = os.path.join(root, 'slot-0.json') atomic_write_private_json(slot_path, { 'phase': 'assigned', 'assignment': {'reservation': {'reservation_id': 7}}, }) for path in ( archived['record'], archived['artifacts']['body']['path'], ): absolute = os.path.join(root, *path.split('/')) os.utime(absolute, (time.time() - 172800, time.time() - 172800)) bundle = os.path.join(root, 'bundles') work = os.path.join(root, 'work') os.makedirs(bundle) os.makedirs(work) open(os.path.join(bundle, 'ready.trb'), 'wb').write(b'bundle') open(os.path.join(work, 'active.tmp'), 'wb').write(b'work') usage = state.retention_usage({'bundles': bundle, 'work': work}) self.assertEqual(usage['categories']['bundles']['evictable_bytes'], 0) self.assertEqual(usage['categories']['work']['evictable_bytes'], 0) self.assertGreater(usage['categories']['state']['non_evictable_bytes'], 0) self.assertGreater(usage['categories']['diagnostics']['non_evictable_bytes'], 0) state.cleanup_retention() self.assertTrue(os.path.exists(os.path.join(root, *archived['record'].split('/')))) self.assertTrue(os.path.exists(os.path.join(bundle, 'ready.trb'))) self.assertTrue(os.path.exists(os.path.join(work, 'active.tmp'))) def test_retention_preserves_diagnostics_referenced_by_retained_history(self): envelope = build_diagnostic_envelope( occurrence_id='retained-1', reservation_id=7, scan_event_id='a' * 32, slot_id=0, source='dockerhub', phase=WorkerPhase.SCANNING, kind=DiagnosticKind.PROVIDER_HTTP, category=DiagnosticCategory.AUTHORIZATION, code='fixture', summary='fixture', retryable=False, attempt=1, assignment_outcome=AssignmentOutcome.ACCEPTED, scan_outcome=ScanOutcome.ERROR, occurred_at=UTC, captured_at=UTC, http=DiagnosticHTTPContext( operation='fixture', status_code=500, content_type=None, request_id=None, body=make_body_material(b'retained body'), ), ) with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root, retention_days=1) archived = state.archive_diagnostic(envelope) state.append_history({ 'history_id': 'retained-history', 'instance_id': 'instance-1', 'slot_id': 0, 'reservation_id': 7, 'source': 'dockerhub', 'outcome': 'bundle_accepted', 'receipt': {}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': None, 'last_sequence': None, 'diagnostics': [archived], 'timeline': [], 'phase_durations': {}, }) old = time.time() - 172800 paths = [archived['record']] + [ artifact['path'] for artifact in archived['artifacts'].values() if artifact is not None ] for relative in paths: absolute = os.path.join(root, *relative.split('/')) os.utime(absolute, (old, old)) state.cleanup_retention() for relative in paths: self.assertTrue(os.path.exists(os.path.join(root, *relative.split('/')))) def test_segment_rotation_restart_cursor_and_retention_preserve_active_files(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, history_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) for index in range(3): state.append_history({ 'history_id': f'history-{index}', 'instance_id': 'instance-1', 'slot_id': 0, 'reservation_id': 100 + index, 'source': 'gitlab', 'outcome': 'bundle_accepted', 'receipt': {'padding': 'y' * 700}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': None, 'last_sequence': None, 'diagnostics': [], 'timeline': [], 'phase_durations': {}, }) event_segments = state._closed_event_paths() history_segments = state._closed_history_paths() self.assertTrue(event_segments) self.assertTrue(history_segments) self.assertEqual( [item['sequence'] for item in state.events_after(2, 3)], [3, 4, 5], ) restarted = WorkerLocalState( root, event_segment_bytes=1024, history_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) self.assertEqual(restarted.next_sequence, 13) self.assertEqual(len(restarted.history()), 3) old = time.time() - 172800 for _first, _last, path in restarted._closed_event_paths(): os.utime(path, (old, old)) for _index, path in restarted._closed_history_paths(): os.utime(path, (old, old)) atomic_write_private_json( os.path.join(root, 'control', 'progress-outbox.json'), {'schema': 1, 'sequence': 12}, ) before = restarted.retention_usage() self.assertGreater(before['categories']['events']['evictable_bytes'], 0) self.assertEqual(before['categories']['history']['evictable_bytes'], 0) restarted.cleanup_retention() self.assertTrue(os.path.exists(restarted.event_path)) self.assertTrue(os.path.exists(restarted.history_path)) self.assertEqual(len(restarted._closed_history_paths()), len(history_segments)) self.assertEqual(len(restarted.history()), 3) retained_events = restarted.events_after(0, 100) self.assertTrue(retained_events) self.assertGreater(retained_events[0]['sequence'], 1) after_cleanup = WorkerLocalState( root, event_segment_bytes=1024, history_segment_bytes=1024, ) self.assertEqual(after_cleanup.next_sequence, 13) self.assertEqual( after_cleanup.events_after(retained_events[0]['sequence'] - 1, 1)[0]['sequence'], retained_events[0]['sequence'], ) def test_aggressive_retention_preserves_active_timeline_and_terminal_history(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, history_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) state.emit_phase('instance-1', 0, WorkerPhase.IDLE, timestamp=UTC) state.emit_phase( 'instance-1', 0, WorkerPhase.CLAIMING, reservation_id=7, source='gitlab', timestamp=UTC, progress={'padding': 'x' * 400}, ) for index in range(6): state.emit_phase( 'instance-1', 0, WorkerPhase.ASSIGNED, reservation_id=7, source='gitlab', timestamp=UTC, progress={'sample': index, 'padding': 'x' * 400}, ) atomic_write_private_json(os.path.join(root, 'slot-0.json'), { 'phase': 'assigned', 'assignment': {'reservation': {'reservation_id': 7}}, }) for index in range(3): state.append_history({ 'history_id': f'terminal-{index}', 'instance_id': 'instance-1', 'slot_id': 1, 'reservation_id': 100 + index, 'source': 'gitlab', 'outcome': 'bundle_accepted', 'receipt': {'padding': 'y' * 700}, 'started_at': UTC, 'completed_at': UTC, 'duration_seconds': 0.0, 'first_sequence': None, 'last_sequence': None, 'diagnostics': [], 'timeline': [], 'phase_durations': {}, }) event_segments = [path for _first, _last, path in state._closed_event_paths()] history_segments = [path for _index, path in state._closed_history_paths()] self.assertTrue(event_segments) self.assertTrue(history_segments) old = time.time() - 172800 for path in event_segments + history_segments: os.utime(path, (old, old)) usage = state.retention_usage() self.assertEqual(usage['categories']['events']['evictable_bytes'], 0) self.assertEqual(usage['categories']['history']['evictable_bytes'], 0) state.cleanup_retention() self.assertTrue(all(os.path.exists(path) for path in event_segments)) self.assertTrue(all(os.path.exists(path) for path in history_segments)) self.assertEqual(len(state.history()), 3) def test_event_retention_usage_stops_at_recent_earlier_segment(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) segments = [path for _first, _last, path in state._closed_event_paths()] self.assertGreaterEqual(len(segments), 2) old = time.time() - 172800 for path in segments: os.utime(path, (old, old)) os.utime(segments[0], None) usage = state.retention_usage() self.assertEqual(usage['categories']['events']['evictable_bytes'], 0) state.cleanup_retention() self.assertTrue(all(os.path.exists(path) for path in segments)) def test_event_retention_usage_stops_at_in_use_earlier_segment(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) segments = [path for _first, _last, path in state._closed_event_paths()] old = time.time() - 172800 for path in segments: os.utime(path, (old, old)) state._retain_segment_paths([segments[0]]) try: usage = state.retention_usage() self.assertEqual(usage['categories']['events']['evictable_bytes'], 0) state.cleanup_retention() self.assertTrue(all(os.path.exists(path) for path in segments)) finally: state._release_segment_paths([segments[0]]) def test_event_retention_byte_pressure_overrides_recent_age_until_cap(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=30, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) segments = [path for _first, _last, path in state._closed_event_paths()] work = os.path.join(root, 'work-retention') os.makedirs(work) atomic_write_private_json( os.path.join(root, 'control', 'progress-outbox.json'), {'schema': 1, 'sequence': 12}, ) initial = state.retention_usage()['total_bytes'] first_size = os.path.getsize(segments[0]) filler_size = state.retention_bytes - initial + max(1, first_size // 2) with open(os.path.join(work, 'active.bin'), 'wb') as handle: handle.write(b'w' * filler_size) usage = state.retention_usage({'work': work}) self.assertEqual( usage['categories']['events']['evictable_bytes'], first_size, ) state.cleanup_retention({'work': work}) self.assertFalse(os.path.exists(segments[0])) self.assertTrue(all(os.path.exists(path) for path in segments[1:])) self.assertTrue(os.path.exists(os.path.join(work, 'active.bin'))) def test_retention_never_deletes_past_a_protected_event_segment(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) segments = state._closed_event_paths() self.assertGreaterEqual(len(segments), 2) protected = os.path.abspath(segments[0][2]) old = time.time() - 172800 for _first, _last, path in segments: os.utime(path, (old, old)) def evictable(path, *_args): return os.path.abspath(path) != protected with mock.patch.object( state, '_event_segment_is_evictable', side_effect=evictable, ): usage = state.retention_usage() self.assertEqual(usage['categories']['events']['evictable_bytes'], 0) state.cleanup_retention() self.assertTrue(all(os.path.exists(path) for _first, _last, path in segments)) restarted = WorkerLocalState(root, event_segment_bytes=1024) self.assertEqual(restarted.next_sequence, 13) def test_progress_outbox_cursor_protects_newer_segments_in_usage_and_cleanup(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(18): state.emit_phase( 'instance-1', 0, WorkerPhase.ASSIGNED, reservation_id=7, source='gitlab', timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) segments = state._closed_event_paths() self.assertGreaterEqual(len(segments), 3) old = time.time() - 172800 for _first, _last, path in segments: os.utime(path, (old, old)) cursor = segments[0][1] atomic_write_private_json( os.path.join(root, 'control', 'progress-outbox.json'), {'schema': 1, 'sequence': cursor}, ) with mock.patch.object( state, '_event_segment_is_evictable', return_value=True, ): usage = state.retention_usage() self.assertEqual( usage['categories']['events']['evictable_files'], 1, ) self.assertEqual( usage['progress_outbox']['cursor_sequence'], cursor, ) state.cleanup_retention() self.assertFalse(os.path.exists(segments[0][2])) self.assertTrue(all( os.path.exists(path) for _first, _last, path in segments[1:] )) def test_old_root_progress_cursor_migrates_before_retention(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState(root, event_segment_bytes=1024) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) old_path = os.path.join(root, 'progress-outbox.json') atomic_write_private_json(old_path, {'schema': 1, 'sequence': 7}) path, sequence = worker_local_state.prepare_progress_outbox_cursor(root) self.assertEqual(sequence, 7) self.assertFalse(os.path.exists(old_path)) self.assertEqual( json.load(open(path, encoding='utf-8'))['sequence'], 7, ) self.assertEqual( state.retention_usage()['progress_outbox']['cursor_sequence'], 7, ) def test_missing_or_conflicting_progress_cursor_protects_journal(self): with tempfile.TemporaryDirectory() as root: state = WorkerLocalState( root, event_segment_bytes=1024, retention_days=1, retention_bytes=1024 * 1024, ) for index in range(12): state.emit_phase( 'instance-1', 0, WorkerPhase.IDLE, timestamp=UTC, progress={'sample': index, 'padding': 'x' * 300}, ) old = time.time() - 172800 for _first, _last, path in state._closed_event_paths(): os.utime(path, (old, old)) self.assertEqual( state.retention_usage()['categories']['events']['evictable_files'], 0, ) atomic_write_private_json( os.path.join(root, 'progress-outbox.json'), {'schema': 1, 'sequence': 3}, ) atomic_write_private_json( os.path.join(root, 'control', 'progress-outbox.json'), {'schema': 1, 'sequence': 9}, ) with self.assertRaisesRegex( WorkerLocalStateError, 'conflicting progress outbox cursors', ): worker_local_state.prepare_progress_outbox_cursor(root) self.assertEqual( state.retention_usage()['progress_outbox']['cursor_sequence'], 0, ) self.assertEqual( state.retention_usage()['categories']['events']['evictable_files'], 0, ) if __name__ == '__main__': unittest.main()