import hashlib import json import os import re import struct from dataclasses import dataclass from runtime_security import ( PrivatePathState, durable_publish, ensure_private_directory, harden_private_file, inspect_private_relative_path, private_file_ready, reject_reparse_components, require_private_directory, ) from worker_contracts import ( AssignmentOutcome, MAX_DIAGNOSTIC_AGGREGATE_BYTES, MAX_DIAGNOSTICS_PER_ASSIGNMENT, decode_diagnostic_envelope, build_legacy_error_frame_diagnostics, encode_diagnostic_envelope, ) MAGIC = b'TRUF-RB2\n' FORMAT_VERSION = 2 FRAME_HEADER = struct.Struct('!cI') FRAME_TYPES = frozenset((b'H', b'F', b'E', b'D', b'K', b'M', b'C')) FRAME_ORDER = {name: index for index, name in enumerate((b'H', b'F', b'E', b'D', b'K', b'M', b'C'))} ID_RE = re.compile(r'^[a-f0-9]{32,64}$') DEFAULT_MAX_EVENT_BYTES = 64 * 1024 * 1024 DEFAULT_MAX_FRAME_BYTES = 16 * 1024 * 1024 FRAME_BOUNDS = { b'H': 1024 * 1024, b'F': 16 * 1024 * 1024, b'E': 1024 * 1024, b'D': 64 * 1024, b'K': 2 * 1024 * 1024, b'M': 16 * 1024 * 1024, b'C': 1024 * 1024, } MAX_FINDING_FRAMES = 20000 MAX_ERROR_FRAMES = 2000 MAX_CANDIDATE_FRAMES = 2000 MAX_TOTAL_FRAMES = 24003 + MAX_DIAGNOSTICS_PER_ASSIGNMENT class ResultBundleError(ValueError): pass class ResultBundleConflictError(ResultBundleError): pass class ResultBundleUnavailableError(OSError): pass def canonical_json_bytes(value): try: return json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('utf-8') except (TypeError, ValueError) as exc: raise ResultBundleError('bundle frame is not canonical JSON data') from exc def _validated_id(value, name): text = str(value or '').lower() if not ID_RE.fullmatch(text): raise ResultBundleError(f'invalid {name}') return text def _relative_ready_path(bundle_id): bundle_id = _validated_id(bundle_id, 'bundle_id') return os.path.join('ready', bundle_id[:2], f'{bundle_id}.trb') def bundle_ready_path(root, bundle_id): return os.path.join(os.path.abspath(root), _relative_ready_path(bundle_id)) def bundle_partial_relative_path(bundle_id, reservation_token): bundle_id = _validated_id(bundle_id, 'bundle_id') producer_token = hashlib.sha256(str(reservation_token).encode('utf-8')).hexdigest()[:24] return os.path.join('tmp', bundle_id[:2], f'{bundle_id}.{producer_token}.partial') def bundle_partial_path(root, bundle_id, reservation_token): return os.path.join( os.path.abspath(root), bundle_partial_relative_path(bundle_id, reservation_token), ) def ensure_bundle_reservation_paths(root, reservation): root = require_private_directory(os.path.abspath(root), create=False) reservation = ( reservation if isinstance(reservation, BundleReservation) else BundleReservation.from_mapping(reservation) ) for name in ('tmp', 'ready', 'quarantine'): ensure_private_directory( os.path.join(root, name, reservation.bundle_id[:2]), reject_reparse=True, ) return reservation @dataclass(frozen=True) class BundleReservation: reservation_id: int reservation_token: str bundle_id: str scan_event_id: str queue_id: int claim_lease_token: str declared_bytes: int ready_path: str source: str = '' platform: str = '' query: str = '' target: str = '' normalized_target: str = '' run_id: int | None = None cycle_id: int | None = None producer_instance_id: str = '' producer_pid: int = 0 producer_creation_time: str = '' producer_executable: str = '' @classmethod def from_mapping(cls, value): data = dict(value or {}) ready_path = data.get('ready_path') or data.get('ready_relative_path') or '' return cls( reservation_id=int(data['reservation_id'] if 'reservation_id' in data else data['id']), reservation_token=str(data['reservation_token']), bundle_id=_validated_id(data['bundle_id'], 'bundle_id'), scan_event_id=_validated_id(data['scan_event_id'], 'scan_event_id'), queue_id=int(data['queue_id']), claim_lease_token=str(data.get('claim_lease_token') or data.get('claim_lease_token_value') or ''), declared_bytes=int(data.get('declared_bytes') or data.get('declared_bundle_bytes') or 0), ready_path=str(ready_path), source=str(data.get('source') or ''), platform=str(data.get('platform') or ''), query=str(data.get('query') or ''), target=str(data.get('target') or ''), normalized_target=str(data.get('normalized_target') or ''), run_id=data.get('run_id'), cycle_id=data.get('cycle_id'), producer_instance_id=str(data.get('producer_instance_id') or ''), producer_pid=int(data.get('producer_pid') or 0), producer_creation_time=str(data.get('producer_creation_time') or ''), producer_executable=str(data.get('producer_executable') or ''), ) def header(self): return { 'format_version': FORMAT_VERSION, 'reservation_id': self.reservation_id, 'reservation_token': self.reservation_token, 'bundle_id': self.bundle_id, 'scan_event_id': self.scan_event_id, 'queue_id': self.queue_id, 'claim_lease_token': self.claim_lease_token, 'declared_bytes': self.declared_bytes, 'ready_relative_path': self.ready_path.replace('\\', '/'), 'source': self.source, 'platform': self.platform, 'query': self.query, 'target': self.target, 'normalized_target': self.normalized_target, 'run_id': self.run_id, 'cycle_id': self.cycle_id, 'producer_instance_id': self.producer_instance_id, 'producer_pid': self.producer_pid, 'producer_creation_time': self.producer_creation_time, 'producer_executable': self.producer_executable, } @dataclass(frozen=True) class BundleCommit: reservation_id: int bundle_id: str scan_event_id: str scan_event_hash: str relative_path: str actual_bytes: int frame_count: int finding_count: int error_count: int candidate_count: int def as_dict(self): return dict(self.__dict__) @dataclass(frozen=True) class BundleMetadata(BundleCommit): header: dict result_metadata: dict diagnostic_count: int diagnostic_bytes: int class ResultBundleWriter: def __init__(self, root, reservation, handle, partial_path, ready_path, fault=None): self.root = root self.reservation = reservation self.handle = handle self.partial_path = partial_path self.ready_path = ready_path self.fault = fault self.digest = hashlib.sha256() self.bytes_written = 0 self.frame_count = 0 self.finding_count = 0 self.error_count = 0 self.candidate_count = 0 self.diagnostic_count = 0 self.diagnostic_bytes = 0 self.diagnostic_aggregate_bytes = 0 self._diagnostic_uids = set() self._last_frame_order = FRAME_ORDER[b'H'] self.finished = False self.ready_published = False self._write_bytes(MAGIC) self._write_frame(b'H', reservation.header()) @classmethod def open(cls, root, reservation, fault=None, require_s_drive=False): root = require_private_directory(os.path.abspath(root), create=False) drive = os.path.splitdrive(root)[0].upper() if require_s_drive and drive != 'S:': raise ResultBundleError('production result bundle root must be on S:') reservation = ensure_bundle_reservation_paths(root, reservation) if reservation.declared_bytes <= len(MAGIC) or reservation.declared_bytes > DEFAULT_MAX_EVENT_BYTES: raise ResultBundleError('declared bundle byte bound is invalid') expected_relative = _relative_ready_path(reservation.bundle_id) supplied_relative = str(reservation.ready_path or expected_relative).replace('/', os.sep) if os.path.normcase(os.path.normpath(supplied_relative)) != os.path.normcase(os.path.normpath(expected_relative)): raise ResultBundleError('reservation ready path is not deterministic for its bundle ID') tmp_dir = os.path.join(root, 'tmp', reservation.bundle_id[:2]) ready_dir = os.path.join(root, 'ready', reservation.bundle_id[:2]) partial_path = bundle_partial_path( root, reservation.bundle_id, reservation.reservation_token, ) ready_path = os.path.join(ready_dir, f'{reservation.bundle_id}.trb') if os.path.lexists(ready_path): raise ResultBundleConflictError('deterministic ready bundle path already exists') descriptor = os.open( partial_path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) try: os.close(descriptor) descriptor = None harden_private_file(partial_path) handle = open(partial_path, 'w+b', buffering=0) return cls(root, reservation, handle, partial_path, ready_path, fault=fault) except BaseException: if descriptor is not None: os.close(descriptor) try: os.remove(partial_path) except OSError: pass raise def _inject(self, stage): if self.fault is not None: self.fault(stage, self) def _write_bytes(self, payload): if self.bytes_written + len(payload) > self.reservation.declared_bytes: raise ResultBundleError('bundle exceeded its pre-reserved byte bound') self.handle.write(payload) self.digest.update(payload) self.bytes_written += len(payload) def _write_frame(self, frame_type, value): if self.finished or frame_type not in FRAME_TYPES or frame_type == b'C': raise ResultBundleError('invalid bundle frame write') if self.frame_count >= MAX_TOTAL_FRAMES - 1: raise ResultBundleError('bundle frame count exceeds its bound') if FRAME_ORDER[frame_type] < self._last_frame_order: raise ResultBundleError('bundle frame order is not deterministic') payload = canonical_json_bytes(value) bound = min(FRAME_BOUNDS[frame_type], self.reservation.declared_bytes) if len(payload) > bound: raise ResultBundleError(f'{frame_type.decode()} frame exceeds its byte bound') framed = FRAME_HEADER.pack(frame_type, len(payload)) + payload self._inject(f'before_frame_{frame_type.decode()}') self._write_bytes(framed) self.frame_count += 1 self._last_frame_order = FRAME_ORDER[frame_type] self._inject(f'after_frame_{frame_type.decode()}') return len(payload) def write_finding(self, finding): if self.finding_count >= MAX_FINDING_FRAMES: raise ResultBundleError('bundle finding count exceeds its bound') self._write_frame(b'F', finding) self.finding_count += 1 def write_error(self, error): if self.error_count >= MAX_ERROR_FRAMES: raise ResultBundleError('bundle error count exceeds its bound') self._write_frame(b'E', {'error': str(error)}) self.error_count += 1 def write_diagnostic(self, diagnostic): if self.diagnostic_count >= MAX_DIAGNOSTICS_PER_ASSIGNMENT: raise ResultBundleError('bundle diagnostic count exceeds its bound') try: if isinstance(diagnostic, dict): envelope = decode_diagnostic_envelope(canonical_json_bytes(diagnostic)) else: envelope = decode_diagnostic_envelope( encode_diagnostic_envelope(diagnostic) ) payload = encode_diagnostic_envelope(envelope) value = json.loads(payload.decode('ascii')) except (TypeError, ValueError, UnicodeError) as exc: raise ResultBundleError('bundle diagnostic frame is invalid') from exc if envelope.diagnostic_uid in self._diagnostic_uids: raise ResultBundleError('bundle diagnostic identity is duplicated') if ( envelope.reservation_id != self.reservation.reservation_id or envelope.scan_event_id != self.reservation.scan_event_id or envelope.source != self.reservation.source ): raise ResultBundleError('bundle diagnostic identity conflicts with its reservation') if ( self.diagnostic_aggregate_bytes + len(payload) + 1 > MAX_DIAGNOSTIC_AGGREGATE_BYTES ): raise ResultBundleError('bundle diagnostic aggregate exceeds its byte bound') self._write_frame(b'D', value) self._diagnostic_uids.add(envelope.diagnostic_uid) self.diagnostic_count += 1 self.diagnostic_bytes += FRAME_HEADER.size + len(payload) self.diagnostic_aggregate_bytes += len(payload) + 1 def write_candidate(self, candidate): if self.candidate_count >= MAX_CANDIDATE_FRAMES: raise ResultBundleError('bundle candidate count exceeds its bound') self._write_frame(b'K', candidate) self.candidate_count += 1 def finish(self, metadata): if self.finished: raise ResultBundleError('bundle writer is already finished') self._write_frame(b'M', metadata) content_hash = self.digest.hexdigest() footer = { 'format_version': FORMAT_VERSION, 'content_sha256': content_hash, 'content_bytes': self.bytes_written, 'byte_count': 0, 'frame_count': self.frame_count + 1, 'finding_count': self.finding_count, 'error_count': self.error_count, 'candidate_count': self.candidate_count, } if self.diagnostic_count: footer.update({ 'diagnostic_count': self.diagnostic_count, 'diagnostic_bytes': self.diagnostic_bytes, }) while True: payload = canonical_json_bytes(footer) framed = FRAME_HEADER.pack(b'C', len(payload)) + payload total = self.bytes_written + len(framed) if footer['byte_count'] == total: break footer['byte_count'] = total if total > self.reservation.declared_bytes: raise ResultBundleError('bundle footer exceeds its pre-reserved byte bound') self._inject('before_footer') self.handle.write(framed) self.bytes_written = total self.frame_count += 1 self._inject('after_footer') self._inject('before_fsync') self.handle.flush() os.fsync(self.handle.fileno()) self._inject('after_fsync') self.handle.close() self.handle = None if not private_file_ready(self.partial_path): raise ResultBundleError('private bundle ACL verification failed before publication') if os.path.getsize(self.partial_path) != self.bytes_written: raise ResultBundleError('bundle size changed before publication') self._inject('before_rename') durable_publish(self.partial_path, self.ready_path) self.ready_published = True self._inject('after_rename') if not private_file_ready(self.ready_path): inspection = inspect_private_relative_path( self.root, os.path.relpath(self.ready_path, self.root), ) if inspection.state == PrivatePathState.UNKNOWN: raise ResultBundleUnavailableError( 'ready bundle state is unavailable after publication' ) if inspection.state == PrivatePathState.PRESENT: raise ResultBundleError('ready bundle ACL verification failed') self.finished = True relative = os.path.relpath(self.ready_path, self.root) return BundleCommit( reservation_id=self.reservation.reservation_id, bundle_id=self.reservation.bundle_id, scan_event_id=self.reservation.scan_event_id, scan_event_hash=content_hash, relative_path=relative.replace(os.sep, '/'), actual_bytes=self.bytes_written, frame_count=self.frame_count, finding_count=self.finding_count, error_count=self.error_count, candidate_count=self.candidate_count, ) def abort(self): if self.handle is not None: self.handle.close() self.handle = None if not self.ready_published: try: os.remove(self.partial_path) except FileNotFoundError: pass def __enter__(self): return self def __exit__(self, exc_type, value, traceback): if not self.finished: self.abort() class ResultBundleReader: def __init__(self, path, max_event_bytes=DEFAULT_MAX_EVENT_BYTES): self.path = os.path.abspath(path) self.max_event_bytes = max(1, int(max_event_bytes)) self._validated = None self._validated_fingerprint = None @classmethod def from_reservation(cls, root, reservation, max_event_bytes=DEFAULT_MAX_EVENT_BYTES): reservation = reservation if isinstance(reservation, BundleReservation) else BundleReservation.from_mapping(reservation) return cls(bundle_ready_path(root, reservation.bundle_id), max_event_bytes=max_event_bytes) @staticmethod def _stat_fingerprint(value): return ( int(value.st_dev), int(value.st_ino), int(value.st_mode), int(value.st_size), int(value.st_mtime_ns), ) def _file_fingerprint(self): reject_reparse_components(self.path) if not private_file_ready(self.path): raise ResultBundleUnavailableError('bundle path is not currently available as an exact private regular file') return self._stat_fingerprint(os.stat(self.path, follow_symlinks=False)) def _iter_frames(self, expected_fingerprint=None): fingerprint = self._file_fingerprint() if expected_fingerprint is not None and fingerprint != expected_fingerprint: raise ResultBundleError('bundle changed after validation') size = fingerprint[3] if size <= len(MAGIC) or size > self.max_event_bytes: raise ResultBundleError('bundle aggregate byte bound is invalid') with open(self.path, 'rb', buffering=0) as handle: opened_fingerprint = self._stat_fingerprint(os.fstat(handle.fileno())) if opened_fingerprint != fingerprint: raise ResultBundleUnavailableError('bundle identity changed while it was opened') magic = handle.read(len(MAGIC)) if magic != MAGIC: raise ResultBundleError('bundle magic/version mismatch') offset = len(MAGIC) while offset < size: header = handle.read(FRAME_HEADER.size) if len(header) != FRAME_HEADER.size: raise ResultBundleError('truncated bundle frame header') frame_type, payload_length = FRAME_HEADER.unpack(header) if frame_type not in FRAME_TYPES: raise ResultBundleError('unknown bundle frame type') bound = min(FRAME_BOUNDS[frame_type], self.max_event_bytes) if payload_length > bound or offset + FRAME_HEADER.size + payload_length > size: raise ResultBundleError('bundle frame length exceeds its bound') payload = handle.read(payload_length) if len(payload) != payload_length: raise ResultBundleError('truncated bundle frame payload') try: value = json.loads(payload.decode('utf-8', errors='strict')) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise ResultBundleError('bundle frame contains invalid UTF-8 JSON') from exc if canonical_json_bytes(value) != payload: raise ResultBundleError('bundle frame JSON is not canonical') offset += FRAME_HEADER.size + payload_length yield frame_type, value, header + payload, offset if offset != size: raise ResultBundleError('bundle byte count is inconsistent') if self._stat_fingerprint(os.fstat(handle.fileno())) != opened_fingerprint: raise ResultBundleError('bundle changed while it was read') if self._file_fingerprint() != fingerprint: raise ResultBundleError('bundle path changed while it was read') def validate(self): if self._validated is not None: return self._validated digest = hashlib.sha256(MAGIC) header_value = None metadata_value = None footer = None counts = {b'F': 0, b'E': 0, b'D': 0, b'K': 0} diagnostic_bytes = 0 diagnostic_aggregate_bytes = 0 diagnostic_uids = set() diagnostic_scan_outcomes = [] frame_count = 0 last_frame_order = -1 final_offset = len(MAGIC) content_bytes = None fingerprint = self._file_fingerprint() for frame_type, value, framed, offset in self._iter_frames(fingerprint): frame_count += 1 if frame_count > MAX_TOTAL_FRAMES: raise ResultBundleError('bundle frame count exceeds its bound') if FRAME_ORDER[frame_type] < last_frame_order: raise ResultBundleError('bundle frame order is not deterministic') last_frame_order = FRAME_ORDER[frame_type] final_offset = offset if footer is not None: raise ResultBundleError('commit footer is not the final frame') if metadata_value is not None and frame_type != b'C': raise ResultBundleError('result metadata is not immediately before the commit footer') if frame_count == 1 and frame_type != b'H': raise ResultBundleError('bundle header is not the first frame') if frame_type in FRAME_TYPES and not isinstance(value, dict): raise ResultBundleError('bundle typed frame must contain a JSON object') if frame_type == b'E' and not isinstance(value.get('error'), str): raise ResultBundleError('bundle error frame is invalid') if frame_type == b'H': if header_value is not None: raise ResultBundleError('bundle contains duplicate headers') header_value = value elif frame_type == b'M': if metadata_value is not None: raise ResultBundleError('bundle contains duplicate metadata') metadata_value = value elif frame_type == b'C': footer = value content_bytes = offset - len(framed) continue elif frame_type == b'D': try: envelope = decode_diagnostic_envelope(canonical_json_bytes(value)) except (TypeError, ValueError, UnicodeError) as exc: raise ResultBundleError('bundle diagnostic frame is invalid') from exc if envelope.diagnostic_uid in diagnostic_uids: raise ResultBundleError('bundle diagnostic identity is duplicated') if header_value is None or ( envelope.reservation_id != int(header_value.get('reservation_id') or 0) or envelope.scan_event_id != str(header_value.get('scan_event_id') or '') or envelope.source != str(header_value.get('source') or '') ): raise ResultBundleError('bundle diagnostic identity conflicts with its header') diagnostic_uids.add(envelope.diagnostic_uid) if envelope.assignment_outcome is not AssignmentOutcome.ACCEPTED: raise ResultBundleError( 'bundle diagnostic assignment outcome is invalid' ) diagnostic_scan_outcomes.append(envelope.scan_outcome.value) counts[b'D'] += 1 if counts[b'D'] > MAX_DIAGNOSTICS_PER_ASSIGNMENT: raise ResultBundleError('bundle typed frame count exceeds its bound') envelope_bytes = len(encode_diagnostic_envelope(envelope)) diagnostic_bytes += len(framed) diagnostic_aggregate_bytes += envelope_bytes + 1 if ( diagnostic_aggregate_bytes > MAX_DIAGNOSTIC_AGGREGATE_BYTES ): raise ResultBundleError('bundle diagnostic aggregate exceeds its byte bound') elif frame_type in counts: counts[frame_type] += 1 limit = { b'F': MAX_FINDING_FRAMES, b'E': MAX_ERROR_FRAMES, b'D': MAX_DIAGNOSTICS_PER_ASSIGNMENT, b'K': MAX_CANDIDATE_FRAMES, }[frame_type] if counts[frame_type] > limit: raise ResultBundleError('bundle typed frame count exceeds its bound') digest.update(framed) if not isinstance(header_value, dict) or not isinstance(metadata_value, dict) or not isinstance(footer, dict): raise ResultBundleError('bundle is missing required header, metadata, or footer') if frame_count < 3 or footer.get('format_version') != FORMAT_VERSION: raise ResultBundleError('bundle footer version is invalid') expected_scan_outcome = { 'clean': 'clean', 'found': 'found', 'degraded': 'degraded', 'error': 'error', 'skipped': 'skipped', }.get(str(metadata_value.get('status') or 'clean'), 'error') if any( outcome != expected_scan_outcome for outcome in diagnostic_scan_outcomes ): raise ResultBundleError( 'bundle diagnostic scan outcome conflicts with result metadata' ) base_footer_fields = { 'format_version', 'content_sha256', 'content_bytes', 'byte_count', 'frame_count', 'finding_count', 'error_count', 'candidate_count', } expected_footer_fields = ( base_footer_fields | {'diagnostic_count', 'diagnostic_bytes'} if counts[b'D'] else base_footer_fields ) if set(footer) != expected_footer_fields: raise ResultBundleError('bundle footer shape is invalid') expected = { 'content_sha256': digest.hexdigest(), 'content_bytes': content_bytes, 'byte_count': final_offset, 'frame_count': frame_count, 'finding_count': counts[b'F'], 'error_count': counts[b'E'], 'candidate_count': counts[b'K'], } if counts[b'D']: expected.update({ 'diagnostic_count': counts[b'D'], 'diagnostic_bytes': diagnostic_bytes, }) if footer.get('content_sha256') != expected['content_sha256']: raise ResultBundleError('bundle content hash mismatch') for key in ( 'content_bytes', 'byte_count', 'frame_count', 'finding_count', 'error_count', 'candidate_count', ): value = footer.get(key) if isinstance(value, bool) or not isinstance(value, int) or value != expected[key]: raise ResultBundleError(f'bundle footer {key} mismatch') diagnostic_footer_fields = {'diagnostic_count', 'diagnostic_bytes'} & set(footer) if counts[b'D']: if diagnostic_footer_fields != {'diagnostic_count', 'diagnostic_bytes'}: raise ResultBundleError('bundle footer diagnostic accounting is missing') for key in ('diagnostic_count', 'diagnostic_bytes'): value = footer.get(key) if isinstance(value, bool) or not isinstance(value, int) or value != expected[key]: raise ResultBundleError(f'bundle footer {key} mismatch') elif diagnostic_footer_fields: raise ResultBundleError('D-less bundle has unexpected diagnostic accounting') bundle_id = _validated_id(header_value.get('bundle_id'), 'bundle_id') event_id = _validated_id(header_value.get('scan_event_id'), 'scan_event_id') reservation_id = header_value.get('reservation_id') if isinstance(reservation_id, bool) or not isinstance(reservation_id, int) or reservation_id <= 0: raise ResultBundleError('invalid reservation_id') self._validated = BundleMetadata( reservation_id=reservation_id, bundle_id=bundle_id, scan_event_id=event_id, scan_event_hash=expected['content_sha256'], relative_path='', actual_bytes=final_offset, frame_count=frame_count, finding_count=counts[b'F'], error_count=counts[b'E'], candidate_count=counts[b'K'], header=header_value, result_metadata=metadata_value, diagnostic_count=counts[b'D'], diagnostic_bytes=diagnostic_bytes, ) self._validated_fingerprint = fingerprint return self._validated def _values(self, wanted): validated = self.validate() digest = hashlib.sha256(MAGIC) footer = None for frame_type, value, framed, _ in self._iter_frames( self._validated_fingerprint ): if frame_type == b'C': footer = value else: digest.update(framed) if frame_type == wanted: yield value if ( digest.hexdigest() != validated.scan_event_hash or not isinstance(footer, dict) or footer.get('content_sha256') != validated.scan_event_hash ): raise ResultBundleError('bundle content changed after validation') def iter_findings(self): return self._values(b'F') def iter_errors(self): for value in self._values(b'E'): yield value.get('error') if isinstance(value, dict) else value def iter_candidates(self): return self._values(b'K') def iter_diagnostics(self): for value in self._values(b'D'): envelope = decode_diagnostic_envelope(canonical_json_bytes(value)) yield json.loads(encode_diagnostic_envelope(envelope).decode('ascii')) def effective_diagnostics(self): validated = self.validate() if validated.diagnostic_count: return tuple(self.iter_diagnostics()) metadata = self.metadata() header = self.header() envelopes = build_legacy_error_frame_diagnostics( reservation_id=validated.reservation_id, scan_event_id=validated.scan_event_id, slot_id=0, source=str(header.get('source') or ''), timestamp=( metadata.get('timestamp') or metadata.get('scan_started_at') ), errors=tuple(self.iter_errors()), retryable=bool(metadata.get('retryable', False)), attempt=max(1, int(metadata.get('attempt') or 1)), ) return tuple( json.loads(encode_diagnostic_envelope(envelope).decode('ascii')) for envelope in envelopes ) def metadata(self): self.validate() if self._file_fingerprint() != self._validated_fingerprint: raise ResultBundleError('bundle changed after validation') return dict(self.validate().result_metadata) def header(self): self.validate() if self._file_fingerprint() != self._validated_fingerprint: raise ResultBundleError('bundle changed after validation') return dict(self.validate().header)