Files
2026-09-30 20:30:56 +03:00

758 lines
32 KiB
Python

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)