850 lines
34 KiB
Python
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())
|