1167 lines
52 KiB
Python
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
|