Files
truf-server/app/scanner_error_policy_smoke.py
2026-09-30 20:30:56 +03:00

424 lines
20 KiB
Python

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