import sys sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('Docker shadow runner could not disable bytecode writes') import argparse import copy import hashlib import math import os import re import secrets import socket import time import scanner import console_runner from keycheck_candidates import extract_candidates from lifecycle_authority import require_active_supervisor_child from scanner_db import ( DOCKER_ADAPTIVE_GATE_MAX_CONTROLS, DOCKER_ADAPTIVE_GATE_MIN_CONTROLS, DOCKER_ADAPTIVE_LAYER_CLASS_ORDER, DOCKER_ADAPTIVE_SELECTOR_VERSION, DOCKER_ADAPTIVE_SHADOW_SELECTION_METRIC_KEYS, ScannerDB, canonical_docker_layer_plan_bytes, docker_content_media_class, docker_layer_execution_policy_sha256, docker_layer_selection_policy_sha256, finding_identity, select_docker_adaptive_payload, validate_docker_adaptive_checkpoint, validate_docker_adaptive_shadow_selection_metrics, validate_docker_layer_execution, validate_docker_layer_limits, validate_docker_layer_plan, validate_docker_layer_resolution, ) _SHA256_RE = re.compile(r'[a-f0-9]{64}') _ROUTED_SERVICE_RE = re.compile(r'[a-z0-9][a-z0-9_.-]{0,63}') SHADOW_FAILURE_METRIC_KEYS = ( 'diagnostic_full_incomplete', 'diagnostic_adaptive_incomplete', 'diagnostic_blob_chunk_processing', 'diagnostic_blob_detector_timeout', 'diagnostic_blob_network', 'diagnostic_blob_timeout', 'diagnostic_blob_mixed', 'diagnostic_blob_other', ) _SHADOW_BLOB_FAILURE_METRICS = { 'chunk_processing': 'diagnostic_blob_chunk_processing', 'detector_timeout': 'diagnostic_blob_detector_timeout', 'network': 'diagnostic_blob_network', 'transfer_timeout': 'diagnostic_blob_timeout', 'timeout': 'diagnostic_blob_timeout', 'mixed': 'diagnostic_blob_mixed', } class DockerShadowPrivacyError(RuntimeError): pass def empty_selection_metrics(): return {name: 0 for name in DOCKER_ADAPTIVE_SHADOW_SELECTION_METRIC_KEYS} def empty_failure_metrics(): return {name: 0 for name in SHADOW_FAILURE_METRIC_KEYS} def _blob_failure_metric(error_code): return _SHADOW_BLOB_FAILURE_METRICS.get( str(error_code or ''), 'diagnostic_blob_other', ) def _canonical_private_plan(plan): plan = validate_docker_layer_plan(plan) payload = canonical_docker_layer_plan_bytes(plan) return plan, hashlib.sha256(payload).hexdigest() def _validate_duplicate_descriptors(descriptors): identities = {} for descriptor in descriptors: identity = ( descriptor['kind'], descriptor['size'], docker_content_media_class(descriptor['kind'], descriptor['media_type']), ) previous = identities.setdefault(descriptor['digest'], identity) if previous != identity: raise ValueError('Docker shadow duplicate digest metadata conflicts') def build_private_adaptive_plan( resolved, payload_classes, limits, checkpoint, scan_policy_sha256, ): resolved = validate_docker_layer_resolution(resolved) limits = validate_docker_layer_limits(limits) checkpoint = validate_docker_adaptive_checkpoint(checkpoint) scan_policy_sha256 = str(scan_policy_sha256 or '').lower() if not _SHA256_RE.fullmatch(scan_policy_sha256): raise ValueError('Docker shadow scan policy must be a lowercase SHA-256') descriptors = [resolved['config'], *resolved['layers']] _validate_duplicate_descriptors(descriptors) decisions = select_docker_adaptive_payload( descriptors, payload_classes, limits, covered_digests=(), ) metrics = empty_selection_metrics() planned = [] omitted = 0 for descriptor, selected, reason, payload_class in decisions: if reason == 'config_selected': metrics['selected_config'] += 1 elif reason.startswith('selected_'): metrics[reason] += 1 elif reason == 'already_covered': metrics['reuse_already_covered'] += 1 elif reason == 'duplicate_digest': metrics['reuse_duplicate_digest'] += 1 elif not selected: metric_name = f'omitted_{reason}' if metric_name not in metrics: raise ValueError('Docker shadow selection reason is not aggregate-safe') metrics[metric_name] += 1 omitted += 1 planned.append({ **descriptor, 'payload_class': payload_class, 'selected': bool(selected), 'selection_reason': reason, 'coverage_state': 'selected' if selected else 'skipped', 'lease_token': None, 'attempt': 0, 'max_attempts': limits['blob_max_attempts'], }) plan, plan_sha256 = _canonical_private_plan({ 'version': 2, 'image': resolved['image'], 'repository': resolved['repository'], 'manifest_digest': resolved['manifest_digest'], 'platform_os': resolved['platform_os'], 'platform_arch': resolved['platform_arch'], 'manifest_media_type': resolved['manifest_media_type'], 'limits': limits, 'selector_version': DOCKER_ADAPTIVE_SELECTOR_VERSION, 'selection_policy_sha256': docker_layer_selection_policy_sha256(limits), 'scan_policy_sha256': scan_policy_sha256, 'execution_policy_sha256': docker_layer_execution_policy_sha256( scan_policy_sha256, limits, ), 'checkpoint': checkpoint, 'descriptors': planned, }) return { 'plan': plan, 'plan_sha256': plan_sha256, 'selection_metrics': validate_docker_adaptive_shadow_selection_metrics(metrics), 'omitted_descriptor_count': omitted, } def _checkpoint_order(descriptor): if descriptor['kind'] == 'config': return (0, 0, -descriptor['position'], descriptor['size'], descriptor['digest']) try: class_rank = DOCKER_ADAPTIVE_LAYER_CLASS_ORDER.index(descriptor['payload_class']) except ValueError as exc: raise ValueError('Docker shadow descriptor class is not schedulable') from exc return (1, class_rank, -descriptor['position'], descriptor['size'], descriptor['digest']) def lease_private_adaptive_checkpoint(plan, token_factory=None): plan = validate_docker_layer_plan(plan) if plan['version'] != 2: raise ValueError('Docker shadow checkpoints require a version-two plan') if any( descriptor['coverage_state'] in ('leased', 'shared_pending') for descriptor in plan['descriptors'] ): raise ValueError('Docker shadow plan already contains active leases') pending = {} for descriptor in plan['descriptors']: if descriptor['coverage_state'] == 'selected': pending.setdefault(descriptor['digest'], []).append(descriptor) if not pending: return None groups = [] for digest, descriptors in pending.items(): attempts = {item['attempt'] for item in descriptors} maximums = {item['max_attempts'] for item in descriptors} if len(attempts) != 1 or len(maximums) != 1: raise ValueError('Docker shadow duplicate attempts conflict') if next(iter(attempts)) >= next(iter(maximums)): raise ValueError('Docker shadow plan exceeded its private attempt budget') representative = min(descriptors, key=_checkpoint_order) groups.append((representative, digest)) groups.sort(key=lambda item: _checkpoint_order(item[0])) leased_digests = set() leased_bytes = 0 max_blobs = plan['checkpoint']['max_blobs'] max_bytes = plan['checkpoint']['max_bytes'] for descriptor, digest in groups: if len(leased_digests) >= max_blobs: continue if leased_digests and leased_bytes + descriptor['size'] > max_bytes: continue leased_digests.add(digest) leased_bytes += descriptor['size'] if not leased_digests: raise ValueError('Docker shadow checkpoint made no bounded progress') token_factory = token_factory or (lambda: secrets.token_urlsafe(32)) tokens = {digest: str(token_factory()) for digest in leased_digests} leased_plan = copy.deepcopy(plan) for descriptor in leased_plan['descriptors']: if descriptor['digest'] not in leased_digests: continue descriptor['coverage_state'] = 'leased' descriptor['lease_token'] = tokens[descriptor['digest']] descriptor['attempt'] += 1 leased_plan, plan_sha256 = _canonical_private_plan(leased_plan) return {'plan': leased_plan, 'plan_sha256': plan_sha256} def apply_private_adaptive_execution(plan, execution): plan, plan_sha256 = _canonical_private_plan(plan) execution = validate_docker_layer_execution(execution, plan, plan_sha256) records = {record['digest']: record for record in execution['blobs']} next_plan = copy.deepcopy(plan) for descriptor in next_plan['descriptors']: if descriptor['coverage_state'] != 'leased': continue status = records[descriptor['digest']]['status'] if status == 'covered': next_state = 'covered' elif status == 'retryable_failed' and descriptor['attempt'] < descriptor['max_attempts']: next_state = 'selected' else: next_state = 'terminal_failed' descriptor['coverage_state'] = next_state descriptor['lease_token'] = None return _canonical_private_plan(next_plan)[0] def private_result_identities(result, normalized_target): if not isinstance(result, dict): raise ValueError('Docker shadow private result must be an object') routed = set() detectors = set() try: scanner.strip_nearby_context_for_persistence(result) findings = result.get('findings') or [] if not isinstance(findings, list): raise ValueError('Docker shadow private findings must be a list') for finding in findings: if not isinstance(finding, dict): raise ValueError('Docker shadow private finding must be an object') detector_hash = str( finding_identity('dockerhub', normalized_target, finding)[2] or '' ).lower() if not _SHA256_RE.fullmatch(detector_hash): raise DockerShadowPrivacyError('Docker shadow detector identity is invalid') detectors.add(detector_hash) for candidate in extract_candidates(finding): service = str(candidate.service or '') provider_key_hash = str(candidate.provider_key_hash or '').lower() if ( not _ROUTED_SERVICE_RE.fullmatch(service) or not _SHA256_RE.fullmatch(provider_key_hash) ): raise DockerShadowPrivacyError('Docker shadow routed identity is invalid') routed.add((service, provider_key_hash)) finally: findings = result.get('findings') if isinstance(result, dict) else None if isinstance(findings, list): for finding in findings: if isinstance(finding, dict): finding.clear() findings.clear() result.clear() return frozenset(routed), frozenset(detectors) def run_timed_private_side( side, operation, durable_checkpoint, *, timeout_sec, monotonic_ns=time.monotonic_ns, ): if side not in ('full', 'adaptive'): raise ValueError('Docker shadow side is invalid') with scanner.scan_slot_scope(['docker-shadow', side], timeout_sec=timeout_sec): started_ns = monotonic_ns() metrics = operation() if not isinstance(metrics, dict) or any( not isinstance(name, str) or isinstance(value, bool) or not isinstance(value, int) or value < 0 for name, value in metrics.items() ): raise ValueError('Docker shadow private sink accepts only non-negative aggregates') durable_checkpoint() elapsed_ns = max(1, monotonic_ns() - started_ns) return metrics, max(1, math.ceil(elapsed_ns / 1_000_000)) def _clear_private_result(result): if not isinstance(result, dict): return findings = result.get('findings') if isinstance(findings, list): for finding in findings: if isinstance(finding, dict): finding.clear() findings.clear() result.clear() def _full_scan_incomplete_reason(result): if result.get('source_failure'): return 'full_scan_source_failure' if result.get('skipped'): return 'full_scan_skipped' return 'full_scan_error' def _adaptive_scan_incomplete_reason(exc): if isinstance(exc, scanner.DockerLayerInfrastructureError): return 'adaptive_scan_infrastructure' if isinstance(exc, scanner.DockerContentTransferError): return 'adaptive_scan_transfer' if isinstance(exc, TimeoutError): return 'adaptive_scan_timeout' if isinstance(exc, scanner.DockerRemoteAccessError): return 'adaptive_scan_remote_access' if isinstance(exc, ValueError): return 'adaptive_scan_contract' return 'adaptive_scan_error' def shadow_failure_reason_code(exc): if isinstance(exc, DockerShadowPrivacyError): return 'privacy_violation' if isinstance(exc, (KeyboardInterrupt, SystemExit)): return 'operator_interrupt' if exc.__class__.__name__ == 'ScanSlotFatalError': return 'scan_slot_fatal' if isinstance(exc, TimeoutError): return 'timeout' if isinstance(exc, ValueError): return 'invalid_contract' if isinstance(exc, RuntimeError): message = str(exc).lower() if 'checkpoint' in message: return 'checkpoint_failure' if 'fence' in message or 'lease' in message: return 'report_fence_failure' if 'cohort' in message or 'control' in message: return 'control_failure' return 'runtime_failure' return 'unexpected_failure' def execute_private_adaptive_scan( normalized_target, *, limits, checkpoint, scan_policy_sha256, timeout_sec, platform_os='linux', platform_arch='amd64', min_free_bytes=0, scan_kwargs=None, ): timeout_sec = max(1, int(timeout_sec)) deadline = time.monotonic() + timeout_sec scan_kwargs = dict(scan_kwargs or {}) resolved, bearer_auth = scanner.resolve_docker_content_manifest( normalized_target, platform_os=str(platform_os or 'linux'), platform_arch=str(platform_arch or 'amd64'), deadline=deadline, ) payload_classes, bearer_auth = scanner.fetch_docker_config_payload_classes( resolved, bearer_auth, deadline=deadline, min_free_bytes=max(0, int(min_free_bytes or 0)), ) built = build_private_adaptive_plan( resolved, payload_classes, limits, checkpoint, scan_policy_sha256, ) plan = built['plan'] selection_metrics = dict(built['selection_metrics']) failure_metrics = empty_failure_metrics() routed = set() detectors = set() selected_unique = { item['digest']: item['max_attempts'] for item in plan['descriptors'] if item['selected'] } max_checkpoints = max(1, sum(selected_unique.values())) checkpoint_count = 0 while True: leased = lease_private_adaptive_checkpoint(plan) if leased is None: break checkpoint_count += 1 if checkpoint_count > max_checkpoints or time.monotonic() >= deadline: raise RuntimeError('Docker shadow adaptive checkpoint budget was exhausted') work = { 'plan': leased['plan'], 'plan_sha256': leased['plan_sha256'], 'bearer_auth': bearer_auth, 'min_free_bytes': max(0, int(min_free_bytes or 0)), 'deadline': deadline, } result = scanner.scan_docker_layer_plan( normalized_target, work, timeout_sec=max(1, math.ceil(deadline - time.monotonic())), detectors=scan_kwargs.get('detectors'), exclude_detectors=scan_kwargs.get('exclude_detectors'), no_verification=bool(scan_kwargs.get('no_verification', False)), trufflehog_config=scan_kwargs.get('trufflehog_config'), log_target=False, ) try: execution = result.get('docker_layer_execution') next_plan = apply_private_adaptive_execution( leased['plan'], result.get('docker_layer_execution'), ) records = { record['digest']: record for record in execution['blobs'] } newly_terminal = { item['digest'] for item in next_plan['descriptors'] if item['coverage_state'] == 'terminal_failed' and item['digest'] in records } for digest in newly_terminal: failure_metrics[ _blob_failure_metric(records[digest]['error_code']) ] += 1 plan = next_plan checkpoint_routed, checkpoint_detectors = private_result_identities( result, normalized_target, ) except Exception: _clear_private_result(result) raise routed.update(checkpoint_routed) detectors.update(checkpoint_detectors) del checkpoint_routed, checkpoint_detectors if any( item['coverage_state'] in ('selected', 'leased', 'shared_pending') for item in plan['descriptors'] ): raise RuntimeError('Docker shadow adaptive plan did not reach a terminal state') selection_metrics['adaptive_checkpoints'] += checkpoint_count terminal_digests = { item['digest'] for item in plan['descriptors'] if item['coverage_state'] == 'terminal_failed' } if sum(failure_metrics.values()) != len(terminal_digests): raise RuntimeError('Docker shadow terminal failure accounting is inconsistent') return { 'routed': frozenset(routed), 'detectors': frozenset(detectors), 'selection_metrics': validate_docker_adaptive_shadow_selection_metrics( selection_metrics, ), 'omitted_descriptor_count': int(built['omitted_descriptor_count']), 'failure_count': len(terminal_digests), 'failure_metrics': failure_metrics, } def private_full_side_metrics(db, control, scan_kwargs): reference_routed, reference_detectors = db.docker_adaptive_shadow_control_identities( control['target_scan_id'], ) reference_routed = set(reference_routed) reference_detectors = set(reference_detectors) rerun_routed = set() rerun_detectors = set() result = None try: max_attempts = min(10, max(1, int(scan_kwargs.get('max_attempts', 3) or 3))) incomplete_reason = 'full_scan_error' attempts = 0 for attempts in range(1, max_attempts + 1): result = scanner.scan_docker_image( control['normalized_target'], timeout_sec=int(scan_kwargs['timeout_sec']), detectors=scan_kwargs.get('detectors'), exclude_detectors=scan_kwargs.get('exclude_detectors'), no_verification=bool(scan_kwargs.get('no_verification', False)), trufflehog_config=scan_kwargs.get('trufflehog_config'), config_dir=scanner.docker_token_manager.get_next_config(), trufflehog_concurrency=int(scan_kwargs.get('trufflehog_concurrency', 0) or 0), docker_recovery_limits=scan_kwargs.get('docker_recovery_limits'), docker_recovery_min_free_bytes=int(scan_kwargs.get('docker_recovery_min_free_bytes', 20 << 30)), log_target=False, ) if not isinstance(result, dict): raise ValueError('Docker shadow full result must be an object') if not ( result.get('errors') or result.get('skipped') or result.get('source_failure') or result.get('warnings') or result.get('degraded') or ((result.get('scan_meta') or {}).get('docker_layer_scope') or {}).get('coverage_complete') is False or ((result.get('scan_meta') or {}).get('docker_full_recovery') or {}).get('coverage_complete') is False ): rerun_routed, rerun_detectors = map( set, private_result_identities(result, control['normalized_target']), ) result = None return { 'full_routed_count': len(reference_routed), 'full_detector_count': len(reference_detectors), 'failure_count': 0, 'safety_regression_count': int( rerun_routed != reference_routed or rerun_detectors != reference_detectors ), **empty_failure_metrics(), } incomplete_reason = _full_scan_incomplete_reason(result) retryable = bool(result.get('retryable', True)) _clear_private_result(result) result = None if not retryable: break print( 'Docker adaptive shadow full side incomplete: ' f'reason_code={incomplete_reason} attempts={attempts}' ) metrics = empty_failure_metrics() metrics['diagnostic_full_incomplete'] = 1 metrics.update({ 'full_routed_count': len(reference_routed), 'full_detector_count': len(reference_detectors), 'failure_count': 1, 'safety_regression_count': 0, }) return metrics finally: _clear_private_result(result) reference_routed.clear() reference_detectors.clear() rerun_routed.clear() rerun_detectors.clear() def private_adaptive_side_metrics( db, control, *, limits, checkpoint, scan_policy_sha256, scan_kwargs, platform_os='linux', platform_arch='amd64', min_free_bytes=0, ): reference_routed, reference_detectors = db.docker_adaptive_shadow_control_identities( control['target_scan_id'], ) reference_routed = set(reference_routed) reference_detectors = set(reference_detectors) adaptive_routed = set() adaptive_detectors = set() outcome = None try: max_attempts = min(10, max(1, int(scan_kwargs.get('max_attempts', 3) or 3))) incomplete_reason = 'adaptive_scan_error' attempts = 0 for attempts in range(1, max_attempts + 1): try: outcome = execute_private_adaptive_scan( control['normalized_target'], limits=limits, checkpoint=checkpoint, scan_policy_sha256=scan_policy_sha256, timeout_sec=int(scan_kwargs['timeout_sec']), platform_os=platform_os, platform_arch=platform_arch, min_free_bytes=min_free_bytes, scan_kwargs=scan_kwargs, ) adaptive_routed = set(outcome.pop('routed')) adaptive_detectors = set(outcome.pop('detectors')) metrics = { 'adaptive_routed_count': len(adaptive_routed), 'routed_intersection_count': len(reference_routed & adaptive_routed), 'adaptive_detector_count': len(adaptive_detectors), 'detector_intersection_count': len(reference_detectors & adaptive_detectors), 'omitted_descriptor_count': int(outcome['omitted_descriptor_count']), 'failure_count': int(outcome['failure_count']), } metrics.update(outcome['selection_metrics']) metrics.update(outcome['failure_metrics']) return metrics except DockerShadowPrivacyError: raise except scanner.ScanSlotFatalError: raise except MemoryError: raise except Exception as exc: incomplete_reason = _adaptive_scan_incomplete_reason(exc) retryable = bool(getattr(exc, 'retryable', not isinstance(exc, ValueError))) if not retryable: break finally: if isinstance(outcome, dict): outcome.clear() outcome = None adaptive_routed.clear() adaptive_detectors.clear() print( 'Docker adaptive shadow adaptive side incomplete: ' f'reason_code={incomplete_reason} attempts={attempts}' ) metrics = empty_selection_metrics() metrics.update(empty_failure_metrics()) metrics['diagnostic_adaptive_incomplete'] = 1 metrics.update({ 'adaptive_routed_count': 0, 'routed_intersection_count': 0, 'adaptive_detector_count': 0, 'detector_intersection_count': 0, 'omitted_descriptor_count': 0, 'failure_count': 1, }) return metrics finally: if isinstance(outcome, dict): outcome.clear() reference_routed.clear() reference_detectors.clear() adaptive_routed.clear() adaptive_detectors.clear() def _shadow_config(config): supervisor = config.get('supervisor') if isinstance(config, dict) else None supervisor = supervisor if isinstance(supervisor, dict) else {} value = supervisor.get('docker_shadow') if not isinstance(value, dict): raise ValueError('Docker shadow supervisor configuration is required') allowed = {'enabled', 'cohort_size', 'lease_seconds'} if set(value) - allowed: raise ValueError('Docker shadow supervisor configuration has unsupported fields') if not console_runner.bool_config(value.get('enabled'), False): raise ValueError('Docker shadow operator command is disabled') cohort_size = int(value.get('cohort_size', DOCKER_ADAPTIVE_GATE_MIN_CONTROLS)) lease_seconds = int(value.get('lease_seconds', 3600)) if not DOCKER_ADAPTIVE_GATE_MIN_CONTROLS <= cohort_size <= DOCKER_ADAPTIVE_GATE_MAX_CONTROLS: raise ValueError('Docker shadow cohort size must be between 50 and 100') if not 60 <= lease_seconds <= 86400: raise ValueError('Docker shadow lease duration is outside the supported range') return {'cohort_size': cohort_size, 'lease_seconds': lease_seconds} def _scan_kwargs(source_args): return { 'timeout_sec': max(1, int(source_args.timeout)), 'detectors': source_args.detectors, 'exclude_detectors': source_args.exclude_detectors, 'no_verification': bool(source_args.no_verification), 'trufflehog_config': source_args.trufflehog_config, 'trufflehog_concurrency': max(0, int(source_args.trufflehog_concurrency or 0)), 'max_attempts': max( 1, int(getattr(source_args, 'target_retry_max_attempts', 3) or 3), ), } def run_shadow(config_path): require_active_supervisor_child( config_path, child_kind='docker-shadow', require_dsn=True, ) config = console_runner.load_config(config_path) settings = _shadow_config(config) source_config = (config.get('sources') or {}).get('dockerhub') if not isinstance(source_config, dict): raise ValueError('Docker shadow requires the DockerHub source configuration') global_config = config.get('global') or {} if not isinstance(global_config, dict): raise ValueError('Docker shadow global configuration must be a mapping') console_runner.apply_global_config(global_config) scanner.initialize_scanner_runtime(preflight_complete=True) secrets_config = console_runner.load_secrets(config, config_path) state = { 'version': 1, 'sources': {'dockerhub': console_runner.default_source_state()}, } console_runner.configure_source_auth( 'dockerhub', source_config, state=state, secrets=secrets_config, ) source_args = console_runner.build_args_from_source_config( 'dockerhub', source_config, global_config, '', auth_entry=None, ) scanner.scan_config.drop_detectors = scanner.csv_items(source_args.drop_detectors) scanner.scan_config.trufflehog_job_memory_limit_bytes = int( source_args.trufflehog_job_memory_limit_bytes ) scan_kwargs = _scan_kwargs(source_args) limits = console_runner.docker_layer_limits(source_args) scan_kwargs['docker_recovery_limits'] = limits scan_kwargs['docker_recovery_min_free_bytes'] = int(getattr(source_args, 'docker_layer_min_free_bytes', 20 << 30)) checkpoint = console_runner.docker_adaptive_checkpoint(source_args) scan_policy_sha256 = console_runner.docker_layer_scan_policy_sha256( source_args, scan_kwargs, ) execution_policy_sha256 = docker_layer_execution_policy_sha256( scan_policy_sha256, limits, ) selection_policy_sha256 = docker_layer_selection_policy_sha256(limits) selection_salt = hashlib.sha256( ('docker-shadow-controls-v1:' + scan_policy_sha256 + ':' + execution_policy_sha256 + ':' + selection_policy_sha256).encode('ascii') ).hexdigest() db_url = str(os.getenv('TRUF_MANAGED_POSTGRES_DSN') or '') if not db_url: raise RuntimeError('Docker shadow canonical PostgreSQL authority is unavailable') db = ScannerDB(db_url=db_url, initialize=False) controls = [] report = None owner = 'docker-shadow:{}:{}'.format( os.getpid(), hashlib.sha256(socket.gethostname().encode('utf-8')).hexdigest()[:16], ) try: if not db.enabled: raise RuntimeError('Docker shadow PostgreSQL connection is unavailable') db.set_application_name('truf-docker-adaptive-shadow') controls = db.docker_adaptive_shadow_controls( scan_policy_sha256, settings['cohort_size'], selection_salt, ) report = db.start_docker_adaptive_shadow_report( scan_policy_sha256, execution_policy_sha256, selection_policy_sha256, settings['cohort_size'], owner, lease_seconds=settings['lease_seconds'], ) aggregate = { 'completed_pairs': 0, 'full_routed_count': 0, 'adaptive_routed_count': 0, 'routed_intersection_count': 0, 'full_detector_count': 0, 'adaptive_detector_count': 0, 'detector_intersection_count': 0, 'full_slot_ms': 0, 'adaptive_slot_ms': 0, 'omitted_descriptor_count': 0, 'failure_count': 0, 'privacy_violation_count': 0, 'safety_regression_count': 0, } selection_metrics = empty_selection_metrics() failure_metrics = empty_failure_metrics() def durable_checkpoint(): db.checkpoint_docker_adaptive_shadow_report( report['report_token'], owner, report['lease_token'], lease_seconds=settings['lease_seconds'], ) for index, control in enumerate(controls): sides = ('full', 'adaptive') if index % 2 == 0 else ('adaptive', 'full') for side in sides: if side == 'full': operation = lambda control=control: private_full_side_metrics( db, control, scan_kwargs, ) else: operation = lambda control=control: private_adaptive_side_metrics( db, control, limits=limits, checkpoint=checkpoint, scan_policy_sha256=scan_policy_sha256, scan_kwargs=scan_kwargs, platform_os=source_args.docker_platform_os, platform_arch=source_args.docker_platform_arch, min_free_bytes=source_args.docker_layer_min_free_bytes, ) side_metrics, elapsed_ms = run_timed_private_side( side, operation, durable_checkpoint, timeout_sec=scan_kwargs['timeout_sec'], ) aggregate[f'{side}_slot_ms'] += elapsed_ms for name, value in side_metrics.items(): if name in selection_metrics: selection_metrics[name] += value elif name in failure_metrics: failure_metrics[name] += value else: aggregate[name] += value side_metrics.clear() aggregate['completed_pairs'] += 1 control.clear() completed = db.finish_docker_adaptive_shadow_report( report['report_token'], owner, report['lease_token'], selection_metrics=selection_metrics, **aggregate, ) print( 'Docker adaptive shadow report: ' f'id={completed["report_id"]} passed={str(completed["passed"]).lower()} ' f'controls={completed["completed_pairs"]} ' f'routed_recall_ppm={completed["routed_recall_ppm"]} ' f'slot_ratio_ppm={completed["slot_ratio_ppm"]} ' 'failure_categories=' + ','.join( f'{name.removeprefix("diagnostic_")}:{failure_metrics[name]}' for name in SHADOW_FAILURE_METRIC_KEYS if failure_metrics[name] ) ) return 0 if completed['passed'] else 2 except BaseException as exc: if report is not None: try: db.fail_docker_adaptive_shadow_report( report['report_token'], owner, report['lease_token'], privacy_violation_count=int(isinstance(exc, DockerShadowPrivacyError)), safety_regression_count=0, ) except Exception: pass print( 'Docker adaptive shadow report failed safely: ' f'reason_code={shadow_failure_reason_code(exc)}' ) return 1 finally: for control in controls: if isinstance(control, dict): control.clear() controls.clear() db.close() scanner.docker_token_manager.cleanup() def parse_args(argv=None): parser = argparse.ArgumentParser(description='Run one private Docker adaptive shadow report') parser.add_argument('--config', required=True) return parser.parse_args(argv) def main(argv=None): args = parse_args(argv) return run_shadow(args.config) if __name__ == '__main__': raise SystemExit(main())