"""Strict file protocol and package-local process for one remote assignment.""" import hashlib import json import os import re import time from datetime import datetime, timezone from janitor import JanitorBudget, bounded_remove_tree, run_janitor_pass from process_identity import current_process_identity, serialize_process_identity from result_bundle import ( BundleReservation, ResultBundleReader, bundle_ready_path, ensure_bundle_reservation_paths, ) from runtime_security import ( atomic_write_private_json, canonical_path, durable_publish_directory, durable_publish, ensure_private_directory, fsync_directory, harden_private_directory, harden_private_file, private_file_ready, read_private_json, reject_reparse_components, require_private_directory, write_private_json_exclusive, ) from scan_execution import ( WorkerBuildCompatibility, execute_protocol2_remote_claim, stage_scan_result_in_scope, validate_protocol2_remote_assignment, ) import scanner from worker_contracts import WorkerPhase, validate_phase_transition from worker_package import PACKAGE_DETECTOR_POLICY RUNNER_PROTOCOL_SCHEMA = 1 RUNNER_INPUT_NAME = 'input.json' RUNNER_EVENTS_NAME = 'events.jsonl' RUNNER_START_NAME = 'start.json' RUNNER_TERMINAL_NAME = 'terminal.json' RUNNER_ROOT_PREFIX = 'worker-assignment-' MAX_RUNNER_INPUT_BYTES = 4 * 1024 * 1024 MAX_RUNNER_EVENT_BYTES = 64 * 1024 MAX_RUNNER_EVENTS_BYTES = 4 * 1024 * 1024 MAX_RUNNER_OUTCOME_BYTES = 1024 * 1024 _GENERATION_RE = re.compile(r'^[a-f0-9]{32}$') _ROOT_RE = re.compile(r'^worker-assignment-(0|[1-9][0-9]*)-([1-9][0-9]*)-([a-f0-9]{32})$') _DIGEST_RE = re.compile(r'^[a-f0-9]{64}$') _ASSIGNMENT_FIELDS = { 'reservation', 'deadlines', 'compatibility', 'scan_kwargs', 'event_scan_options', 'queue_policy', 'limits', 'scan_policy', 'execution_snapshot', 'execution_snapshot_sha256', 'execution_plan', } _COMMIT_FIELDS = { 'target', 'scan_event_id', 'bundle_id', 'reservation_id', 'scan_event_hash', 'actual_bytes', 'relative_path', 'frame_count', 'finding_count', 'error_count', 'candidate_count', 'queue_status', 'source_failure', 'source_failure_category', 'source_failure_auth_related', 'first_error', } class RunnerProtocolError(RuntimeError): pass class RunnerFencedError(RunnerProtocolError): pass class RunnerDeadlineElapsed(RunnerProtocolError): pass def utc_now(): return datetime.now(timezone.utc).isoformat(timespec='milliseconds').replace('+00:00', 'Z') def canonical_json_bytes(value, *, newline=False): try: payload = json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), allow_nan=False, ).encode('utf-8') except (TypeError, ValueError) as exc: raise RunnerProtocolError('runner protocol value is not canonical JSON') from exc return payload + (b'\n' if newline else b'') def runner_root_name(slot_id, reservation_id, generation): slot_id = int(slot_id) reservation_id = int(reservation_id) generation = str(generation or '') if slot_id < 0 or reservation_id <= 0 or not _GENERATION_RE.fullmatch(generation): raise RunnerProtocolError('runner root identity is invalid') return f'{RUNNER_ROOT_PREFIX}{slot_id}-{reservation_id}-{generation}' def runner_paths(work_root, root_name): work_root = require_private_directory(os.path.abspath(work_root), create=False) match = _ROOT_RE.fullmatch(str(root_name or '')) if match is None: raise RunnerProtocolError('runner root name is invalid') root = os.path.join(work_root, root_name) if os.path.dirname(os.path.abspath(root)) != os.path.abspath(work_root): raise RunnerProtocolError('runner root escapes worker work storage') return { 'root': root, 'input': os.path.join(root, RUNNER_INPUT_NAME), 'events': os.path.join(root, RUNNER_EVENTS_NAME), 'start': os.path.join(root, RUNNER_START_NAME), 'terminal': os.path.join(root, RUNNER_TERMINAL_NAME), 'bundle_root': os.path.join(root, 'bundle'), 'scanner_work': os.path.join(root, 'scanner-work'), } def _timestamp(value, field): if type(value) is not str: raise RunnerProtocolError(f'runner {field} is invalid') try: parsed = datetime.fromisoformat(value.replace('Z', '+00:00')) except ValueError as exc: raise RunnerProtocolError(f'runner {field} is invalid') from exc if parsed.tzinfo is None: raise RunnerProtocolError(f'runner {field} is invalid') return parsed.astimezone(timezone.utc) def validate_runner_input(value): if not isinstance(value, dict) or set(value) != { 'schema', 'generation', 'slot_id', 'created_at', 'scan_started_at', 'scan_deadline_at', 'watchdog_deadline_at', 'operation', 'timeout_phase', 'assignment', }: raise RunnerProtocolError('runner input shape is invalid') if value.get('schema') != RUNNER_PROTOCOL_SCHEMA: raise RunnerProtocolError('runner input schema is invalid') generation = str(value.get('generation') or '') if _GENERATION_RE.fullmatch(generation) is None: raise RunnerProtocolError('runner generation is invalid') if type(value.get('slot_id')) is not int or value['slot_id'] < 0: raise RunnerProtocolError('runner slot identity is invalid') created = _timestamp(value.get('created_at'), 'created timestamp') started = _timestamp(value.get('scan_started_at'), 'scan start timestamp') deadline = _timestamp(value.get('scan_deadline_at'), 'scan deadline timestamp') watchdog_deadline = _timestamp( value.get('watchdog_deadline_at'), 'watchdog deadline timestamp', ) if deadline <= started or created < started: raise RunnerProtocolError('runner scan deadline ordering is invalid') if watchdog_deadline <= created: raise RunnerProtocolError('runner watchdog deadline ordering is invalid') operation = value.get('operation') timeout_phase = value.get('timeout_phase') if operation not in {'execute', 'timeout_bundle'}: raise RunnerProtocolError('runner operation is invalid') if operation == 'execute' and timeout_phase is not None: raise RunnerProtocolError('scan runner cannot carry a timeout phase') if operation == 'timeout_bundle': try: timeout_phase = WorkerPhase(timeout_phase).value except ValueError as exc: raise RunnerProtocolError('timeout runner phase is invalid') from exc if timeout_phase in { WorkerPhase.IDLE.value, WorkerPhase.CLAIMING.value, WorkerPhase.UPLOADING.value, WorkerPhase.AWAITING_RECEIPT.value, WorkerPhase.BACKOFF.value, WorkerPhase.DRAINING.value, WorkerPhase.STOPPED.value, }: raise RunnerProtocolError('timeout runner phase is outside the scan stage') if ( not isinstance(value.get('assignment'), dict) or set(value['assignment']) != _ASSIGNMENT_FIELDS ): raise RunnerProtocolError('runner assignment is invalid') reservation = dict(value['assignment'].get('reservation') or {}) if int(reservation.get('reservation_id') or 0) <= 0: raise RunnerProtocolError('runner reservation identity is invalid') runner_root_name( value['slot_id'], reservation['reservation_id'], generation, ) return dict(value) def build_runner_input( assignment, *, generation, slot_id, scan_started_at, scan_deadline_at, watchdog_deadline_at=None, created_at=None, operation='execute', timeout_phase=None, ): return validate_runner_input({ 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': str(generation), 'slot_id': int(slot_id), 'created_at': created_at or utc_now(), 'scan_started_at': str(scan_started_at), 'scan_deadline_at': str(scan_deadline_at), 'watchdog_deadline_at': str(watchdog_deadline_at or scan_deadline_at), 'operation': str(operation), 'timeout_phase': timeout_phase, 'assignment': dict(assignment), }) def _owner_marker(work_root, root, owner_identity, parent_identity): relative = os.path.relpath(root, work_root) if relative.startswith('..' + os.sep) or os.path.isabs(relative): raise RunnerProtocolError('runner root escapes worker work storage') owner = dict(owner_identity) parent = dict(parent_identity) fields = ('pid', 'creation_time', 'executable') if any(not owner.get(field) or not parent.get(field) for field in fields): raise RunnerProtocolError('runner process identity is incomplete') return { 'schema': 2, **{f'owner_{field}': owner[field] for field in fields}, **{f'parent_{field}': parent[field] for field in fields}, 'created_at': datetime.now(timezone.utc).isoformat(timespec='seconds'), 'root_kind': 'work', 'relative_path': relative.replace(os.sep, '/'), 'command': ['worker-assignment-runner'], } def create_runner_root(work_root, root_name, runner_input): paths = runner_paths(work_root, root_name) if os.path.lexists(paths['root']): raise RunnerProtocolError('runner root already exists') os.mkdir(paths['root'], 0o700) harden_private_directory(paths['root']) try: for name in ('bundle', 'scanner-work'): ensure_private_directory(os.path.join(paths['root'], name), reject_reparse=True) for name in ('tmp', 'ready', 'quarantine'): ensure_private_directory(os.path.join(paths['bundle_root'], name), reject_reparse=True) current = serialize_process_identity(current_process_identity()) atomic_write_private_json( os.path.join(paths['root'], '.scanner-owner.json'), _owner_marker(work_root, paths['root'], current, current), ) atomic_write_private_json( paths['input'], validate_runner_input(runner_input), max_bytes=MAX_RUNNER_INPUT_BYTES, ) with open(paths['input'], 'rb') as handle: input_sha256 = hashlib.sha256(handle.read(MAX_RUNNER_INPUT_BYTES + 1)).hexdigest() return paths, input_sha256 except BaseException: try: bounded_remove_tree(paths['root'], JanitorBudget( max_candidates=1, max_entries=4000, max_bytes=512 * 1024 * 1024, max_seconds=2.0, max_depth=64, )) except OSError: pass raise def bind_runner_owner(work_root, root_name, payload_identity): paths = runner_paths(work_root, root_name) parent = serialize_process_identity(current_process_identity()) atomic_write_private_json( os.path.join(paths['root'], '.scanner-owner.json'), _owner_marker(work_root, paths['root'], payload_identity, parent), ) def bind_transferred_runner_owner(work_root, root_name, payload_identity): work_root = require_private_directory(os.path.abspath(work_root), create=False) root = os.path.join(work_root, 'abandoned', root_name) if not os.path.isdir(root): return False parent = serialize_process_identity(current_process_identity()) atomic_write_private_json( os.path.join(root, '.scanner-owner.json'), _owner_marker(work_root, root, payload_identity, parent), ) return True def _load_canonical_object(path, maximum, label): reject_reparse_components(path) if not private_file_ready(path): raise RunnerProtocolError(f'runner {label} is not an exact private file') with open(path, 'rb') as handle: payload = handle.read(maximum + 1) if len(payload) > maximum: raise RunnerProtocolError(f'runner {label} exceeds its byte bound') try: value = json.loads(payload.decode('utf-8', errors='strict')) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise RunnerProtocolError(f'runner {label} is invalid JSON') from exc if canonical_json_bytes(value, newline=True) != payload: raise RunnerProtocolError(f'runner {label} is not canonical JSON') return value, hashlib.sha256(payload).hexdigest() def load_runner_input(path): value, digest = _load_canonical_object(path, MAX_RUNNER_INPUT_BYTES, 'input') return validate_runner_input(value), digest def validate_runner_event(value, *, generation=None, input_sha256=None): if not isinstance(value, dict) or set(value) != { 'schema', 'generation', 'input_sha256', 'sequence', 'timestamp', 'phase', 'phase_started_at', 'progress', } or value.get('schema') != RUNNER_PROTOCOL_SCHEMA: raise RunnerProtocolError('runner event shape is invalid') if _GENERATION_RE.fullmatch(str(value.get('generation') or '')) is None: raise RunnerProtocolError('runner event generation is invalid') if generation is not None and value['generation'] != generation: raise RunnerProtocolError('runner event generation conflicts with slot authority') if _DIGEST_RE.fullmatch(str(value.get('input_sha256') or '')) is None: raise RunnerProtocolError('runner event input hash is invalid') if input_sha256 is not None and value['input_sha256'] != input_sha256: raise RunnerProtocolError('runner event input hash conflicts with slot authority') if type(value.get('sequence')) is not int or value['sequence'] <= 0: raise RunnerProtocolError('runner event sequence is invalid') _timestamp(value.get('timestamp'), 'event timestamp') _timestamp(value.get('phase_started_at'), 'phase start timestamp') try: WorkerPhase(value.get('phase')) except ValueError as exc: raise RunnerProtocolError('runner event phase is invalid') from exc if not isinstance(value.get('progress'), dict): raise RunnerProtocolError('runner event progress is invalid') return dict(value) def read_runner_events( path, *, generation, input_sha256, operation=None, after_sequence=0, ): if not os.path.exists(path): return [] reject_reparse_components(path) if not private_file_ready(path): raise RunnerProtocolError('runner event journal is not an exact private file') if os.path.getsize(path) > MAX_RUNNER_EVENTS_BYTES: raise RunnerProtocolError('runner event journal exceeds its byte bound') events = [] previous = None with open(path, 'rb') as handle: while True: payload = handle.readline(MAX_RUNNER_EVENT_BYTES + 2) if not payload: break if len(payload) > MAX_RUNNER_EVENT_BYTES + 1: raise RunnerProtocolError('runner event exceeds its byte bound') if not payload.endswith(b'\n'): break try: value = json.loads(payload.decode('utf-8', errors='strict')) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise RunnerProtocolError('runner event journal contains invalid JSON') from exc if canonical_json_bytes(value, newline=True) != payload: raise RunnerProtocolError('runner event is not canonical JSON') event = validate_runner_event( value, generation=generation, input_sha256=input_sha256, ) if previous is not None: if event['sequence'] != previous['sequence'] + 1: raise RunnerProtocolError('runner event sequence is not contiguous') validate_phase_transition(previous['phase'], event['phase']) previous_time = _timestamp(previous['timestamp'], 'event timestamp') current_time = _timestamp(event['timestamp'], 'event timestamp') if current_time < previous_time: raise RunnerProtocolError('runner event timestamps are not monotonic') expected_phase_start = ( previous['phase_started_at'] if event['phase'] == previous['phase'] else event['timestamp'] ) if event['phase_started_at'] != expected_phase_start: raise RunnerProtocolError('runner phase start timestamp is inconsistent') elif event['sequence'] != 1: raise RunnerProtocolError('runner event journal does not begin at sequence one') elif operation is not None and event['phase'] != ( WorkerPhase.PREPARING.value if operation == 'execute' else WorkerPhase.BUNDLING.value ): raise RunnerProtocolError('runner event journal begins with the wrong operation phase') elif event['phase_started_at'] != event['timestamp']: raise RunnerProtocolError('runner first phase start timestamp is inconsistent') previous = event if event['sequence'] > int(after_sequence or 0): events.append(event) return events def phase_durations_from_events(events, completed_at): completed = _timestamp(completed_at, 'outcome timestamp') durations = {} if not events: return durations for index, event in enumerate(events): started = _timestamp(event['timestamp'], 'event timestamp') ended = ( _timestamp(events[index + 1]['timestamp'], 'event timestamp') if index + 1 < len(events) else completed ) if ended < started: raise RunnerProtocolError('runner duration timestamps are inconsistent') name = event['phase'] durations[name] = durations.get(name, 0.0) + (ended - started).total_seconds() return {name: round(value, 6) for name, value in durations.items()} def validate_terminal_against_journal(terminal, events, runner_input): terminal = validate_generation_terminal( terminal, generation=runner_input['generation'], ) if terminal['input_sha256'] != hashlib.sha256( canonical_json_bytes(runner_input, newline=True), ).hexdigest(): raise RunnerProtocolError('runner terminal does not match canonical input') if terminal['decision'] != 'completed': return terminal outcome = terminal['outcome'] if not events: raise RunnerProtocolError('completed runner has no durable phase events') expected_sequence = events[-1]['sequence'] if events else 0 expected_phase = events[-1]['phase'] if events else None expected_durations = phase_durations_from_events( events, outcome['completed_at'], ) if ( outcome['last_event_sequence'] != expected_sequence or outcome['final_phase'] != expected_phase or outcome['phase_durations'] != expected_durations ): raise RunnerProtocolError('runner outcome conflicts with its durable event journal') expected_first = ( WorkerPhase.PREPARING.value if runner_input['operation'] == 'execute' else WorkerPhase.BUNDLING.value ) if events and events[0]['phase'] != expected_first: raise RunnerProtocolError('runner operation began with an invalid phase') return terminal def _duration_map(value): if not isinstance(value, dict) or set(value) - {phase.value for phase in WorkerPhase}: raise RunnerProtocolError('runner phase durations are invalid') result = {} for name, duration in value.items(): if isinstance(duration, bool) or not isinstance(duration, (int, float)) or not 0 <= duration < 10 ** 9: raise RunnerProtocolError('runner phase duration is invalid') result[name] = round(float(duration), 6) return result def validate_runner_outcome(value, *, generation=None, input_sha256=None): if not isinstance(value, dict) or set(value) != { 'schema', 'generation', 'input_sha256', 'identity', 'status', 'completed_at', 'last_event_sequence', 'final_phase', 'phase_durations', 'bundle', 'error', } or value.get('schema') != RUNNER_PROTOCOL_SCHEMA: raise RunnerProtocolError('runner outcome shape is invalid') if _GENERATION_RE.fullmatch(str(value.get('generation') or '')) is None: raise RunnerProtocolError('runner outcome generation is invalid') if generation is not None and value['generation'] != generation: raise RunnerProtocolError('runner outcome generation conflicts with slot authority') if _DIGEST_RE.fullmatch(str(value.get('input_sha256') or '')) is None: raise RunnerProtocolError('runner outcome input hash is invalid') if input_sha256 is not None and value['input_sha256'] != input_sha256: raise RunnerProtocolError('runner outcome input hash conflicts with slot authority') identity = value.get('identity') if not isinstance(identity, dict) or set(identity) != { 'slot_id', 'reservation_id', 'bundle_id', 'scan_event_id', 'execution_snapshot_sha256', }: raise RunnerProtocolError('runner outcome identity shape is invalid') if type(identity.get('slot_id')) is not int or identity['slot_id'] < 0 or type(identity.get('reservation_id')) is not int or identity['reservation_id'] <= 0: raise RunnerProtocolError('runner outcome identity is invalid') if not re.fullmatch(r'[a-f0-9]{32,64}', str(identity.get('bundle_id') or '')) or not re.fullmatch(r'[a-f0-9]{32,64}', str(identity.get('scan_event_id') or '')) or _DIGEST_RE.fullmatch(str(identity.get('execution_snapshot_sha256') or '')) is None: raise RunnerProtocolError('runner outcome immutable identity is invalid') if value.get('status') not in {'succeeded', 'failed'}: raise RunnerProtocolError('runner outcome status is invalid') _timestamp(value.get('completed_at'), 'outcome timestamp') if type(value.get('last_event_sequence')) is not int or value['last_event_sequence'] < 0: raise RunnerProtocolError('runner outcome event sequence is invalid') if value.get('final_phase') is not None: try: WorkerPhase(value['final_phase']) except ValueError as exc: raise RunnerProtocolError('runner outcome final phase is invalid') from exc value = dict(value) value['phase_durations'] = _duration_map(value.get('phase_durations')) if value['status'] == 'succeeded': bundle = value.get('bundle') if not isinstance(bundle, dict) or set(bundle) != { 'commit', 'payload_sha256', 'ready_relative_path', } or value.get('error') is not None: raise RunnerProtocolError('successful runner outcome is invalid') commit = bundle.get('commit') if ( not isinstance(commit, dict) or set(commit) != _COMMIT_FIELDS or _DIGEST_RE.fullmatch(str(bundle.get('payload_sha256') or '')) is None ): raise RunnerProtocolError('runner bundle outcome is invalid') if ( type(commit.get('reservation_id')) is not int or commit['reservation_id'] <= 0 or not re.fullmatch(r'[a-f0-9]{32,64}', str(commit.get('bundle_id') or '')) or not re.fullmatch(r'[a-f0-9]{32,64}', str(commit.get('scan_event_id') or '')) or _DIGEST_RE.fullmatch(str(commit.get('scan_event_hash') or '')) is None or any( type(commit.get(name)) is not int or commit[name] < 0 for name in ( 'actual_bytes', 'frame_count', 'finding_count', 'error_count', 'candidate_count', ) ) or type(commit.get('source_failure')) is not bool or type(commit.get('source_failure_auth_related')) is not bool or any( type(commit.get(name)) is not str for name in ( 'target', 'relative_path', 'queue_status', 'source_failure_category', 'first_error', ) ) ): raise RunnerProtocolError('runner bundle commit is invalid') relative = str(bundle.get('ready_relative_path') or '') if not relative.startswith('ready/') or '\\' in relative or '..' in relative.split('/'): raise RunnerProtocolError('runner bundle reference is invalid') else: error = value.get('error') if value.get('bundle') is not None or not isinstance(error, dict) or set(error) != { 'code', 'type', 'summary', }: raise RunnerProtocolError('failed runner outcome is invalid') if not re.fullmatch(r'[a-z0-9_]{1,64}', str(error.get('code') or '')) or any(type(error.get(name)) is not str for name in ('type', 'summary')): raise RunnerProtocolError('runner error outcome is invalid') error['summary'] = error['summary'][:1000] return value def validate_generation_terminal(value, *, generation=None, input_sha256=None): if not isinstance(value, dict) or set(value) != { 'schema', 'generation', 'input_sha256', 'decision', 'decided_at', 'reason', 'outcome', } or value.get('schema') != RUNNER_PROTOCOL_SCHEMA: raise RunnerProtocolError('runner terminal-generation record shape is invalid') if _GENERATION_RE.fullmatch(str(value.get('generation') or '')) is None: raise RunnerProtocolError('runner terminal generation is invalid') if generation is not None and value['generation'] != generation: raise RunnerProtocolError('runner terminal generation conflicts with slot authority') if _DIGEST_RE.fullmatch(str(value.get('input_sha256') or '')) is None: raise RunnerProtocolError('runner terminal input hash is invalid') if input_sha256 is not None and value['input_sha256'] != input_sha256: raise RunnerProtocolError('runner terminal input hash conflicts with slot authority') _timestamp(value.get('decided_at'), 'terminal decision timestamp') decision = value.get('decision') if decision not in {'completed', 'fenced', 'timed_out'}: raise RunnerProtocolError('runner terminal decision is invalid') if decision == 'completed': if value.get('reason') is not None: raise RunnerProtocolError('completed runner terminal reason is invalid') value = dict(value) value['outcome'] = validate_runner_outcome( value.get('outcome'), generation=value['generation'], input_sha256=value['input_sha256'], ) elif ( value.get('outcome') is not None or type(value.get('reason')) is not str or not re.fullmatch(r'[a-z0-9_]{1,64}', value['reason']) ): raise RunnerProtocolError('controller runner terminal decision is invalid') return dict(value) def load_generation_terminal(path, *, generation, input_sha256): value, digest = _load_canonical_object( path, MAX_RUNNER_OUTCOME_BYTES, 'terminal-generation record', ) return validate_generation_terminal( value, generation=generation, input_sha256=input_sha256, ), digest def publish_generation_terminal( path, *, generation, input_sha256, decision, reason=None, outcome=None, ): value = validate_generation_terminal({ 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': str(generation), 'input_sha256': str(input_sha256), 'decision': str(decision), 'decided_at': utc_now(), 'reason': reason, 'outcome': outcome, }, generation=generation, input_sha256=input_sha256) try: write_private_json_exclusive( path, value, max_bytes=MAX_RUNNER_OUTCOME_BYTES, ) return value, True except FileExistsError: existing, _digest = load_generation_terminal( path, generation=generation, input_sha256=input_sha256, ) return existing, False def publish_start_gate(path, *, generation, input_sha256, host, payload): fields = ('pid', 'creation_time', 'executable') value = { 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': str(generation), 'input_sha256': str(input_sha256), 'host': {field: host[field] for field in fields}, 'payload': {field: payload[field] for field in fields}, 'released_at': utc_now(), } write_private_json_exclusive(path, value) return value def wait_for_start_gate(paths, runner_input, input_sha256): deadline = _timestamp( runner_input['watchdog_deadline_at'], 'watchdog deadline timestamp', ) fields = {'pid', 'creation_time', 'executable'} current = serialize_process_identity(current_process_identity()) while True: if os.path.exists(paths['terminal']): terminal, _digest = load_generation_terminal( paths['terminal'], generation=runner_input['generation'], input_sha256=input_sha256, ) raise RunnerFencedError( f"runner startup was closed by {terminal['decision']}" ) if os.path.exists(paths['start']): value, _digest = _load_canonical_object( paths['start'], MAX_RUNNER_EVENT_BYTES, 'startup gate', ) if not isinstance(value, dict) or set(value) != { 'schema', 'generation', 'input_sha256', 'host', 'payload', 'released_at', } or value.get('schema') != RUNNER_PROTOCOL_SCHEMA: raise RunnerProtocolError('runner startup gate shape is invalid') if ( value.get('generation') != runner_input['generation'] or value.get('input_sha256') != input_sha256 or not isinstance(value.get('host'), dict) or set(value['host']) != fields or not isinstance(value.get('payload'), dict) or set(value['payload']) != fields or any(value['payload'][field] != current[field] for field in fields) ): raise RunnerProtocolError('runner startup gate identity is invalid') _timestamp(value.get('released_at'), 'startup gate timestamp') return value if datetime.now(timezone.utc) >= deadline: raise TimeoutError('runner startup gate was not released before its deadline') time.sleep(0.01) def adopt_runner_bundle(work_root, root_name, outcome, assignment, bundle_root): outcome = validate_runner_outcome(outcome) if outcome['status'] != 'succeeded': raise RunnerProtocolError('failed runner outcome has no adoptable bundle') paths = runner_paths(work_root, root_name) reservation = BundleReservation.from_mapping(assignment['reservation']) identity = outcome['identity'] root_match = _ROOT_RE.fullmatch(root_name) if outcome['generation'] != root_match.group(3): raise RunnerProtocolError('runner outcome generation conflicts with isolated root') if identity != { 'slot_id': int(root_match.group(1)), 'reservation_id': reservation.reservation_id, 'bundle_id': reservation.bundle_id, 'scan_event_id': reservation.scan_event_id, 'execution_snapshot_sha256': assignment['execution_snapshot_sha256'], }: raise RunnerProtocolError('runner outcome identity conflicts with assignment') relative = outcome['bundle']['ready_relative_path'].replace('/', os.sep) source = os.path.abspath(os.path.join(paths['bundle_root'], relative)) if os.path.commonpath((os.path.abspath(paths['bundle_root']), source)) != os.path.abspath(paths['bundle_root']): raise RunnerProtocolError('runner bundle reference escapes its isolated root') metadata = ResultBundleReader( source, max_event_bytes=reservation.declared_bytes, ).validate() if metadata.header != reservation.header(): raise RunnerProtocolError('runner bundle header conflicts with assignment') commit = outcome['bundle']['commit'] for name in ( 'reservation_id', 'bundle_id', 'scan_event_id', 'scan_event_hash', 'actual_bytes', 'frame_count', 'finding_count', 'error_count', 'candidate_count', ): if getattr(metadata, name) != commit.get(name): raise RunnerProtocolError('runner bundle commit conflicts with durable bundle') digest = hashlib.sha256() with open(source, 'rb', buffering=0) as handle: for block in iter(lambda: handle.read(1024 * 1024), b''): digest.update(block) if digest.hexdigest() != outcome['bundle']['payload_sha256']: raise RunnerProtocolError('runner bundle payload hash is invalid') ensure_bundle_reservation_paths(bundle_root, reservation) destination = bundle_ready_path(bundle_root, reservation.bundle_id) if os.path.lexists(destination): adopted = ResultBundleReader( destination, max_event_bytes=reservation.declared_bytes, ).validate() if adopted.header != reservation.header(): raise RunnerProtocolError('canonical ready bundle conflicts with runner outcome') return adopted generation = outcome['generation'] temporary = os.path.join( os.path.dirname(os.path.dirname(destination)), '..', 'tmp', reservation.bundle_id[:2], f'{reservation.bundle_id}.{generation}.adopt.partial', ) temporary = os.path.normpath(temporary) descriptor = None try: if not os.path.lexists(temporary): descriptor = os.open( temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) harden_private_file(temporary) with os.fdopen(descriptor, 'wb', buffering=0) as target: descriptor = None with open(source, 'rb', buffering=0) as source_handle: for block in iter(lambda: source_handle.read(1024 * 1024), b''): target.write(block) target.flush() os.fsync(target.fileno()) copied = ResultBundleReader( temporary, max_event_bytes=reservation.declared_bytes, ).validate() if copied.header != reservation.header() or copied.scan_event_hash != metadata.scan_event_hash: raise RunnerProtocolError('copied runner bundle failed canonical validation') durable_publish(temporary, destination) finally: if descriptor is not None: os.close(descriptor) if os.path.exists(temporary): os.remove(temporary) return ResultBundleReader( destination, max_event_bytes=reservation.declared_bytes, ).validate() def _check_terminal(path, generation, input_sha256): if not os.path.exists(path): return value, _digest = load_generation_terminal( path, generation=generation, input_sha256=input_sha256, ) raise RunnerFencedError( f"runner generation was closed by {value['decision']}" ) class RunnerEventJournal: def __init__(self, path, terminal_path, generation, input_sha256, fault=None): self.path = path self.terminal_path = terminal_path self.generation = generation self.input_sha256 = input_sha256 self.fault = fault self.sequence = 0 self.phase = None self.phase_started_at = None self.phase_started_monotonic = None self.durations = {} self.events = [] def check_fence(self): _check_terminal(self.terminal_path, self.generation, self.input_sha256) def current_timestamp(self): now = utc_now() if ( self.events and _timestamp(now, 'event timestamp') < _timestamp(self.events[-1]['timestamp'], 'event timestamp') ): return self.events[-1]['timestamp'] return now def emit(self, phase, progress=None): self.check_fence() phase = WorkerPhase(phase) if self.phase is not None: validate_phase_transition(self.phase, phase) now = self.current_timestamp() monotonic_now = time.monotonic() measured = dict(progress or {}) if self.phase is not None and phase != self.phase: duration = max(0.0, monotonic_now - self.phase_started_monotonic) self.durations[self.phase.value] = self.durations.get(self.phase.value, 0.0) + duration measured.update({ 'previous_phase': self.phase.value, 'previous_duration_seconds': round(duration, 6), }) if phase != self.phase: self.phase_started_at = now self.phase_started_monotonic = monotonic_now event = validate_runner_event({ 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': self.generation, 'input_sha256': self.input_sha256, 'sequence': self.sequence + 1, 'timestamp': now, 'phase': phase.value, 'phase_started_at': self.phase_started_at, 'progress': measured, }, generation=self.generation, input_sha256=self.input_sha256) payload = canonical_json_bytes(event, newline=True) if len(payload) > MAX_RUNNER_EVENT_BYTES: raise RunnerProtocolError('runner event exceeds its byte bound') flags = os.O_WRONLY | os.O_APPEND | os.O_CREAT | getattr(os, 'O_BINARY', 0) created = not os.path.lexists(self.path) descriptor = os.open(self.path, flags, 0o600) try: harden_private_file(self.path) remaining = memoryview(payload) while remaining: written = os.write(descriptor, remaining) if written <= 0: raise OSError('runner event append made no progress') remaining = remaining[written:] os.fsync(descriptor) finally: os.close(descriptor) if created: fsync_directory(os.path.dirname(self.path)) self.sequence = event['sequence'] self.phase = phase self.events.append(event) if self.fault is not None: self.fault(phase.value) return event def finish_durations(self, completed_at): return phase_durations_from_events(self.events, completed_at) def _outcome_identity(runner_input): assignment = runner_input['assignment'] reservation = assignment['reservation'] return { 'slot_id': runner_input['slot_id'], 'reservation_id': int(reservation['reservation_id']), 'bundle_id': str(reservation['bundle_id']), 'scan_event_id': str(reservation['scan_event_id']), 'execution_snapshot_sha256': str(assignment['execution_snapshot_sha256']), } def run_assignment(root, package_runtime, *, fault=None, bundle_fault=None): root = require_private_directory(os.path.abspath(root), create=False) root_name = os.path.basename(root) work_root = os.path.dirname(root) paths = runner_paths(work_root, root_name) runner_input, input_sha256 = load_runner_input(paths['input']) reservation_id = int(runner_input['assignment']['reservation']['reservation_id']) if root_name != runner_root_name( runner_input['slot_id'], reservation_id, runner_input['generation'], ): raise RunnerProtocolError('runner root identity conflicts with input') journal = RunnerEventJournal( paths['events'], paths['terminal'], runner_input['generation'], input_sha256, fault=fault, ) identity = _outcome_identity(runner_input) outcome = None wait_for_start_gate(paths, runner_input, input_sha256) try: journal.emit( 'preparing' if runner_input['operation'] == 'execute' else 'bundling', {'operation': runner_input['operation']}, ) if set(package_runtime) != { 'build_compatibility', 'code_manifest', 'code_manifest_sha256', 'trufflehog_path', 'git_path', 'detector_policy_path', 'capabilities', }: raise RunnerProtocolError('verified runner package runtime is invalid') local_build = WorkerBuildCompatibility.from_mapping( package_runtime['build_compatibility'], ) validated = validate_protocol2_remote_assignment( runner_input['assignment'], local_build, package_runtime['capabilities'], ) scan_kwargs = dict(runner_input['assignment']['scan_kwargs']) if scan_kwargs.get('trufflehog_config') != PACKAGE_DETECTOR_POLICY: raise RunnerProtocolError('runner detector policy reference is invalid') def phase(phase_value, progress=None): journal.emit(phase_value, progress) def publication_fault(stage, writer): journal.check_fence() if bundle_fault is not None: bundle_fault(stage, writer) limits = dict(runner_input['assignment']['limits']) if runner_input['operation'] == 'timeout_bundle': phase('bundling', { 'reason': 'scan_stage_timeout', 'timed_out_phase': runner_input['timeout_phase'], }) reservation = validated['reservation'] started = _timestamp(runner_input['scan_started_at'], 'scan start timestamp') result = { 'findings': [], 'errors': [ f"Scan-stage deadline exceeded during {runner_input['timeout_phase']}" ], 'target': reservation.target, 'scan_type': reservation.platform, 'scan_event_id': reservation.scan_event_id, 'scan_started_at': runner_input['scan_started_at'], 'duration_sec': max( 0.0, (datetime.now(timezone.utc) - started).total_seconds(), ), 'timestamp': datetime.now(timezone.utc).isoformat(), 'error_class': 'timeout', 'retryable': True, 'scan_meta': { 'command_timed_out': True, 'full_stage_timeout': True, 'timed_out_phase': runner_input['timeout_phase'], 'scan_deadline_at': runner_input['scan_deadline_at'], }, } staged = stage_scan_result_in_scope( result, reservation, paths['bundle_root'], runner_input['assignment']['event_scan_options'], runner_input['assignment']['queue_policy'], attempts=int(runner_input['assignment']['reservation'].get('attempts') or 0), candidate_max_items=int(limits.get('candidate_max_items') or 2000), candidate_max_bytes=int(limits.get('candidate_max_bytes') or 2 * 1024 * 1024), require_s_drive=False, fault=publication_fault, diagnostic_slot_id=runner_input['slot_id'], ) else: remaining = ( _timestamp(runner_input['scan_deadline_at'], 'scan deadline timestamp') - datetime.now(timezone.utc) ).total_seconds() if remaining < 1: raise RunnerDeadlineElapsed( 'scan-stage deadline has less than one executable second remaining', ) scan_kwargs['trufflehog_config'] = package_runtime['detector_policy_path'] scanner.scan_config.trufflehog_path = package_runtime['trufflehog_path'] scanner.scan_config.trufflehog_config = package_runtime['detector_policy_path'] scanner.scan_config.work_dir = paths['scanner_work'] scanner.initialize_scanner_runtime(preflight_complete=True, register_cleanup=False) with scanner.client_scan_launch_authority( package_runtime['code_manifest'], expected_sha256=package_runtime['code_manifest_sha256'], ): staged = execute_protocol2_remote_claim( validated, paths['bundle_root'], scan_kwargs, runner_input['assignment']['event_scan_options'], runner_input['assignment']['queue_policy'], runner_input['assignment']['scan_policy'], attempts=int(runner_input['assignment']['reservation'].get('attempts') or 0), candidate_max_items=int(limits.get('candidate_max_items') or 2000), candidate_max_bytes=int(limits.get('candidate_max_bytes') or 2 * 1024 * 1024), require_s_drive=False, phase_callback=phase, bundle_fault=publication_fault, diagnostic_slot_id=runner_input['slot_id'], ) journal.check_fence() ready = bundle_ready_path( paths['bundle_root'], runner_input['assignment']['reservation']['bundle_id'], ) metadata = ResultBundleReader( ready, max_event_bytes=int(runner_input['assignment']['reservation']['declared_bundle_bytes']), ).validate() staged_value = staged.as_dict() for name in ( 'reservation_id', 'bundle_id', 'scan_event_id', 'scan_event_hash', 'actual_bytes', 'frame_count', 'finding_count', 'error_count', 'candidate_count', ): if getattr(metadata, name) != staged_value[name]: raise RunnerProtocolError('runner staged bundle commit is inconsistent') digest = hashlib.sha256() with open(ready, 'rb', buffering=0) as handle: for block in iter(lambda: handle.read(1024 * 1024), b''): digest.update(block) relative = os.path.relpath(ready, paths['bundle_root']).replace(os.sep, '/') completed_at = journal.current_timestamp() outcome = { 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': runner_input['generation'], 'input_sha256': input_sha256, 'identity': identity, 'status': 'succeeded', 'completed_at': completed_at, 'last_event_sequence': journal.sequence, 'final_phase': journal.phase.value if journal.phase is not None else None, 'phase_durations': journal.finish_durations(completed_at), 'bundle': { 'commit': staged.as_dict(), 'payload_sha256': digest.hexdigest(), 'ready_relative_path': relative, }, 'error': None, } except RunnerDeadlineElapsed: publish_generation_terminal( paths['terminal'], generation=runner_input['generation'], input_sha256=input_sha256, decision='timed_out', reason='scan_stage_deadline', ) return 1 except BaseException as exc: code = 'runner_fenced' if isinstance(exc, RunnerFencedError) else 'runner_timeout' if isinstance(exc, TimeoutError) else 'runner_failed' completed_at = journal.current_timestamp() outcome = { 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': runner_input['generation'], 'input_sha256': input_sha256, 'identity': identity, 'status': 'failed', 'completed_at': completed_at, 'last_event_sequence': journal.sequence, 'final_phase': journal.phase.value if journal.phase is not None else None, 'phase_durations': journal.finish_durations(completed_at), 'bundle': None, 'error': { 'code': code, 'type': f'{type(exc).__module__}.{type(exc).__qualname__}', 'summary': str(exc)[:1000], }, } outcome = validate_runner_outcome( outcome, generation=runner_input['generation'], input_sha256=input_sha256, ) terminal = validate_terminal_against_journal({ 'schema': RUNNER_PROTOCOL_SCHEMA, 'generation': runner_input['generation'], 'input_sha256': input_sha256, 'decision': 'completed', 'decided_at': outcome['completed_at'], 'reason': None, 'outcome': outcome, }, journal.events, runner_input) decided, won = publish_generation_terminal( paths['terminal'], generation=runner_input['generation'], input_sha256=input_sha256, decision='completed', outcome=outcome, ) if not won or decided['decision'] != 'completed': return 1 return 0 if terminal['outcome']['status'] == 'succeeded' else 1 def cleanup_abandoned_runner_roots( work_root, *, minimum_age_sec=60, budget=None, active_root_names=(), ): work_root = require_private_directory(os.path.abspath(work_root), create=False) allowed = {canonical_path(os.path.abspath(os.sys.executable))} if getattr(os.sys, '_base_executable', None): allowed.add(canonical_path(os.path.abspath(os.sys._base_executable))) return run_janitor_pass( work_root, allowed, minimum_age_sec=minimum_age_sec, excluded_relative_paths=tuple(active_root_names), budget=budget or JanitorBudget( max_candidates=8, max_entries=4000, max_bytes=512 * 1024 * 1024, max_seconds=2.0, max_depth=64, max_enumerated=256, ), ) def cleanup_runner_root(work_root, root_name): paths = runner_paths(work_root, root_name) if not os.path.exists(paths['root']): return True return bounded_remove_tree(paths['root'], JanitorBudget( max_candidates=1, max_entries=4000, max_bytes=512 * 1024 * 1024, max_seconds=2.0, max_depth=64, )) def transfer_runner_to_janitor(work_root, root_name, generation): paths = runner_paths(work_root, root_name) abandoned = ensure_private_directory( os.path.join(work_root, 'abandoned'), reject_reparse=True, ) destination = os.path.join(abandoned, root_name) destination_relative = f'abandoned/{root_name}' if not os.path.exists(paths['root']): if os.path.isdir(destination): marker = read_private_json( os.path.join(destination, '.scanner-owner.json'), max_bytes=64 * 1024, ) intent_path = os.path.join(destination, '.janitor-transfer.json') intent = read_private_json(intent_path, max_bytes=64 * 1024) if ( not isinstance(marker, dict) or marker.get('schema') != 2 or marker.get('root_kind') != 'work' or marker.get('relative_path') != destination_relative or not isinstance(intent, dict) or set(intent) != { 'schema', 'generation', 'source_relative_path', 'destination_relative_path', 'status', 'prepared_at', } or intent.get('schema') != 1 or intent.get('generation') != str(generation) or intent.get('source_relative_path') != root_name or intent.get('destination_relative_path') != destination_relative or intent.get('status') not in {'prepared', 'transferred'} ): raise RunnerProtocolError('janitor runner transfer evidence is invalid') if intent['status'] == 'prepared': intent['status'] = 'transferred' atomic_write_private_json(intent_path, intent) return destination_relative raise RunnerProtocolError('runner work tree is unavailable for janitor transfer') if os.path.lexists(destination): raise RunnerProtocolError('janitor runner destination already exists') marker_path = os.path.join(paths['root'], '.scanner-owner.json') marker = read_private_json(marker_path, max_bytes=64 * 1024) if ( not isinstance(marker, dict) or marker.get('schema') != 2 or marker.get('root_kind') != 'work' or marker.get('relative_path') not in {root_name, destination_relative} ): raise RunnerProtocolError('runner owner marker is invalid for janitor transfer') marker = dict(marker) marker['relative_path'] = destination_relative atomic_write_private_json(marker_path, marker) intent_path = os.path.join(paths['root'], '.janitor-transfer.json') intent = { 'schema': 1, 'generation': str(generation), 'source_relative_path': root_name, 'destination_relative_path': destination_relative, 'status': 'prepared', 'prepared_at': utc_now(), } atomic_write_private_json(intent_path, intent) durable_publish_directory(paths['root'], destination) intent['status'] = 'transferred' atomic_write_private_json( os.path.join(destination, '.janitor-transfer.json'), intent, ) return destination_relative