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

1167 lines
52 KiB
Python

"""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