774 lines
37 KiB
Python
774 lines
37 KiB
Python
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()
|