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

850 lines
34 KiB
Python

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())