1834 lines
78 KiB
Python
1834 lines
78 KiB
Python
import contextlib
|
|
import hashlib
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import shutil
|
|
import socket
|
|
import sys
|
|
import tempfile
|
|
import time
|
|
from types import SimpleNamespace
|
|
import unittest
|
|
from unittest import mock
|
|
|
|
|
|
ROOT = Path(__file__).resolve().parents[1]
|
|
APP_DIR = ROOT / 'app'
|
|
sys.path.insert(0, str(APP_DIR))
|
|
|
|
import console_runner
|
|
import docker_shadow
|
|
import scan_execution
|
|
import scanner
|
|
import scanner_db
|
|
from result_bundle import BundleReservation, ResultBundleReader
|
|
from worker_contracts import WorkerPhase, validate_phase_transition
|
|
|
|
|
|
class StreamResponse:
|
|
def __init__(self, status_code=200, chunks=(), headers=None, url=''):
|
|
self.status_code = status_code
|
|
self._chunks = list(chunks)
|
|
self.headers = dict(headers or {})
|
|
self.url = url
|
|
self.closed = False
|
|
|
|
@property
|
|
def content(self):
|
|
raise AssertionError('streamed response content was materialized')
|
|
|
|
@property
|
|
def text(self):
|
|
raise AssertionError('streamed response text was materialized')
|
|
|
|
def json(self):
|
|
raise AssertionError('streamed response JSON parser was used')
|
|
|
|
def iter_content(self, chunk_size=1):
|
|
del chunk_size
|
|
yield from self._chunks
|
|
|
|
def close(self):
|
|
self.closed = True
|
|
|
|
|
|
def docker_limits():
|
|
return {
|
|
'config_max_bytes': 1024,
|
|
'layer_max_bytes': 1024,
|
|
'image_max_bytes': 2048,
|
|
'max_layers': 2,
|
|
'archive_max_size_bytes': 1024,
|
|
'archive_max_depth': 2,
|
|
'archive_timeout_sec': 5,
|
|
'blob_timeout_sec': 10,
|
|
'filesystem_concurrency': 1,
|
|
'blob_max_attempts': 3,
|
|
}
|
|
|
|
|
|
def docker_plan(payload=b'{}', coverage_state='leased', extra_descriptors=()):
|
|
limits = docker_limits()
|
|
descriptor = {
|
|
'digest': 'sha256:' + hashlib.sha256(payload).hexdigest(),
|
|
'size': len(payload),
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
'kind': 'config',
|
|
'position': 0,
|
|
'selected': coverage_state != 'skipped',
|
|
'selection_reason': (
|
|
'config_selected' if coverage_state != 'shared_pending' else 'shared_active'
|
|
),
|
|
'coverage_state': coverage_state,
|
|
'lease_token': 'l' * 43 if coverage_state == 'leased' else None,
|
|
'attempt': 1 if coverage_state == 'leased' else 0,
|
|
'max_attempts': 3,
|
|
}
|
|
descriptors = [descriptor, *extra_descriptors]
|
|
plan = {
|
|
'version': 1,
|
|
'image': 'owner/repo@sha256:' + ('a' * 64),
|
|
'repository': 'owner/repo',
|
|
'manifest_digest': 'sha256:' + ('a' * 64),
|
|
'platform_os': 'linux',
|
|
'platform_arch': 'amd64',
|
|
'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json',
|
|
'limits': limits,
|
|
'selection_policy_sha256': hashlib.sha256(json.dumps(
|
|
limits, ensure_ascii=True, sort_keys=True, separators=(',', ':'),
|
|
).encode('ascii')).hexdigest(),
|
|
'scan_policy_sha256': 'b' * 64,
|
|
'descriptors': descriptors,
|
|
}
|
|
try:
|
|
scanner_db.validate_docker_layer_plan(plan)
|
|
except ValueError:
|
|
coverage_hash = getattr(scanner_db, 'docker_layer_coverage_policy_sha256', None)
|
|
if not callable(coverage_hash):
|
|
raise
|
|
plan['coverage_policy_sha256'] = coverage_hash('b' * 64, limits)
|
|
scanner_db.validate_docker_layer_plan(plan)
|
|
return plan
|
|
|
|
|
|
def docker_v2_plan(payload=b'{}', coverage_state='leased', extra_descriptors=()):
|
|
limits = docker_limits()
|
|
descriptor = {
|
|
'digest': 'sha256:' + hashlib.sha256(payload).hexdigest(),
|
|
'size': len(payload),
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
'kind': 'config',
|
|
'position': 0,
|
|
'payload_class': 'config',
|
|
'selected': coverage_state != 'skipped',
|
|
'selection_reason': (
|
|
'config_too_large' if coverage_state == 'skipped' else 'config_selected'
|
|
),
|
|
'coverage_state': coverage_state,
|
|
'lease_token': 'v' * 43 if coverage_state == 'leased' else None,
|
|
'attempt': 1 if coverage_state == 'leased' else 0,
|
|
'max_attempts': 3,
|
|
}
|
|
plan = {
|
|
'version': 2,
|
|
'image': 'owner/repo@sha256:' + ('a' * 64),
|
|
'repository': 'owner/repo',
|
|
'manifest_digest': 'sha256:' + ('a' * 64),
|
|
'platform_os': 'linux',
|
|
'platform_arch': 'amd64',
|
|
'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json',
|
|
'limits': limits,
|
|
'selector_version': scanner_db.DOCKER_ADAPTIVE_SELECTOR_VERSION,
|
|
'selection_policy_sha256': scanner_db.docker_layer_selection_policy_sha256(limits),
|
|
'scan_policy_sha256': 'b' * 64,
|
|
'execution_policy_sha256': scanner_db.docker_layer_execution_policy_sha256(
|
|
'b' * 64, limits,
|
|
),
|
|
'checkpoint': {'max_blobs': 2, 'max_bytes': 2048},
|
|
'descriptors': [descriptor, *extra_descriptors],
|
|
}
|
|
return scanner_db.validate_docker_layer_plan(plan)
|
|
|
|
|
|
def layer_work(plan, bearer='opaque-bearer', deadline=None):
|
|
return {
|
|
'plan': plan,
|
|
'plan_sha256': hashlib.sha256(
|
|
scanner_db.canonical_docker_layer_plan_bytes(plan)
|
|
).hexdigest(),
|
|
'bearer_auth': scanner.DockerRegistryAuth(token=bearer),
|
|
'min_free_bytes': 0,
|
|
'deadline': deadline if deadline is not None else time.monotonic() + 30,
|
|
}
|
|
|
|
|
|
class DockerDBFreeParityTests(unittest.TestCase):
|
|
def test_already_covered_plan_has_identical_local_and_db_free_evidence(self):
|
|
plan = docker_v2_plan(coverage_state='covered')
|
|
work = layer_work(plan)
|
|
bundle_id = 'd' * 32
|
|
reservation = BundleReservation(
|
|
reservation_id=17, reservation_token='request-token',
|
|
bundle_id=bundle_id, scan_event_id='e' * 32, queue_id=19,
|
|
claim_lease_token='lease-token', declared_bytes=1024 * 1024,
|
|
ready_path=f'ready/{bundle_id[:2]}/{bundle_id}.trb',
|
|
source='docker', platform='docker', query='fixture',
|
|
target=plan['image'],
|
|
normalized_target=scanner.normalize_target(plan['image'], 'docker'),
|
|
)
|
|
kwargs = {'timeout_sec': 60, 'docker_layer_work': work}
|
|
queue_policy = scan_execution.QueueDispositionPolicy()
|
|
with tempfile.TemporaryDirectory() as temp_dir:
|
|
local_root = os.path.join(temp_dir, 'local')
|
|
remote_root = os.path.join(temp_dir, 'remote')
|
|
scanner.ensure_private_directory(local_root, reject_reparse=True)
|
|
scanner.ensure_private_directory(remote_root, reject_reparse=True)
|
|
with mock.patch.object(
|
|
scanner, 'run_command_streamed',
|
|
side_effect=AssertionError('covered Docker plan must not launch a scanner'),
|
|
), mock.patch.object(
|
|
scanner, 'stream_docker_registry_blob',
|
|
side_effect=AssertionError('covered Docker plan must not download a blob'),
|
|
):
|
|
local_result = scanner.scan_target_result(
|
|
reservation.target, 'docker', reservation.scan_event_id, kwargs,
|
|
)
|
|
local = scan_execution.stage_scan_result_in_scope(
|
|
local_result, reservation, local_root, {}, queue_policy, attempts=1,
|
|
)
|
|
remote = scan_execution.execute_planned_result_in_scope(
|
|
reservation, remote_root, kwargs, {}, queue_policy, attempts=1,
|
|
)
|
|
local_metadata = ResultBundleReader(
|
|
os.path.join(local_root, local.relative_path),
|
|
).metadata()
|
|
remote_metadata = ResultBundleReader(
|
|
os.path.join(remote_root, remote.relative_path),
|
|
).metadata()
|
|
|
|
self.assertEqual(local.queue_status, 'done')
|
|
self.assertEqual(remote.queue_status, 'done')
|
|
self.assertEqual(
|
|
local_metadata['docker_layer_plan'], remote_metadata['docker_layer_plan'],
|
|
)
|
|
self.assertEqual(
|
|
local_metadata['docker_layer_execution'],
|
|
remote_metadata['docker_layer_execution'],
|
|
)
|
|
self.assertEqual(
|
|
local_metadata['scan_meta']['docker_layer_scope'],
|
|
remote_metadata['scan_meta']['docker_layer_scope'],
|
|
)
|
|
execution = remote_metadata['docker_layer_execution']
|
|
self.assertEqual(execution['blobs'], [])
|
|
self.assertEqual(
|
|
execution['plan_sha256'],
|
|
hashlib.sha256(scanner_db.canonical_docker_layer_plan_bytes(plan)).hexdigest(),
|
|
)
|
|
self.assertTrue(remote_metadata['scan_meta']['docker_layer_scope']['coverage_complete'])
|
|
|
|
|
|
def unlink(path):
|
|
try:
|
|
os.unlink(path)
|
|
except FileNotFoundError:
|
|
pass
|
|
|
|
|
|
class DockerRegistryBoundTests(unittest.TestCase):
|
|
def test_manifest_and_token_bodies_stream_under_one_deadline(self):
|
|
manifest = json.dumps({
|
|
'schemaVersion': 2,
|
|
'config': {
|
|
'digest': 'sha256:' + ('c' * 64),
|
|
'size': 2,
|
|
'mediaType': 'application/vnd.oci.image.config.v1+json',
|
|
},
|
|
'layers': [],
|
|
}, separators=(',', ':')).encode('utf-8')
|
|
digest = 'sha256:' + hashlib.sha256(manifest).hexdigest()
|
|
challenge = (
|
|
'Bearer realm="https://auth.docker.io/token",'
|
|
'service="registry.docker.io",scope="repository:owner/repo:pull"'
|
|
)
|
|
unauthorized = StreamResponse(
|
|
401, headers={'WWW-Authenticate': challenge},
|
|
url='https://registry-1.docker.io/v2/owner/repo/manifests/' + digest,
|
|
)
|
|
token_body = b'{"token":"opaque-registry-token"}'
|
|
token = StreamResponse(
|
|
200, [token_body], {'Content-Length': str(len(token_body))},
|
|
'https://auth.docker.io/token',
|
|
)
|
|
success = StreamResponse(200, [manifest[:17], manifest[17:]], {
|
|
'Content-Length': str(len(manifest)),
|
|
'Docker-Content-Digest': digest,
|
|
})
|
|
deadline = time.monotonic() + 30
|
|
manager = scanner.DockerTokenManager()
|
|
|
|
with mock.patch.object(scanner, 'docker_token_manager', manager), \
|
|
mock.patch.object(
|
|
scanner, 'api_request', side_effect=[unauthorized, token, success],
|
|
) as request:
|
|
payload, auth = scanner.docker_registry_manifest(
|
|
'owner/repo', digest, verify_content_digest=True, deadline=deadline,
|
|
)
|
|
|
|
self.assertEqual(payload['schemaVersion'], 2)
|
|
self.assertEqual(auth.token, 'opaque-registry-token')
|
|
self.assertTrue(all(response.closed for response in (unauthorized, token, success)))
|
|
self.assertTrue(all(call.kwargs['stream'] for call in request.call_args_list))
|
|
self.assertTrue(all(call.kwargs['deadline'] == deadline for call in request.call_args_list))
|
|
self.assertTrue(all(call.kwargs.get('use_proxy') is None for call in request.call_args_list))
|
|
|
|
def test_manifest_bearer_401_is_target_scoped(self):
|
|
digest = 'sha256:' + ('a' * 64)
|
|
challenge = (
|
|
'Bearer realm="https://auth.docker.io/token",'
|
|
'service="registry.docker.io",scope="repository:owner/repo:pull"'
|
|
)
|
|
responses = [
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
StreamResponse(200, [b'{"token":"bearer-a"}']),
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
StreamResponse(200, [b'{"token":"bearer-b"}']),
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
]
|
|
manager = scanner.DockerTokenManager()
|
|
manager.accounts = [
|
|
scanner.DockerAccount('account-a', 'user-a', 'secret-a', ''),
|
|
scanner.DockerAccount('account-b', 'user-b', 'secret-b', ''),
|
|
]
|
|
manager.explicit_pool = True
|
|
|
|
with mock.patch.object(scanner, 'docker_token_manager', manager), \
|
|
mock.patch.object(scanner, 'api_request', side_effect=responses):
|
|
with self.assertRaises(scanner.DockerRemoteAccessError) as raised:
|
|
scanner.docker_registry_manifest('owner/repo', digest)
|
|
|
|
self.assertEqual(raised.exception.status, 'target_forbidden')
|
|
self.assertEqual(manager.invalid_accounts, set())
|
|
self.assertNotIn(
|
|
'auth_invalid',
|
|
[event['category'] for event in manager.drain_status_events()],
|
|
)
|
|
|
|
def test_registry_token_401_still_invalidates_rejected_account(self):
|
|
manifest = b'{"schemaVersion":2,"layers":[]}'
|
|
digest = 'sha256:' + hashlib.sha256(manifest).hexdigest()
|
|
challenge = (
|
|
'Bearer realm="https://auth.docker.io/token",'
|
|
'service="registry.docker.io",scope="repository:owner/repo:pull"'
|
|
)
|
|
responses = [
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
StreamResponse(401),
|
|
StreamResponse(200, [b'{"token":"bearer-b"}']),
|
|
StreamResponse(200, [manifest], {'Docker-Content-Digest': digest}),
|
|
]
|
|
manager = scanner.DockerTokenManager()
|
|
manager.accounts = [
|
|
scanner.DockerAccount('account-a', 'user-a', 'secret-a', ''),
|
|
scanner.DockerAccount('account-b', 'user-b', 'secret-b', ''),
|
|
]
|
|
manager.explicit_pool = True
|
|
|
|
with mock.patch.object(scanner, 'docker_token_manager', manager), \
|
|
mock.patch.object(scanner, 'api_request', side_effect=responses):
|
|
payload, auth = scanner.docker_registry_manifest('owner/repo', digest)
|
|
|
|
self.assertEqual(payload['schemaVersion'], 2)
|
|
self.assertEqual(auth.account_name, 'account-b')
|
|
self.assertEqual(manager.invalid_accounts, {'account-a'})
|
|
events = manager.drain_status_events()
|
|
self.assertEqual([
|
|
event['name'] for event in events if event['category'] == 'auth_invalid'
|
|
], ['account-a'])
|
|
|
|
def test_blob_bearer_401_is_target_scoped(self):
|
|
digest = 'sha256:' + ('b' * 64)
|
|
challenge = (
|
|
'Bearer realm="https://auth.docker.io/token",'
|
|
'service="registry.docker.io",scope="repository:owner/repo:pull"'
|
|
)
|
|
responses = [
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
StreamResponse(200, [b'{"token":"bearer-b"}']),
|
|
StreamResponse(401, headers={'WWW-Authenticate': challenge}),
|
|
]
|
|
manager = scanner.DockerTokenManager()
|
|
manager.accounts = [
|
|
scanner.DockerAccount('account-a', 'user-a', 'secret-a', ''),
|
|
scanner.DockerAccount('account-b', 'user-b', 'secret-b', ''),
|
|
]
|
|
manager.explicit_pool = True
|
|
|
|
with mock.patch.object(scanner, 'docker_token_manager', manager), \
|
|
mock.patch.object(scanner, 'api_request', side_effect=responses):
|
|
with self.assertRaises(scanner.DockerContentTransferError) as raised:
|
|
scanner._docker_registry_blob_response(
|
|
'owner/repo', digest,
|
|
scanner.DockerRegistryAuth('bearer-a', 'account-a', challenge),
|
|
time.monotonic() + 30,
|
|
)
|
|
|
|
self.assertEqual(raised.exception.error_code, 'target_forbidden')
|
|
self.assertEqual(manager.invalid_accounts, set())
|
|
self.assertNotIn(
|
|
'auth_invalid',
|
|
[event['category'] for event in manager.drain_status_events()],
|
|
)
|
|
|
|
|
|
def test_index_resolution_keeps_claimed_immutable_root_and_deadline(self):
|
|
index_digest = 'sha256:' + ('1' * 64)
|
|
child_digest = 'sha256:' + ('2' * 64)
|
|
config_digest = 'sha256:' + ('3' * 64)
|
|
layer_digest = 'sha256:' + ('4' * 64)
|
|
index = {
|
|
'manifests': [{
|
|
'digest': child_digest,
|
|
'platform': {'os': 'linux', 'architecture': 'amd64'},
|
|
}],
|
|
}
|
|
child = {
|
|
'mediaType': 'application/vnd.oci.image.manifest.v1+json',
|
|
'config': {
|
|
'digest': config_digest, 'size': 2,
|
|
'mediaType': 'application/vnd.oci.image.config.v1+json',
|
|
},
|
|
'layers': [{
|
|
'digest': layer_digest, 'size': 10,
|
|
'mediaType': 'application/vnd.oci.image.layer.v1.tar+gzip',
|
|
}],
|
|
}
|
|
deadline = time.monotonic() + 30
|
|
auth = scanner.DockerRegistryAuth(token='opaque')
|
|
with mock.patch.object(
|
|
scanner, 'docker_registry_manifest',
|
|
side_effect=[(index, auth), (child, auth)],
|
|
) as manifest:
|
|
resolved, _ = scanner.resolve_docker_content_manifest(
|
|
f'owner/repo@{index_digest}', deadline=deadline,
|
|
)
|
|
|
|
self.assertEqual(resolved['manifest_digest'], index_digest)
|
|
self.assertEqual(
|
|
set(resolved['config']), {'digest', 'size', 'media_type'},
|
|
)
|
|
self.assertEqual(
|
|
set(resolved['layers'][0]), {'digest', 'size', 'media_type'},
|
|
)
|
|
self.assertEqual(resolved['layers'][0]['digest'], layer_digest)
|
|
self.assertEqual(
|
|
scanner_db.validate_docker_layer_resolution(resolved)['image'],
|
|
f'owner/repo@{index_digest}',
|
|
)
|
|
self.assertTrue(all(
|
|
call.kwargs['deadline'] == deadline for call in manifest.call_args_list
|
|
))
|
|
|
|
def test_streamed_manifest_rejects_body_over_limit(self):
|
|
response = StreamResponse(200, [b'1234', b'5'])
|
|
with self.assertRaisesRegex(
|
|
scanner.DockerRegistryResolutionError, 'size limit',
|
|
):
|
|
scanner._bounded_docker_registry_json(response, 'fixture', max_bytes=4)
|
|
|
|
def test_streamed_request_error_does_not_echo_signed_query(self):
|
|
secret = 'signed-query-secret'
|
|
error = scanner.requests.exceptions.ConnectionError(
|
|
f'failed https://layers.cloudfront.net/blob?token={secret}'
|
|
)
|
|
with mock.patch.object(scanner.scan_config, 'api_proxy_enabled', False), \
|
|
mock.patch.object(scanner.requests, 'request', side_effect=error), \
|
|
mock.patch.object(scanner.logger, 'warning') as warning:
|
|
with self.assertRaises(scanner.ApiRequestError) as raised:
|
|
scanner.api_request(
|
|
'GET', f'https://layers.cloudfront.net/blob?token={secret}',
|
|
stream=True, max_retries=1,
|
|
)
|
|
self.assertNotIn(secret, str(raised.exception))
|
|
self.assertNotIn(secret, repr(warning.call_args_list))
|
|
|
|
def test_streamed_http_error_body_is_never_materialized(self):
|
|
response = StreamResponse(503, [b'server detail that must not be read'])
|
|
with mock.patch.object(scanner.scan_config, 'api_proxy_enabled', False), \
|
|
mock.patch.object(scanner.requests, 'request', return_value=response):
|
|
with self.assertRaises(scanner.ApiRequestError):
|
|
scanner.api_request(
|
|
'GET', 'https://registry-1.docker.io/v2/fixture',
|
|
stream=True, max_retries=1, retry_statuses={503},
|
|
)
|
|
self.assertTrue(response.closed)
|
|
|
|
def test_redirect_dns_must_be_entirely_global(self):
|
|
private_answer = [
|
|
(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('127.0.0.1', 443)),
|
|
]
|
|
mixed_answers = private_answer + [
|
|
(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('93.184.216.34', 443)),
|
|
]
|
|
with mock.patch.object(scanner.socket, 'getaddrinfo', return_value=mixed_answers):
|
|
self.assertFalse(scanner._docker_blob_url_allowed(
|
|
'https://layers.cloudfront.net/content',
|
|
))
|
|
|
|
def test_cross_host_redirect_strips_registry_authorization(self):
|
|
payload = b'{"history":[]}'
|
|
digest = 'sha256:' + hashlib.sha256(payload).hexdigest()
|
|
descriptor = {
|
|
'digest': digest,
|
|
'size': len(payload),
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
'kind': 'config',
|
|
}
|
|
redirect = StreamResponse(307, headers={
|
|
'Location': 'https://layers.cloudfront.net/blob?signature=sensitive',
|
|
}, url='https://registry-1.docker.io/v2/owner/repo/blobs/' + digest)
|
|
success = StreamResponse(200, [payload], {
|
|
'Content-Length': str(len(payload)),
|
|
}, url='https://layers.cloudfront.net/blob?signature=sensitive')
|
|
answers = [
|
|
(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, '', ('93.184.216.34', 443)),
|
|
]
|
|
|
|
with tempfile.TemporaryDirectory() as temp_dir:
|
|
destination = os.path.join(temp_dir, 'blob')
|
|
with contextlib.ExitStack() as stack:
|
|
stack.enter_context(mock.patch.object(
|
|
scanner, '_docker_registry_blob_response',
|
|
return_value=(redirect, scanner.DockerRegistryAuth(token='registry-secret')),
|
|
))
|
|
request = stack.enter_context(mock.patch.object(
|
|
scanner, 'api_request', return_value=success,
|
|
))
|
|
stack.enter_context(mock.patch.object(
|
|
scanner.socket, 'getaddrinfo', return_value=answers,
|
|
))
|
|
stack.enter_context(mock.patch.object(
|
|
scanner, 'require_private_directory', side_effect=lambda path, create=False: path,
|
|
))
|
|
stack.enter_context(mock.patch.object(scanner, 'reject_reparse_components'))
|
|
stack.enter_context(mock.patch.object(scanner, 'harden_private_file'))
|
|
stack.enter_context(mock.patch.object(scanner, 'private_file_ready', return_value=True))
|
|
stack.enter_context(mock.patch.object(scanner, 'durable_replace', side_effect=os.replace))
|
|
stack.enter_context(mock.patch.object(scanner, 'durable_unlink', side_effect=unlink))
|
|
outcome = scanner.stream_docker_registry_blob(
|
|
'owner/repo', descriptor, destination,
|
|
scanner.DockerRegistryAuth(token='registry-secret'),
|
|
deadline=time.monotonic() + 30,
|
|
)
|
|
|
|
self.assertEqual(outcome.verified_bytes, len(payload))
|
|
self.assertNotIn('Authorization', request.call_args.kwargs['headers'])
|
|
self.assertNotIn('registry-secret', repr(request.call_args.kwargs))
|
|
self.assertIs(request.call_args.kwargs['use_proxy'], False)
|
|
|
|
def test_digest_mismatch_removes_partial_and_destination(self):
|
|
payload = b'{"history":[]}'
|
|
descriptor = {
|
|
'digest': 'sha256:' + ('0' * 64),
|
|
'size': len(payload),
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
'kind': 'config',
|
|
}
|
|
response = StreamResponse(200, [payload], {
|
|
'Content-Length': str(len(payload)),
|
|
})
|
|
|
|
with tempfile.TemporaryDirectory() as temp_dir:
|
|
destination = os.path.join(temp_dir, 'blob')
|
|
with contextlib.ExitStack() as stack:
|
|
stack.enter_context(mock.patch.object(
|
|
scanner, '_docker_registry_blob_response',
|
|
return_value=(response, scanner.DockerRegistryAuth(token='opaque')),
|
|
))
|
|
stack.enter_context(mock.patch.object(
|
|
scanner, 'require_private_directory', side_effect=lambda path, create=False: path,
|
|
))
|
|
stack.enter_context(mock.patch.object(scanner, 'reject_reparse_components'))
|
|
stack.enter_context(mock.patch.object(scanner, 'harden_private_file'))
|
|
stack.enter_context(mock.patch.object(scanner, 'durable_replace', side_effect=os.replace))
|
|
stack.enter_context(mock.patch.object(scanner, 'durable_unlink', side_effect=unlink))
|
|
with self.assertRaisesRegex(
|
|
scanner.DockerContentTransferError, 'SHA-256',
|
|
) as raised:
|
|
scanner.stream_docker_registry_blob(
|
|
'owner/repo', descriptor, destination,
|
|
deadline=time.monotonic() + 30,
|
|
)
|
|
self.assertEqual(raised.exception.transfer_bytes, len(payload))
|
|
self.assertEqual(list(Path(temp_dir).iterdir()), [])
|
|
|
|
def test_blob_http_failures_keep_target_and_source_scope_distinct(self):
|
|
descriptor = {
|
|
'digest': 'sha256:' + ('0' * 64),
|
|
'size': 2,
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
'kind': 'config',
|
|
}
|
|
cases = {
|
|
401: (scanner.DockerContentTransferError, 'target_forbidden'),
|
|
403: (scanner.DockerContentTransferError, 'target_forbidden'),
|
|
429: (scanner.DockerLayerInfrastructureError, 'remote_rate_limit'),
|
|
503: (scanner.DockerLayerInfrastructureError, 'remote_transient'),
|
|
}
|
|
for status, (error_type, error_code) in cases.items():
|
|
with self.subTest(status=status), tempfile.TemporaryDirectory() as temp_dir:
|
|
response = StreamResponse(status)
|
|
destination = os.path.join(temp_dir, 'blob')
|
|
with mock.patch.object(
|
|
scanner, '_docker_registry_blob_response',
|
|
return_value=(response, scanner.DockerRegistryAuth(token='opaque')),
|
|
), mock.patch.object(
|
|
scanner, 'require_private_directory', side_effect=lambda path, create=False: path,
|
|
), mock.patch.object(scanner, 'reject_reparse_components'):
|
|
with self.assertRaises(error_type) as raised:
|
|
scanner.stream_docker_registry_blob(
|
|
'owner/repo', descriptor, destination,
|
|
deadline=time.monotonic() + 30,
|
|
)
|
|
self.assertEqual(raised.exception.error_code, error_code)
|
|
if status == 403:
|
|
self.assertNotIsInstance(
|
|
raised.exception, scanner.DockerLayerInfrastructureError,
|
|
)
|
|
|
|
|
|
class DockerAdaptivePolicyTests(unittest.TestCase):
|
|
def test_private_shadow_checkpoints_are_canonical_frozen_and_db_free(self):
|
|
limits = dict(docker_limits(), image_max_bytes=1500)
|
|
resolved = {
|
|
'version': 1,
|
|
'image': 'owner/repo@sha256:' + ('a' * 64),
|
|
'repository': 'owner/repo',
|
|
'manifest_digest': 'sha256:' + ('a' * 64),
|
|
'platform_os': 'linux',
|
|
'platform_arch': 'amd64',
|
|
'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json',
|
|
'config': {
|
|
'digest': 'sha256:' + ('0' * 64), 'size': 100,
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
},
|
|
'layers': [{
|
|
'digest': 'sha256:' + (char * 64), 'size': size,
|
|
'media_type': 'application/vnd.oci.image.layer.v1.tar+gzip',
|
|
} for char, size in (('1', 900), ('2', 800), ('3', 700))],
|
|
}
|
|
|
|
built = docker_shadow.build_private_adaptive_plan(
|
|
resolved, ['bulk_data', 'copy_add', 'app_config_run'], limits,
|
|
{'max_blobs': 2, 'max_bytes': 1024}, 'b' * 64,
|
|
)
|
|
|
|
plan = built['plan']
|
|
self.assertEqual(
|
|
[(item['position'], item['coverage_state'], item['selection_reason'])
|
|
for item in plan['descriptors']],
|
|
[
|
|
(0, 'selected', 'config_selected'),
|
|
(1, 'skipped', 'layer_limit_exhausted'),
|
|
(2, 'selected', 'selected_copy_add'),
|
|
(3, 'selected', 'selected_app_config_run'),
|
|
],
|
|
)
|
|
self.assertEqual(built['omitted_descriptor_count'], 1)
|
|
self.assertEqual(built['selection_metrics']['selected_config'], 1)
|
|
self.assertEqual(built['selection_metrics']['selected_copy_add'], 1)
|
|
self.assertEqual(built['selection_metrics']['selected_app_config_run'], 1)
|
|
self.assertEqual(built['selection_metrics']['omitted_layer_limit_exhausted'], 1)
|
|
self.assertEqual(
|
|
built['plan_sha256'],
|
|
hashlib.sha256(scanner_db.canonical_docker_layer_plan_bytes(plan)).hexdigest(),
|
|
)
|
|
|
|
first = docker_shadow.lease_private_adaptive_checkpoint(
|
|
plan, token_factory=lambda: 'p' * 43,
|
|
)
|
|
first_plan = first['plan']
|
|
self.assertEqual(
|
|
[item['position'] for item in first_plan['descriptors']
|
|
if item['coverage_state'] == 'leased'],
|
|
[0, 2],
|
|
)
|
|
execution = {
|
|
'version': 2,
|
|
'plan_sha256': first['plan_sha256'],
|
|
'blobs': [{
|
|
'digest': item['digest'],
|
|
'lease_token': item['lease_token'],
|
|
'status': 'covered',
|
|
'verified_bytes': item['size'],
|
|
'transfer_bytes': item['size'],
|
|
'transfer_duration_ms': 1,
|
|
'scan_duration_ms': 1,
|
|
'finding_count': 0,
|
|
'error_code': None,
|
|
} for item in first_plan['descriptors']
|
|
if item['coverage_state'] == 'leased'],
|
|
}
|
|
resumed = docker_shadow.apply_private_adaptive_execution(first_plan, execution)
|
|
second = docker_shadow.lease_private_adaptive_checkpoint(
|
|
resumed, token_factory=lambda: 'q' * 43,
|
|
)
|
|
self.assertEqual(
|
|
[item['position'] for item in second['plan']['descriptors']
|
|
if item['coverage_state'] == 'leased'],
|
|
[3],
|
|
)
|
|
self.assertEqual(
|
|
[item['selection_reason'] for item in second['plan']['descriptors']],
|
|
[item['selection_reason'] for item in plan['descriptors']],
|
|
)
|
|
|
|
def test_private_shadow_identity_reduction_clears_raw_material(self):
|
|
result = {'findings': [{'Raw': 'private-value', 'DetectorName': 'Fixture'}]}
|
|
candidate = SimpleNamespace(service='openai', provider_key_hash='c' * 64)
|
|
with mock.patch.object(
|
|
docker_shadow, 'extract_candidates', return_value=[candidate],
|
|
), mock.patch.object(
|
|
docker_shadow, 'finding_identity',
|
|
return_value=('private-value', 'd' * 64, 'e' * 64, 'f' * 64, {}),
|
|
):
|
|
routed, detectors = docker_shadow.private_result_identities(
|
|
result, 'owner/repo@sha256:' + ('a' * 64),
|
|
)
|
|
|
|
self.assertEqual(routed, frozenset({('openai', 'c' * 64)}))
|
|
self.assertEqual(detectors, frozenset({'e' * 64}))
|
|
self.assertEqual(result, {})
|
|
|
|
def test_private_shadow_timer_includes_durable_sink_inside_slot(self):
|
|
events = []
|
|
|
|
@contextlib.contextmanager
|
|
def slot(command, timeout_sec=None):
|
|
events.append(('enter', command, timeout_sec))
|
|
try:
|
|
yield object()
|
|
finally:
|
|
events.append(('release',))
|
|
|
|
ticks = iter((1_000_000, 4_000_001))
|
|
with mock.patch.object(scanner, 'scan_slot_scope', side_effect=slot):
|
|
metrics, elapsed_ms = docker_shadow.run_timed_private_side(
|
|
'adaptive',
|
|
lambda: events.append(('operation',)) or {'findings': 1},
|
|
lambda: events.append(('checkpoint',)),
|
|
timeout_sec=30,
|
|
monotonic_ns=lambda: next(ticks),
|
|
)
|
|
|
|
self.assertEqual(metrics, {'findings': 1})
|
|
self.assertEqual(elapsed_ms, 4)
|
|
self.assertEqual(events, [
|
|
('enter', ['docker-shadow', 'adaptive'], 30),
|
|
('operation',),
|
|
('checkpoint',),
|
|
('release',),
|
|
])
|
|
|
|
def test_private_full_side_retries_transient_incomplete_scan(self):
|
|
class IdentityDB:
|
|
@staticmethod
|
|
def docker_adaptive_shadow_control_identities(_target_scan_id):
|
|
return frozenset(), frozenset()
|
|
|
|
incomplete = {
|
|
'findings': [{'Raw': 'private-value'}],
|
|
'errors': ['private-error'],
|
|
'retryable': True,
|
|
}
|
|
complete = {'findings': [], 'errors': []}
|
|
control = {
|
|
'target_scan_id': 7,
|
|
'normalized_target': 'owner/repo@sha256:' + ('a' * 64),
|
|
}
|
|
with mock.patch.object(
|
|
scanner, 'scan_docker_image', side_effect=(incomplete, complete),
|
|
) as scan:
|
|
metrics = docker_shadow.private_full_side_metrics(
|
|
IdentityDB(), control,
|
|
{'timeout_sec': 30, 'max_attempts': 3},
|
|
)
|
|
|
|
self.assertEqual(scan.call_count, 2)
|
|
self.assertEqual(incomplete, {})
|
|
self.assertEqual(complete, {})
|
|
self.assertEqual(metrics, {
|
|
'full_routed_count': 0,
|
|
'full_detector_count': 0,
|
|
'failure_count': 0,
|
|
'safety_regression_count': 0,
|
|
**docker_shadow.empty_failure_metrics(),
|
|
})
|
|
self.assertTrue(all(
|
|
call.kwargs['log_target'] is False for call in scan.call_args_list
|
|
))
|
|
|
|
def test_private_full_side_records_terminal_incomplete_and_continues(self):
|
|
class IdentityDB:
|
|
@staticmethod
|
|
def docker_adaptive_shadow_control_identities(_target_scan_id):
|
|
return frozenset({('openai', 'b' * 64)}), frozenset({'c' * 64})
|
|
|
|
retained = []
|
|
|
|
def incomplete(*_args, **_kwargs):
|
|
result = {
|
|
'findings': [{'Raw': 'private-value'}],
|
|
'errors': ['private-error'],
|
|
'retryable': True,
|
|
}
|
|
retained.append(result)
|
|
return result
|
|
|
|
control = {
|
|
'target_scan_id': 8,
|
|
'normalized_target': 'owner/private@sha256:' + ('d' * 64),
|
|
}
|
|
with mock.patch.object(
|
|
scanner, 'scan_docker_image', side_effect=incomplete,
|
|
) as scan, mock.patch('builtins.print') as output:
|
|
metrics = docker_shadow.private_full_side_metrics(
|
|
IdentityDB(), control,
|
|
{'timeout_sec': 30, 'max_attempts': 3},
|
|
)
|
|
|
|
self.assertEqual(scan.call_count, 3)
|
|
self.assertTrue(all(result == {} for result in retained))
|
|
self.assertEqual(metrics, {
|
|
'full_routed_count': 1,
|
|
'full_detector_count': 1,
|
|
'failure_count': 1,
|
|
'safety_regression_count': 0,
|
|
**dict(
|
|
docker_shadow.empty_failure_metrics(),
|
|
diagnostic_full_incomplete=1,
|
|
),
|
|
})
|
|
rendered = ' '.join(str(call) for call in output.call_args_list)
|
|
self.assertIn('reason_code=full_scan_error attempts=3', rendered)
|
|
self.assertNotIn('owner/private', rendered)
|
|
self.assertNotIn('private-error', rendered)
|
|
|
|
def test_shadow_failure_reason_codes_do_not_render_exception_details(self):
|
|
self.assertEqual(
|
|
docker_shadow.shadow_failure_reason_code(
|
|
docker_shadow.DockerShadowPrivacyError('private-value')
|
|
),
|
|
'privacy_violation',
|
|
)
|
|
self.assertEqual(
|
|
docker_shadow.shadow_failure_reason_code(
|
|
RuntimeError('Docker shadow checkpoint contains private-value')
|
|
),
|
|
'checkpoint_failure',
|
|
)
|
|
|
|
def test_private_shadow_executor_retries_without_authoritative_state(self):
|
|
resolved = {
|
|
'version': 1,
|
|
'image': 'owner/repo@sha256:' + ('a' * 64),
|
|
'repository': 'owner/repo',
|
|
'manifest_digest': 'sha256:' + ('a' * 64),
|
|
'platform_os': 'linux',
|
|
'platform_arch': 'amd64',
|
|
'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json',
|
|
'config': {
|
|
'digest': 'sha256:' + ('0' * 64), 'size': 100,
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
},
|
|
'layers': [],
|
|
}
|
|
executions = []
|
|
target_logging = []
|
|
|
|
def execute(_target, work, **kwargs):
|
|
target_logging.append(kwargs['log_target'])
|
|
descriptor = next(
|
|
item for item in work['plan']['descriptors']
|
|
if item['coverage_state'] == 'leased'
|
|
)
|
|
executions.append(descriptor['attempt'])
|
|
covered = len(executions) == 2
|
|
return {
|
|
'findings': [],
|
|
'errors': [],
|
|
'docker_layer_execution': {
|
|
'version': 2,
|
|
'plan_sha256': work['plan_sha256'],
|
|
'blobs': [{
|
|
'digest': descriptor['digest'],
|
|
'lease_token': descriptor['lease_token'],
|
|
'status': 'covered' if covered else 'retryable_failed',
|
|
'verified_bytes': descriptor['size'] if covered else 0,
|
|
'transfer_bytes': descriptor['size'],
|
|
'transfer_duration_ms': 1,
|
|
'scan_duration_ms': 1,
|
|
'finding_count': 0,
|
|
'error_code': None if covered else 'transfer_timeout',
|
|
}],
|
|
},
|
|
}
|
|
|
|
with mock.patch.object(
|
|
scanner, 'resolve_docker_content_manifest', return_value=(resolved, None),
|
|
), mock.patch.object(
|
|
scanner, 'fetch_docker_config_payload_classes', return_value=([], None),
|
|
), mock.patch.object(
|
|
scanner, 'scan_docker_layer_plan', side_effect=execute,
|
|
):
|
|
outcome = docker_shadow.execute_private_adaptive_scan(
|
|
resolved['image'], limits=docker_limits(),
|
|
checkpoint={'max_blobs': 1, 'max_bytes': 1024},
|
|
scan_policy_sha256='b' * 64, timeout_sec=30,
|
|
)
|
|
|
|
self.assertEqual(executions, [1, 2])
|
|
self.assertEqual(outcome['failure_count'], 0)
|
|
self.assertEqual(outcome['omitted_descriptor_count'], 0)
|
|
self.assertEqual(outcome['selection_metrics']['adaptive_checkpoints'], 2)
|
|
self.assertEqual(outcome['failure_metrics'], docker_shadow.empty_failure_metrics())
|
|
self.assertEqual(outcome['routed'], frozenset())
|
|
self.assertEqual(outcome['detectors'], frozenset())
|
|
self.assertEqual(target_logging, [False, False])
|
|
|
|
def test_private_shadow_executor_classifies_terminal_blob_failure(self):
|
|
resolved = {
|
|
'version': 1,
|
|
'image': 'owner/repo@sha256:' + ('a' * 64),
|
|
'repository': 'owner/repo',
|
|
'manifest_digest': 'sha256:' + ('a' * 64),
|
|
'platform_os': 'linux',
|
|
'platform_arch': 'amd64',
|
|
'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json',
|
|
'config': {
|
|
'digest': 'sha256:' + ('0' * 64), 'size': 100,
|
|
'media_type': 'application/vnd.oci.image.config.v1+json',
|
|
},
|
|
'layers': [],
|
|
}
|
|
|
|
def execute(_target, work, **_kwargs):
|
|
descriptor = next(
|
|
item for item in work['plan']['descriptors']
|
|
if item['coverage_state'] == 'leased'
|
|
)
|
|
return {
|
|
'findings': [],
|
|
'errors': [],
|
|
'docker_layer_execution': {
|
|
'version': 2,
|
|
'plan_sha256': work['plan_sha256'],
|
|
'blobs': [{
|
|
'digest': descriptor['digest'],
|
|
'lease_token': descriptor['lease_token'],
|
|
'status': 'terminal_failed',
|
|
'verified_bytes': descriptor['size'],
|
|
'transfer_bytes': descriptor['size'],
|
|
'transfer_duration_ms': 1,
|
|
'scan_duration_ms': 1,
|
|
'finding_count': 0,
|
|
'error_code': 'chunk_processing',
|
|
}],
|
|
},
|
|
}
|
|
|
|
with mock.patch.object(
|
|
scanner, 'resolve_docker_content_manifest', return_value=(resolved, None),
|
|
), mock.patch.object(
|
|
scanner, 'fetch_docker_config_payload_classes', return_value=([], None),
|
|
), mock.patch.object(
|
|
scanner, 'scan_docker_layer_plan', side_effect=execute,
|
|
):
|
|
outcome = docker_shadow.execute_private_adaptive_scan(
|
|
resolved['image'], limits=docker_limits(),
|
|
checkpoint={'max_blobs': 1, 'max_bytes': 1024},
|
|
scan_policy_sha256='b' * 64, timeout_sec=30,
|
|
)
|
|
|
|
self.assertEqual(outcome['failure_count'], 1)
|
|
self.assertEqual(
|
|
outcome['failure_metrics']['diagnostic_blob_chunk_processing'], 1,
|
|
)
|
|
self.assertEqual(sum(outcome['failure_metrics'].values()), 1)
|
|
|
|
def test_private_shadow_side_uses_only_hashed_identity_reader(self):
|
|
class ReadOnlyIdentityDB:
|
|
def docker_adaptive_shadow_control_identities(self, target_scan_id):
|
|
self.target_scan_id = target_scan_id
|
|
return frozenset({('openai', 'a' * 64)}), frozenset({'b' * 64})
|
|
|
|
def __getattr__(self, name):
|
|
raise AssertionError(f'unexpected authoritative database access: {name}')
|
|
|
|
db = ReadOnlyIdentityDB()
|
|
outcome = {
|
|
'routed': frozenset({('openai', 'a' * 64)}),
|
|
'detectors': frozenset({'c' * 64}),
|
|
'selection_metrics': docker_shadow.empty_selection_metrics(),
|
|
'omitted_descriptor_count': 0,
|
|
'failure_count': 0,
|
|
'failure_metrics': docker_shadow.empty_failure_metrics(),
|
|
}
|
|
with mock.patch.object(
|
|
docker_shadow, 'execute_private_adaptive_scan', return_value=outcome,
|
|
):
|
|
metrics = docker_shadow.private_adaptive_side_metrics(
|
|
db, {'target_scan_id': 7, 'normalized_target': 'private'},
|
|
limits=docker_limits(), checkpoint={'max_blobs': 1, 'max_bytes': 1024},
|
|
scan_policy_sha256='b' * 64,
|
|
scan_kwargs={'timeout_sec': 30},
|
|
)
|
|
|
|
self.assertEqual(db.target_scan_id, 7)
|
|
self.assertEqual(metrics['adaptive_routed_count'], 1)
|
|
self.assertEqual(metrics['routed_intersection_count'], 1)
|
|
self.assertEqual(metrics['adaptive_detector_count'], 1)
|
|
self.assertEqual(metrics['detector_intersection_count'], 0)
|
|
self.assertEqual(outcome, {})
|
|
|
|
def test_private_adaptive_side_retries_transient_infrastructure_failure(self):
|
|
class IdentityDB:
|
|
@staticmethod
|
|
def docker_adaptive_shadow_control_identities(_target_scan_id):
|
|
return frozenset(), frozenset()
|
|
|
|
outcome = {
|
|
'routed': frozenset(),
|
|
'detectors': frozenset(),
|
|
'selection_metrics': docker_shadow.empty_selection_metrics(),
|
|
'omitted_descriptor_count': 0,
|
|
'failure_count': 0,
|
|
'failure_metrics': docker_shadow.empty_failure_metrics(),
|
|
}
|
|
failure = scanner.DockerLayerInfrastructureError(
|
|
'transfer_timeout', 'private remote detail',
|
|
)
|
|
with mock.patch.object(
|
|
docker_shadow, 'execute_private_adaptive_scan',
|
|
side_effect=(failure, outcome),
|
|
) as execute:
|
|
metrics = docker_shadow.private_adaptive_side_metrics(
|
|
IdentityDB(), {
|
|
'target_scan_id': 8,
|
|
'normalized_target': 'owner/private@sha256:' + ('d' * 64),
|
|
},
|
|
limits=docker_limits(), checkpoint={'max_blobs': 1, 'max_bytes': 1024},
|
|
scan_policy_sha256='b' * 64,
|
|
scan_kwargs={'timeout_sec': 30, 'max_attempts': 3},
|
|
)
|
|
|
|
self.assertEqual(execute.call_count, 2)
|
|
self.assertEqual(metrics['failure_count'], 0)
|
|
self.assertEqual(outcome, {})
|
|
|
|
def test_private_adaptive_side_records_terminal_error_and_continues(self):
|
|
class IdentityDB:
|
|
@staticmethod
|
|
def docker_adaptive_shadow_control_identities(_target_scan_id):
|
|
return frozenset(), frozenset()
|
|
|
|
with mock.patch.object(
|
|
docker_shadow, 'execute_private_adaptive_scan',
|
|
side_effect=RuntimeError('private remote detail'),
|
|
) as execute, mock.patch('builtins.print') as output:
|
|
metrics = docker_shadow.private_adaptive_side_metrics(
|
|
IdentityDB(), {
|
|
'target_scan_id': 9,
|
|
'normalized_target': 'owner/private@sha256:' + ('e' * 64),
|
|
},
|
|
limits=docker_limits(), checkpoint={'max_blobs': 1, 'max_bytes': 1024},
|
|
scan_policy_sha256='b' * 64,
|
|
scan_kwargs={'timeout_sec': 30, 'max_attempts': 3},
|
|
)
|
|
|
|
self.assertEqual(execute.call_count, 3)
|
|
self.assertEqual(metrics['failure_count'], 1)
|
|
self.assertEqual(metrics['adaptive_routed_count'], 0)
|
|
self.assertEqual(metrics['adaptive_checkpoints'], 0)
|
|
rendered = ' '.join(str(call) for call in output.call_args_list)
|
|
self.assertIn('reason_code=adaptive_scan_error attempts=3', rendered)
|
|
self.assertNotIn('owner/private', rendered)
|
|
self.assertNotIn('private remote detail', rendered)
|
|
|
|
def test_scan_policy_includes_full_docker_execution_concurrency(self):
|
|
args = SimpleNamespace(drop_detectors=[])
|
|
base = {
|
|
'detectors': None,
|
|
'exclude_detectors': None,
|
|
'no_verification': False,
|
|
'trufflehog_config': None,
|
|
'trufflehog_concurrency': 0,
|
|
}
|
|
with mock.patch.object(console_runner.shutil, 'which', return_value='trufflehog'), \
|
|
mock.patch.object(console_runner, 'hash_file', return_value='a' * 64):
|
|
serial = console_runner.docker_layer_scan_policy_sha256(args, base)
|
|
parallel = console_runner.docker_layer_scan_policy_sha256(
|
|
args, dict(base, trufflehog_concurrency=3),
|
|
)
|
|
|
|
self.assertNotEqual(serial, parallel)
|
|
|
|
def test_shadow_runner_compares_exact_fifty_without_authoritative_writes(self):
|
|
events = []
|
|
|
|
class AggregateOnlyDB:
|
|
enabled = True
|
|
|
|
def __init__(self, **kwargs):
|
|
events.append(('database', kwargs))
|
|
self.checkpoints = 0
|
|
self.finished = None
|
|
|
|
def set_application_name(self, value):
|
|
events.append(('application', value))
|
|
|
|
def docker_adaptive_shadow_controls(self, policy, size, salt):
|
|
events.append(('controls', policy, size, salt))
|
|
return [{
|
|
'target_scan_id': value,
|
|
'normalized_target': f'private-control-{value}',
|
|
} for value in range(size)]
|
|
|
|
def start_docker_adaptive_shadow_report(self, *args, **kwargs):
|
|
events.append(('start-report', args, kwargs))
|
|
return {'report_token': 'report', 'lease_token': 'lease'}
|
|
|
|
def checkpoint_docker_adaptive_shadow_report(self, *args, **kwargs):
|
|
self.checkpoints += 1
|
|
|
|
def finish_docker_adaptive_shadow_report(self, *args, **kwargs):
|
|
self.finished = kwargs
|
|
return {
|
|
'report_id': 9, 'passed': True, 'completed_pairs': 50,
|
|
'routed_recall_ppm': 900000, 'slot_ratio_ppm': 400000,
|
|
}
|
|
|
|
def close(self):
|
|
events.append(('close',))
|
|
|
|
def __getattr__(self, name):
|
|
raise AssertionError(f'unexpected authoritative database access: {name}')
|
|
|
|
source_args = SimpleNamespace(
|
|
drop_detectors=[], trufflehog_job_memory_limit_bytes=1024,
|
|
timeout=30, detectors=None, exclude_detectors=None,
|
|
no_verification=False, trufflehog_config=None,
|
|
trufflehog_concurrency=1, docker_platform_os='linux',
|
|
docker_platform_arch='amd64', docker_layer_min_free_bytes=0,
|
|
)
|
|
config = {
|
|
'supervisor': {'docker_shadow': {
|
|
'enabled': True, 'cohort_size': 50, 'lease_seconds': 3600,
|
|
}},
|
|
'sources': {'dockerhub': {}},
|
|
'global': {},
|
|
}
|
|
db_holder = {}
|
|
|
|
def database_factory(**kwargs):
|
|
db_holder['db'] = AggregateOnlyDB(**kwargs)
|
|
return db_holder['db']
|
|
|
|
def timed(side, operation, checkpoint, **_kwargs):
|
|
events.append(('side', side))
|
|
metrics = operation()
|
|
checkpoint()
|
|
return metrics, 10 if side == 'full' else 4
|
|
|
|
def full_metrics(_db, control, _scan_kwargs):
|
|
self.assertIn('target_scan_id', control)
|
|
return {
|
|
'full_routed_count': 10,
|
|
'full_detector_count': 5,
|
|
'safety_regression_count': 0,
|
|
}
|
|
|
|
def adaptive_metrics(_db, control, **_kwargs):
|
|
self.assertIn('target_scan_id', control)
|
|
values = docker_shadow.empty_selection_metrics()
|
|
values.update({
|
|
'adaptive_routed_count': 9,
|
|
'routed_intersection_count': 9,
|
|
'adaptive_detector_count': 4,
|
|
'detector_intersection_count': 4,
|
|
'omitted_descriptor_count': 0,
|
|
'failure_count': 0,
|
|
})
|
|
values['selected_config'] = 1
|
|
values['adaptive_checkpoints'] = 1
|
|
return values
|
|
|
|
with contextlib.ExitStack() as patches:
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'require_active_supervisor_child',
|
|
side_effect=lambda *_args, **_kwargs: events.append(('authority',)),
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'load_config',
|
|
side_effect=lambda _path: events.append(('config',)) or config,
|
|
))
|
|
patches.enter_context(mock.patch.object(console_runner, 'apply_global_config'))
|
|
patches.enter_context(mock.patch.object(scanner, 'initialize_scanner_runtime'))
|
|
patches.enter_context(mock.patch.object(console_runner, 'load_secrets', return_value={}))
|
|
patches.enter_context(mock.patch.object(console_runner, 'configure_source_auth'))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'default_source_state', return_value={},
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'build_args_from_source_config', return_value=source_args,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'docker_layer_limits', return_value=docker_limits(),
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'docker_adaptive_checkpoint',
|
|
return_value={'max_blobs': 2, 'max_bytes': 2048},
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
console_runner, 'docker_layer_scan_policy_sha256', return_value='a' * 64,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'docker_layer_execution_policy_sha256', return_value='b' * 64,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'docker_layer_selection_policy_sha256', return_value='c' * 64,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'ScannerDB', side_effect=database_factory,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'run_timed_private_side', side_effect=timed,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'private_full_side_metrics', side_effect=full_metrics,
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
docker_shadow, 'private_adaptive_side_metrics', side_effect=adaptive_metrics,
|
|
))
|
|
patches.enter_context(mock.patch.object(scanner.docker_token_manager, 'cleanup'))
|
|
patches.enter_context(mock.patch.dict(
|
|
os.environ, {'TRUF_MANAGED_POSTGRES_DSN': 'postgresql://managed'}, clear=False,
|
|
))
|
|
output = patches.enter_context(mock.patch('builtins.print'))
|
|
exit_code = docker_shadow.run_shadow('config.yaml')
|
|
|
|
db = db_holder['db']
|
|
self.assertEqual(exit_code, 0)
|
|
self.assertEqual(events[:2], [('authority',), ('config',)])
|
|
self.assertEqual(db.checkpoints, 100)
|
|
self.assertEqual(db.finished['completed_pairs'], 50)
|
|
self.assertEqual(db.finished['full_routed_count'], 500)
|
|
self.assertEqual(db.finished['adaptive_routed_count'], 450)
|
|
self.assertEqual(db.finished['routed_intersection_count'], 450)
|
|
self.assertEqual(db.finished['full_slot_ms'], 500)
|
|
self.assertEqual(db.finished['adaptive_slot_ms'], 200)
|
|
self.assertEqual(db.finished['selection_metrics']['selected_config'], 50)
|
|
self.assertEqual(db.finished['selection_metrics']['adaptive_checkpoints'], 50)
|
|
sides = [event[1] for event in events if event[0] == 'side']
|
|
self.assertEqual(sides[:4], ['full', 'adaptive', 'adaptive', 'full'])
|
|
rendered = ' '.join(str(call) for call in output.call_args_list)
|
|
self.assertNotIn('private-control', rendered)
|
|
|
|
def test_adaptive_payload_selection_is_deterministic_bounded_and_honest(self):
|
|
def descriptor(position, char, size, kind='layer'):
|
|
return {
|
|
'digest': 'sha256:' + (char * 64),
|
|
'size': size,
|
|
'media_type': (
|
|
'application/vnd.oci.image.config.v1+json'
|
|
if kind == 'config'
|
|
else 'application/vnd.oci.image.layer.v1.tar+gzip'
|
|
),
|
|
'kind': kind,
|
|
'position': position,
|
|
}
|
|
|
|
descriptors = [
|
|
descriptor(0, '0', 100, 'config'),
|
|
descriptor(1, '1', 800),
|
|
descriptor(2, '2', 900),
|
|
descriptor(3, '3', 700),
|
|
descriptor(4, '3', 700),
|
|
descriptor(5, '5', 100),
|
|
]
|
|
classes = [
|
|
'copy_add', 'bulk_data', 'app_config_run', 'unknown', 'other_run',
|
|
]
|
|
limits = dict(docker_limits(), image_max_bytes=1600)
|
|
|
|
selected = scanner_db.select_docker_adaptive_payload(
|
|
descriptors, classes, limits, {descriptors[2]['digest']},
|
|
)
|
|
self.assertEqual(
|
|
[(entry[0]['position'], entry[1], entry[2], entry[3]) for entry in selected],
|
|
[
|
|
(0, True, 'config_selected', 'config'),
|
|
(1, True, 'selected_copy_add', 'copy_add'),
|
|
(2, True, 'already_covered', 'bulk_data'),
|
|
(3, True, 'selected_app_config_run', 'app_config_run'),
|
|
(4, True, 'duplicate_digest', 'unknown'),
|
|
(5, False, 'layer_limit_exhausted', 'other_run'),
|
|
],
|
|
)
|
|
self.assertEqual(
|
|
selected,
|
|
scanner_db.select_docker_adaptive_payload(
|
|
descriptors, classes, limits, {descriptors[2]['digest']},
|
|
),
|
|
)
|
|
|
|
all_fit = scanner_db.select_docker_adaptive_payload(
|
|
descriptors, classes,
|
|
dict(limits, image_max_bytes=4096, max_layers=5),
|
|
)
|
|
self.assertTrue(all(entry[1] for entry in all_fit))
|
|
|
|
def test_history_classification_is_bounded_aligned_and_secret_free(self):
|
|
secret = 'must-not-survive-classification'
|
|
config = {
|
|
'rootfs': {'diff_ids': ['sha256:' + char * 64 for char in '12345']},
|
|
'history': [
|
|
{'created_by': 'RUN apt-get install -y curl'},
|
|
{'created_by': 'metadata only', 'empty_layer': True},
|
|
{'created_by': '#(nop) COPY dir:abc in /app'},
|
|
{'created_by': f'RUN printf {secret} > /app/config.json'},
|
|
{'created_by': 'RUN huggingface-cli download x/model.safetensors'},
|
|
{'created_by': 'RUN echo complete'},
|
|
],
|
|
}
|
|
|
|
classes = scanner.docker_config_payload_classes(config, 5)
|
|
|
|
self.assertEqual(classes, [
|
|
'package_run', 'copy_add', 'app_config_run', 'bulk_data', 'other_run',
|
|
])
|
|
self.assertNotIn(secret, json.dumps(classes))
|
|
self.assertEqual(
|
|
scanner.docker_config_payload_classes(config, 4),
|
|
['unknown'] * 4,
|
|
)
|
|
oversized = 'RUN ' + ('x' * scanner.DOCKER_CONFIG_HISTORY_COMMAND_MAX_CHARS)
|
|
self.assertEqual(scanner.docker_history_payload_class(oversized), 'unknown')
|
|
|
|
def test_adaptive_policy_hashes_separate_execution_selection_and_scheduling(self):
|
|
limits = docker_limits()
|
|
selection_changed = dict(limits, image_max_bytes=4096)
|
|
scheduling_changed = dict(limits, blob_timeout_sec=20, blob_max_attempts=4)
|
|
execution_changed = dict(limits, archive_max_depth=3)
|
|
scan_policy = 'b' * 64
|
|
|
|
execution = scanner_db.docker_layer_execution_policy_sha256(scan_policy, limits)
|
|
selection = scanner_db.docker_layer_selection_policy_sha256(limits)
|
|
|
|
self.assertEqual(
|
|
execution,
|
|
scanner_db.docker_layer_execution_policy_sha256(
|
|
scan_policy, selection_changed,
|
|
),
|
|
)
|
|
self.assertEqual(
|
|
execution,
|
|
scanner_db.docker_layer_execution_policy_sha256(
|
|
scan_policy, scheduling_changed,
|
|
),
|
|
)
|
|
self.assertNotEqual(
|
|
execution,
|
|
scanner_db.docker_layer_execution_policy_sha256(
|
|
scan_policy, execution_changed,
|
|
),
|
|
)
|
|
self.assertNotEqual(
|
|
selection,
|
|
scanner_db.docker_layer_selection_policy_sha256(selection_changed),
|
|
)
|
|
self.assertEqual(
|
|
selection,
|
|
scanner_db.docker_layer_selection_policy_sha256(scheduling_changed),
|
|
)
|
|
|
|
def test_version_two_plan_and_assignment_are_exact_and_stable(self):
|
|
plan = docker_v2_plan()
|
|
plan_sha256 = hashlib.sha256(
|
|
scanner_db.canonical_docker_layer_plan_bytes(plan)
|
|
).hexdigest()
|
|
execution = scanner_db.validate_docker_layer_execution({
|
|
'version': 2,
|
|
'plan_sha256': plan_sha256,
|
|
'blobs': [{
|
|
'digest': plan['descriptors'][0]['digest'],
|
|
'lease_token': plan['descriptors'][0]['lease_token'],
|
|
'status': 'covered',
|
|
'verified_bytes': plan['descriptors'][0]['size'],
|
|
'transfer_bytes': plan['descriptors'][0]['size'],
|
|
'transfer_duration_ms': 1,
|
|
'scan_duration_ms': 1,
|
|
'finding_count': 0,
|
|
'error_code': None,
|
|
}],
|
|
}, plan, plan_sha256)
|
|
self.assertEqual(execution['version'], 2)
|
|
|
|
target = plan['image']
|
|
args = SimpleNamespace(
|
|
docker_content_scan_mode='adaptive-canary',
|
|
docker_adaptive_canary_basis_points=4173,
|
|
)
|
|
assignment = console_runner.docker_layer_effective_mode(
|
|
args, target, adaptive_gate_passed=True,
|
|
selection_policy_sha256=plan['selection_policy_sha256'],
|
|
)
|
|
self.assertEqual(
|
|
assignment,
|
|
console_runner.docker_layer_effective_mode(
|
|
args, target, adaptive_gate_passed=True,
|
|
selection_policy_sha256=plan['selection_policy_sha256'],
|
|
),
|
|
)
|
|
self.assertEqual(
|
|
console_runner.docker_layer_effective_mode(args, target)[0], 'full',
|
|
)
|
|
|
|
malformed = dict(plan, unexpected='value')
|
|
with self.assertRaises(ValueError):
|
|
scanner_db.validate_docker_layer_plan(malformed)
|
|
|
|
|
|
class DockerLayerExecutionTests(unittest.TestCase):
|
|
def run_config_scan(
|
|
self, payload=b'{}', stderr=None, stdout='', stream_error=None, plan=None,
|
|
):
|
|
plan = plan or docker_plan(payload)
|
|
deadline = time.monotonic() + 30
|
|
work = layer_work(plan, deadline=deadline)
|
|
stderr = stderr if stderr is not None else json.dumps({
|
|
'level': 'info-0', 'msg': 'finished scanning',
|
|
})
|
|
temporary = tempfile.TemporaryDirectory()
|
|
self.addCleanup(temporary.cleanup)
|
|
|
|
def download(_repository, _descriptor, destination, bearer_auth, **_kwargs):
|
|
if stream_error is not None:
|
|
raise stream_error
|
|
with open(destination, 'wb') as output:
|
|
output.write(payload)
|
|
return scanner.DockerBlobDownloadOutcome(
|
|
destination, len(payload), len(payload), 1, bearer_auth,
|
|
)
|
|
|
|
patches = contextlib.ExitStack()
|
|
self.addCleanup(patches.close)
|
|
patches.enter_context(mock.patch.object(scanner, 'get_work_dir', return_value=temporary.name))
|
|
patches.enter_context(mock.patch.object(scanner, 'harden_private_directory'))
|
|
patches.enter_context(mock.patch.object(scanner, 'write_temp_owner'))
|
|
patches.enter_context(mock.patch.object(scanner, 'durable_unlink', side_effect=unlink))
|
|
patches.enter_context(mock.patch.object(
|
|
scanner, 'cleanup_command_work_dir', side_effect=lambda path: shutil.rmtree(path),
|
|
))
|
|
patches.enter_context(mock.patch.object(
|
|
scanner, 'stream_docker_registry_blob', side_effect=download,
|
|
))
|
|
command = patches.enter_context(mock.patch.object(
|
|
scanner, 'run_command_streamed',
|
|
side_effect=lambda *_args, **_kwargs: scanner.streamed_output_from_text(
|
|
stdout=stdout, stderr=stderr, returncode=0,
|
|
),
|
|
))
|
|
return plan, work, deadline, command
|
|
|
|
def test_warning_preserves_findings_but_never_marks_blob_covered(self):
|
|
for message, error_code in (
|
|
('a detector ignored the context timeout', 'detector_timeout'),
|
|
('non-critical error processing chunk', 'chunk_processing'),
|
|
):
|
|
with self.subTest(error_code=error_code):
|
|
stderr = '\n'.join((
|
|
json.dumps({'level': 'error', 'msg': message}),
|
|
json.dumps({'level': 'info-0', 'msg': 'finished scanning'}),
|
|
))
|
|
finding = json.dumps({
|
|
'DetectorName': 'Fixture', 'SourceMetadata': {'Data': {}},
|
|
})
|
|
plan, work, _, _ = self.run_config_scan(
|
|
stderr=stderr, stdout=finding,
|
|
)
|
|
|
|
result = scanner.scan_docker_layer_plan(plan['image'], work)
|
|
|
|
record = result['docker_layer_execution']['blobs'][0]
|
|
self.assertEqual(record['status'], 'terminal_failed')
|
|
self.assertEqual(record['error_code'], error_code)
|
|
self.assertEqual(len(result['findings']), 1)
|
|
self.assertEqual(record['finding_count'], len(result['findings']))
|
|
self.assertFalse(result['retryable'])
|
|
self.assertFalse(
|
|
result['scan_meta']['docker_layer_scope']['coverage_complete']
|
|
)
|
|
|
|
def test_finding_filter_target_logging_can_be_suppressed(self):
|
|
private_target = 'owner/private@sha256:' + ('f' * 64)
|
|
with mock.patch.object(
|
|
scanner, 'filter_dropped_detectors',
|
|
return_value=([], 1, {'Fixture': 1}),
|
|
), mock.patch.object(
|
|
scanner, 'filter_noisy_findings', return_value=([], 0),
|
|
), mock.patch.object(scanner.logger, 'info') as log:
|
|
scanner.apply_finding_filters(
|
|
{'findings': [{}]}, private_target, log_target=False,
|
|
)
|
|
private_message = str(log.call_args.args[0])
|
|
scanner.apply_finding_filters({'findings': [{}]}, private_target)
|
|
production_message = str(log.call_args.args[0])
|
|
|
|
self.assertNotIn(private_target, private_message)
|
|
self.assertIn(private_target, production_message)
|
|
|
|
def test_invalid_configuration_is_terminal_and_not_executed(self):
|
|
plan, work, _, command = self.run_config_scan(payload=b'not-json')
|
|
|
|
result = scanner.scan_docker_layer_plan(plan['image'], work)
|
|
|
|
record = result['docker_layer_execution']['blobs'][0]
|
|
self.assertEqual(record['status'], 'terminal_failed')
|
|
self.assertEqual(record['error_code'], 'invalid_config_json')
|
|
command.assert_not_called()
|
|
|
|
def test_layer_archive_media_is_prevalidated_before_execution(self):
|
|
payload = b'not-a-gzip-layer'
|
|
layer = {
|
|
'digest': 'sha256:' + hashlib.sha256(payload).hexdigest(),
|
|
'size': len(payload),
|
|
'media_type': 'application/vnd.oci.image.layer.v1.tar+gzip',
|
|
'kind': 'layer',
|
|
'position': 1,
|
|
'selected': True,
|
|
'selection_reason': 'layer_selected',
|
|
'coverage_state': 'leased',
|
|
'lease_token': 'm' * 43,
|
|
'attempt': 1,
|
|
'max_attempts': 3,
|
|
}
|
|
plan = docker_plan(b'{}', coverage_state='skipped', extra_descriptors=(layer,))
|
|
plan, work, _, command = self.run_config_scan(payload=payload, plan=plan)
|
|
|
|
result = scanner.scan_docker_layer_plan(plan['image'], work)
|
|
|
|
record = result['docker_layer_execution']['blobs'][0]
|
|
self.assertEqual(record['status'], 'terminal_failed')
|
|
self.assertEqual(record['error_code'], 'invalid_layer_archive')
|
|
command.assert_not_called()
|
|
|
|
def test_process_deadline_never_exceeds_original_absolute_deadline(self):
|
|
plan, work, deadline, command = self.run_config_scan()
|
|
|
|
result = scanner.scan_docker_layer_plan(plan['image'], work)
|
|
|
|
self.assertEqual(result['docker_layer_execution']['blobs'][0]['status'], 'covered')
|
|
self.assertLessEqual(command.call_args.kwargs['deadline'], deadline)
|
|
|
|
def test_each_separate_blob_transfer_emits_downloading_then_scanning(self):
|
|
config_payload = b'{}'
|
|
layer_payload = b'layer payload'
|
|
layer = {
|
|
'digest': 'sha256:' + hashlib.sha256(layer_payload).hexdigest(),
|
|
'size': len(layer_payload),
|
|
'media_type': 'application/vnd.oci.image.layer.v1.tar',
|
|
'kind': 'layer', 'position': 1, 'selected': True,
|
|
'selection_reason': 'layer_selected', 'coverage_state': 'leased',
|
|
'lease_token': 'm' * 43, 'attempt': 1, 'max_attempts': 3,
|
|
}
|
|
plan = docker_plan(
|
|
config_payload, coverage_state='leased', extra_descriptors=(layer,),
|
|
)
|
|
plan, work, _, _ = self.run_config_scan(
|
|
payload=config_payload, plan=plan,
|
|
)
|
|
payloads = {
|
|
plan['descriptors'][0]['digest']: config_payload,
|
|
layer['digest']: layer_payload,
|
|
}
|
|
phases = []
|
|
|
|
def download(_repository, descriptor, destination, bearer_auth, **_kwargs):
|
|
payload = payloads[descriptor['digest']]
|
|
with open(destination, 'wb') as output:
|
|
output.write(payload)
|
|
return scanner.DockerBlobDownloadOutcome(
|
|
destination, len(payload), len(payload), 1, bearer_auth,
|
|
)
|
|
|
|
with scanner.client_scan_phase_events(
|
|
lambda phase, progress=None: phases.append(
|
|
(phase, dict(progress or {}), time.monotonic()),
|
|
),
|
|
), mock.patch.object(
|
|
scanner, 'stream_docker_registry_blob', side_effect=download,
|
|
), mock.patch.object(
|
|
scanner, '_scan_docker_content_file',
|
|
return_value={'findings': [], 'errors': []},
|
|
):
|
|
scanner.scan_docker_layer_plan(plan['image'], work)
|
|
names = [item[0] for item in phases]
|
|
self.assertEqual(
|
|
names[:4],
|
|
['downloading', 'scanning', 'downloading', 'scanning'],
|
|
)
|
|
previous = WorkerPhase.WAITING_PERMIT
|
|
for name in names:
|
|
current = WorkerPhase(name)
|
|
validate_phase_transition(previous, current)
|
|
previous = current
|
|
measured = [
|
|
later[2] - earlier[2]
|
|
for earlier, later in zip(phases, phases[1:])
|
|
]
|
|
self.assertTrue(all(duration >= 0 for duration in measured))
|
|
|
|
def test_version_two_execution_and_finding_provenance_remain_versioned(self):
|
|
finding = json.dumps({
|
|
'DetectorName': 'OpenAI',
|
|
'Raw': 'sk-' + ('a' * 32),
|
|
'SourceMetadata': {'Data': {'Filesystem': {'file': 'config.json'}}},
|
|
})
|
|
plan = docker_v2_plan()
|
|
plan, work, _, _ = self.run_config_scan(plan=plan, stdout=finding)
|
|
|
|
result = scanner.scan_docker_layer_plan(plan['image'], work)
|
|
|
|
self.assertEqual(result['docker_layer_execution']['version'], 2)
|
|
docker_content = result['findings'][0]['SourceMetadata']['Data']['DockerContent']
|
|
self.assertEqual(docker_content['payload_class'], 'config')
|
|
self.assertEqual(docker_content['descriptor_kind'], 'config')
|
|
|
|
def test_full_docker_invocation_shape_is_exact_and_target_log_is_optional(self):
|
|
target = 'owner/repo@sha256:' + ('a' * 64)
|
|
completion = json.dumps({'level': 'info-0', 'msg': 'finished scanning'})
|
|
with mock.patch.object(
|
|
scanner, 'run_command_streamed',
|
|
return_value=scanner.streamed_output_from_text(stderr=completion),
|
|
) as command, mock.patch.object(scanner.logger, 'info') as log:
|
|
scanner.scan_docker_image(
|
|
target, detectors='OpenAI,Github', exclude_detectors='AWS',
|
|
no_verification=True, trufflehog_config='policy.yaml',
|
|
config_dir='private-config', trufflehog_concurrency=3,
|
|
log_target=False,
|
|
)
|
|
|
|
self.assertEqual(command.call_args.args[0], [
|
|
scanner.get_trufflehog_cmd(), 'docker', '--image', target,
|
|
'--json', '--no-update', '--local-dev', '--log-level', '2', '--concurrency', '3',
|
|
'--config', 'policy.yaml', '--include-detectors', 'OpenAI,Github',
|
|
'--exclude-detectors', 'AWS', '--no-verification',
|
|
])
|
|
self.assertEqual(command.call_args.args[2]['DOCKER_CONFIG'], 'private-config')
|
|
log.assert_not_called()
|
|
|
|
def test_source_outage_propagates_for_fenced_refund(self):
|
|
outage = scanner.DockerLayerInfrastructureError(
|
|
'remote_rate_limit', 'registry unavailable', category='docker_rate_limit',
|
|
)
|
|
plan, work, _, _ = self.run_config_scan(stream_error=outage)
|
|
|
|
with self.assertRaises(scanner.DockerLayerInfrastructureError) as raised:
|
|
scanner.scan_target_result(plan['image'], 'docker', 'event-id', {
|
|
'docker_layer_work': work,
|
|
})
|
|
self.assertEqual(raised.exception.category, 'docker_rate_limit')
|
|
|
|
def test_process_launch_failure_cannot_stage_missing_execution(self):
|
|
plan, work, _, command = self.run_config_scan()
|
|
command.side_effect = OSError('launch detail must not be staged')
|
|
|
|
with self.assertRaises(scanner.DockerLayerInfrastructureError):
|
|
scanner.scan_target_result(plan['image'], 'docker', 'event-id', {
|
|
'docker_layer_work': work,
|
|
})
|
|
|
|
|
|
class DockerRunnerRolloutTests(unittest.TestCase):
|
|
def test_parent_attempt_resets_only_without_failed_blob_work(self):
|
|
pending = {
|
|
'docker_layer_plan': {
|
|
'descriptors': [{'coverage_state': 'shared_pending'}],
|
|
},
|
|
'docker_layer_execution': {'blobs': []},
|
|
'errors': ['shared wait'],
|
|
'retryable': True,
|
|
}
|
|
timed_out = {
|
|
'docker_layer_plan': {
|
|
'descriptors': [{'coverage_state': 'leased'}],
|
|
},
|
|
'docker_layer_execution': {
|
|
'blobs': [{'status': 'retryable_failed', 'error_code': 'transfer_timeout'}],
|
|
},
|
|
'errors': ['timeout'],
|
|
'retryable': True,
|
|
}
|
|
|
|
self.assertTrue(console_runner.queue_result_resets_attempts(pending))
|
|
self.assertFalse(console_runner.queue_result_resets_attempts(timed_out))
|
|
args = SimpleNamespace(
|
|
docker_layer_checkpoint_delay_sec=1,
|
|
target_retry_max_attempts=3,
|
|
)
|
|
disposition = console_runner.docker_layer_queue_disposition(
|
|
timed_out, args, attempts=2,
|
|
)
|
|
self.assertEqual((disposition[0], disposition[2]), ('deferred', False))
|
|
disposition = console_runner.docker_layer_queue_disposition(
|
|
timed_out, args, attempts=3,
|
|
)
|
|
self.assertEqual(disposition, ('failed', None, False))
|
|
disposition = console_runner.docker_layer_queue_disposition(
|
|
pending, args, attempts=3,
|
|
)
|
|
self.assertEqual((disposition[0], disposition[2]), ('deferred', True))
|
|
|
|
def test_remote_failure_classification_is_target_or_source_scoped(self):
|
|
args = SimpleNamespace(timeout=10, docker_platform_os='linux', docker_platform_arch='amd64')
|
|
claim = {'target': 'owner/repo@sha256:' + ('a' * 64)}
|
|
cases = {
|
|
'auth_failed': (True, True, 'docker_auth'),
|
|
'rate_limited': (True, True, 'docker_rate_limit'),
|
|
'remote_transient': (True, True, 'remote_transient'),
|
|
'target_forbidden': (False, False, 'docker_target_forbidden'),
|
|
}
|
|
for status, expected in cases.items():
|
|
with self.subTest(status=status), mock.patch.object(
|
|
console_runner, 'resolve_docker_content_manifest',
|
|
side_effect=scanner.DockerRemoteAccessError('safe', status=status),
|
|
):
|
|
with self.assertRaises(console_runner._DockerPlanningFailure) as raised:
|
|
console_runner.resolve_and_bind_docker_claim(
|
|
args, 'unused', 'dockerhub', claim, {}, 'b' * 64,
|
|
)
|
|
failure = raised.exception
|
|
self.assertEqual(
|
|
(failure.retryable, failure.source_failure, failure.category), expected,
|
|
)
|
|
|
|
def test_runner_preserves_deadline_across_resolution_binding_and_work(self):
|
|
plan = docker_plan()
|
|
deadline = time.monotonic() + 30
|
|
resolved = {'fixture': True}
|
|
|
|
class PlanningDB:
|
|
enabled = True
|
|
|
|
def set_application_name(self, _name):
|
|
return None
|
|
|
|
def bind_docker_layer_plan(self, *_args, **_kwargs):
|
|
return plan
|
|
|
|
def close(self):
|
|
return None
|
|
|
|
args = SimpleNamespace(timeout=10)
|
|
claim = {
|
|
'target': plan['image'],
|
|
'reservation_id': 1,
|
|
'claim_lease_token': 'claim-token',
|
|
}
|
|
with mock.patch.object(
|
|
console_runner, 'resolve_docker_content_manifest',
|
|
return_value=(resolved, scanner.DockerRegistryAuth(token='opaque')),
|
|
) as resolve, mock.patch.object(
|
|
console_runner, 'ScannerDB', return_value=PlanningDB(),
|
|
):
|
|
work = console_runner.resolve_and_bind_docker_claim(
|
|
args, 'postgresql://unused', 'dockerhub', claim, {}, 'b' * 64,
|
|
deadline=deadline,
|
|
)
|
|
|
|
self.assertEqual(resolve.call_args.kwargs['deadline'], deadline)
|
|
self.assertEqual(work['deadline'], deadline)
|
|
|
|
def test_canary_assignment_is_stable_and_zero_falls_back_to_full(self):
|
|
target = 'owner/repo@sha256:' + ('a' * 64)
|
|
zero = SimpleNamespace(
|
|
docker_content_scan_mode='canary', docker_layer_canary_basis_points=0,
|
|
)
|
|
sampled = SimpleNamespace(
|
|
docker_content_scan_mode='canary', docker_layer_canary_basis_points=4173,
|
|
)
|
|
all_eligible = SimpleNamespace(
|
|
docker_content_scan_mode='canary', docker_layer_canary_basis_points=10000,
|
|
)
|
|
|
|
self.assertEqual(console_runner.docker_layer_effective_mode(zero, 'not-a-digest')[0], 'full')
|
|
self.assertEqual(
|
|
console_runner.docker_layer_effective_mode(all_eligible, target, False)[0],
|
|
'full',
|
|
)
|
|
self.assertEqual(
|
|
console_runner.docker_layer_effective_mode(all_eligible, target, True)[0],
|
|
'layer',
|
|
)
|
|
self.assertEqual(
|
|
console_runner.docker_layer_effective_mode(sampled, target, True),
|
|
console_runner.docker_layer_effective_mode(sampled, target, True),
|
|
)
|
|
config = console_runner.load_config(str(APP_DIR / 'config.yaml'))['sources']['dockerhub']
|
|
self.assertEqual(config['docker_content_scan_mode'], 'canary')
|
|
self.assertEqual(config['docker_layer_canary_basis_points'], 10000)
|
|
|
|
def test_canary_eligibility_uses_dedicated_fenced_database_check(self):
|
|
calls = []
|
|
|
|
class EligibilityDB:
|
|
enabled = True
|
|
|
|
def set_application_name(self, name):
|
|
calls.append(('application', name))
|
|
|
|
def docker_layer_timeout_canary_eligible(self, reservation_id, token):
|
|
calls.append(('eligible', reservation_id, token))
|
|
return True
|
|
|
|
def close(self):
|
|
calls.append(('close',))
|
|
|
|
claim = {'reservation_id': 17, 'claim_lease_token': 'claim-token'}
|
|
with mock.patch.object(console_runner, 'ScannerDB', return_value=EligibilityDB()):
|
|
eligible = console_runner.docker_layer_timeout_canary_eligible(
|
|
'postgresql://unused', 'dockerhub', claim,
|
|
)
|
|
self.assertTrue(eligible)
|
|
self.assertEqual(calls, [
|
|
('application', 'truf-docker-layer-canary:dockerhub'),
|
|
('eligible', 17, 'claim-token'),
|
|
('close',),
|
|
])
|
|
|
|
def test_planning_diagnostics_redact_tokens(self):
|
|
secret = 'diagnostic-secret-token'
|
|
failure = console_runner._DockerPlanningFailure(
|
|
RuntimeError('failed with ' + secret),
|
|
retryable=True, source_failure=True, category='docker_auth',
|
|
)
|
|
args = SimpleNamespace(
|
|
token=secret, docker_token=secret, docker_layer_canary_basis_points=10,
|
|
)
|
|
claim = {
|
|
'target': 'owner/repo@sha256:' + ('a' * 64),
|
|
'scan_event_id': 'event-id',
|
|
}
|
|
|
|
result = console_runner.docker_planning_failure_result(
|
|
args, claim, {'token': secret}, failure, 'canary', 'layer',
|
|
)
|
|
|
|
self.assertNotIn(secret, json.dumps(result, sort_keys=True))
|
|
|
|
|
|
if __name__ == '__main__':
|
|
unittest.main()
|