import sys sys.dont_write_bytecode = True import json import gzip import os import tarfile import tempfile import uuid import zipfile from datetime import datetime, timedelta, timezone from types import SimpleNamespace import scanner as scanner_module from console_runner import queue_error_disposition, target_retry_delay_sec from db_backend import redact_database_url from scanner import ( apply_trufflehog_diagnostics, build_authenticated_git_url, cached_gharchive_hour, convert_package_git_unavailable_to_skip, docker_tag_platform_support, fetch_dockerhub_images, normalize_git_repo_candidate, redact_command_args, run_command, safe_extract_package_zip, safe_extract_tar, ) from scanner_db import ScannerDB, normalize_target, target_status from runtime_security import ensure_private_directory, harden_private_file def diagnostic(message, error, **extra): return json.dumps({'level': 'error', 'msg': message, 'error': error, **extra}) def assert_diagnostic_policy(): result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('non-critical error processing chunk', 'error reading chunk: brotli: PADDING_2'), 0, 'pypi', ) assert not result.get('errors') assert result.get('warnings') and result.get('degraded') assert target_status(result) == 'degraded' result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics(result, diagnostic('Skipping result: invalid', 'empty raw'), 0, 'npm') assert not result.get('errors') and result.get('warnings') result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('error reading chunk', 'read tcp: connection reset by peer'), 0, 'pypi', ) assert result.get('errors') and result.get('retryable') is True assert result.get('error_class') == 'network' for detail in ('unexpected EOF', 'permission denied', 'no space left on device'): result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics(result, diagnostic('error reading chunk', detail), 0, 'pypi') assert result.get('errors'), detail result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics(result, diagnostic('error reading chunk', 'brotli: PADDING_2'), 0, 'git') assert result.get('errors') result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics(result, diagnostic('space', 'no repo found for repo'), 0, 'huggingface') assert result.get('skipped') and not result.get('errors') assert target_status(result) == 'skipped' result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('error processing image', 'no child with platform linux/amd64 in index image:tag'), 1, 'docker', ) assert result.get('skipped') and not result.get('errors') result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('error processing layer', 'gzip: invalid header') + '\n' + json.dumps({'level': 'info-0', 'msg': 'finished scanning'}), 0, 'docker', ) assert result.get('degraded') and not result.get('errors') result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('a detector ignored the context timeout', 'context deadline exceeded') + '\n' + json.dumps({'level': 'info-0', 'msg': 'finished scanning'}), 0, 'docker', ) assert result.get('degraded') and not result.get('errors') assert result.get('warning_classes') == ['detector_timeout'] result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics( result, diagnostic('error processing layer', 'gzip: invalid header'), 1, 'docker', ) assert result.get('errors') and not result.get('degraded') result = {'findings': [], 'errors': []} apply_trufflehog_diagnostics(result, '', 2, 'filesystem') assert result.get('errors') and result.get('retryable') is True result['findings'] = [{'DetectorName': 'Example'}] assert target_status(result) == 'error' result = {'errors': [diagnostic('error running scan', 'remote: Repository not found.')], 'findings': []} convert_package_git_unavailable_to_skip(result) assert result.get('skipped') and not result.get('errors') def assert_docker_platform_policy(): assert docker_tag_platform_support({'images': [{'os': 'linux', 'architecture': 'amd64'}]}) is True assert docker_tag_platform_support({'images': [{'os': 'linux', 'architecture': 'arm64'}]}) is False assert docker_tag_platform_support({'images': []}) is None assert docker_tag_platform_support({'images': [{'os': 'linux'}]}) is None assert docker_tag_platform_support({'images': [{'os': 'unknown', 'architecture': 'unknown'}]}) is None assert docker_tag_platform_support({'images': [{'os': 'linux', 'architecture': 'arm64'}, {'os': 'unknown', 'architecture': 'unknown'}]}) is None assert docker_tag_platform_support({}) is None assert normalize_target('Owner/Repo:Prod', 'docker') == 'owner/repo:prod' assert normalize_git_repo_candidate('https://[invalid url, do not cite]/repo') is None def assert_docker_partial_pagination_policy(): class Response: def __init__(self, page): self.page = page def raise_for_status(self): if self.page > 1: raise RuntimeError('page outside result set') def json(self): return {'count': 1, 'results': [{'repo_name': 'owner/repo'}]} original = scanner_module.api_request try: def fake_request(method, url, **kwargs): page = int(url.split('page=', 1)[1].split('&', 1)[0]) return Response(page) scanner_module.api_request = fake_request assert fetch_dockerhub_images('test', 3, per_page=10, resolve_tags=False) == ['owner/repo'] finally: scanner_module.api_request = original def assert_retry_policy(): assert target_retry_delay_sec(1, 60, 3600) == 60 assert target_retry_delay_sec(2, 60, 3600) == 120 assert target_retry_delay_sec(10, 60, 3600) == 3600 args = SimpleNamespace(target_retry_max_attempts=3, target_timeout_retry_delay_sec=21600) status, available_after, attempts, max_attempts = queue_error_disposition( None, 'github', 'github', 'https://github.com/o/r', {'errors': ['timeout'], 'error_class': 'timeout', 'scan_meta': {'command_timed_out': True}}, args, {'attempts': 3}, ) assert status == 'failed' and available_after is None and attempts == 3 and max_attempts == 3 def assert_timeout_output_policy(): marker = '{"DetectorName":"PartialFinding"}' old_work_dir = scanner_module.scan_config.work_dir old_min_free = scanner_module.scan_config.min_free_gb old_authority_check = scanner_module.require_trufflehog_launch_authority with tempfile.TemporaryDirectory() as temp_dir: work_dir = os.path.join(temp_dir, 'work') ensure_private_directory(work_dir, reject_reparse=True) scanner_module.initialize_scanner_runtime(preflight_complete=True, register_cleanup=False) scanner_module.scan_config.work_dir = work_dir scanner_module.scan_config.min_free_gb = 0 scanner_module.require_trufflehog_launch_authority = lambda _command: None try: stdout, stderr, returncode = run_command( [sys.executable, '-c', f'import time; print({marker!r}, flush=True); time.sleep(5)'], 1, ) finally: scanner_module.scan_config.work_dir = old_work_dir scanner_module.scan_config.min_free_gb = old_min_free scanner_module.require_trufflehog_launch_authority = old_authority_check assert returncode == -1 and marker in stdout and 'timed out' in stderr.lower() def assert_archive_and_secret_safety(): url, secrets = build_authenticated_git_url('https://attacker.example/repo.git', 'github', 'sentinel-token') assert url == 'https://attacker.example/repo.git' and not secrets and 'sentinel-token' not in url redacted = redact_command_args(['git', 'https://x-access-token:sentinel-token@github.com/org/repo.git', '--token', 'sentinel-token']) assert all('sentinel-token' not in value for value in redacted) with tempfile.TemporaryDirectory() as temp_dir: tar_path = os.path.join(temp_dir, 'unsafe.tar') with tarfile.open(tar_path, 'w') as archive: link = tarfile.TarInfo('link') link.type = tarfile.SYMTYPE link.linkname = '..' archive.addfile(link) try: safe_extract_tar(tar_path, os.path.join(temp_dir, 'tar-out')) raise AssertionError('unsafe tar link was accepted') except ValueError: pass zip_path = os.path.join(temp_dir, 'large.zip') with zipfile.ZipFile(zip_path, 'w') as archive: archive.writestr('large.txt', b'x' * 2048) try: safe_extract_package_zip( zip_path, os.path.join(temp_dir, 'zip-out'), max_total_size_mb=0, max_file_size_mb=0, ) except ValueError: raise AssertionError('disabled package zip bounds rejected a valid archive') crowded_tar = os.path.join(temp_dir, 'crowded.tar') with tarfile.open(crowded_tar, 'w') as archive: archive.addfile(tarfile.TarInfo('one')) archive.addfile(tarfile.TarInfo('two')) try: safe_extract_tar(crowded_tar, os.path.join(temp_dir, 'crowded-out'), max_files=1) raise AssertionError('tar member bound was not enforced') except ValueError: pass assert 'query-secret' not in redact_database_url('postgresql://u@localhost/db?password=query-secret') malformed = ScannerDB(db_path=os.path.join(tempfile.gettempdir(), 'must-not-open.db'), db_url='postgreql://bad') assert malformed.postgres_required and not malformed.enabled and malformed.path is None with tempfile.TemporaryDirectory() as temp_dir: runtime_dir = os.path.join(temp_dir, 'runtime') cache_dir = os.path.join(runtime_dir, 'gharchive') ensure_private_directory(runtime_dir) ensure_private_directory(cache_dir) hour = datetime(2026, 7, 11, 12, tzinfo=timezone.utc) cache_path = os.path.join(cache_dir, '2026-07-11-12.json.gz') with gzip.open(cache_path, 'wb') as archive: archive.write(b'{"type":"PushEvent"}\n') harden_private_file(cache_path) old_runtime = scanner_module.scan_config.runtime_dir old_free = scanner_module.scan_config.gharchive_cache_min_free_bytes try: scanner_module.scan_config.runtime_dir = runtime_dir scanner_module.scan_config.gharchive_cache_min_free_bytes = 0 assert scanner_module.canonical_path( cached_gharchive_hour(hour, cache_dir, request_timeout=1, retries=1) ) == scanner_module.canonical_path(cache_path) finally: scanner_module.scan_config.runtime_dir = old_runtime scanner_module.scan_config.gharchive_cache_min_free_bytes = old_free def assert_queue_policy(): old_urls = {key: os.environ.pop(key, None) for key in ('SCANNER_DB_URL', 'DATABASE_URL')} try: with tempfile.TemporaryDirectory() as temp_dir: db = ScannerDB(db_path=os.path.join(temp_dir, 'scanner.db')) try: assert db.enabled target = 'owner/repo:tag' db.enqueue_targets('dockerhub', 'docker', 'q', [target]) claimed = db.claim_targets('dockerhub', 'docker', 1, 'owner-a', 3600, return_rows=True) assert [row['target'] for row in claimed] == [target] first_claim = claimed[0] row = db.target_queue_item('dockerhub', 'docker', target) assert row['status'] == 'in_progress' and int(row['attempts']) == 1 assert db.reclaim_target_leases('dockerhub', 'docker', 'owner-b') == 0 future = (datetime.now(timezone.utc) + timedelta(hours=1)).isoformat(timespec='seconds') db.complete_target_queue_item( 'dockerhub', 'docker', target, None, 'deferred', 'network', future, queue_id=first_claim['id'], lease_token=first_claim['lease_token'], ) db.sync_target_queue_from_files('dockerhub', 'docker', [target], [], 'q') row = db.target_queue_item('dockerhub', 'docker', target) assert row['status'] == 'deferred' and row['available_after'] == future updated_at = db.conn.execute( 'SELECT updated_at FROM target_queue WHERE source = ? AND normalized_target = ?', ('dockerhub', normalize_target(target, 'docker')), ).fetchone()['updated_at'] db.sync_target_queue_from_files('dockerhub', 'docker', [target], [], 'q') assert db.conn.execute( 'SELECT updated_at FROM target_queue WHERE source = ? AND normalized_target = ?', ('dockerhub', normalize_target(target, 'docker')), ).fetchone()['updated_at'] == updated_at assert db.claim_targets('dockerhub', 'docker', 1, 'owner-b', 3600) == [] db.conn.execute( "UPDATE target_queue SET available_after = ? WHERE source = ?", ('2000-01-01T00:00:00+00:00', 'dockerhub'), ) db.conn.commit() reclaimed = db.claim_targets('dockerhub', 'docker', 1, 'owner-b', 3600, return_rows=True) assert [item['target'] for item in reclaimed] == [target] db.complete_target_queue_item( 'dockerhub', 'docker', target, None, 'failed', 'permanent', queue_id=reclaimed[0]['id'], lease_token=reclaimed[0]['lease_token'], ) db.sync_target_queue_from_files('dockerhub', 'docker', [], [target], 'q') row = db.target_queue_item('dockerhub', 'docker', target) assert row['status'] == 'failed' and row['available_after'] is None done_target = 'owner/other:tag' db.enqueue_targets('dockerhub', 'docker', 'q', [done_target]) done_claim = db.claim_targets('dockerhub', 'docker', 1, 'owner-a', 3600, return_rows=True)[0] db.complete_target_queue_item( 'dockerhub', 'docker', done_target, None, 'done', queue_id=done_claim['id'], lease_token=done_claim['lease_token'], ) db.enqueue_targets('dockerhub', 'docker', 'q', [done_target], requeue_done=True) row = db.target_queue_item('dockerhub', 'docker', done_target) assert row['status'] == 'pending' and int(row['attempts']) == 0 done_claim = db.claim_targets('dockerhub', 'docker', 1, 'owner-a', 3600, return_rows=True)[0] db.complete_target_queue_item( 'dockerhub', 'docker', done_target, None, 'done', queue_id=done_claim['id'], lease_token=done_claim['lease_token'], ) legacy_target = 'owner/legacy:tag' db.enqueue_targets('dockerhub', 'docker', 'q', [legacy_target]) db.conn.execute( "UPDATE target_queue SET status = 'pending', available_after = ? WHERE normalized_target = ?", (future, normalize_target(legacy_target, 'docker')), ) db.conn.commit() db.sync_target_queue_from_files('dockerhub', 'docker', [legacy_target], [], 'q') row = db.target_queue_item('dockerhub', 'docker', legacy_target) assert row['status'] == 'deferred' and row['available_after'] == future stale_target = 'owner/stale:tag' db.enqueue_targets('dockerhub', 'docker', 'q', [stale_target]) stale_claim = db.claim_targets('dockerhub', 'docker', 1, 'owner-a', 60, max_attempts=3, return_rows=True)[0] assert stale_claim['target'] == stale_target db.conn.execute( "UPDATE target_queue SET lease_expires_at = ? WHERE normalized_target = ?", ('2000-01-01T00:00:00+00:00', normalize_target(stale_target, 'docker')), ) db.conn.commit() newer_claim = db.claim_targets('dockerhub', 'docker', 1, 'owner-b', 60, max_attempts=3, return_rows=True)[0] assert newer_claim['target'] == stale_target assert newer_claim['lease_token'] != stale_claim['lease_token'] assert not db.complete_target_queue_item( 'dockerhub', 'docker', stale_target, None, 'done', queue_id=stale_claim['id'], lease_token=stale_claim['lease_token'], ) row = db.target_queue_item('dockerhub', 'docker', stale_target) assert row['status'] == 'in_progress' and row['lease_owner'] == 'owner-b' run_id = db.start_run('smoke', ['smoke']) cycle_id = db.start_source_cycle(run_id, 'dockerhub', 'docker', 'search', 'q', 1, 1, None, {}, {}) result = { 'target': stale_target, 'scan_type': 'docker', 'findings': [], 'errors': [], 'scan_event_id': str(uuid.uuid4()), 'timestamp': datetime.now(timezone.utc).isoformat(timespec='seconds'), } before = db.conn.execute('SELECT COUNT(*) AS count FROM target_scans').fetchone()['count'] stale = db.record_and_complete_target_result( run_id, cycle_id, 'dockerhub', 'q', stale_target, result, {}, row['id'], 'owner-a', 'done', lease_token=stale_claim['lease_token'], ) assert stale and stale['stale'] assert db.conn.execute('SELECT COUNT(*) AS count FROM target_scans').fetchone()['count'] == before + 1 row = db.target_queue_item('dockerhub', 'docker', stale_target) assert row['status'] == 'in_progress' and row['lease_token'] == newer_claim['lease_token'] result = dict(result, scan_event_id=str(uuid.uuid4())) owned = db.record_and_complete_target_result( run_id, cycle_id, 'dockerhub', 'q', stale_target, result, {}, row['id'], 'owner-b', 'done', lease_token=newer_claim['lease_token'], ) assert owned and not owned['stale'] assert db.conn.execute('SELECT COUNT(*) AS count FROM target_scans').fetchone()['count'] == before + 2 exhausted_target = 'owner/exhausted:tag' db.enqueue_targets('dockerhub', 'docker', 'q', [exhausted_target]) assert db.claim_targets('dockerhub', 'docker', 1, 'owner-x', 60, max_attempts=1) == [exhausted_target] db.conn.execute( "UPDATE target_queue SET lease_expires_at = ? WHERE normalized_target = ?", ('2000-01-01T00:00:00+00:00', normalize_target(exhausted_target, 'docker')), ) db.conn.commit() assert db.claim_targets('dockerhub', 'docker', 1, 'owner-y', 60, max_attempts=1) == [] assert db.target_queue_item('dockerhub', 'docker', exhausted_target)['status'] == 'failed' finally: db.close() finally: for key, value in old_urls.items(): if value is not None: os.environ[key] = value def main(): for key in ('SCANNER_DB_URL', 'DATABASE_URL', 'TRUF_MANAGED_POSTGRES_DSN', 'KEYCHECK_DB_URL'): os.environ.pop(key, None) assert_diagnostic_policy() assert_docker_platform_policy() assert_docker_partial_pagination_policy() assert_retry_policy() assert_timeout_output_policy() assert_archive_and_secret_safety() assert_queue_policy() print('scanner error policy smoke: OK') if __name__ == '__main__': main()