import sys sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('console runner could not disable bytecode writes') import argparse import concurrent.futures import copy import contextlib import hashlib import json import logging import os import secrets import shutil import threading import time from dataclasses import dataclass from datetime import datetime, timedelta, timezone from types import SimpleNamespace if os.name == 'nt': import ctypes from ctypes import wintypes class _PROCESS_MEMORY_COUNTERS_EX(ctypes.Structure): _fields_ = [ ('cb', wintypes.DWORD), ('PageFaultCount', wintypes.DWORD), ('PeakWorkingSetSize', ctypes.c_size_t), ('WorkingSetSize', ctypes.c_size_t), ('QuotaPeakPagedPoolUsage', ctypes.c_size_t), ('QuotaPagedPoolUsage', ctypes.c_size_t), ('QuotaPeakNonPagedPoolUsage', ctypes.c_size_t), ('QuotaNonPagedPoolUsage', ctypes.c_size_t), ('PagefileUsage', ctypes.c_size_t), ('PeakPagefileUsage', ctypes.c_size_t), ('PrivateUsage', ctypes.c_size_t), ] _P_PROCESS_MEMORY_COUNTERS_EX = ctypes.POINTER(_PROCESS_MEMORY_COUNTERS_EX) _METRIC_KERNEL32 = ctypes.WinDLL('kernel32', use_last_error=True) _METRIC_PSAPI = ctypes.WinDLL('psapi', use_last_error=True) _METRIC_GET_CURRENT_PROCESS = _METRIC_KERNEL32.GetCurrentProcess _METRIC_GET_CURRENT_PROCESS.argtypes = [] _METRIC_GET_CURRENT_PROCESS.restype = wintypes.HANDLE _GET_PROCESS_MEMORY_INFO = _METRIC_PSAPI.GetProcessMemoryInfo _GET_PROCESS_MEMORY_INFO.argtypes = [ wintypes.HANDLE, _P_PROCESS_MEMORY_COUNTERS_EX, wintypes.DWORD, ] _GET_PROCESS_MEMORY_INFO.restype = wintypes.BOOL else: _PROCESS_MEMORY_COUNTERS_EX = None from scanner import ( ApiRequestError, DockerContentTransferError, DockerLayerInfrastructureError, DockerRegistryResolutionError, DockerResolverLeaseLostError, DockerRemoteAccessError, DockerHubDiscoveryTransportError, GitLabDiscoveryTransportError, RateLimitError, ResultSinkError, StagedResult, acquire_scan_slot, acquire_scan_slot_leases, assign_finding_uids, check_dependencies, cleanup_pending_command_work_dirs, cleanup_stale_temp_dirs, configure_docker_accounts, configure_docker_discovery_accounts, configure_docker_discovery_tokens, configure_docker_tokens, csv_items, docker_images_per_repository_limit, docker_tag_resolution_is_conclusive, drain_docker_auth_events, fetch_github_archive_repos, fetch_github_archive_file_targets, fetch_github_gist_targets, fetch_dockerhub_tags, fetch_dockerhub_images, fetch_dockerhub_search_page, fetch_github_repo_items, fetch_github_postman_targets, fetch_github_repos, fetch_gitlab_repo_items, fetch_gitlab_repos, fetch_huggingface_spaces, fetch_npm_packages, fetch_npm_package_git_repos, fetch_pypi_packages, fetch_pypi_package_git_repos, fetch_recent_dockerhub_images, fetch_recent_github_repo_items, fetch_recent_github_repos, fetch_docker_config_payload_classes, fetch_recent_gitlab_repo_items, fetch_recent_gitlab_repos, restore_docker_endpoint_cooldowns, get_trufflehog_cmd, normalize_git_scan_resolution_target, parse_github_repo_target, parse_gitlab_project_target, redact_secrets, resolve_git_scan_target, resolve_docker_content_manifest, scan_config, scan_limiter_enabled, scan_slot_scope, scan_target_result, scan_targets_batch, stage_result_bundle, save_scan_result, strip_nearby_context_for_persistence, validate_postman_cache_artifact, initialize_scanner_runtime, write_structured_keycheck_candidates, ) from scan_execution import QueueDispositionPolicy, stage_scan_result_in_scope from scanner_db import ( CI_SOFT_SKIP_REASONS, DiscoveryPausedError, DiscoveryRetryLeaseError, ScanEventConflictError, ScannerDB, canonical_docker_layer_plan_bytes, docker_adaptive_canary_selected, docker_layer_canary_selected, docker_layer_execution_policy_sha256, docker_layer_selection_policy_sha256, hash_file, normalize_target as normalize_db_target, queue_counts, summarize_results, validate_docker_adaptive_checkpoint, validate_docker_layer_limits, ) from result_spool import ( ResultSpool, SpoolBackpressureError, SpoolContentionError, SpoolTransientCapacityError, prepare_scan_event, ) from paths import apply_path_config, default_project_paths, resolve_optional_path from docker_depth_experiment import ( DOCKERHUB_DISCOVERY_ALGORITHM_VERSION, DOCKERHUB_DISCOVERY_MAX_PAGES, DOCKERHUB_DISCOVERY_MAX_PER_PAGE, canonical_dockerhub_discovery_policy, docker_depth_resolver_authority, validate_docker_depth_config, validate_dockerhub_discovery_policies, ) from lifecycle_authority import ( CHILD_KIND_ENV, DISCOVERY_PRODUCER_ROLE, DISCOVERY_PRODUCER_SOURCES, LifecycleAuthorityError, require_active_supervisor_child, ) from process_identity import current_process_identity from result_bundle import ( ResultBundleReader, bundle_partial_relative_path, ensure_bundle_reservation_paths, ) from target_identity import parse_docker_target from runtime_security import ( durable_unlink, PrivatePathState, PrivateFileLock, ensure_private_directory, harden_private_file, inspect_private_relative_path, preflight_lifecycle_paths, read_private_json, private_file_ready, require_private_directory, require_sensitive_runtime_paths, ) logger = logging.getLogger(__name__) SOURCE_INFRASTRUCTURE_HOLD_EXIT = 75 DOCKERHUB_DEEP_INTERVAL = timedelta(hours=72) DOCKERHUB_DISCOVERY_STATE_KEY = 'dockerhub_discovery' DOCKERHUB_DISCOVERY_RETRY_LEASE_SECONDS = 300 DOCKERHUB_DISCOVERY_ERROR_CATEGORIES = frozenset(( 'account_pool_exhausted', 'auth_forbidden', 'auth_invalid', 'auth_unavailable', 'invalid_payload', 'network', 'page_unavailable', 'policy_mismatch', 'provider_cooldown', 'provider_unavailable', 'query_removed', 'rate_limit', 'remote_transient', 'request_failed', 'tail_unavailable', 'transport', )) class UnresolvedHandoffInfrastructureError(RuntimeError): pass def ci_seed_statement_timeout(exc): if isinstance(exc, TimeoutError): return True return bool( str(getattr(exc, 'sqlstate', '') or '') == '57014' and 'statement timeout' in str(exc).lower() ) def commit_if_postgres(db): try: conn = getattr(db, 'conn', None) if conn and getattr(conn, 'is_postgres', False): conn.commit() except Exception: pass def load_set_from_file(path): if not os.path.exists(path): return set() with open(path, 'r', encoding='utf-8') as f: return {line.strip().lstrip('\ufeff') for line in f if line.strip()} @contextlib.contextmanager def projection_file_lock(path, timeout_sec=30): lock_path = f'{path}.projection.lock' parent = os.path.dirname(os.path.abspath(lock_path)) ensure_private_directory(parent, reject_reparse=True) deadline = time.monotonic() + max(0.1, float(timeout_sec)) while True: lock = PrivateFileLock(lock_path) try: lock.acquire() break except BlockingIOError: if time.monotonic() >= deadline: raise TimeoutError(f'timed out acquiring projection lock {lock_path}') time.sleep(0.05) try: yield finally: lock.release() def _write_lines_unlocked(path, lines): parent = os.path.dirname(path) if parent: ensure_private_directory(parent, reject_reparse=True) tmp_path = f'{path}.{os.getpid()}.{threading.get_ident()}.{time.time_ns()}.tmp' with open(tmp_path, 'w', encoding='utf-8') as f: for line in lines: f.write(f"{line}\n") f.flush() try: os.fsync(f.fileno()) except OSError: pass os.replace(tmp_path, path) harden_private_file(path) def write_lines(path, lines): with projection_file_lock(path): _write_lines_unlocked(path, lines) def _append_lines_unlocked(path, lines): parent = os.path.dirname(path) if parent: ensure_private_directory(parent, reject_reparse=True) with open(path, 'a', encoding='utf-8') as f: for line in lines: f.write(f"{line}\n") f.flush() try: os.fsync(f.fileno()) except OSError: pass harden_private_file(path) def append_lines(path, lines): with projection_file_lock(path): _append_lines_unlocked(path, lines) def skip_startup_cleanup(): return str(os.getenv('SCANNER_SKIP_STARTUP_CLEANUP', '')).strip().lower() in ('1', 'true', 'yes', 'on') def console_safe_text(value): text = str(value) encoding = getattr(sys.stdout, 'encoding', None) or 'utf-8' return text.encode(encoding, errors='replace').decode(encoding, errors='replace') def safe_print(value='', **kwargs): print(console_safe_text(value), **kwargs) def bool_config(value, default=False): if value is None: return default if isinstance(value, bool): return value return str(value).strip().lower() in ('1', 'true', 'yes', 'on') def normalize_target(target, platform): return normalize_db_target(target, platform) def read_custom_targets(path): if not path: raise ValueError('--target-file is required for custom mode') if not os.path.exists(path): raise FileNotFoundError(path) with open(path, 'r', encoding='utf-8') as f: return [line.strip().lstrip('\ufeff') for line in f if line.strip()] def read_queries(args): queries = [] if args.query_file: if not os.path.exists(args.query_file): raise FileNotFoundError(args.query_file) with open(args.query_file, 'r', encoding='utf-8') as f: queries.extend(line.strip() for line in f if line.strip()) if args.query: queries.extend(query.strip() for query in args.query.split(',') if query.strip()) return list(dict.fromkeys(queries)) def parse_datetime(value): if not value: return None try: normalized = str(value).replace('Z', '+00:00') parsed = datetime.fromisoformat(normalized) if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=timezone.utc) return parsed.astimezone(timezone.utc) except ValueError: return None def updated_target_rescan_enabled(args): return bool(getattr(args, 'updated_target_rescan_enabled', False)) and args.platform in { 'github', 'gitlab', 'huggingface', } def repo_item_discovery(item, args): target = item.get('url') if not target: return None if not updated_target_rescan_enabled(args) and args.platform != 'huggingface': return target remote_field = { 'github': 'pushed_at', 'gitlab': 'last_activity_at', 'huggingface': 'updated_at', }.get(args.platform) remote_time = parse_datetime(item.get(remote_field)) if remote_field else None discovery = { 'target': target, 'remote_modified_at': remote_time.isoformat(timespec='seconds') if remote_time else None, } if args.platform == 'huggingface': for field in ('private', 'protected', 'gated', 'disabled'): if field in item: discovery[field] = item[field] return discovery def huggingface_discovery_target_is_restricted(item): if not isinstance(item, dict): return False return any( field in item and item[field] is not False and item[field] is not None for field in ('private', 'protected', 'gated', 'disabled') ) def filter_repo_items_by_age(items, args): max_days = int(getattr(args, 'max_repo_age_days', 0) or 0) if max_days <= 0: return [ discovery for item in items if (discovery := repo_item_discovery(item, args)) is not None ] field = getattr(args, 'repo_age_field', None) or 'created_at' cutoff = datetime.now(timezone.utc) - timedelta(days=max_days) kept = [] skipped_old = 0 skipped_missing = 0 for item in items: url = item.get('url') if not url: continue value = item.get(field) parsed = parse_datetime(value) if not parsed: skipped_missing += 1 if updated_target_rescan_enabled(args): discovery = repo_item_discovery(item, args) if discovery is not None: kept.append(discovery) continue if parsed >= cutoff: discovery = repo_item_discovery(item, args) if discovery is not None: kept.append(discovery) else: skipped_old += 1 print( f"Repo age filter: kept {len(kept)}, skipped {skipped_old} older than " f"{max_days} days by {field}, skipped {skipped_missing} missing dates" ) return kept def make_repo_candidate_callback(db=None, run_id=None, cycle_id=None, query=None): if not db or not run_id: return None def callback(candidates): return db.record_package_repo_candidates(run_id, cycle_id, query, candidates) return callback def fetch_package_git_targets_from_db(args, db=None): if not db: return [] queries = read_queries(args) package_sources = [item.strip().lower() for item in str(getattr(args, 'package_sources', 'npm,pypi') or '').split(',') if item.strip()] targets = [] seen = set() for query in queries: candidates = db.get_package_repo_candidates(query, package_sources, limit=max(1000, int(args.pages or 1) * int(args.per_page or 50) * 10)) for candidate in candidates: target = json.dumps(candidate, separators=(',', ':'), ensure_ascii=False) normalized = normalize_target(target, 'package_git') if normalized not in seen: targets.append(target) seen.add(normalized) return targets def source_list(value, default): if value is None: value = default if isinstance(value, str): return [item.strip() for item in value.split(',') if item.strip()] return [str(item).strip() for item in value if str(item).strip()] def ci_metadata_candidate_values(value): if not value: return [] try: data = json.loads(value) if isinstance(value, str) else value except Exception: return [] output = [] interesting = {'repository', 'repo', 'repo_url', 'url', 'link', 'project', 'path_with_namespace'} def visit(item): if isinstance(item, dict): for key, child in item.items(): if key in interesting and isinstance(child, str): output.append(child) visit(child) elif isinstance(item, list): for child in item: visit(child) visit(data) return output def fetch_ci_repo_targets_from_db(args, db=None, provider='github'): owned_db = None if not db or not getattr(db, 'conn', None): db_path = getattr(args, 'database_path', None) or os.getenv('SCANNER_DB_PATH') or os.getenv('SCAN_DB_PATH') db_url = getattr(args, 'database_url', None) or os.getenv('SCANNER_DB_URL') or os.getenv('DATABASE_URL') if not db_path: db_path = os.path.join(getattr(args, 'save_dir', '') or '', 'scanner_active.db') try: owned_db = ScannerDB(db_path=db_path, db_url=db_url, initialize=False) if not owned_db.enabled: raise RuntimeError('scanner DB disabled') db = owned_db print(f'CI target discovery ({provider}): opened scanner DB {owned_db.db_display}') except Exception as e: print(f'CI target discovery ({provider}): scanner DB unavailable: {e}') return [] seed_sources = source_list(getattr(args, 'ci_seed_sources', None), 'github,gitlab,package_git') limit = max(1, int(getattr(args, 'ci_max_repos_per_cycle', 100) or 100)) scan_limit = max(1, int(getattr(args, 'ci_seed_scan_limit', 5000) or 5000)) query_batch = min(500, max(25, int(getattr(args, 'ci_seed_query_batch_size', 250) or 250))) platform = 'github_actions' if provider == 'github' else 'gitlab_ci' soft_skip_reasons = CI_SOFT_SKIP_REASONS.get(platform, set()) cooldown_days = int(getattr(args, 'ci_soft_cooldown_days', 7) or 0) cooldown_keys = set() cooldown_rows_loaded = 0 cooldown_truncated = False if cooldown_days > 0 and soft_skip_reasons: cutoff = (datetime.now(timezone.utc) - timedelta(days=cooldown_days)).isoformat(timespec='seconds') reason_predicate = ' OR '.join( f"skipped_reason = '{reason}'" for reason in sorted(soft_skip_reasons) ) cooldown_scope = f"source = '{platform}' AND ({reason_predicate})" cursor_ended_at = None cursor_id = None while cooldown_rows_loaded < scan_limit: page_limit = min(query_batch, scan_limit - cooldown_rows_loaded) cursor_clause = '' params = [cutoff] if cursor_ended_at is not None and cursor_id is not None: cursor_clause = 'AND (ended_at, id) < (?, ?)' params.extend((cursor_ended_at, cursor_id)) params.append(page_limit) try: page = db.conn.execute(''' SELECT id, target, skipped_reason, ended_at FROM target_scans WHERE status = 'skipped' AND {} AND ended_at >= ? {} ORDER BY ended_at DESC, id DESC LIMIT ? '''.format(cooldown_scope, cursor_clause), params).fetchall() except Exception: if getattr(db.conn, 'is_postgres', False): db.conn.rollback() raise if not page: break for row in page: if provider == 'github': repo, _ = parse_github_repo_target(row['target']) else: repo, _ = parse_gitlab_project_target(row['target']) if repo: cooldown_keys.add(repo.lower()) cooldown_rows_loaded += len(page) cursor_ended_at = page[-1]['ended_at'] cursor_id = page[-1]['id'] if len(page) < page_limit: break if cooldown_rows_loaded >= scan_limit: probe_params = [cutoff, cursor_ended_at, cursor_id] try: cooldown_truncated = bool(db.conn.execute(''' SELECT 1 FROM target_scans WHERE status = 'skipped' AND {} AND ended_at >= ? AND (ended_at, id) < (?, ?) ORDER BY ended_at DESC, id DESC LIMIT 1 '''.format(cooldown_scope), probe_params).fetchone()) except Exception: if getattr(db.conn, 'is_postgres', False): db.conn.rollback() raise break if cooldown_truncated: safe_print( f'Warning: CI cooldown history reached its bounded {scan_limit}-row limit; ' 'older rows in the configured cooldown window were not loaded', flush=True, ) if getattr(db.conn, 'is_postgres', False): # Keyset pages keep each HDD transaction bounded. The per-source branch # limit still preserves the exact global newest-first ordering. rows = [] cursor_ended_at = None cursor_id = None while len(rows) < scan_limit: page_limit = min(query_batch, scan_limit - len(rows)) cursor_clause = '' params = [*seed_sources] if cursor_ended_at is not None and cursor_id is not None: cursor_clause = 'AND (ended_at < ? OR (ended_at = ? AND id < ?))' params.extend((cursor_ended_at, cursor_ended_at, cursor_id)) params.extend((page_limit, page_limit)) try: page = db.conn.execute(''' SELECT ranked.id, ranked.source, ranked.query, ranked.target, ranked.normalized_target, ranked.ended_at AS latest_seen FROM (VALUES {}) AS seed(source) CROSS JOIN LATERAL ( SELECT id, source, query, target, normalized_target, ended_at FROM target_scans WHERE source = seed.source AND ended_at IS NOT NULL {} ORDER BY ended_at DESC, id DESC LIMIT ? ) AS ranked ORDER BY ranked.ended_at DESC, ranked.id DESC LIMIT ? '''.format(','.join('(?)' for _ in seed_sources), cursor_clause), params).fetchall() except Exception as exc: if not ci_seed_statement_timeout(exc): raise try: db.conn.rollback() except Exception as rollback_exc: reset = getattr(db, '_reset_connection', None) if not reset or not reset(): raise exc from rollback_exc safe_print( f'Warning: CI target discovery ({provider}) degraded after a bounded database ' 'statement timeout; no seed targets from this cycle were used', flush=True, ) if owned_db: owned_db.close() return [] if not page: break rows.extend(page) cursor_ended_at = page[-1]['latest_seen'] cursor_id = page[-1]['id'] if len(page) < page_limit: break else: rows = db.conn.execute(''' SELECT source, query, target, normalized_target, ended_at AS latest_seen FROM target_scans WHERE source IN ({}) ORDER BY ended_at DESC NULLS LAST LIMIT ? '''.format(','.join('?' for _ in seed_sources)), [*seed_sources, scan_limit]).fetchall() seed_records = [] seen_scan_targets = set() for row in rows: scan_key = ( row['source'] or '', row['query'] or '', row['target'] or '', row['normalized_target'] or '', ) if scan_key in seen_scan_targets: continue seen_scan_targets.add(scan_key) candidates = [row['target'], row['normalized_target']] try: parsed_json = json.loads(str(row['target'] or '')) except Exception: parsed_json = None if isinstance(parsed_json, dict): candidates.append(parsed_json) if parsed_json.get('repo_url'): candidates.append(parsed_json.get('repo_url')) seed_records.append({ 'source': row['source'], 'query': row['query'], 'candidates': candidates, }) # Lower-priority sources may fill only the remaining aggregate DB-row # budget; they must not multiply ci_seed_scan_limit. if 'package_git' in seed_sources and len(seed_records) < scan_limit: remaining_seed_rows = scan_limit - len(seed_records) package_rows = db.conn.execute(''' SELECT package_source, package_name, package_version, query, repo_url, provider, last_seen_at FROM package_repo_candidates WHERE LOWER(provider) = ? ORDER BY last_seen_at DESC LIMIT ? ''', [provider, remaining_seed_rows]).fetchall() for row in package_rows: seed_records.append({ 'source': 'package_git', 'query': row['query'], 'candidates': [ row['repo_url'], { 'repo_url': row['repo_url'], 'provider': row['provider'], 'package_source': row['package_source'], 'name': row['package_name'], 'version': row['package_version'], }, ], }) use_finding_seeds = bool(getattr(args, 'ci_use_finding_seeds', False)) finding_sources = [source for source in seed_sources if source in ('github', 'gitlab', 'package_git')] if use_finding_seeds else [] if finding_sources and len(seed_records) < scan_limit: per_source_finding_limit = max(limit * 50, scan_limit // max(1, len(finding_sources))) for source_name in finding_sources: remaining_seed_rows = scan_limit - len(seed_records) if remaining_seed_rows <= 0: break finding_rows = db.conn.execute(''' SELECT source, query, target, normalized_target, source_metadata_json, raw_finding_json, created_at FROM findings WHERE source = ? ORDER BY id DESC LIMIT ? ''', [source_name, min(per_source_finding_limit, remaining_seed_rows)]).fetchall() for row in finding_rows: candidates = [row['target'], row['normalized_target']] candidates.extend(ci_metadata_candidate_values(row['source_metadata_json'])) candidates.extend(ci_metadata_candidate_values(row['raw_finding_json'])) seed_records.append({ 'source': row['source'], 'query': row['query'], 'candidates': candidates, }) offered = [] seen = set() parsed = 0 skipped_known = 0 skipped_duplicate = 0 skipped_unparseable = 0 skipped_cooldown = 0 for record in seed_records: for candidate in record['candidates']: if provider == 'github': repo, url = parse_github_repo_target(candidate) platform = 'github_actions' else: repo, url = parse_gitlab_project_target(candidate) platform = 'gitlab_ci' if not repo: skipped_unparseable += 1 continue key = repo.lower() parsed += 1 if key in cooldown_keys: skipped_cooldown += 1 continue target_payload = json.dumps({ 'repo' if provider == 'github' else 'project': repo, 'url': url, 'seed_source': record['source'], 'seed_query': record['query'], }, separators=(',', ':'), ensure_ascii=False) if key in seen: skipped_duplicate += 1 continue seen.add(key) offered.append((target_payload, normalize_target(target_payload, platform))) if len(offered) >= scan_limit: break if len(offered) >= scan_limit: break known = set() lookup = getattr(db, 'known_target_normalizations_for', None) if lookup: for start in range(0, len(offered), 64): batch = [target for target, _ in offered[start:start + 64]] try: known.update(lookup(platform, platform, batch) or ()) except Exception as exc: safe_print(f'Warning: CI known-target lookup failed open: {exc}') eligible = [] for target, normalized in offered: if normalized in known: skipped_known += 1 continue eligible.append(target) targets = eligible[:limit] print( f'CI target discovery ({provider}): selected {len(targets)} target(s); ' f'db_rows={len(seed_records)}, parsed={parsed}, known={len(known)}, ' f'skipped_known={skipped_known}, skipped_duplicate={skipped_duplicate}, ' f'skipped_unparseable={skipped_unparseable}, skipped_cooldown={skipped_cooldown}, ' f'cooldown_rows={cooldown_rows_loaded}, cooldown_truncated={str(cooldown_truncated).lower()}, ' f'seed_sources={",".join(seed_sources)}. ' f'Increase ci_seed_scan_limit or add fresh seed repos if db_rows equals scan limit.' ) if owned_db: owned_db.close() else: commit_if_postgres(db) return targets def _dockerhub_discovery_policy_values(pages, per_page, sort_by, sort_order): try: effective_pages = max(1, min(DOCKERHUB_DISCOVERY_MAX_PAGES, int(pages))) effective_per_page = max( 1, min(DOCKERHUB_DISCOVERY_MAX_PER_PAGE, int(per_page)), ) except (TypeError, ValueError, OverflowError): raise ValueError('DockerHub discovery page policy is invalid') from None return canonical_dockerhub_discovery_policy( effective_pages, effective_per_page, str(sort_by or 'updated_at'), str(sort_order or 'desc'), ) def dockerhub_discovery_policy(args): return _dockerhub_discovery_policy_values( getattr(args, 'pages', 1), getattr(args, 'per_page', DOCKERHUB_DISCOVERY_MAX_PER_PAGE), getattr(args, 'docker_sort_by', 'updated_at'), getattr(args, 'sort_order', 'desc'), ) def configured_dockerhub_discovery_policies(source_config): return { effective['query']: { key: value for key, value in effective.items() if key not in ('query', 'max_targets') } for effective in validate_dockerhub_discovery_policies(source_config) } def managed_dockerhub_discovery(args, db, source_name): return bool( source_name == 'dockerhub' and getattr(args, 'platform', None) == 'docker' and getattr(args, 'mode', None) == 'search' and db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False) and callable(getattr(db, 'persist_dockerhub_discovery_page', None)) and callable(getattr(db, 'enqueue_discovery_retry', None)) ) def _require_discovery_provider_enabled(db): control = db.runtime_control_state() if control['effective_discovery_paused']: raise DiscoveryPausedError(control) return control def _safe_dockerhub_discovery_category(error): category = str(getattr(error, 'category', '') or '') return category if category in DOCKERHUB_DISCOVERY_ERROR_CATEGORIES else 'page_unavailable' def _validated_dockerhub_page(query, requested_page, args, policy=None): policy = policy or dockerhub_discovery_policy(args) result = fetch_dockerhub_search_page( query, requested_page, per_page=policy['per_page'], sort_by=policy['sort_by'], sort_order=policy['sort_order'], request_timeout=getattr(args, 'fetch_timeout', 15), ) valid = isinstance(result, dict) and result.get('page') == requested_page repositories = result.get('repositories') if valid else None total_count = result.get('total_count') if valid else None if ( not isinstance(repositories, list) or isinstance(total_count, bool) or (isinstance(total_count, float) and not total_count.is_integer()) ): valid = False else: try: total_count = int(total_count) except (TypeError, ValueError, OverflowError): valid = False result_count = len(repositories) if isinstance(repositories, list) else -1 absolute_start = (requested_page - 1) * int(policy['per_page']) absolute_bound = absolute_start + max(0, result_count) if ( not valid or total_count < absolute_bound or result_count > int(policy['per_page']) or (result_count == 0 and total_count > absolute_start) ): raise DockerHubDiscoveryTransportError( 'Docker Hub search page returned an invalid payload', category='invalid_payload', remote_attempted=True, retryable=False, ) return result, repositories, total_count def annotate_dockerhub_experiment_observation_args(args, experiment, cycle_id=None): args.docker_depth_collection_only = bool( experiment is not None and not experiment.enabled ) if experiment is None: return args args.dockerhub_discovery_ordered_queries = tuple(experiment.queries) args.dockerhub_discovery_ordered_query_hash = str( experiment.ordered_query_hash ) args.dockerhub_discovery_query_count = len(experiment.queries) args.dockerhub_discovery_collection_generation = str( experiment.collection_generation ) args.dockerhub_discovery_cycle_id = cycle_id return args def annotate_dockerhub_runtime_args( args, experiment, discovery_policies, cycle_id=None, ): annotate_dockerhub_experiment_observation_args(args, experiment, cycle_id) if experiment is None: return args policy_hashes = { policy['policy_sha256'] for policy in discovery_policies.values() } if len(policy_hashes) != 1: raise RuntimeError('DockerHub experiment resolver policy authority is ambiguous') args.docker_depth_experiment_authority = docker_depth_resolver_authority( experiment, next(iter(policy_hashes)), ) return args def _dockerhub_discovery_observation( args, query, page_number, per_page, policy_sha256, pass_kind, total_count=None, query_complete=False, ): ordered_queries = getattr(args, 'dockerhub_discovery_ordered_queries', None) ordered_query_hash = getattr( args, 'dockerhub_discovery_ordered_query_hash', None, ) if ordered_queries is None and ordered_query_hash is None: return None if not isinstance(ordered_queries, tuple) or not ordered_queries: raise RuntimeError('DockerHub experiment observation queries are unavailable') if ( not isinstance(ordered_query_hash, str) or not _valid_dockerhub_policy_sha256(ordered_query_hash) or len(set(ordered_queries)) != len(ordered_queries) or query not in ordered_queries or int(getattr(args, 'dockerhub_discovery_query_count', 0) or 0) != len(ordered_queries) ): raise RuntimeError('DockerHub experiment observation identity is invalid') return { 'cycle_id': getattr(args, 'dockerhub_discovery_cycle_id', None), 'query_ordinal': ordered_queries.index(query), 'query_count': len(ordered_queries), 'page_number': page_number, 'page_limit': int(getattr(args, 'pages', 1)), 'per_page': per_page, 'total_count': total_count, 'policy_sha256': policy_sha256, 'pass_kind': pass_kind, 'collection_generation': str( getattr(args, 'dockerhub_discovery_collection_generation', '') or '' ), 'ordered_query_hash': ordered_query_hash, 'query_complete': bool(query_complete), } def _confirmed_discovery_retry(db, source, query, policy_sha256, pass_kind, work_kind, page_start, page_end, error, observation=None): retry_at = getattr(error, 'retry_at', None) try: observation_kwargs = ( {'observation': observation} if observation is not None else {} ) result = db.enqueue_discovery_retry( source, query, policy_sha256, pass_kind, work_kind, page_start=page_start, page_end=page_end, available_after=retry_at, error_category=_safe_dockerhub_discovery_category(error), **observation_kwargs, ) confirmed = bool( isinstance(result, dict) and not isinstance(result.get('id'), bool) and int(result.get('id') or 0) >= 1 and result.get('status') in ('pending', 'leased') ) except Exception: raise RuntimeError('DockerHub discovery retry delegation failed') from None if not confirmed: raise RuntimeError('DockerHub discovery retry delegation was not confirmed') return result def run_dockerhub_incremental_discovery( args, db, source='dockerhub', *, partial_metrics=None, ): policy = dockerhub_discovery_policy(args) annotated_policy = str( getattr(args, 'dockerhub_discovery_policy_sha256', '') or '' ) if annotated_policy and annotated_policy != policy['policy_sha256']: raise RuntimeError('DockerHub discovery policy annotation is stale') policy_sha256 = policy['policy_sha256'] deep = bool(getattr(args, 'dockerhub_discovery_deep', False)) pass_kind = 'deep' if deep else 'ordinary' metrics = { 'fetched_count': 0, 'queued_new_count': 0, 'queued_updated_count': 0, 'discovery_pages_fetched': 0, 'discovery_retry_enqueued_count': 0, 'discovery_retry_inserted_count': 0, 'discovery_retry_coalesced_count': 0, 'discovery_preexisting_count': 0, 'discovery_known_page_count': 0, 'discovery_stopped_on_preexisting': False, 'discovery_deep': deep, 'discovery_pass_kind': pass_kind, 'deep_dispatch_durable': False, 'cycle_status': 'completed', } query = str(getattr(args, 'query', '') or '') known_streak = 0 new_this_pass = set() knownness_disabled = False def publish_metrics(): if isinstance(partial_metrics, dict): partial_metrics.update(metrics) publish_metrics() def delegate(error, work_kind, page_start, page_end): observation = _dockerhub_discovery_observation( args, query, page_start, policy['per_page'], policy_sha256, pass_kind, None, ) report = _confirmed_discovery_retry( db, source, query, policy_sha256, pass_kind, work_kind, page_start, page_end, error, observation, ) metrics['discovery_retry_enqueued_count'] += 1 metrics['discovery_retry_inserted_count'] += int( report.get('inserted_count', 0) or 0 ) metrics['discovery_retry_coalesced_count'] += int( report.get('coalesced_count', 0) or 0 ) metrics['cycle_status'] = 'completed_with_retries' publish_metrics() def admit(repositories, page_number, total_count, query_complete=False): try: observation = _dockerhub_discovery_observation( args, query, page_number, policy['per_page'], policy_sha256, pass_kind, total_count, query_complete, ) observation_kwargs = ( {'observation': observation} if observation is not None else {} ) experiment_authority = getattr( args, 'docker_depth_experiment_authority', None, ) if isinstance(experiment_authority, dict): observation_kwargs.update({ 'experiment_authority': experiment_authority, 'final_cutover': True, }) report = db.persist_dockerhub_discovery_page( source, query, repositories, **observation_kwargs, ) except DiscoveryPausedError: raise except Exception: raise RuntimeError('DockerHub discovery page admission failed') from None if not isinstance(report, dict): raise RuntimeError('DockerHub discovery page admission was not confirmed') counts = {} for key in ( 'attempted_count', 'normalized_count', 'preexisting_count', 'inserted_count', ): value = report.get(key) if isinstance(value, bool): raise RuntimeError('DockerHub discovery page admission report is invalid') try: counts[key] = int(value) except (TypeError, ValueError, OverflowError): raise RuntimeError( 'DockerHub discovery page admission report is invalid' ) from None if counts[key] < 0: raise RuntimeError('DockerHub discovery page admission report is invalid') if ( counts['attempted_count'] != len(repositories) or counts['normalized_count'] > counts['attempted_count'] or counts['preexisting_count'] > counts['normalized_count'] or counts['inserted_count'] > counts['normalized_count'] ): raise RuntimeError('DockerHub discovery page admission report is invalid') metrics['fetched_count'] += counts['attempted_count'] metrics['queued_new_count'] += counts['inserted_count'] metrics['discovery_preexisting_count'] += counts['preexisting_count'] metrics['discovery_pages_fetched'] += 1 publish_metrics() return report, counts def observe_knownness(report, counts): nonlocal known_streak, knownness_disabled normalized = report.get('normalized_repositories') preexisting = report.get('preexisting_repositories') certain = bool( isinstance(normalized, (set, frozenset)) and isinstance(preexisting, (set, frozenset)) and len(normalized) == counts['normalized_count'] and len(preexisting) == counts['preexisting_count'] and preexisting.issubset(normalized) ) if not certain: known_streak = 0 if isinstance(normalized, (set, frozenset)): new_this_pass.update(normalized) else: knownness_disabled = True return all_preexisting = bool( normalized and normalized.issubset(preexisting) and normalized.isdisjoint(new_this_pass) ) new_this_pass.update(normalized - preexisting) known_streak = known_streak + 1 if all_preexisting else 0 if all_preexisting: metrics['discovery_known_page_count'] += 1 try: _require_discovery_provider_enabled(db) _, repositories, total_count = _validated_dockerhub_page(query, 1, args) except DockerHubDiscoveryTransportError as error: if not getattr(error, 'retryable', True): raise delegate(error, 'query', 1, policy['pages']) metrics['deep_dispatch_durable'] = deep publish_metrics() return metrics expected_pages = max( 1, min(policy['pages'], (total_count + policy['per_page'] - 1) // policy['per_page']), ) report, counts = admit( repositories, 1, total_count, query_complete=not repositories or expected_pages == 1, ) if not repositories: metrics['deep_dispatch_durable'] = deep publish_metrics() return metrics observe_knownness(report, counts) if not deep and not knownness_disabled and known_streak >= 2: metrics['discovery_stopped_on_preexisting'] = True metrics['deep_dispatch_durable'] = deep publish_metrics() return metrics for page in range(2, expected_pages + 1): try: _require_discovery_provider_enabled(db) _, repositories, page_total_count = _validated_dockerhub_page( query, page, args, ) except DockerHubDiscoveryTransportError as error: if not getattr(error, 'retryable', True): raise known_streak = 0 if getattr(error, 'remote_attempted', True): delegate(error, 'page', page, page) continue delegate(error, 'range', page, expected_pages) break report, counts = admit( repositories, page, page_total_count, query_complete=not repositories or page >= expected_pages, ) if not repositories: known_streak = 0 break observe_knownness(report, counts) if not deep and not knownness_disabled and known_streak >= 2: metrics['discovery_stopped_on_preexisting'] = True break metrics['deep_dispatch_durable'] = deep publish_metrics() return metrics def process_dockerhub_discovery_retry(args, db, source, configured_policies): metrics = { 'discovery_retry_claimed_count': 0, 'discovery_retry_pages_fetched': 0, 'discovery_retry_queued_new_count': 0, 'discovery_retry_completed_count': 0, 'discovery_retry_deferred_count': 0, 'discovery_retry_held_count': 0, 'discovery_retry_error_count': 0, } required = ( 'claim_discovery_retries', 'persist_dockerhub_discovery_page', 'finish_discovery_retry', 'renew_discovery_retry_lease', 'update_discovery_retry', 'hold_discovery_retry', ) if ( source != 'dockerhub' or getattr(args, 'platform', None) != 'docker' or getattr(args, 'mode', None) != 'search' or not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False) or any(not callable(getattr(db, name, None)) for name in required) ): return metrics policy_hashes = { query: policy.get('policy_sha256') for query, policy in configured_policies.items() if isinstance(policy, dict) } lease_owner = f'dockerhub-retry:{os.getpid()}:{secrets.token_hex(8)}'[:128] try: claims = db.claim_discovery_retries( source, policy_hashes, lease_owner, limit=1, lease_seconds=DOCKERHUB_DISCOVERY_RETRY_LEASE_SECONDS, ) except Exception: metrics['discovery_retry_error_count'] = 1 logger.warning('DockerHub discovery retry claim failed') return metrics if not claims: return metrics claim = claims[0] metrics['discovery_retry_claimed_count'] = 1 query = claim.get('query') if isinstance(claim, dict) else None policy = configured_policies.get(query) if isinstance(query, str) else None def hold(category): try: db.hold_discovery_retry( claim.get('id'), claim.get('lease_owner'), claim.get('lease_token'), category, ) metrics['discovery_retry_held_count'] = 1 except Exception: metrics['discovery_retry_error_count'] = 1 logger.warning('DockerHub discovery retry hold failed') if policy is None: hold('query_removed') return metrics if claim.get('policy_sha256') != policy.get('policy_sha256'): hold('policy_mismatch') return metrics try: work_kind = str(claim.get('work_kind') or '') page_start = int(claim.get('page_start')) stored_page_end = int(claim.get('page_end')) next_page = int(claim.get('next_page')) effective_end = min(stored_page_end, int(policy['pages'])) valid_work = work_kind in ('query', 'page', 'range') valid_work = valid_work and 1 <= page_start <= next_page <= stored_page_end <= 30 valid_work = valid_work and (work_kind != 'query' or page_start == 1) valid_work = valid_work and (work_kind != 'page' or page_start == stored_page_end) except (TypeError, ValueError, OverflowError, KeyError): valid_work = False if not valid_work: hold('invalid_payload') return metrics if next_page > effective_end: hold('invalid_payload') return metrics page = next_page remote_attempted_in_claim = False def defer(category, *, retry_at=None, refund_attempt=False): try: db.update_discovery_retry( claim['id'], claim['lease_owner'], claim['lease_token'], category, retry_at=retry_at, refund_attempt=refund_attempt, next_page=page, ) metrics['discovery_retry_deferred_count'] = 1 return True except Exception: metrics['discovery_retry_error_count'] = 1 logger.warning('DockerHub discovery retry deferral failed') return False while page <= effective_end: try: db.renew_discovery_retry_lease( claim['id'], claim['lease_owner'], claim['lease_token'], DOCKERHUB_DISCOVERY_RETRY_LEASE_SECONDS, ) except Exception: metrics['discovery_retry_error_count'] = 1 logger.warning('DockerHub discovery retry lease renewal failed') return metrics try: _require_discovery_provider_enabled(db) _, repositories, total_count = _validated_dockerhub_page( query, page, args, policy=policy, ) except DiscoveryPausedError: defer('provider_cooldown', refund_attempt=True) return metrics except DockerHubDiscoveryTransportError as error: if not getattr(error, 'retryable', True): hold('invalid_payload') return metrics retry_at = getattr(error, 'retry_at', None) refund_attempt = bool( not remote_attempted_in_claim and not getattr(error, 'remote_attempted', True) and retry_at is not None ) defer( _safe_dockerhub_discovery_category(error), retry_at=retry_at if refund_attempt else None, refund_attempt=refund_attempt, ) return metrics except Exception: logger.warning('DockerHub discovery retry acquisition failed') defer('request_failed') return metrics remote_attempted_in_claim = True metrics['discovery_retry_pages_fetched'] += 1 if work_kind == 'query' and page == 1: effective_end = min( effective_end, max(1, (total_count + int(policy['per_page']) - 1) // int(policy['per_page'])), ) complete = not repositories or page >= effective_end query_page_end = min( int(policy['pages']), max(1, (total_count + int(policy['per_page']) - 1) // int(policy['per_page'])), ) observation = _dockerhub_discovery_observation( args, query, page, int(policy['per_page']), str(claim['policy_sha256']), str(claim['pass_kind']), total_count, query_complete=not repositories or page >= query_page_end, ) if observation is not None: observation['cycle_id'] = claim.get('source_cycle_id') observation_kwargs = ( {'observation': observation} if observation is not None else {} ) experiment_authority = getattr( args, 'docker_depth_experiment_authority', None, ) if isinstance(experiment_authority, dict): observation_kwargs.update({ 'experiment_authority': experiment_authority, 'final_cutover': True, }) try: report = db.persist_dockerhub_discovery_page( source, query, repositories, retry_id=claim['id'], lease_owner=claim['lease_owner'], lease_token=claim['lease_token'], next_page=None if complete else page + 1, complete=complete, **observation_kwargs, ) if not isinstance(report, dict): raise RuntimeError('discovery retry page admission was not confirmed') inserted_count = report.get('inserted_count') if isinstance(inserted_count, bool) or int(inserted_count) < 0: raise RuntimeError('discovery retry page admission report is invalid') progress_key = 'retry_completed_count' if complete else 'retry_progress_count' if int(report.get(progress_key) or 0) != 1: raise RuntimeError('discovery retry fence transition was not confirmed') except DiscoveryPausedError: raise except ValueError: hold('invalid_payload') return metrics except DiscoveryRetryLeaseError: metrics['discovery_retry_error_count'] = 1 logger.warning('DockerHub discovery retry page admission lost its lease') return metrics except Exception: logger.warning('DockerHub discovery retry page admission failed') defer('request_failed') return metrics metrics['discovery_retry_queued_new_count'] += int(inserted_count) if complete: metrics['discovery_retry_completed_count'] = 1 return metrics page += 1 return metrics def fetch_targets(args, db=None, run_id=None, cycle_id=None, source_name=None): if args.platform == 'docker' and args.docker_platform_filter_enabled and ( str(args.docker_platform_os).lower(), str(args.docker_platform_arch).lower() ) != ('linux', 'amd64'): raise ValueError('This TruffleHog Docker runner supports platform filtering only for linux/amd64') if args.mode == 'custom': return read_custom_targets(args.target_file) stop_options = fetch_stop_options(args, db, source_name) if args.platform == 'github': token = args.token or os.getenv('GITHUB_TOKEN') use_metadata = ( int(getattr(args, 'max_repo_age_days', 0) or 0) > 0 or updated_target_rescan_enabled(args) ) raise_rate_limit = bool(getattr(args, 'raise_rate_limit', False)) if args.mode == 'recent': since = datetime.now() - timedelta(hours=args.recent_hours) if use_metadata: return filter_repo_items_by_age(fetch_recent_github_repo_items(args.query, since, token, raise_rate_limit, args.pages, args.per_page, **stop_options), args) return fetch_recent_github_repos(args.query, since, token, raise_rate_limit, args.pages, args.per_page, **stop_options) if use_metadata: return filter_repo_items_by_age(fetch_github_repo_items( args.query, args.pages, args.per_page, token, args.sort_by, args.sort_order, args.created_filter, raise_rate_limit, **stop_options, ), args) return fetch_github_repos( args.query, args.pages, args.per_page, token, args.sort_by, args.sort_order, args.created_filter, raise_rate_limit, **stop_options, ) if args.platform == 'github_archive': targets = fetch_github_archive_repos( args.archive_hours_back, args.archive_max_repos_per_cycle, normalize_archive_event_types(args.archive_event_types), args.fetch_timeout, getattr(args, 'gharchive_cache_dir', None), ) return filter_github_archive_targets_by_cooldown(targets, args, db) if args.platform == 'github_archive_files': targets = fetch_github_archive_file_targets( args.archive_hours_back, args.archive_max_files_per_cycle, normalize_archive_event_types(args.archive_event_types), args.fetch_timeout, getattr(args, 'postman_cache_dir', None), getattr(args, 'max_artifact_size_mb', 2), args.token or os.getenv('GITHUB_TOKEN'), getattr(args, 'archive_max_commit_lookups', 200), getattr(args, 'gharchive_cache_dir', None), discovery_max_artifacts=getattr(args, 'postman_discovery_max_artifacts_per_cycle', None), discovery_max_artifacts_per_page=getattr(args, 'postman_discovery_max_artifacts_per_page', None), discovery_max_bytes=getattr(args, 'postman_discovery_max_bytes_per_cycle', None), discovery_max_elapsed_sec=getattr(args, 'postman_discovery_max_elapsed_sec', None), ) return targets if args.platform == 'gitlab': token = args.token or os.getenv('GITLAB_TOKEN') use_metadata = ( int(getattr(args, 'max_repo_age_days', 0) or 0) > 0 or updated_target_rescan_enabled(args) ) raise_rate_limit = bool(getattr(args, 'raise_rate_limit', False)) gitlab_options = { **stop_options, 'request_attempts': max(1, int(getattr(args, 'gitlab_discovery_request_attempts', 1) or 1)), 'retry_delay': max(0, int(getattr(args, 'gitlab_discovery_retry_delay', 0) or 0)), } if args.mode == 'recent': since = datetime.now() - timedelta(hours=args.recent_hours) if use_metadata: return filter_repo_items_by_age(fetch_recent_gitlab_repo_items(args.query, since, token, raise_rate_limit, args.pages, args.per_page, args.gitlab_visibility, **gitlab_options), args) return fetch_recent_gitlab_repos(args.query, since, token, raise_rate_limit, args.pages, args.per_page, args.gitlab_visibility, **gitlab_options) if use_metadata: return filter_repo_items_by_age(fetch_gitlab_repo_items( args.query, args.pages, args.per_page, token, args.gitlab_sort_by, args.sort_order, args.gitlab_visibility, raise_rate_limit=raise_rate_limit, **gitlab_options, ), args) return fetch_gitlab_repos( args.query, args.pages, args.per_page, token, args.gitlab_sort_by, args.sort_order, args.gitlab_visibility, raise_rate_limit, **gitlab_options, ) if args.platform == 'huggingface': token = args.token or os.getenv('HF_TOKEN') or os.getenv('HUGGINGFACE_TOKEN') preserve_metadata = True spaces = fetch_huggingface_spaces( args.pages, token, args.fetch_timeout, return_metadata=preserve_metadata, request_attempts=max( 1, int(getattr(args, 'huggingface_discovery_request_attempts', 1) or 1), ), retry_delay=max( 0, int(getattr(args, 'huggingface_discovery_retry_delay', 0) or 0), ), **stop_options, ) return [ discovery for item in spaces if (discovery := repo_item_discovery(item, args)) is not None and not huggingface_discovery_target_is_restricted(discovery) ] if args.platform == 'docker': queries = read_queries(args) if not queries: print('Docker Hub search requires --query. Empty query returns no results from the Docker Hub API.') return [] images = [] seen = set() queue_resolver = bool( db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False) and hasattr(db, 'claim_docker_resolutions') ) for query in queries: print(f"Fetching Docker Hub targets for query: {query}") if args.mode == 'recent': since = datetime.now() - timedelta(days=args.recent_days) query_images = fetch_recent_dockerhub_images( query, since, args.per_page, args.pages, args.docker_platform_filter_enabled, args.docker_platform_os, args.docker_platform_arch, args.docker_platform_candidate_tags, docker_images_per_repository_limit( getattr(args, 'docker_images_per_repository', 1) ), not queue_resolver, ) else: query_images = fetch_dockerhub_images( query, args.pages, args.per_page, args.docker_sort_by, args.sort_order, args.fetch_workers, args.fetch_timeout, not queue_resolver, args.tag_fetch_workers, args.tag_retry_count, args.tag_retry_delay, args.docker_platform_filter_enabled, args.docker_platform_os, args.docker_platform_arch, args.docker_platform_candidate_tags, docker_images_per_repository_limit( getattr(args, 'docker_images_per_repository', 1) ), ) for image in query_images: if image not in seen: images.append(image) seen.add(image) return images if args.platform == 'npm': queries = read_queries(args) if not queries: print('npm search requires --query') return [] targets = [] seen = set() for query in queries: print(f"Fetching npm targets for query: {query}") query_targets = fetch_npm_packages( query, args.pages, args.per_page, args.max_version_age_days, args.fetch_timeout, args.versions_per_package, make_repo_candidate_callback(db, run_id, cycle_id, query), ) for target in query_targets: normalized = normalize_target(target, 'npm') if normalized not in seen: targets.append(target) seen.add(normalized) return targets if args.platform == 'pypi': queries = read_queries(args) if not queries: print('PyPI search requires --query') return [] targets = [] seen = set() for query in queries: print(f"Fetching PyPI targets for query: {query}") query_targets = fetch_pypi_packages( query, args.pages, args.per_page, args.max_version_age_days, args.fetch_timeout, args.versions_per_package, make_repo_candidate_callback(db, run_id, cycle_id, query), ) for target in query_targets: normalized = normalize_target(target, 'pypi') if normalized not in seen: targets.append(target) seen.add(normalized) return targets if args.platform == 'package_git': queries = read_queries(args) if not queries: print('package_git search requires --query') return [] cached_targets = fetch_package_git_targets_from_db(args, db) refresh_registry = bool(getattr(args, 'refresh_registry', False)) if cached_targets and not refresh_registry: print(f'Using {len(cached_targets)} package git repo targets from scanner.db cache') return cached_targets if cached_targets: print(f'Using {len(cached_targets)} cached package git repo target(s) and refreshing registry discovery') else: print('No package git repo candidates found in scanner.db cache; falling back to direct registry discovery') package_sources = [item.strip().lower() for item in str(getattr(args, 'package_sources', 'npm,pypi') or '').split(',') if item.strip()] targets = [] seen = set() for target in cached_targets: normalized = normalize_target(target, 'package_git') if normalized not in seen: targets.append(target) seen.add(normalized) for query in queries: print(f"Fetching package git repo targets for query: {query}") query_targets = [] if 'npm' in package_sources: query_targets.extend(fetch_npm_package_git_repos( query, args.pages, args.per_page, args.max_version_age_days, args.fetch_timeout, args.versions_per_package, )) if 'pypi' in package_sources: query_targets.extend(fetch_pypi_package_git_repos( query, args.pages, args.per_page, args.max_version_age_days, args.fetch_timeout, args.versions_per_package, )) for target in query_targets: normalized = normalize_target(target, 'package_git') if normalized not in seen: targets.append(target) seen.add(normalized) return targets if args.platform == 'postman': queries = read_queries(args) if not queries: print('Postman GitHub code search requires --query') return [] targets = [] seen = set() for query in queries: print(f"Fetching Postman targets from GitHub code search for query: {query}") query_targets = fetch_github_postman_targets( query, args.pages, args.per_page, token_entries=getattr(args, 'github_tokens', None), token=args.token or os.getenv('GITHUB_TOKEN'), search_kinds=getattr(args, 'search_kinds', 'collection,environment'), cache_dir=getattr(args, 'postman_cache_dir', None), max_file_age_days=getattr(args, 'max_file_age_days', 365), max_artifact_size_mb=getattr(args, 'max_artifact_size_mb', 20), request_timeout=getattr(args, 'fetch_timeout', 20), code_search_rpm_per_token=getattr(args, 'github_code_search_rpm', 8), all_tokens_cooldown=getattr(args, 'all_tokens_cooldown', 1800), auth_status=getattr(args, 'github_auth_status', None), discovery_max_artifacts=getattr(args, 'postman_discovery_max_artifacts_per_cycle', None), discovery_max_artifacts_per_page=getattr(args, 'postman_discovery_max_artifacts_per_page', None), discovery_max_bytes=getattr(args, 'postman_discovery_max_bytes_per_cycle', None), discovery_max_elapsed_sec=getattr(args, 'postman_discovery_max_elapsed_sec', None), **stop_options, ) for target in query_targets: normalized = normalize_target(target, 'postman') if normalized not in seen: targets.append(target) seen.add(normalized) return targets if args.platform == 'github_gists': token = args.token or os.getenv('GITHUB_TOKEN') return fetch_github_gist_targets( args.pages, args.per_page, getattr(args, 'gist_since', '') or None, token, getattr(args, 'postman_cache_dir', None), getattr(args, 'max_artifact_size_mb', 2), getattr(args, 'fetch_timeout', 20), discovery_max_artifacts=getattr(args, 'postman_discovery_max_artifacts_per_cycle', None), discovery_max_artifacts_per_page=getattr(args, 'postman_discovery_max_artifacts_per_page', None), discovery_max_bytes=getattr(args, 'postman_discovery_max_bytes_per_cycle', None), discovery_max_elapsed_sec=getattr(args, 'postman_discovery_max_elapsed_sec', None), **stop_options, ) if args.platform == 'github_actions': return fetch_ci_repo_targets_from_db(args, db, 'github') if args.platform == 'gitlab_ci': return fetch_ci_repo_targets_from_db(args, db, 'gitlab') raise ValueError(f'Unsupported platform: {args.platform}') def known_targets_for_args(args, db=None, source_name=None): postgres_configured = bool( (db and getattr(db, 'postgres_required', False)) or str(getattr(args, 'database_url', '') or '').lower().startswith(('postgresql://', 'postgres://')) ) if postgres_configured: return set() todo_file, checked_file = queue_files_for_args(args) known = set() paths = (todo_file,) if args.platform == 'github_archive' else (todo_file, checked_file) for path in paths: for target in load_set_from_file(path): known.add(normalize_target(target, args.platform)) return known def fetch_stop_options(args, db=None, source_name=None): if updated_target_rescan_enabled(args): return {} if not getattr(args, 'stop_on_seen_pages', False): return {} common = { 'normalize_target': lambda target: normalize_target(target, args.platform), 'stop_on_seen_pages': True, 'seen_page_threshold': int(getattr(args, 'seen_page_threshold', 2) or 2), 'min_pages_before_stop': int(getattr(args, 'min_pages_before_stop', 1) or 1), } if db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False): lookup = getattr(db, 'known_target_normalizations_for', None) if not lookup: return {} source_key = source_name or args.platform common['known_target_lookup'] = lambda targets: lookup(source_key, args.platform, targets) return common known = known_targets_for_args(args, db, source_name) if not known: return {} common['known_targets'] = known return common def get_platform_token(args): if args.platform in ('github', 'github_archive'): return args.token or os.getenv('GITHUB_TOKEN') if args.platform == 'gitlab': return args.token or os.getenv('GITLAB_TOKEN') if args.platform == 'package_git': return args.token if args.platform == 'huggingface': return args.token or os.getenv('HF_TOKEN') or os.getenv('HUGGINGFACE_TOKEN') if args.platform == 'postman': return args.token or os.getenv('GITHUB_TOKEN') if args.platform == 'github_gists': return args.token or os.getenv('GITHUB_TOKEN') if args.platform == 'github_actions': return args.token or os.getenv('GITHUB_TOKEN') if args.platform == 'gitlab_ci': return args.token or os.getenv('GITLAB_TOKEN') return None def queue_files_for_args(args): queue_dir = getattr(args, 'queue_dir', None) or args.save_dir return ( os.path.join(queue_dir, f'todo_{args.platform}.txt'), os.path.join(queue_dir, f'checked_{args.platform}.txt'), ) def github_archive_recently_scanned(db, target, cooldown_hours): if not db or not getattr(db, 'conn', None) or int(cooldown_hours or 0) <= 0: return False normalized = normalize_target(target, 'github_archive') normalized_values = [normalized] if '#' in normalized: normalized_values.append(normalized.split('#', 1)[0]) cutoff = (datetime.now() - timedelta(hours=int(cooldown_hours or 0))).isoformat(timespec='seconds') row = db.conn.execute(''' SELECT 1 FROM target_scans WHERE normalized_target IN ({}) AND source = 'github_archive' AND COALESCE(ended_at, started_at, created_at) >= ? LIMIT 1 '''.format(','.join('?' for _ in normalized_values)), [*normalized_values, cutoff]).fetchone() commit_if_postgres(db) return bool(row) def filter_github_archive_targets_by_cooldown(targets, args, db): cooldown_hours = int(getattr(args, 'archive_rescan_cooldown_hours', 0) or 0) if cooldown_hours <= 0 or not db or not getattr(db, 'conn', None): return targets output = [] skipped = 0 for target in targets: if github_archive_recently_scanned(db, target, cooldown_hours): skipped += 1 continue output.append(target) print(f'GHArchive cooldown filter: kept={len(output)} skipped_recent={skipped} cooldown_hours={cooldown_hours}') return output def queue_counts_for_args(args): todo_file, checked_file = queue_files_for_args(args) return queue_counts(todo_file, checked_file) def refund_claim_setup(db, spool, reservation_id, claims, reason): claims = [claim for claim in (claims or []) if claim is not None] if not claims: if reservation_id: if spool.release_reservation(reservation_id) is False: raise RuntimeError('unable to release exact result-spool reservation') return True if hasattr(db, 'refund_target_claims'): refunded = db.refund_target_claims(claims, reason) else: refunded = all( db.refund_target_claim( claim.get('id') if isinstance(claim, dict) else claim['id'], claim.get('lease_token') if isinstance(claim, dict) else claim['lease_token'], reason, ) for claim in claims ) if not refunded: raise RuntimeError('unable to atomically refund fenced target claims after infrastructure setup failure') if reservation_id: if spool.release_reservation(reservation_id) is False: raise RuntimeError('claims were refunded but the exact result-spool reservation was not released') return True class PostClaimRefundGuard: def __init__( self, db, spool, reservation_id, claims, active_tokens, active_lock, dispatch_leases=None, ): self.db = db self.spool = spool self.reservation_id = reservation_id self.claims = [dict(claim) for claim in claims] self.active_tokens = active_tokens self.active_lock = active_lock self.dispatch_leases = list(dispatch_leases or []) def release_unstarted_dispatch_leases(self): for lease in self.dispatch_leases: if lease.heartbeat_thread is None: lease.release() def refund_undurable(self, reason): if not self.reservation_id: return True with self.active_lock: tokens = set(self.active_tokens) claims = [ claim for claim in self.claims if claim.get('lease_token') in tokens ] if claims: refund_claim_setup( self.db, self.spool, self.reservation_id, claims, reason, ) else: release = getattr(self.spool, 'release_reservation', None) if release: try: release(self.reservation_id) except Exception: logger.exception('Unable to release consumed result-spool reservation') self.release_unstarted_dispatch_leases() self.reservation_id = None return True def call(self, reason, function, *args, **kwargs): try: return function(*args, **kwargs) except Exception as exc: self.refund_undurable(f'{reason}: {exc}') raise def _claim_value(claim, name): if isinstance(claim, dict): return claim.get(name) return claim[name] def validate_recovered_claim_batch(db, rows, reservation_id, lease_owner, claim_limit): rows = list(rows or []) identities = [] targets = set() for row in rows: queue_id = _claim_value(row, 'id') lease_token = _claim_value(row, 'lease_token') claim_batch = _claim_value(row, 'claim_batch') owner = _claim_value(row, 'lease_owner') target = str(_claim_value(row, 'target') or '') if ( queue_id is None or not lease_token or not target or str(claim_batch or '') != str(reservation_id) or str(owner or '') != str(lease_owner) ): raise RuntimeError('claim batch recovery returned inconsistent fenced row metadata') identity = (int(queue_id), str(lease_token)) if identity in identities or target in targets: raise RuntimeError('claim batch recovery returned duplicate fenced rows') identities.append(identity) targets.add(target) if len(rows) > int(claim_limit): raise RuntimeError('claim batch recovery exceeded its reserved row limit') expectation_loader = getattr(db, 'claim_recovery_expectation', None) expected = expectation_loader(reservation_id, lease_owner) if expectation_loader else None if expected is not None: expected_identities = { (int(_claim_value(claim, 'id')), str(_claim_value(claim, 'lease_token'))) for claim in expected } if set(identities) != expected_identities: raise RuntimeError('claim batch recovery did not return the exact pre-commit fenced row set') return rows def enqueue_discovered_targets( args, fetched_targets, db, run_id, cycle_id, source_name=None, partial_metrics=None, ): source_key = source_name or args.platform if not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False): raise RuntimeError('Discovery admission requires PostgreSQL') db.require_runtime_safety_schema() if not run_id or not cycle_id: raise RuntimeError('Postgres discovery admission requires valid run_id and cycle_id') discovery_by_normalized = {} for item in fetched_targets or []: if ( args.platform == 'huggingface' and huggingface_discovery_target_is_restricted(item) ): continue target = item.get('target') or item.get('url') if isinstance(item, dict) else item normalized = normalize_target(target, args.platform) if not normalized: continue remote_value = item.get('remote_modified_at') if isinstance(item, dict) else None remote_time = parse_datetime(remote_value) existing = discovery_by_normalized.get(normalized) if existing is None or ( remote_time and (existing['_remote_time'] is None or remote_time > existing['_remote_time']) ): record = { 'target': target, 'remote_modified_at': ( remote_time.isoformat(timespec='seconds') if remote_time else None ), '_remote_time': remote_time, } discovery_by_normalized[normalized] = record discovery_records = [ {key: value for key, value in item.items() if key != '_remote_time'} for item in discovery_by_normalized.values() ] discoveries = [item['target'] for item in discovery_records] eligible = discoveries unresolved = [] if args.platform == 'docker' and discoveries: eligible = [] visible = [] for target in discoveries: try: eligible.append(parse_docker_target(target)['target']) except (TypeError, ValueError): visible.append(target) eligible_seen = set() deduped_eligible = [] for target in eligible: normalized = normalize_target(target, args.platform) if normalized in eligible_seen: continue eligible_seen.add(normalized) deduped_eligible.append(target) eligible = deduped_eligible eligible_normalized = { normalize_target(target, args.platform) for target in eligible } unresolved_seen = set() for target in visible: normalized = normalize_target(target, args.platform) if normalized in eligible_normalized or normalized in unresolved_seen: continue unresolved.append(target) unresolved_seen.add(normalized) projected_targets = [*eligible, *unresolved] expected = len(projected_targets) queued_updated_count = 0 if updated_target_rescan_enabled(args): observer = getattr(db, 'observe_discovered_targets', None) if not observer: raise RuntimeError( 'Updated-target rescan requires atomic Postgres discovery observation' ) observation = observer( source_key, args.platform, args.query, discovery_records, rescan_limit=max( 0, int(getattr(args, 'updated_target_rescan_max_per_cycle', 0) or 0), ), cooldown_seconds=max( 0, int(getattr(args, 'updated_target_rescan_cooldown_hours', 0) or 0) * 3600, ), ) queued_new_count = int(observation.get('queued_new_count', 0) or 0) queued_updated_count = int(observation.get('queued_updated_count', 0) or 0) if int(observation.get('attempted_count', -1)) != expected: raise RuntimeError( f'Unable to observe all Postgres discoveries for {source_key}: ' f'{observation.get("attempted_count")}/{expected}' ) else: if hasattr(db, 'known_target_normalizations_for'): known_before = db.known_target_normalizations_for( source_key, args.platform, projected_targets, ) else: known_before = set() queued_new_count = sum( 1 for target in projected_targets if normalize_target(target, args.platform) not in known_before ) enqueued = db.enqueue_targets( source_key, args.platform, args.query, eligible, requeue_done=(args.platform == 'github_archive'), unresolved_targets=unresolved, discovery_admission=True, ) if enqueued != expected: raise RuntimeError( f'Unable to enqueue all Postgres discoveries for {source_key}: ' f'{enqueued}/{expected}' ) discovery_info = { 'fetched_count': len(fetched_targets or []), 'queued_new_count': queued_new_count, 'queued_updated_count': queued_updated_count, } if isinstance(partial_metrics, dict): partial_metrics.update(discovery_info) return projected_targets, discovery_info def prepare_targets( args, fetched_targets, db=None, run_id=None, cycle_id=None, source_name=None, spool=None, claim_limit_override=None, dispatch_leases=None, enqueue_only=False, partial_metrics=None, experiment_resolver_already_run=False, ): todo_file, checked_file = queue_files_for_args(args) source_key = source_name or args.platform if db and getattr(db, 'postgres_required', False) and not getattr(db, 'conn', None): raise RuntimeError('Required Postgres connection is unavailable') postgres_db = bool(db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False)) project_files = not postgres_db or bool(getattr(args, 'sync_file_queues', True)) if postgres_db: projected_targets, admission = enqueue_discovered_targets( args, fetched_targets, db, run_id, cycle_id, source_key, partial_metrics=partial_metrics, ) queued_new_count = admission['queued_new_count'] queued_updated_count = admission['queued_updated_count'] projection_error = '' if project_files and projected_targets: try: os.makedirs(getattr(args, 'queue_dir', None) or args.save_dir, exist_ok=True) with projection_file_lock(todo_file): existing_projection = { normalize_target(target, args.platform) for target in load_set_from_file(todo_file) } projection_rows = [ target for target in projected_targets if normalize_target(target, args.platform) not in existing_projection ] if projection_rows: _append_lines_unlocked(todo_file, projection_rows) except Exception as exc: projection_error = str(exc) safe_print(f'Warning: Postgres targets committed but todo projection failed: {exc}') if ( args.platform == 'docker' and not bool(getattr(args, 'docker_depth_collection_only', False)) and hasattr(db, 'claim_docker_resolutions') ): resolve_due_docker_queue_targets_if_scan_queue_empty( db, source_key, args, include_experiment=not experiment_resolver_already_run, ) if enqueue_only: queue_info = queue_counts(todo_file, checked_file) if project_files else { 'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None, } queue_info['db_queue'] = db.target_queue_counts(source_key) queue_info.update({ 'fetched_count': len(fetched_targets or []), 'queued_new_count': queued_new_count, 'queued_updated_count': queued_updated_count, 'scan_requested_count': 0, 'lease_owner': None, 'lease_tokens': [], 'lease_seconds': 0, 'queue_claims': {}, 'spool_reservation_id': None, 'scan_slot_leases': [], 'projection_error': projection_error, }) if queued_new_count: print(f'Committed {queued_new_count} new Postgres target(s).') if queued_updated_count: print(f'Committed {queued_updated_count} updated Postgres target rescan(s).') return [], todo_file, checked_file, queue_info worker_count = max(1, int(getattr(args, 'workers', 1) or 1)) configured_batch = int(getattr(args, 'target_claim_batch_size', 0) or 0) claim_limit = ( max(1, int(claim_limit_override)) if claim_limit_override is not None else configured_batch if configured_batch > 0 else worker_count ) if args.max_targets: claim_limit = min(claim_limit, int(args.max_targets)) target_timeout = max(60, int(getattr(args, 'timeout', 0) or 0)) lease_seconds = max(target_timeout + 900, 1800) lease_owner = f'{source_key}:{os.getpid()}:{cycle_id or "cycle"}' if spool is None: spool = result_spool_for_args(args) reservation_id = reserve_result_spool_claims( spool, db, lease_owner, claim_limit, lease_seconds, stop_event=getattr(args, 'result_spool_stop_event', None), wait_seconds=float(getattr(args, 'result_spool_wait_sec', 1.0) or 1.0), diagnostic_interval=float(getattr(args, 'result_spool_diagnostic_interval_sec', 30.0) or 30.0), ) claim_error = None claim_rows_validated = False try: claimed_rows = db.claim_targets( source_key, args.platform, claim_limit, lease_owner, lease_seconds, max_attempts=int(getattr(args, 'target_retry_max_attempts', 3) or 3), return_rows=True, claim_batch=reservation_id, ) except Exception as exc: claimed_rows = None claim_error = exc if claimed_rows is None: original_error = claim_error or RuntimeError( getattr(db, 'last_error', '') or f'claim outcome is ambiguous for {source_key}' ) try: claimed_rows = db.recover_claim_batch(reservation_id, lease_owner) except Exception as recovery_error: raise RuntimeError( f'Unable to recover ambiguous claim batch for {source_key}; ' f'original claim error: {original_error}; recovery error: {recovery_error}' ) from original_error claimed_rows = validate_recovered_claim_batch( db, claimed_rows, reservation_id, lease_owner, claim_limit, ) claim_rows_validated = True if not claimed_rows: try: released = spool.release_reservation(reservation_id) if released is not True: raise RuntimeError('exact result-spool reservation release was not confirmed') except Exception as release_error: raise RuntimeError( f'FATAL durability error after zero-row claim recovery for {source_key}; ' f'original claim error: {original_error}; reservation rollback error: {release_error}' ) from original_error raise original_error if not claim_rows_validated: claimed_rows = validate_recovered_claim_batch( db, claimed_rows, reservation_id, lease_owner, claim_limit, ) setup_tokens = { str(_claim_value(row, 'lease_token')) for row in claimed_rows if _claim_value(row, 'lease_token') } setup_guard = PostClaimRefundGuard( db, spool, reservation_id, claimed_rows, setup_tokens, threading.Lock(), dispatch_leases=dispatch_leases, ) try: reservation_id = spool.bind_claims(reservation_id, claimed_rows) targets_to_scan = [row['target'] for row in claimed_rows] leases = list(dispatch_leases or []) scan_slot_leases = leases[:len(targets_to_scan)] for extra_lease in leases[len(targets_to_scan):]: extra_lease.release() queue_claims = { str(row['target']): { 'id': row['id'], 'attempts': int(row['attempts'] or 0), 'lease_owner': row['lease_owner'], 'lease_token': row['lease_token'], 'claim_batch': row['claim_batch'], } for row in claimed_rows } queue_info = queue_counts(todo_file, checked_file) if project_files else { 'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None, } queue_info['db_queue'] = db.target_queue_counts(source_key) queue_info.update({ 'fetched_count': len(fetched_targets or []), 'queued_new_count': queued_new_count, 'queued_updated_count': queued_updated_count, 'scan_requested_count': len(targets_to_scan), 'lease_owner': lease_owner, 'lease_tokens': [row['lease_token'] for row in claimed_rows], 'lease_seconds': lease_seconds, 'queue_claims': queue_claims, 'spool_reservation_id': reservation_id, 'scan_slot_leases': scan_slot_leases, 'projection_error': projection_error, }) except Exception as exc: setup_guard.refund_undurable(f'infrastructure claim setup failure: {exc}') raise if queued_new_count: print(f'Committed {queued_new_count} new Postgres target(s).') if queued_updated_count: print(f'Committed {queued_updated_count} updated Postgres target rescan(s).') return targets_to_scan, todo_file, checked_file, queue_info os.makedirs(args.save_dir, exist_ok=True) os.makedirs(getattr(args, 'queue_dir', None) or args.save_dir, exist_ok=True) with projection_file_lock(todo_file): checked = load_set_from_file(checked_file) todo = load_set_from_file(todo_file) checked_normalized = set() if args.platform == 'github_archive' else {normalize_target(target, args.platform) for target in checked} todo_normalized = {normalize_target(target, args.platform) for target in todo} new_targets = [] for item in fetched_targets: target = item.get('target') or item.get('url') if isinstance(item, dict) else item normalized = normalize_target(target, args.platform) if normalized not in checked_normalized and normalized not in todo_normalized: new_targets.append(target) todo.add(target) todo_normalized.add(normalized) if new_targets: _append_lines_unlocked(todo_file, new_targets) print(f'Queued {len(new_targets)} new targets.') targets_to_scan = [ target for target in sorted(todo) if normalize_target(target, args.platform) not in checked_normalized ] if args.platform == 'docker': targets_to_scan, todo_entries = resolve_docker_targets_for_scan(targets_to_scan, args) _write_lines_unlocked(todo_file, sorted(set(todo_entries))) queue_info = queue_counts(todo_file, checked_file) queue_info.update({ 'fetched_count': len(fetched_targets), 'queued_new_count': len(new_targets), 'queued_updated_count': 0, 'scan_requested_count': len(targets_to_scan), 'lease_owner': None, 'lease_tokens': [], 'lease_seconds': 0, 'queue_claims': {}, 'spool_reservation_id': None, 'scan_slot_leases': [], }) return targets_to_scan, todo_file, checked_file, queue_info def resolve_docker_targets_for_scan(targets, args, resolve_all=False): def has_tag(target): text = str(target or '').strip() return '@' in text or ':' in text.rsplit('/', 1)[-1] tagged_targets = [target for target in targets if has_tag(target)] bare_targets = [target for target in targets if not has_tag(target)] if not bare_targets: return tagged_targets, tagged_targets resolve_limit = 0 if resolve_all else int(getattr(args, 'tag_resolve_limit', 100) or 0) if resolve_limit > 0: bare_to_resolve = bare_targets[:resolve_limit] bare_to_keep = bare_targets[resolve_limit:] else: bare_to_resolve = bare_targets bare_to_keep = [] print( f'Resolving tags for {len(bare_to_resolve)} queued Docker repositories ' f'({len(bare_to_keep)} deferred)...' ) import concurrent.futures max_workers = max(1, min(args.tag_fetch_workers, len(bare_to_resolve))) resolved_targets = list(tagged_targets) unresolved_targets = [] with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: futures = { executor.submit( fetch_dockerhub_tags, target, None, docker_images_per_repository_limit( getattr(args, 'docker_images_per_repository', 1) ), args.tag_retry_count, args.tag_retry_delay, args.docker_platform_filter_enabled, args.docker_platform_os, args.docker_platform_arch, args.docker_platform_candidate_tags, True, ): target for target in bare_to_resolve } for future in concurrent.futures.as_completed(futures): target = futures[future] tags, status = future.result() if tags: resolved_targets.extend(tags) if status != 'ok': print(f'Preserving unresolved Docker repository {target}: tag resolution status={status}') unresolved_targets.append(target) todo_entries = resolved_targets + unresolved_targets + bare_to_keep return resolved_targets, todo_entries def resolve_due_docker_queue_targets(db, source, args): limit = max(1, min(100, int(getattr(args, 'tag_resolve_limit', 100) or 100))) owner = f'docker-resolver:{source}:{os.getpid()}:{threading.get_ident()}' periodic_limit = max( 0, min(1, int(getattr(args, 'docker_repository_refresh_max_per_cycle', 0) or 0)), ) refresh_interval = max( 0, int(getattr(args, 'docker_repository_refresh_interval_sec', 0) or 0), ) allowed_queries = getattr(args, 'configured_queries', None) processed = 0 stopped = False def process(row): nonlocal processed paused_error = None try: _require_discovery_provider_enabled(db) outcome = fetch_dockerhub_tags( row['target'], None, docker_images_per_repository_limit( getattr(args, 'docker_images_per_repository', 1) ), int(getattr(args, 'tag_retry_count', 2) or 2), int(getattr(args, 'tag_retry_delay', 5) or 5), bool(getattr(args, 'docker_platform_filter_enabled', True)), str(getattr(args, 'docker_platform_os', 'linux')), str(getattr(args, 'docker_platform_arch', 'amd64')), int(getattr(args, 'docker_platform_candidate_tags', 20) or 20), return_outcome=True, ) except DiscoveryPausedError as exc: paused_error = exc outcome = SimpleNamespace( tags=(), status='global_cooldown', remote_attempted=False, retry_at=None, error='Discovery paused by runtime control.', ) except Exception as exc: outcome = SimpleNamespace( tags=(), status='unknown', remote_attempted=True, retry_at=None, error=str(exc), ) resolver_status = outcome.status resolver_retry_at = outcome.retry_at resolver_error = outcome.error or f'Docker tag resolution status={outcome.status}' resolver_remote_attempted = outcome.remote_attempted if not db.finish_docker_resolution( source, row['id'], row['resolver_token'], list(outcome.tags), resolver_error, complete=docker_tag_resolution_is_conclusive(outcome.status), retry_at=resolver_retry_at, claim_attempt_consumed=( False if paused_error is not None else not ( resolver_status == 'global_cooldown' and not resolver_remote_attempted and bool(resolver_retry_at) ) ), refresh_interval_sec=refresh_interval, ): raise RuntimeError(f'Unable to persist Docker resolver outcome for queue row {row["id"]}') processed += 1 if paused_error is not None: raise paused_error return bool( resolver_retry_at or resolver_status in ( 'global_cooldown', 'rate_limited', 'auth_failed', 'remote_transient', ) ) retry_budget = limit - 1 if periodic_limit and limit > 1 else limit for _ in range(retry_budget): rows = db.claim_docker_resolutions( source, 1, owner, lease_seconds=300, allowed_queries=allowed_queries, ) if not rows: break if process(rows[0]): stopped = True break if not stopped and periodic_limit: rows = db.claim_docker_resolutions( source, 1, owner, lease_seconds=300, periodic_limit=1, periodic_only=True, allowed_queries=allowed_queries, ) if rows: stopped = process(rows[0]) while not stopped and processed < limit: rows = db.claim_docker_resolutions( source, 1, owner, lease_seconds=300, allowed_queries=allowed_queries, ) if not rows: break stopped = process(rows[0]) if processed: print(f'Processed {processed} due PostgreSQL Docker resolver row(s).') return processed def resolve_due_docker_experiment_targets(db, source, args): authority = getattr(args, 'docker_depth_experiment_authority', None) if ( source != 'dockerhub' or not isinstance(authority, dict) or authority.get('enabled') is not True or bool(getattr(args, 'sync_file_queues', True)) or not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False) or not callable(getattr(db, 'claim_docker_depth_experiment_resolutions', None)) or not callable(getattr(db, 'renew_docker_depth_experiment_resolution', None)) or not callable(getattr(db, 'finish_docker_depth_experiment_resolution', None)) ): return 0, False limit = max(1, min(100, int(getattr(args, 'tag_resolve_limit', 100) or 100))) owner = f'docker-depth-resolver:{source}:{os.getpid()}:{threading.get_ident()}' processed = 0 stopped = False for _ in range(limit): rows = db.claim_docker_depth_experiment_resolutions( source, 1, owner, lease_seconds=300, authority=authority, final_cutover=True, ) if not rows: break row = rows[0] paused_error = None def renew(): _require_discovery_provider_enabled(db) renewal = db.renew_docker_depth_experiment_resolution( source, row['id'], row['resolver_generation'], row['resolver_token'], resolver_owner=row['resolver_owner'], lease_seconds=300, authority=authority, final_cutover=True, ) return bool(renewal.get('renewed')) try: if not renew(): raise DockerResolverLeaseLostError( 'Docker depth resolver lease is no longer owned' ) outcome = fetch_dockerhub_tags( row['target'], None, int(row['selection_limit']), int(getattr(args, 'tag_retry_count', 2) or 2), int(getattr(args, 'tag_retry_delay', 5) or 5), bool(getattr(args, 'docker_platform_filter_enabled', True)), str(getattr(args, 'docker_platform_os', 'linux')), str(getattr(args, 'docker_platform_arch', 'amd64')), int(getattr(args, 'docker_platform_candidate_tags', 20) or 20), return_outcome=True, fresh_graph_evidence=True, lease_renewal_callback=renew, selector_version=authority['selector_version'], ) if not renew(): raise DockerResolverLeaseLostError( 'Docker depth resolver lease expired during remote resolution' ) except DockerResolverLeaseLostError: stopped = True break except DiscoveryPausedError as exc: paused_error = exc outcome = SimpleNamespace( tags=(), status='global_cooldown', remote_attempted=False, retry_at=None, error='Discovery paused by runtime control.', selection_records=(), selector_version='', selector_hash='', candidate_distinct_graph_count=0, fresh_graph_evidence=False, cache_bypassed=True, ) except Exception as exc: outcome = SimpleNamespace( tags=(), status='unknown', remote_attempted=True, retry_at=None, error=str(exc), selection_records=(), selector_version='', selector_hash='', candidate_distinct_graph_count=0, fresh_graph_evidence=False, cache_bypassed=True, ) result = db.finish_docker_depth_experiment_resolution( source, row['id'], row['resolver_generation'], row['resolver_token'], outcome, outcome.error or f'Docker tag resolution status={outcome.status}', complete=docker_tag_resolution_is_conclusive(outcome.status), resolver_owner=row['resolver_owner'], authority=authority, retry_at=outcome.retry_at, claim_attempt_consumed=( False if paused_error is not None else not ( outcome.status == 'global_cooldown' and not outcome.remote_attempted and bool(outcome.retry_at) ) ), final_cutover=True, ) if result.get('committed'): processed += 1 if result.get('status') in ('stale', 'unavailable', 'invalid', 'error'): raise RuntimeError( f'Unable to persist Docker experiment resolver outcome for member {row["id"]}: ' f'{result.get("status")}' ) if paused_error is not None: raise paused_error stopped = bool( result.get('status') in ('conflict_retry', 'held') or outcome.retry_at or outcome.status in ('global_cooldown', 'rate_limited', 'auth_failed') ) if stopped: break if processed: print(f'Processed {processed} fenced Docker depth experiment resolver row(s).') return processed, stopped def resolve_due_docker_queue_targets_if_scan_queue_empty( db, source, args, *, include_experiment=True, ): experiment_processed, experiment_stopped = (0, False) if include_experiment: experiment_processed, experiment_stopped = resolve_due_docker_experiment_targets( db, source, args, ) if experiment_stopped: return experiment_processed backlog_probe = getattr(db, 'has_claimable_targets_v2', None) if callable(backlog_probe): has_scan_targets = backlog_probe( source, 'docker', max_attempts=max(0, int(getattr(args, 'target_retry_max_attempts', 3) or 3)), ) if getattr(db, 'last_error', None): raise RuntimeError( f'Unable to inspect Docker scan backlog before resolver work: {db.last_error}' ) if has_scan_targets: return experiment_processed return experiment_processed + resolve_due_docker_queue_targets(db, source, args) def mark_checked(results, todo_file, checked_file, platform): completed_targets = [ result.get('target', '') for result in results if result.get('target') and str(result.get('skipped') or '') not in CI_SOFT_SKIP_REASONS.get(platform, set()) ] with projection_file_lock(todo_file): checked_normalized = {normalize_target(target, platform) for target in load_set_from_file(checked_file)} new_completed_targets = [] for result in results: target = result.get('target', '') if not target: continue if str(result.get('skipped') or '') in CI_SOFT_SKIP_REASONS.get(platform, set()): continue normalized = normalize_target(target, platform) if normalized not in checked_normalized: new_completed_targets.append(target) checked_normalized.add(normalized) if new_completed_targets: _append_lines_unlocked(checked_file, new_completed_targets) completed_normalized = {normalize_target(target, platform) for target in completed_targets} remaining = [ target for target in load_set_from_file(todo_file) if normalize_target(target, platform) not in completed_normalized ] _write_lines_unlocked(todo_file, sorted(remaining)) def enqueue_targets_for_platform(queue_dir, platform, targets): if not targets: return 0 os.makedirs(queue_dir, exist_ok=True) todo_file = os.path.join(queue_dir, f'todo_{platform}.txt') checked_file = os.path.join(queue_dir, f'checked_{platform}.txt') with projection_file_lock(todo_file): todo = load_set_from_file(todo_file) checked = load_set_from_file(checked_file) known = {normalize_target(target, platform) for target in todo} known.update(normalize_target(target, platform) for target in checked) new_targets = [] for target in targets: normalized = normalize_target(target, platform) if normalized in known: continue new_targets.append(target) known.add(normalized) if new_targets: _append_lines_unlocked(todo_file, new_targets) return len(new_targets) def collect_postman_targets(results): targets = [] for result in results or []: for target in result.get('postman_targets') or []: if target: targets.append(target) return targets def publish_scan_payload(result): try: publication_result = copy.deepcopy(result) candidate_result = copy.deepcopy(result) if candidate_result.get('structured_keycheck_pending') and isinstance(candidate_result.get('postman'), dict): postman_data = candidate_result['postman'] candidate_path, _ = validate_postman_cache_artifact( postman_data, int(candidate_result.get('postman_max_artifact_size_mb') or 20), expected_size=candidate_result.get('bytes'), ) write_structured_keycheck_candidates(candidate_path, postman_data) if not save_scan_result(publication_result): raise RuntimeError('one or more scanner JSONL/keycheck publications failed') return True, '' except Exception as exc: return False, str(exc) def drain_scan_publication_outbox(db, limit=100, require_v2_schema=True): if not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False): return 0 lease_seconds = 300 heartbeat_interval = min(60.0, lease_seconds / 3.0) heartbeat_join_timeout = 12.0 db_path = getattr(db, 'path', None) db_url = getattr(db, 'url', None) # Every enabled production ScannerDB has one of these identities. Their # absence is tolerated for lightweight test doubles only. canonical_db = bool(db_path or db_url) def fenced_finish(row_id, owner, delivered, error): try: return bool(db.finish_scan_publication(row_id, owner, delivered, error)) except Exception as exc: logger.error('Publication outbox fenced finish failed for row %s: %s', row_id, exc) return False def stop_heartbeat(stop_event, thread, heartbeat_db): if stop_event is not None: stop_event.set() alive = False if thread is not None: try: alive = thread.is_alive() if alive: thread.join(heartbeat_join_timeout) alive = thread.is_alive() except Exception as exc: alive = True logger.error('Unable to join publication heartbeat thread: %s', exc) if alive: logger.error('Publication heartbeat did not stop within %.1f seconds', heartbeat_join_timeout) elif heartbeat_db is not None: try: heartbeat_db.close() except Exception as exc: logger.error('Unable to close publication heartbeat DB: %s', exc) return not alive delivered = 0 for sequence in range(max(1, int(limit or 100))): lease_owner = ( f'publisher:{os.getpid()}:{threading.get_ident()}:{sequence}:' f'{secrets.token_urlsafe(18)}' ) rows = db.claim_scan_publications(lease_owner, 1) if not rows: break row = rows[0] try: result = json.loads(row['payload_json']) except (TypeError, ValueError) as exc: if not fenced_finish(row['id'], lease_owner, False, f'invalid publication payload: {exc}'): if canonical_db: logger.warning('Publication ownership handoff detected for row %s; stopping this drain pass', row['id']) break logger.critical('Publication outbox lease acknowledgement failed for row %s', row['id']) raise RuntimeError(f'publication outbox lease acknowledgement failed for row {row["id"]}') continue heartbeat_db = None heartbeat_stop = None heartbeat_thread = None heartbeat_lost = None heartbeat_failed = None if canonical_db: try: heartbeat_stop = threading.Event() heartbeat_lost = threading.Event() heartbeat_failed = threading.Event() heartbeat_db = ScannerDB(db_path=db_path, db_url=db_url, initialize=False) if not heartbeat_db.enabled: raise RuntimeError('dedicated publication heartbeat DB is unavailable') if getattr(heartbeat_db.conn, 'is_postgres', False): heartbeat_db.conn.execute("SELECT set_config('statement_timeout', '10000ms', false)") heartbeat_db.conn.execute("SELECT set_config('lock_timeout', '5000ms', false)") heartbeat_db.conn.commit() if require_v2_schema: heartbeat_db.require_runtime_safety_schema() if not heartbeat_db.renew_scan_publication(row['id'], lease_owner, lease_seconds): detail = getattr(heartbeat_db, 'last_error', '') raise RuntimeError(detail or 'initial publication lease renewal did not confirm ownership') def renew_publication_lease(connection=heartbeat_db): try: while not heartbeat_stop.wait(heartbeat_interval): if connection.renew_scan_publication(row['id'], lease_owner, lease_seconds): continue if getattr(connection, 'last_error', ''): heartbeat_failed.set() logger.error( 'Publication lease heartbeat failed for row %s: %s', row['id'], connection.last_error, ) else: heartbeat_lost.set() logger.warning('Publication lease ownership was lost for row %s', row['id']) return except Exception as exc: heartbeat_failed.set() logger.error('Publication lease heartbeat failed for row %s: %s', row['id'], exc) finally: try: connection.close() except Exception as exc: logger.error('Unable to close publication heartbeat DB: %s', exc) heartbeat_thread = threading.Thread( target=renew_publication_lease, name=f'scan-publication-heartbeat-{row["id"]}', daemon=True, ) heartbeat_thread.start() except Exception as exc: stopped = stop_heartbeat(heartbeat_stop, heartbeat_thread, heartbeat_db) setup_error = f'publication lease heartbeat setup failed: {exc}' acknowledged = fenced_finish(row['id'], lease_owner, False, setup_error) if acknowledged: logger.error('Publication row %s was requeued after heartbeat setup failure: %s', row['id'], exc) else: logger.warning( 'Publication ownership handoff detected while handling setup failure for row %s', row['id'], ) if not stopped: logger.error('Publication row %s heartbeat cleanup remains in progress', row['id']) break try: ok, error = publish_scan_payload(result) except Exception as exc: ok, error = False, str(exc) heartbeat_stopped = True if canonical_db: heartbeat_stopped = stop_heartbeat(heartbeat_stop, heartbeat_thread, heartbeat_db) acknowledged = fenced_finish(row['id'], lease_owner, ok, error) if not acknowledged: if canonical_db: logger.warning('Publication ownership handoff detected for row %s; stopping this drain pass', row['id']) break logger.critical('Publication outbox lease acknowledgement failed for row %s', row['id']) raise RuntimeError(f'publication outbox lease acknowledgement failed for row {row["id"]}') if ok: delivered += 1 if canonical_db and ( not heartbeat_stopped or heartbeat_lost.is_set() or heartbeat_failed.is_set() ): logger.warning('Stopping publication drain after heartbeat ownership uncertainty for row %s', row['id']) break return delivered def require_scan_publication_capacity(db, args, additional_items=1): health_loader = getattr(db, 'scan_publication_backlog_health', None) if health_loader is None: return None health = health_loader( int(getattr(args, 'scan_outbox_max_pending_items', 10000)), int(getattr(args, 'scan_outbox_max_pending_bytes', 1024 * 1024 * 1024)), int(getattr(args, 'scan_outbox_max_pending_age_sec', 24 * 60 * 60)), additional_items=additional_items, ) if not isinstance(health, dict) or not health.get('accepting'): detail = (health or {}).get('reason') if isinstance(health, dict) else 'health query failed' raise RuntimeError(f'scan publication backlog gate is closed: {detail or "configured bound reached"}') return health def result_spool_for_args(args): runtime_dir = os.path.realpath(os.path.abspath(getattr(args, 'runtime_dir', '') or '')) spool_dir = os.path.realpath(os.path.abspath(getattr(args, 'result_spool_dir', '') or '')) if not runtime_dir or not spool_dir: raise RuntimeError('Postgres queue mode requires runtime_dir and result_spool_dir') try: within_runtime = os.path.commonpath((runtime_dir, spool_dir)) == runtime_dir except ValueError: within_runtime = False if not within_runtime or spool_dir == runtime_dir: raise RuntimeError('result_spool_dir must be a dedicated directory under runtime_dir') return ResultSpool( spool_dir, max_event_bytes=int(getattr(args, 'result_spool_max_event_bytes', 192 * 1024 * 1024)), max_events=int(getattr(args, 'result_spool_max_events', 10000)), max_total_bytes=int(getattr(args, 'result_spool_max_total_bytes', 2 * 1024 * 1024 * 1024)), min_free_bytes=int(getattr(args, 'result_spool_min_free_bytes', 1024 * 1024 * 1024)), ) class ResultSpoolPublisherBusy(RuntimeError): def __init__(self, state): self.state = dict(state or {}) super().__init__('result-spool publisher lease is held by another live database session') class ResultSpoolPublisherHandoff(RuntimeError): pass def supervised_spool_stop_requested(): instance_file = str(os.getenv('TRUF_SUPERVISOR_INSTANCE_FILE') or '') inherited_id = str(os.getenv('TRUF_SUPERVISOR_INSTANCE_ID') or '') if not instance_file or not inherited_id: return False try: metadata = read_private_json(instance_file) except (OSError, ValueError): return False if str(metadata.get('instance_id') or '') != inherited_id: raise RuntimeError('supervisor instance identity changed while waiting on result-spool backpressure') return str(metadata.get('activation_state') or '').upper() in ('STOPPING', 'FAILED_HOLD', 'INACTIVE') def drain_result_spool(spool, db): if not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False): raise RuntimeError('result spool ingestion requires PostgreSQL') db.require_runtime_safety_schema() acquire = getattr(db, 'try_acquire_result_spool_publisher', None) release = getattr(db, 'release_result_spool_publisher', None) acquired = False if acquire: acquired = bool(acquire()) if not acquired: inspect_owner = getattr(db, 'result_spool_publisher_state', None) state = inspect_owner() if inspect_owner else { 'status': 'held', 'authenticated': True, 'application_name': 'test-publisher', 'holder_identity': 'test', } if not isinstance(state, dict): raise RuntimeError('result-spool publisher ownership lookup returned an invalid state') if state.get('status') == 'free' and state.get('authenticated') is True: raise ResultSpoolPublisherHandoff( 'result-spool publisher released ownership before inspection' ) if ( state.get('status') != 'held' or state.get('authenticated') is not True or not state.get('application_name') or not state.get('holder_identity') ): raise RuntimeError('result-spool publisher ownership is unknown or unauthenticated') raise ResultSpoolPublisherBusy(state) outcomes = {} legacy_records = None try: while True: if hasattr(spool, 'next_pending_event'): record = spool.next_pending_event() else: if legacy_records is None: legacy_records = iter(spool.pending_events()) record = next(legacy_records, None) if record is None: break try: outcome = db.ingest_scan_event(record.envelope) except ScanEventConflictError as exc: spool.quarantine_event(record.event_id, f'database scan-event hash conflict: {exc}') raise RuntimeError(f'scan event {record.event_id} was quarantined after a database hash conflict') from exc if not scan_outcome_matches(record, outcome): raise RuntimeError(f'database did not confirm ingestion of scan event {record.event_id}') if spool.acknowledge(record.event_id, record.event_hash) is not True: raise RuntimeError(f'result-spool acknowledgement was not confirmed for event {record.event_id}') outcomes[record.event_id] = outcome spool.assert_claims_allowed() return outcomes finally: if acquired and (not release or release() is not True): raise RuntimeError('result-spool publisher lease release was not confirmed') def wait_for_result_spool_ready( spool, db, outcomes=None, stop_event=None, wait_seconds=1.0, diagnostic_interval=30.0, ): outcomes = outcomes if outcomes is not None else {} wait_seconds = max(0.05, float(wait_seconds or 1.0)) diagnostic_interval = max(wait_seconds, float(diagnostic_interval or 30.0)) last_diagnostic = 0.0 while True: if (stop_event is not None and stop_event.is_set()) or supervised_spool_stop_requested(): raise KeyboardInterrupt('coordinated shutdown while waiting for result-spool backpressure') try: drained = drain_result_spool(spool, db) outcomes.update(drained) if drained: safe_print(f'Drained {len(drained)} prior durable result-spool event(s).', flush=True) return outcomes except ResultSpoolPublisherBusy as exc: publisher = str(exc.state.get('application_name') or '')[:80] state = None reason = f'publisher={publisher}' except ResultSpoolPublisherHandoff: state = None reason = 'publisher handoff observed; retrying acquisition' except (SpoolBackpressureError, SpoolContentionError) as exc: reason = str(exc)[:200] try: state = spool.backpressure_state() if hasattr(spool, 'backpressure_state') else None except SpoolContentionError: state = None now = time.monotonic() if now - last_diagnostic >= diagnostic_interval: details = '' if state: details = ( f" pending={state.get('pending_count', 0)}" f" bytes={state.get('pending_bytes', 0)}" f" oldest_age_sec={state.get('oldest_pending_age_sec', 0)}" f" reservations={state.get('reservation_count', 0)}" ) safe_print(f'Result-spool backpressure; waiting in-process:{details} {reason}', flush=True) last_diagnostic = now if stop_event is not None: if stop_event.wait(wait_seconds): raise KeyboardInterrupt('coordinated shutdown while waiting for result-spool backpressure') else: time.sleep(wait_seconds) def reserve_result_spool_claims( spool, db, owner, count, lease_seconds, stop_event=None, wait_seconds=1.0, diagnostic_interval=30.0, ): last_capacity_diagnostic = 0.0 while True: if hasattr(spool, 'next_pending_event') or hasattr(spool, 'pending_events'): wait_for_result_spool_ready( spool, db, stop_event=stop_event, wait_seconds=wait_seconds, diagnostic_interval=diagnostic_interval, ) try: return spool.reserve_claims(owner, count, lease_seconds) except (SpoolBackpressureError, SpoolContentionError): continue except SpoolTransientCapacityError as exc: snapshot_loader = getattr(spool, 'reservation_snapshot', None) progress_loader = getattr(db, 'result_spool_reservation_progress', None) if not snapshot_loader or not progress_loader: raise RuntimeError('result-spool temporary capacity ownership cannot be verified') from exc snapshot = snapshot_loader() if not snapshot.get('reservations'): if (stop_event is not None and stop_event.is_set()) or supervised_spool_stop_requested(): raise KeyboardInterrupt('coordinated shutdown during result-spool capacity handoff') if stop_event is not None: if stop_event.wait(wait_seconds): raise KeyboardInterrupt('coordinated shutdown during result-spool capacity handoff') else: time.sleep(max(0.05, wait_seconds)) continue progress = progress_loader(snapshot['reservations']) if not isinstance(progress, dict) or progress.get('safe_progress') is not True: raise RuntimeError('result-spool reservation progress is unsafe or indeterminate') from exc now = time.monotonic() if now - last_capacity_diagnostic >= max(wait_seconds, diagnostic_interval): counts = progress.get('counts') or {} safe_print( 'Result-spool capacity reserved; waiting in-process: ' f"reserved_bytes={snapshot.get('reserved_future_bytes', 0)} " f"reservations={snapshot.get('reservation_count', 0)} " f"exact_live={counts.get('exact_live', 0)} " f"exact_expired={counts.get('exact_expired', 0)} " f"stale={counts.get('stale_or_reassigned', 0)} " f"next_progress_sec={progress.get('earliest_progress_in_sec', 0)}", flush=True, ) last_capacity_diagnostic = now if (stop_event is not None and stop_event.is_set()) or supervised_spool_stop_requested(): raise KeyboardInterrupt('coordinated shutdown while waiting for result-spool capacity') if stop_event is not None: if stop_event.wait(wait_seconds): raise KeyboardInterrupt('coordinated shutdown while waiting for result-spool capacity') else: time.sleep(max(0.05, wait_seconds)) def scan_outcome_matches(record, outcome): return bool( outcome and outcome.get('ingested') and str(outcome.get('scan_event_id') or '') == str(record.event_id) and str(outcome.get('scan_event_hash') or '') == str(record.event_hash) ) def recheck_active_lease_ownership(db, lease_owner, requested_tokens, active_tokens, active_lock): with active_lock: expected_before_query = set(requested_tokens).intersection(active_tokens) if not expected_before_query: return True confirmed = set(db.active_target_lease_tokens(lease_owner, expected_before_query)) with active_lock: expected_after_query = expected_before_query.intersection(active_tokens) return not expected_after_query.difference(confirmed) def renew_active_result_spool_reservation( spool, reservation_id, lease_seconds, reserved_tokens, active_lock, ): with active_lock: if not reservation_id or not reserved_tokens: return True return spool.renew_reservation(reservation_id, lease_seconds) is True def write_reserved_result_spool_event( spool, event, reservation_id, queue_id, lease_token, reserved_tokens, active_tokens, active_lock, ): # Keep the token transition and reservation consumption in one lock order. # A heartbeat can then observe either the live reservation or its completed # durable handoff, never the transient gap between those states. with active_lock: record = spool.write_event( event, reservation_id=reservation_id, queue_id=queue_id, lease_token=lease_token, ) reserved_tokens.discard(lease_token) active_tokens.discard(lease_token) return record def first_error_line(result): for error in result.get('errors', []): for line in str(error).splitlines(): line = line.strip() if line: try: import json payload = json.loads(line) if payload.get('error'): return str(payload.get('error'))[:300] if payload.get('msg'): return str(payload.get('msg'))[:300] except Exception: pass return line[:300] return '' def target_retry_delay_sec(attempt, base_delay_sec=3600, max_delay_sec=86400): attempt = max(1, int(attempt or 1)) base_delay_sec = max(1, int(base_delay_sec or 3600)) max_delay_sec = max(base_delay_sec, int(max_delay_sec or 86400)) return min(max_delay_sec, base_delay_sec * (2 ** max(0, attempt - 1))) def queue_error_disposition(db, source, platform, target, result, args, queue_claim=None): queue_row = None if queue_claim else db.target_queue_item(source, platform, target) if db else None attempts = int((queue_claim or {}).get('attempts') or 0) if queue_claim else int(queue_row['attempts'] or 0) if queue_row else 1 max_attempts = max(1, int(getattr(args, 'target_retry_max_attempts', 3) or 3)) timed_out = bool((result.get('scan_meta') or {}).get('command_timed_out')) or result.get('error_class') == 'timeout' if timed_out: if attempts >= max_attempts: return 'failed', None, attempts, max_attempts delay = max(60, int(getattr(args, 'target_timeout_retry_delay_sec', 21600) or 21600)) available_after = (datetime.now(timezone.utc) + timedelta(seconds=delay)).isoformat(timespec='seconds') return 'deferred', available_after, attempts, max_attempts if result.get('source_failure'): if not bool(result.get('retryable', True)) and attempts >= max_attempts: return 'failed', None, attempts, max_attempts delay = max(1, int(getattr(args, 'target_retry_max_delay_sec', 86400) or 86400)) available_after = (datetime.now(timezone.utc) + timedelta(seconds=delay)).isoformat(timespec='seconds') return 'deferred', available_after, attempts, max_attempts retryable = bool(result.get('retryable', True)) if not retryable or attempts >= max_attempts: return 'failed', None, attempts, max_attempts delay = target_retry_delay_sec( attempts, getattr(args, 'target_retry_base_delay_sec', 3600), getattr(args, 'target_retry_max_delay_sec', 86400), ) available_after = (datetime.now(timezone.utc) + timedelta(seconds=delay)).isoformat(timespec='seconds') return 'deferred', available_after, attempts, max_attempts def queue_result_resets_attempts(result): return bool( (result.get('source_failure') and result.get('retryable', True)) or docker_layer_result_resets_attempts(result) ) def docker_layer_result_resets_attempts(result): if result.get('docker_layer_plan') is None or not result.get('retryable', False): return False execution = result.get('docker_layer_execution') records = execution.get('blobs') if isinstance(execution, dict) else None if not isinstance(records, list): return False descriptors = result['docker_layer_plan'].get('descriptors') if not isinstance(descriptors, list) or not any( item.get('coverage_state') in ('selected', 'shared_pending') for item in descriptors if isinstance(item, dict) ): return False return not any( item.get('status') in ('retryable_failed', 'terminal_failed') for item in records if isinstance(item, dict) ) def docker_layer_queue_disposition(result, args, attempts=None): if result.get('docker_layer_plan') is None: return None if not result.get('errors'): return 'done', None, False if not bool(result.get('retryable', False)): return 'failed', None, False reset_attempts = docker_layer_result_resets_attempts(result) max_attempts = max(1, int(getattr(args, 'target_retry_max_attempts', 3) or 3)) if not reset_attempts and attempts is not None and int(attempts or 0) >= max_attempts: return 'failed', None, False delay = max(1, int(getattr(args, 'docker_layer_checkpoint_delay_sec', 60) or 60)) available_after = ( datetime.now(timezone.utc) + timedelta(seconds=delay) ).isoformat(timespec='seconds') return 'deferred', available_after, reset_attempts def finding_summary(finding): detector = finding.get('DetectorName', 'Unknown') verified = 'verified' if finding.get('Verified', False) else 'unverified' source = finding.get('SourceMetadata', {}).get('Data', {}) git_source = source.get('Git', {}) if isinstance(source, dict) else {} file_name = git_source.get('file', '') line_number = git_source.get('line', '') location = '' if file_name and line_number: location = f' at {file_name}:{line_number}' elif file_name: location = f' at {file_name}' return f'{detector} ({verified}){location}' def print_result_report(results, save_dir): results_file = os.path.join(save_dir, 'scan_results.jsonl') secrets_file = os.path.join(save_dir, 'found_secrets.jsonl') errors_file = os.path.join(save_dir, 'scan_errors.log') for index, result in enumerate(results, 1): target = result.get('target', '') findings = result.get('findings', []) errors = result.get('errors', []) warnings = result.get('warnings', []) skipped = result.get('skipped') safe_print(f'\n[{index}/{len(results)}] Target: {target}') if skipped: safe_print(f' SKIPPED: {skipped}') if findings: safe_print(f' FOUND: {len(findings)} secret(s). Saved: {secrets_file}') for finding in findings[:5]: safe_print(f' - {finding_summary(finding)}') if len(findings) > 5: safe_print(f' - ... {len(findings) - 5} more') if errors: safe_print(f' ERROR: saved to {errors_file}') safe_print(f' First error: {first_error_line(result)}') if warnings: safe_print(f' DEGRADED: {len(warnings)} non-fatal diagnostic(s)') if not skipped and not findings and not errors and not warnings: safe_print(' CLEAN: no secrets found') if findings or errors or warnings: safe_print(f' Full non-clean result saved: {results_file}') def prepare_scan_options(args, target_count, *, quiet=False): if not quiet: print(f'Scanning {target_count} targets with {args.workers} workers...') print(f'TruffleHog timeout per target: {args.timeout}s') print(f'Results directory: {args.save_dir}') scan_config.drop_detectors = csv_items( getattr(args, 'drop_detectors', getattr(scan_config, 'drop_detectors', [])) ) if scan_config.drop_detectors and not quiet: print(f'Dropping detector findings before persistence: {",".join(scan_config.drop_detectors)}') scan_kwargs = { 'timeout_sec': args.timeout, 'detectors': args.detectors, 'exclude_detectors': args.exclude_detectors, 'no_verification': bool(getattr(args, 'no_verification', False)), 'trufflehog_config': getattr(args, 'trufflehog_config', ''), } if args.platform == 'docker': scan_kwargs['trufflehog_concurrency'] = int(getattr(args, 'trufflehog_concurrency', 0) or 0) scan_kwargs['docker_recovery_limits'] = docker_layer_limits(args) scan_kwargs['docker_recovery_min_free_bytes'] = int(getattr(args, 'docker_layer_min_free_bytes', 20 << 30)) if scan_kwargs['trufflehog_concurrency'] > 0 and not quiet: print(f"TruffleHog internal concurrency: {scan_kwargs['trufflehog_concurrency']}") if args.platform == 'gitlab': scan_kwargs['external_trufflehog_lifecycle'] = bool( getattr(args, 'external_trufflehog_lifecycle', False) ) if args.platform in ('github', 'github_archive', 'gitlab', 'package_git') and not getattr(args, 'scan_full_history', False): max_depth = int(getattr(args, 'max_depth', 0) or 0) if max_depth > 0: scan_kwargs['max_depth'] = max_depth if not quiet: print(f'TruffleHog git max depth: {max_depth} commits') if args.platform in ('github', 'github_archive', 'gitlab', 'package_git'): max_commit_age_days = int(getattr(args, 'max_commit_age_days', 0) or 0) if max_commit_age_days > 0: scan_kwargs['max_commit_age_days'] = max_commit_age_days scan_kwargs['commit_lookup_pages'] = int(getattr(args, 'commit_lookup_pages', 3) or 3) scan_kwargs['skip_if_commit_lookup_fails'] = bool(getattr(args, 'skip_if_commit_lookup_fails', True)) if not quiet: print(f'Git scan commit max age: {max_commit_age_days} days') if args.platform in ('npm', 'pypi'): scan_kwargs['max_artifact_size_mb'] = int(getattr(args, 'max_artifact_size_mb', 50) or 50) if args.platform == 'postman': scan_kwargs['max_artifact_size_mb'] = int(getattr(args, 'max_artifact_size_mb', 20) or 20) if args.platform == 'github_actions': scan_kwargs.update({ 'ci_runs_per_repo': int(getattr(args, 'ci_runs_per_repo', 5) or 5), 'ci_lookback_days': int(getattr(args, 'ci_lookback_days', 30) or 30), 'ci_max_log_archive_mb': int(getattr(args, 'ci_max_log_archive_mb', 50) or 50), 'ci_max_log_file_mb': int(getattr(args, 'ci_max_log_file_mb', 20) or 20), 'ci_failed_first': bool(getattr(args, 'ci_failed_first', True)), 'ci_scan_artifacts': bool(getattr(args, 'ci_scan_artifacts', False)), 'ci_max_artifacts_per_run': int(getattr(args, 'ci_max_artifacts_per_run', 3) or 3), 'ci_max_artifact_archive_mb': int(getattr(args, 'ci_max_artifact_archive_mb', 50) or 50), 'ci_max_artifact_file_mb': int(getattr(args, 'ci_max_artifact_file_mb', 10) or 10), 'ci_max_artifact_files': int(getattr(args, 'ci_max_artifact_files', 1000) or 1000), 'ci_target_max_download_mb': int(getattr(args, 'ci_target_max_download_mb', 500) or 500), 'fetch_timeout': int(getattr(args, 'fetch_timeout', 20) or 20), }) if args.platform == 'gitlab_ci': scan_kwargs.update({ 'ci_pipelines_per_project': int(getattr(args, 'ci_pipelines_per_project', 5) or 5), 'ci_jobs_per_pipeline': int(getattr(args, 'ci_jobs_per_pipeline', 20) or 20), 'ci_lookback_days': int(getattr(args, 'ci_lookback_days', 30) or 30), 'ci_max_trace_mb': int(getattr(args, 'ci_max_trace_mb', 20) or 20), 'ci_scan_artifacts': bool(getattr(args, 'ci_scan_artifacts', False)), 'ci_max_artifacts_per_pipeline': int(getattr(args, 'ci_max_artifacts_per_pipeline', 5) or 5), 'ci_max_artifact_archive_mb': int(getattr(args, 'ci_max_artifact_archive_mb', 50) or 50), 'ci_max_artifact_file_mb': int(getattr(args, 'ci_max_artifact_file_mb', 10) or 10), 'ci_max_artifact_files': int(getattr(args, 'ci_max_artifact_files', 1000) or 1000), 'ci_target_max_download_mb': int(getattr(args, 'ci_target_max_download_mb', 500) or 500), 'fetch_timeout': int(getattr(args, 'fetch_timeout', 20) or 20), }) token = get_platform_token(args) if token: scan_kwargs['token'] = token return scan_kwargs def validate_v2_capacity_model( max_active_scans, max_event_bytes, projection_max_bytes, projection_headroom_bytes, ): slots = max(1, int(max_active_scans)) event_bytes = max(1, int(max_event_bytes)) per_scan_projection = event_bytes * 2 headroom = max(per_scan_projection, int(projection_headroom_bytes)) required = slots * per_scan_projection + headroom if int(projection_max_bytes) < required: raise RuntimeError( 'projection capacity cannot cover every physical scan slot plus bounded backlog ' f'headroom: configured={int(projection_max_bytes)} required={required}' ) return { 'physical_slots': slots, 'per_scan_projection_bytes': per_scan_projection, 'headroom_bytes': headroom, 'required_projection_bytes': required, } @dataclass(frozen=True) class V2AdmissionOutcome: claim: object permit_released: bool retry_without_claim: bool = False reason: str | None = None def reserve_v2_admission_with_recovery( db_url, source, platform, producer_mapping, supervisor_instance_id, declared_bundle_bytes, projection_bytes, candidate_items, candidate_bytes, *, lease_seconds, max_attempts, capacity_limits, run_id, cycle_id, reservation_token, bundle_id, scan_event_id, resolution_attempts=8, resolution_seconds=30, retry_delay=0.2, claim_order='oldest', docker_depth_authority=None, final_cutover=False, db_factory=ScannerDB, stop_requested=supervised_spool_stop_requested, sleep=time.sleep, release_permit=lambda: None, remote_assignment=None, reserved_bundle_bytes=None, remote_max_active=50, ): deadline = time.monotonic() + max(0.1, float(resolution_seconds)) last_error = None admission = db_factory(db_url=db_url, initialize=False) try: if not admission.enabled: release_permit() return V2AdmissionOutcome(None, permit_released=True) admission.set_application_name(f'truf-admission:{source}') claim = admission.reserve_and_claim_target( source, platform, producer_mapping, supervisor_instance_id, declared_bundle_bytes, projection_bytes, candidate_items, candidate_bytes, lease_seconds=lease_seconds, max_attempts=max_attempts, capacity_limits=capacity_limits, run_id=run_id, cycle_id=cycle_id, reservation_token=reservation_token, bundle_id=bundle_id, scan_event_id=scan_event_id, claim_order=claim_order, docker_depth_authority=docker_depth_authority, final_cutover=final_cutover, remote_assignment=remote_assignment, reserved_bundle_bytes=reserved_bundle_bytes, remote_max_active=remote_max_active, ) reason = ( admission.admission_intent_resolution(reservation_token) if claim is None else None ) return V2AdmissionOutcome( claim, permit_released=False, reason=reason, ) except (KeyboardInterrupt, SystemExit): raise except Exception as exc: last_error = exc release_permit() logger.warning( 'Admission outcome is ambiguous; released physical permit and entered ' 'bounded idempotent exact-token resolution: %s', type(exc).__name__, ) finally: admission.close() expected = { 'bundle_id': str(bundle_id).lower(), 'scan_event_id': str(scan_event_id).lower(), 'source': str(source), 'platform': str(platform), 'producer_instance_id': str(supervisor_instance_id or ''), 'producer_pid': int(producer_mapping['pid']), 'producer_creation_time': str(producer_mapping['creation_time']), 'producer_executable': str(producer_mapping['executable']), 'declared_bundle_bytes': int(declared_bundle_bytes), 'reserved_bundle_bytes': int( declared_bundle_bytes if reserved_bundle_bytes is None else reserved_bundle_bytes ), 'reserved_projection_items': 1, 'reserved_projection_bytes': int(projection_bytes), 'reserved_candidate_items': int(candidate_items), 'reserved_candidate_bytes': int(candidate_bytes), 'run_id': run_id, 'cycle_id': cycle_id, 'assignment_kind': 'local', 'remote_user_id': None, 'remote_device_id': None, 'remote_effective_config_sha256': None, 'remote_client_compat_sha256': None, 'remote_execution_snapshot_json': None, 'remote_execution_snapshot_sha256': None, } if remote_assignment is not None: remote = ScannerDB._remote_assignment_mapping(remote_assignment) expected.update({ 'assignment_kind': 'remote', 'remote_user_id': remote['user_id'], 'remote_device_id': remote['device_id'], 'remote_effective_config_sha256': remote['effective_config_sha256'], 'remote_client_compat_sha256': remote['client_compat_sha256'], 'remote_execution_snapshot_json': remote['execution_snapshot_json'], 'remote_execution_snapshot_sha256': remote['execution_snapshot_sha256'], }) for attempt in range(max(1, int(resolution_attempts))): if stop_requested(): raise KeyboardInterrupt('coordinated shutdown during exact claim recovery') recovery = db_factory(db_url=db_url, initialize=False) try: if not recovery.enabled: raise RuntimeError('admission PostgreSQL connection is unavailable') recovery.set_application_name(f'truf-admission-recovery:{source}') claim = recovery.recover_result_reservation_claim( reservation_token, expected, ) if claim: logger.info( 'Recovered exact admission reservation after an ambiguous response: %s', claim['reservation_id'], ) return V2AdmissionOutcome(claim, permit_released=True) return V2AdmissionOutcome(None, permit_released=True) except (KeyboardInterrupt, SystemExit): raise except Exception as exc: last_error = exc logger.warning( 'Admission serialized exact-token resolution remains unavailable (%s/%s): %s', attempt + 1, max(1, int(resolution_attempts)), type(exc).__name__, ) finally: recovery.close() remaining = deadline - time.monotonic() if attempt + 1 >= max(1, int(resolution_attempts)) or remaining <= 0: break sleep(min(remaining, min(1.0, max(0.05, float(retry_delay))))) raise UnresolvedHandoffInfrastructureError( 'admission outcome remained ambiguous after bounded exact-token resolution' ) from last_error def reacquire_scan_permit_bounded( command, timeout_sec, wait_seconds, *, acquire=acquire_scan_slot, stop_requested=supervised_spool_stop_requested, sleep=time.sleep, ): if not scan_limiter_enabled(): return None deadline = time.monotonic() + max(0.1, float(wait_seconds)) while time.monotonic() < deadline: if stop_requested(): raise KeyboardInterrupt('coordinated shutdown while reacquiring scan permit') lease = acquire(command, timeout_sec, wait=False, start_heartbeat=False) if lease is not None: return lease sleep(min(0.2, max(0.0, deadline - time.monotonic()))) raise TimeoutError('bounded scan permit reacquisition expired') def refund_v2_claim_after_no_handoff( db, bundle_root, claim, producer_mapping, reason, ): ready_relative = str(claim['ready_relative_path']).replace('\\', '/') ready = inspect_private_relative_path(bundle_root, ready_relative) if ready.state != PrivatePathState.ABSENT: return False partial_relative = bundle_partial_relative_path( claim['bundle_id'], claim['reservation_token'], ).replace(os.sep, '/') partial = inspect_private_relative_path(bundle_root, partial_relative) if partial.state == PrivatePathState.UNKNOWN: return False if partial.state == PrivatePathState.PRESENT: if not private_file_ready(partial.path): return False try: durable_unlink(partial.path) except OSError: return False partial = inspect_private_relative_path(bundle_root, partial_relative) if partial.state != PrivatePathState.ABSENT: return False return bool(db.refund_uncommitted_reservation( claim['reservation_id'], producer_mapping, reason, partial_absence_confirmed=True, )) class _GitResolutionFailure(Exception): def __init__(self, error, invalid_target=False): super().__init__(str(error)) self.error = error self.invalid_target = bool(invalid_target) class _DockerPlanningFailure(Exception): def __init__(self, error, *, retryable, source_failure, category): super().__init__(str(error)) self.error = error self.retryable = bool(retryable) self.source_failure = bool(source_failure) self.category = str(category or 'docker_planning') def docker_layer_limits(args): return validate_docker_layer_limits({ 'config_max_bytes': int(getattr(args, 'docker_layer_config_max_bytes', 1 << 20)), 'layer_max_bytes': int(getattr(args, 'docker_layer_max_bytes', 256 << 20)), 'image_max_bytes': int(getattr(args, 'docker_layer_image_max_bytes', 1 << 30)), 'max_layers': int(getattr(args, 'docker_layer_max_layers', 8)), 'archive_max_size_bytes': int( getattr(args, 'docker_layer_archive_max_size_bytes', 256 << 20) ), 'archive_max_depth': int(getattr(args, 'docker_layer_archive_max_depth', 4)), 'archive_timeout_sec': int(getattr(args, 'docker_layer_archive_timeout_sec', 30)), 'blob_timeout_sec': int(getattr(args, 'docker_layer_blob_timeout_sec', 600)), 'filesystem_concurrency': int( getattr(args, 'docker_layer_filesystem_concurrency', 2) ), 'blob_max_attempts': int(getattr(args, 'docker_layer_blob_max_attempts', 3)), }) def docker_adaptive_checkpoint(args): return validate_docker_adaptive_checkpoint({ 'max_blobs': int(getattr(args, 'docker_adaptive_checkpoint_max_blobs', 4)), 'max_bytes': int( getattr(args, 'docker_adaptive_checkpoint_max_bytes', 512 << 20) ), }) def docker_layer_effective_mode( args, target, canary_eligible=False, *, adaptive_gate_passed=False, selection_policy_sha256=None, ): configured = str( getattr(args, 'docker_content_scan_mode', 'full') or 'full' ).strip().lower() if configured not in ('full', 'canary', 'layer', 'adaptive-canary', 'adaptive'): raise ValueError( 'Docker content scan mode must be full, canary, layer, ' 'adaptive-canary, or adaptive' ) basis_points = max(0, min( 10000, int(getattr(args, 'docker_layer_canary_basis_points', 0) or 0), )) if configured in ('full', 'layer'): return configured, basis_points if configured == 'canary': if basis_points == 0 or not canary_eligible: return 'full', basis_points image = str(parse_docker_target(target)['image']).lower() manifest_digest = image.rsplit('@', 1)[-1] effective = ( 'layer' if docker_layer_canary_selected(manifest_digest, basis_points) else 'full' ) return effective, basis_points adaptive_basis_points = max(0, min( 10000, int(getattr(args, 'docker_adaptive_canary_basis_points', 0) or 0), )) if not adaptive_gate_passed: return 'full', adaptive_basis_points if configured == 'adaptive': return 'adaptive', 10000 if adaptive_basis_points == 0: return 'full', adaptive_basis_points selection_policy_sha256 = selection_policy_sha256 or ( docker_layer_selection_policy_sha256(docker_layer_limits(args)) ) image = str(parse_docker_target(target)['image']).lower() manifest_digest = image.rsplit('@', 1)[-1] effective = ( 'adaptive' if docker_adaptive_canary_selected( manifest_digest, selection_policy_sha256, adaptive_basis_points, ) else 'full' ) return effective, adaptive_basis_points def docker_layer_timeout_canary_eligible(db_url, source, claim): eligibility_db = ScannerDB(db_url=db_url, initialize=False) try: if not eligibility_db.enabled: raise RuntimeError('Docker layer canary eligibility database is unavailable') eligibility_db.set_application_name(f'truf-docker-layer-canary:{source}') return eligibility_db.docker_layer_timeout_canary_eligible( claim['reservation_id'], claim['claim_lease_token'], ) finally: eligibility_db.close() def docker_adaptive_gate_report( db_url, source, claim, scan_policy_sha256, execution_policy_sha256, selection_policy_sha256, max_age_sec, ): gate_db = ScannerDB(db_url=db_url, initialize=False) try: if not gate_db.enabled: raise RuntimeError('Docker adaptive gate database is unavailable') gate_db.set_application_name(f'truf-docker-adaptive-gate:{source}') return gate_db.docker_adaptive_gate_report( claim['reservation_id'], claim['claim_lease_token'], scan_policy_sha256, execution_policy_sha256, selection_policy_sha256, max_age_sec=max_age_sec, ) finally: gate_db.close() def docker_layer_scan_policy_sha256(args, scan_kwargs): executable = shutil.which(get_trufflehog_cmd()) or get_trufflehog_cmd() executable_sha256 = hash_file(executable) if not executable_sha256: raise RuntimeError('Docker layer scanner executable fingerprint is unavailable') config_path = str(scan_kwargs.get('trufflehog_config') or '') config_sha256 = hash_file(config_path) if config_path else None if config_path and not config_sha256: raise RuntimeError('Docker layer detector policy fingerprint is unavailable') candidate_normalizer = os.path.join( os.path.dirname(os.path.abspath(__file__)), 'keycheck_candidates.py', ) candidate_normalizer_sha256 = hash_file(candidate_normalizer) if not candidate_normalizer_sha256: raise RuntimeError('Docker routed identity policy fingerprint is unavailable') recovery_limits = validate_docker_layer_limits( scan_kwargs.get('docker_recovery_limits') or docker_layer_limits(args), ) payload = { 'version': 'docker-content-scan-v3', 'trufflehog_sha256': executable_sha256, 'docker_recovery_execution': {key: recovery_limits[key] for key in ( 'archive_max_size_bytes', 'archive_max_depth', 'archive_timeout_sec', 'filesystem_concurrency', )}, 'detectors': sorted(csv_items(scan_kwargs.get('detectors'))), 'exclude_detectors': sorted(csv_items(scan_kwargs.get('exclude_detectors'))), 'drop_detectors': sorted(csv_items(getattr(args, 'drop_detectors', []))), 'no_verification': bool(scan_kwargs.get('no_verification', False)), 'detector_config_sha256': config_sha256, 'strict_git_provider_token_filter': bool( getattr(scan_config, 'strict_git_provider_token_filter', True) ), 'docker_concurrency': max( 0, min(64, int(scan_kwargs.get('trufflehog_concurrency', 0) or 0)), ), 'candidate_normalizer_sha256': candidate_normalizer_sha256, 'finding_identity_version': 'scanner-db-finding-identity-v1', } encoded = json.dumps( payload, ensure_ascii=True, sort_keys=True, separators=(',', ':'), ).encode('ascii') return hashlib.sha256(encoded).hexdigest() def docker_planning_failure_result( args, claim, scan_kwargs, failure, configured_mode, effective_mode, canary_eligible=False, adaptive_gate=None, ): adaptive_gate = adaptive_gate if isinstance(adaptive_gate, dict) else {} now = datetime.now().isoformat() message = redact_secrets(str(failure.error), [ str(scan_kwargs.get('token') or ''), str(getattr(args, 'token', '') or ''), str(getattr(args, 'docker_token', '') or ''), ]) result = { 'findings': [], 'errors': [f'Docker layer planning failed: {message[:500]}'], 'error_class': failure.category, 'retryable': failure.retryable, 'source_failure': failure.source_failure, 'source_failure_category': failure.category if failure.source_failure else '', 'source_failure_auth_related': failure.category == 'docker_auth', 'scan_meta': { 'docker_scan_assignment': { 'configured_mode': configured_mode, 'effective_mode': effective_mode, 'canary_basis_points': max(0, min( 10000, int(getattr( args, 'docker_adaptive_canary_basis_points' if configured_mode in ('adaptive-canary', 'adaptive') else 'docker_layer_canary_basis_points', 0, ) or 0), )), 'canary_eligible': bool(canary_eligible), 'adaptive_gate_passed': bool(adaptive_gate.get('passed')), 'adaptive_gate_reason': str(adaptive_gate.get('reason') or 'not_configured'), 'adaptive_gate_report_id': adaptive_gate.get('report_id'), 'status': 'planning_failed', }, }, 'target': str(claim['target']), 'scan_type': 'docker', 'scan_event_id': str(claim['scan_event_id']), 'scan_started_at': now, 'duration_sec': 0.0, 'timestamp': now, } diagnostic_http = getattr(failure.error, 'diagnostic_http', None) if isinstance(diagnostic_http, dict): result['_diagnostic_http'] = dict(diagnostic_http) assign_finding_uids(result) return result def resolve_and_bind_docker_claim( args, db_url, source, claim, scan_kwargs, scan_policy_sha256, *, deadline=None, adaptive=False, ): if deadline is None: deadline = time.monotonic() + max( 0.001, float(getattr(args, 'timeout', 600) or 0.001), ) try: resolved, bearer_auth = resolve_docker_content_manifest( claim['target'], getattr(args, 'docker_platform_os', 'linux'), getattr(args, 'docker_platform_arch', 'amd64'), deadline=deadline, ) except DockerRemoteAccessError as exc: retryable, source_failure, category = { 'rate_limited': (True, True, 'docker_rate_limit'), 'auth_failed': (True, True, 'docker_auth'), 'target_forbidden': (False, False, 'docker_target_forbidden'), 'remote_transient': (True, True, 'remote_transient'), }.get(exc.status, (True, True, 'remote_transient')) raise _DockerPlanningFailure( exc, retryable=retryable, source_failure=source_failure, category=category, ) from exc except ApiRequestError as exc: raise _DockerPlanningFailure( exc, retryable=True, source_failure=True, category='remote_transient', ) from exc except (DockerRegistryResolutionError, ValueError) as exc: raise _DockerPlanningFailure( exc, retryable=False, source_failure=False, category='invalid_target', ) from exc payload_classes = None if adaptive: try: payload_classes, bearer_auth = fetch_docker_config_payload_classes( resolved, bearer_auth, deadline=deadline, min_free_bytes=max( 0, int(getattr(args, 'docker_layer_min_free_bytes', 20 << 30) or 0), ), ) except DockerLayerInfrastructureError as exc: raise _DockerPlanningFailure( exc, retryable=True, source_failure=bool(exc.source_failure), category=str(exc.category or 'remote_transient'), ) from exc except DockerContentTransferError as exc: raise _DockerPlanningFailure( exc, retryable=bool(exc.retryable), source_failure=False, category='docker_config_transfer', ) from exc planning_db = ScannerDB(db_url=db_url, initialize=False) try: if not planning_db.enabled: raise RuntimeError('Docker layer planning database is unavailable') planning_db.set_application_name(f'truf-docker-layer-plan:{source}') plan = planning_db.bind_docker_layer_plan( claim['reservation_id'], claim['claim_lease_token'], resolved, docker_layer_limits(args), scan_policy_sha256, blob_lease_seconds=max( int(getattr(args, 'docker_layer_blob_lease_sec', 1800) or 1800), int(getattr(args, 'timeout', 600) or 600) + 300, ), payload_classes=payload_classes, checkpoint=docker_adaptive_checkpoint(args) if adaptive else None, ) finally: planning_db.close() plan_bytes = canonical_docker_layer_plan_bytes(plan) return { 'plan': plan, 'plan_sha256': hashlib.sha256(plan_bytes).hexdigest(), 'bearer_auth': bearer_auth, 'deadline': deadline, 'min_free_bytes': max( 0, int(getattr(args, 'docker_layer_min_free_bytes', 20 << 30) or 0), ), } def git_resolution_failure_result(args, claim, scan_kwargs, failure): error = failure.error token = str(scan_kwargs.get('token') or '') message = redact_secrets(str(error), [token]) category = str(getattr(error, 'category', '') or '') auth_related = bool(getattr(error, 'auth_related', False)) if failure.invalid_target: error_class = 'invalid_target' category = 'invalid_target' retryable = False source_failure = False elif isinstance(error, RateLimitError): retryable = bool(getattr(error, 'retryable', True)) source_failure = category != 'not_found' error_class = 'source_auth' if auth_related else category or 'remote_transient' else: error_class = 'remote_transient' category = category or 'remote_transient' retryable = True source_failure = True auth_related = False now = datetime.now().isoformat() result = { 'findings': [], 'errors': [f'Exact Git ref resolution failed: {message}'], 'error_class': error_class, 'retryable': retryable, 'source_failure': source_failure, 'source_failure_category': category if source_failure else '', 'source_failure_auth_related': auth_related if source_failure else False, 'scan_meta': { 'exact_git_resolution': { 'status': 'failed', 'provider': str(args.platform), 'category': category, 'retryable': retryable, 'auth_related': auth_related, }, }, 'target': str(claim['target']), 'scan_type': str(args.platform), 'scan_event_id': str(claim['scan_event_id']), 'scan_started_at': now, 'duration_sec': 0.0, 'timestamp': now, } diagnostic_http = getattr(error, 'diagnostic_http', None) if isinstance(diagnostic_http, dict): result['_diagnostic_http'] = dict(diagnostic_http) assign_finding_uids(result) return result def resolve_and_bind_git_claim( args, db_url, source, claim, scan_kwargs, *, remote_credential=None, ): try: normalize_git_scan_resolution_target(claim['target'], args.platform) except ValueError as exc: raise _GitResolutionFailure(exc, invalid_target=True) from exc try: resolved = resolve_git_scan_target( claim['target'], args.platform, scan_kwargs.get('token'), request_attempts=max( 1, int(getattr(args, 'git_ref_resolution_attempts', 2) or 2), ), timeout_sec=max( 0.1, float(getattr(args, 'git_ref_resolution_timeout_sec', 10) or 10), ), max_response_bytes=max( 1024, int(getattr(args, 'git_ref_resolution_max_bytes', 1 << 20) or (1 << 20)), ), ) except (RateLimitError, ApiRequestError, ValueError) as exc: raise _GitResolutionFailure(exc) from exc planning_db = ScannerDB(db_url=db_url, initialize=False) try: if not planning_db.enabled: raise RuntimeError('exact Git planning database is unavailable') planning_db.set_application_name(f'truf-git-plan:{source}') binding_options = {} if remote_credential is not None: binding_options['remote_credential'] = remote_credential return planning_db.bind_git_scan_plan( claim['reservation_id'], claim['claim_lease_token'], resolved, max(1, int(getattr(args, 'git_baseline_depth', 100) or 100)), **binding_options, ) finally: planning_db.close() def run_discovery_cycle( args, db, run_id, cycle_id, source_name=None, partial_metrics=None, ): source = source_name or args.platform expected_platforms = { 'gitlab': 'gitlab', 'dockerhub': 'docker', 'huggingface': 'huggingface', } if source not in expected_platforms or args.platform != expected_platforms[source]: raise ValueError('Discovery-only cycles support GitLab, DockerHub, and HuggingFace') if not db or not getattr(db, 'conn', None) or not getattr(db.conn, 'is_postgres', False): raise RuntimeError('Discovery-only cycles require PostgreSQL') if not run_id or not cycle_id: raise RuntimeError('Discovery-only cycles require valid run_id and cycle_id') db.require_runtime_safety_schema() db.require_final_cutover() partial_metrics = partial_metrics if isinstance(partial_metrics, dict) else {} status = 'completed' message = None discovery_info = { 'fetched_count': 0, 'queued_new_count': 0, 'queued_updated_count': 0, } control = db.runtime_control_state() if control['effective_discovery_paused']: status = 'paused' message = 'Discovery paused by runtime control.' else: try: _require_discovery_provider_enabled(db) if managed_dockerhub_discovery(args, db, source): collection_only = bool( getattr(args, 'docker_depth_collection_only', False) ) experiment_stopped = False if not collection_only: _, experiment_stopped = resolve_due_docker_experiment_targets( db, source, args, ) incremental = run_dockerhub_incremental_discovery( args, db, source, partial_metrics=discovery_info, ) discovery_info.update(incremental) status = str(discovery_info.get('cycle_status') or 'completed') if not collection_only and not experiment_stopped: resolve_due_docker_queue_targets(db, source, args) else: fetched_targets = fetch_targets(args, db, run_id, cycle_id, source) _, discovery_info = enqueue_discovered_targets( args, fetched_targets or [], db, run_id, cycle_id, source, partial_metrics=partial_metrics, ) except DiscoveryPausedError: status = 'paused' message = 'Discovery paused before provider work or target admission.' metrics = { **summarize_results([]), **discovery_info, 'scan_requested_count': 0, 'staged_count': 0, 'backlog_only': False, 'docker_depth_collection_only': bool( args.platform == 'docker' and getattr(args, 'docker_depth_collection_only', False) ), 'cycle_status': status, } partial_metrics.update(metrics) db.finish_source_cycle( cycle_id, status, metrics, { 'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None, }, message, ) return metrics def run_cycle_v2(args, db, run_id, cycle_id, source_name, partial_metrics=None): source = source_name or args.platform partial_metrics = partial_metrics if isinstance(partial_metrics, dict) else {} docker_depth_collection_only = bool( args.platform == 'docker' and getattr(args, 'docker_depth_collection_only', False) ) db.require_runtime_safety_schema() db.require_final_cutover() bundle_root = require_private_directory( getattr(args, 'result_bundle_dir', None) or scan_config.result_bundle_dir, create=False, ) minimum_free = max(0, int(getattr(args, 'result_bundle_min_free_bytes', 20 * 1024 * 1024 * 1024))) max_event_bytes = max(1, int(getattr(args, 'result_bundle_max_event_bytes', 64 * 1024 * 1024))) max_targets = max(0, int(getattr(args, 'max_targets', 0) or 0)) configured_slots = max(1, int(getattr(scan_config, 'max_active_scans', 3) or 3)) capacity_slots = configured_slots + max(0, min(1, int(getattr( scan_config, 'opportunistic_scan_slots', 0, ) or 0))) dispatch_limit = min(max(1, int(getattr(args, 'workers', 1) or 1)), configured_slots) if max_targets: dispatch_limit = min(dispatch_limit, max_targets) validate_v2_capacity_model( capacity_slots, max_event_bytes, int(getattr(args, 'projection_backlog_max_bytes', 2 * 1024 * 1024 * 1024)), int(getattr(args, 'projection_backlog_headroom_bytes', max_event_bytes * 2)), ) producer_identity = current_process_identity() producer_mapping = producer_identity.as_dict() supervisor_instance_id = str(os.getenv('TRUF_SUPERVISOR_INSTANCE_ID') or '') if not supervisor_instance_id: raise RuntimeError('v2 source admission requires an authenticated supervisor instance ID') fetched_targets = [] if args.platform == 'docker' and not docker_depth_collection_only: resolve_due_docker_experiment_targets(db, source, args) has_backlog = db.has_claimable_targets_v2( source, args.platform, max_attempts=int(getattr(args, 'target_retry_max_attempts', 3) or 3), ) if db.last_error: raise RuntimeError(f'Unable to inspect target backlog for {source}: {db.last_error}') todo_file, checked_file = queue_files_for_args(args) discovery_info = { 'fetched_count': 0, 'queued_new_count': 0, 'queued_updated_count': 0, 'scan_requested_count': 0, 'projection_error': '', 'cycle_status': 'completed', 'deep_dispatch_durable': False, } refresh_backlog = bool(getattr(args, 'refresh_registry', False)) backlog_only = bool(has_backlog and not refresh_backlog) if not has_backlog or refresh_backlog: incremental_info = None if managed_dockerhub_discovery(args, db, source): incremental_info = run_dockerhub_incremental_discovery(args, db, source) else: fetched_targets = fetch_targets(args, db, run_id, cycle_id, source) _, todo_file, checked_file, discovery_info = prepare_targets( args, fetched_targets or [], db, run_id, cycle_id, source, enqueue_only=True, partial_metrics=partial_metrics, experiment_resolver_already_run=True, ) if incremental_info is not None: discovery_info.update(incremental_info) partial_metrics.update({ key: value for key, value in discovery_info.items() if key in { 'fetched_count', 'queued_new_count', 'queued_updated_count', 'discovery_pages_fetched', 'discovery_retry_enqueued_count', 'discovery_retry_inserted_count', 'discovery_retry_coalesced_count', 'discovery_preexisting_count', 'discovery_known_page_count', 'discovery_stopped_on_preexisting', 'discovery_deep', 'discovery_pass_kind', 'deep_dispatch_durable', 'cycle_status', } }) scan_kwargs = prepare_scan_options(args, dispatch_limit) event_scan_options = {key: value for key, value in scan_kwargs.items() if key != 'token'} docker_configured_mode = str( getattr(args, 'docker_content_scan_mode', 'full') or 'full' ).strip().lower() docker_canary_basis_points = max(0, min( 10000, int(getattr(args, 'docker_layer_canary_basis_points', 0) or 0), )) docker_scan_policy_sha256 = None docker_execution_policy_sha256 = None docker_selection_policy_sha256 = None docker_limits_value = None docker_gate_max_age_sec = 604800 if args.platform == 'docker': if docker_configured_mode not in ( 'full', 'canary', 'layer', 'adaptive-canary', 'adaptive', ): raise ValueError( 'Docker content scan mode must be full, canary, layer, ' 'adaptive-canary, or adaptive' ) docker_scan_policy_sha256 = docker_layer_scan_policy_sha256(args, scan_kwargs) event_scan_options['scan_policy_sha256'] = docker_scan_policy_sha256 if docker_configured_mode == 'layer' or ( docker_configured_mode == 'canary' and docker_canary_basis_points > 0 ) or docker_configured_mode in ('adaptive-canary', 'adaptive'): docker_limits_value = docker_layer_limits(args) if docker_configured_mode in ('adaptive-canary', 'adaptive'): docker_execution_policy_sha256 = docker_layer_execution_policy_sha256( docker_scan_policy_sha256, docker_limits_value, ) docker_selection_policy_sha256 = docker_layer_selection_policy_sha256( docker_limits_value, ) docker_adaptive_checkpoint(args) docker_gate_max_age_sec = int(getattr( args, 'docker_adaptive_gate_max_age_sec', 604800, )) if not 60 <= docker_gate_max_age_sec <= 2592000: raise ValueError( 'Docker adaptive gate freshness must be between 60 and 2592000 seconds' ) capacity_limits = { 'bundle_items': int(getattr(args, 'result_bundle_max_items', 10000)), 'bundle_bytes': int(getattr(args, 'result_bundle_max_total_bytes', 3 * 1024 * 1024 * 1024)), 'projection_items': int(getattr(args, 'projection_backlog_max_items', 10000)), 'projection_bytes': int(getattr(args, 'projection_backlog_max_bytes', 2 * 1024 * 1024 * 1024)), 'keycheck_items': int(getattr(args, 'keycheck_queue_max_items', 100000)), 'keycheck_bytes': int(getattr(args, 'keycheck_queue_max_bytes', 512 * 1024 * 1024)), 'quarantine_items': int(getattr(args, 'pipeline_quarantine_max_items', 10000)), 'quarantine_bytes': int(getattr(args, 'pipeline_quarantine_max_bytes', 1024 * 1024 * 1024)), } candidate_items = int(getattr(args, 'keycheck_candidates_per_event', 2000)) candidate_bytes = int(getattr(args, 'keycheck_candidate_bytes_per_event', 2 * 1024 * 1024)) queue_policy = QueueDispositionPolicy( target_retry_max_attempts=int(getattr(args, 'target_retry_max_attempts', 3) or 3), target_retry_base_delay_sec=int(getattr(args, 'target_retry_base_delay_sec', 3600) or 3600), target_retry_max_delay_sec=int(getattr(args, 'target_retry_max_delay_sec', 86400) or 86400), target_timeout_retry_delay_sec=int(getattr(args, 'target_timeout_retry_delay_sec', 21600) or 21600), docker_layer_checkpoint_delay_sec=int(getattr(args, 'docker_layer_checkpoint_delay_sec', 60) or 60), ci_soft_cooldown_days=int(getattr(args, 'ci_soft_cooldown_days', 7) or 7), soft_skip_reasons=tuple(sorted(CI_SOFT_SKIP_REASONS.get(args.platform, set()))), ) lease_seconds = max(1800, int(getattr(args, 'timeout', 0) or 0) + 900) require_s_drive = bool(os.name == 'nt' and os.getenv('SCANNER_SUPERVISED') == '1') def reserve_with_exact_recovery(permit_box): reservation_token = secrets.token_urlsafe(32) bundle_id = secrets.token_hex(16) scan_event_id = secrets.token_hex(16) def release_permit(): lease = permit_box.get('lease') if lease is not None: lease.release() permit_box['lease'] = None return reserve_v2_admission_with_recovery( db.url, source, args.platform, producer_mapping, supervisor_instance_id, max_event_bytes, max_event_bytes * 2, candidate_items, candidate_bytes, lease_seconds=lease_seconds, max_attempts=int(getattr(args, 'target_retry_max_attempts', 3) or 3), capacity_limits=capacity_limits, run_id=run_id, cycle_id=cycle_id, reservation_token=reservation_token, bundle_id=bundle_id, scan_event_id=scan_event_id, resolution_attempts=int(getattr(args, 'admission_resolution_attempts', 8) or 8), resolution_seconds=float(getattr(args, 'admission_resolution_seconds', 30) or 30), retry_delay=float( getattr(args, 'admission_resolution_retry_delay_sec', 0.2) or 0.2 ), claim_order=str(getattr(args, 'target_claim_order', 'oldest') or 'oldest'), docker_depth_authority=getattr( args, 'docker_depth_experiment_authority', None, ), final_cutover=True, release_permit=release_permit, ) def notify_ready(claim, staged): notification = ScannerDB(db_url=db.url, initialize=False) try: if not notification.enabled: return False notification.set_application_name(f'truf-bundle-notify:{source}') return notification.mark_result_bundle_ready( claim['reservation_id'], { 'bundle_id': staged.bundle_id, 'scan_event_id': staged.scan_event_id, 'scan_event_hash': staged.scan_event_hash, 'relative_path': staged.relative_path, 'actual_bytes': staged.actual_bytes, 'frame_count': staged.frame_count, 'finding_count': staged.finding_count, 'error_count': staged.error_count, 'candidate_count': staged.candidate_count, }, ) except Exception as exc: logger.warning( 'Ready bundle DB notification failed after durable handoff event=%s: %s', staged.scan_event_id, exc, ) return False finally: notification.close() def stage_claim(claim, lease): claim_deadline = time.monotonic() + max( 0.001, float(getattr(args, 'timeout', 600) or 0.001), ) with scan_slot_scope( ['scan-target', args.platform], getattr(args, 'timeout', None), lease=lease, ): claim_scan_kwargs = scan_kwargs if ( args.platform in ('github', 'gitlab') and bool(getattr(args, 'exact_git_planning_enabled', False)) ): try: git_plan = resolve_and_bind_git_claim( args, db.url, source, claim, scan_kwargs, ) except _GitResolutionFailure as failure: result = git_resolution_failure_result( args, claim, scan_kwargs, failure, ) else: claim_scan_kwargs = dict(scan_kwargs) claim_scan_kwargs['git_plan'] = git_plan result = scan_target_result( claim['target'], args.platform, claim['scan_event_id'], claim_scan_kwargs, ) elif args.platform == 'docker': canary_eligible = False adaptive_gate = {'passed': False, 'reason': 'not_configured'} try: if ( docker_configured_mode == 'canary' and docker_canary_basis_points > 0 ): canary_eligible = docker_layer_timeout_canary_eligible( db.url, source, claim, ) if docker_configured_mode in ('adaptive-canary', 'adaptive'): adaptive_gate = docker_adaptive_gate_report( db.url, source, claim, docker_scan_policy_sha256, docker_execution_policy_sha256, docker_selection_policy_sha256, docker_gate_max_age_sec, ) effective_mode, basis_points = docker_layer_effective_mode( args, claim['target'], canary_eligible, adaptive_gate_passed=bool(adaptive_gate.get('passed')), selection_policy_sha256=docker_selection_policy_sha256, ) except ValueError as exc: failure = _DockerPlanningFailure( exc, retryable=False, source_failure=False, category='invalid_target', ) result = docker_planning_failure_result( args, claim, scan_kwargs, failure, docker_configured_mode, 'adaptive' if docker_configured_mode in ('adaptive-canary', 'adaptive') else 'layer', canary_eligible, adaptive_gate, ) effective_mode = ( 'adaptive' if docker_configured_mode in ('adaptive-canary', 'adaptive') else 'layer' ) basis_points = max(0, min(10000, int(getattr( args, 'docker_adaptive_canary_basis_points' if docker_configured_mode in ('adaptive-canary', 'adaptive') else 'docker_layer_canary_basis_points', 0, ) or 0))) else: if effective_mode in ('layer', 'adaptive'): try: layer_work = resolve_and_bind_docker_claim( args, db.url, source, claim, scan_kwargs, docker_scan_policy_sha256, deadline=claim_deadline, adaptive=effective_mode == 'adaptive', ) except _DockerPlanningFailure as failure: result = docker_planning_failure_result( args, claim, scan_kwargs, failure, docker_configured_mode, effective_mode, canary_eligible, adaptive_gate, ) else: claim_scan_kwargs = dict(scan_kwargs) claim_scan_kwargs['docker_layer_work'] = layer_work result = scan_target_result( claim['target'], args.platform, claim['scan_event_id'], claim_scan_kwargs, ) else: result = scan_target_result( claim['target'], args.platform, claim['scan_event_id'], claim_scan_kwargs, ) scan_meta = result.get('scan_meta') if not isinstance(scan_meta, dict): scan_meta = {} result['scan_meta'] = scan_meta scan_meta.setdefault('docker_scan_assignment', { 'configured_mode': docker_configured_mode, 'effective_mode': effective_mode, 'canary_basis_points': basis_points, 'canary_eligible': bool(canary_eligible), 'adaptive_gate_passed': bool(adaptive_gate.get('passed')), 'adaptive_gate_reason': str( adaptive_gate.get('reason') or 'not_configured' ), 'adaptive_gate_report_id': adaptive_gate.get('report_id'), 'selection_policy_sha256': docker_selection_policy_sha256, 'status': 'executed', }) else: result = scan_target_result( claim['target'], args.platform, claim['scan_event_id'], claim_scan_kwargs, ) staged = stage_scan_result_in_scope( result, claim, bundle_root, event_scan_options, queue_policy, attempts=claim.get('attempts'), candidate_max_items=candidate_items, candidate_max_bytes=candidate_bytes, require_s_drive=require_s_drive, ) # Notification is deliberately outside the permit; exact-path recovery handles failure. notify_ready(claim, staged) return staged refunded_claim = object() def recover_worker_failure(claim, exc): ready_relative = str(claim['ready_relative_path']).replace('\\', '/') inspection = inspect_private_relative_path(bundle_root, ready_relative) try: if inspection.state == PrivatePathState.PRESENT: reader = ResultBundleReader( inspection.path, max_event_bytes=max_event_bytes, ) validated = reader.validate() metadata = reader.metadata() staged = StagedResult( target=str(claim['target']), scan_event_id=validated.scan_event_id, bundle_id=validated.bundle_id, reservation_id=validated.reservation_id, scan_event_hash=validated.scan_event_hash, actual_bytes=validated.actual_bytes, relative_path=str(claim['ready_relative_path']).replace('\\', '/'), frame_count=validated.frame_count, finding_count=validated.finding_count, error_count=validated.error_count, candidate_count=validated.candidate_count, queue_status=str(metadata.get('queue_status') or ''), source_failure=bool(metadata.get('source_failure')), source_failure_category=str(metadata.get('source_failure_category') or ''), source_failure_auth_related=bool(metadata.get('source_failure_auth_related')), first_error=str(metadata.get('first_error_summary') or ''), ) notify_ready(claim, staged) return staged if inspection.state == PrivatePathState.UNKNOWN: logger.error( 'Bundle handoff state is unknown; reservation retained id=%s error=%s', claim['reservation_id'], inspection.detail, ) return None refunded = refund_v2_claim_after_no_handoff( db, bundle_root, claim, producer_mapping, f'pre-handoff source failure: {type(exc).__name__}: {exc}', ) if not refunded: logger.error( 'Exact pre-handoff reservation refund was not confirmed; recovery retains it: %s', claim['reservation_id'], ) return None return refunded_claim except (OSError, ValueError) as inspection_error: logger.error( 'Bundle handoff state is unknown; reservation retained for ingester recovery id=%s error=%s', claim['reservation_id'], inspection_error, ) return None staged_count = 0 source_failure_count = 0 first_source_failure = None pending = {} claimed_count = 0 admission_open = not docker_depth_collection_only unresolved_handoff_error = None with concurrent.futures.ThreadPoolExecutor(max_workers=dispatch_limit) as executor: while pending or admission_open: while admission_open and len(pending) < dispatch_limit and (not max_targets or claimed_count < max_targets): if supervised_spool_stop_requested(): admission_open = False break free_bytes = shutil.disk_usage(bundle_root).free if minimum_free and free_bytes < minimum_free: admission_open = False break lease = acquire_scan_slot( ['dispatch-target', args.platform], getattr(args, 'timeout', None), wait=True, start_heartbeat=False, ) permit_box = {'lease': lease} try: outcome = reserve_with_exact_recovery(permit_box) except BaseException: if permit_box.get('lease') is not None: permit_box['lease'].release() raise claim = outcome.claim if not claim: if permit_box.get('lease') is not None: permit_box['lease'].release() if outcome.retry_without_claim: continue admission_open = False break try: ensure_bundle_reservation_paths(bundle_root, claim) except BaseException: if permit_box.get('lease') is not None: permit_box['lease'].release() raise if outcome.permit_released: try: permit_box['lease'] = reacquire_scan_permit_bounded( ['dispatch-target', args.platform], getattr(args, 'timeout', None), float(getattr(args, 'admission_permit_reacquire_seconds', 30) or 30), ) except BaseException as exc: if recover_worker_failure(claim, exc) is not refunded_claim: raise UnresolvedHandoffInfrastructureError( 'confirmed admission claim could not reacquire a permit or refund exactly' ) from exc continue lease = permit_box.get('lease') if lease is None and scan_limiter_enabled(): if recover_worker_failure( claim, RuntimeError('confirmed claim has no physical scan permit'), ) is not refunded_claim: raise UnresolvedHandoffInfrastructureError( 'confirmed admission claim lost its physical scan permit' ) continue claimed_count += 1 future = executor.submit(stage_claim, claim, lease) pending[future] = claim if not pending: break done, _ = concurrent.futures.wait( tuple(pending), timeout=0.2, return_when=concurrent.futures.FIRST_COMPLETED, ) for future in done: claim = pending.pop(future) try: staged = future.result() except Exception as exc: staged = recover_worker_failure(claim, exc) refunded = staged is refunded_claim source_outage_refunded = bool( refunded and isinstance(exc, DockerLayerInfrastructureError) ) if refunded: staged = None if source_outage_refunded: source_failure_count += 1 if first_source_failure is None: first_source_failure = SimpleNamespace( source_failure_category=exc.category, source_failure_auth_related=exc.auth_related, first_error=str(exc)[:300], ) admission_open = False if staged is None: if not source_outage_refunded: logger.error('Target staging failed before durable handoff: %s', exc) unresolved_handoff_error = UnresolvedHandoffInfrastructureError( 'deterministic handoff cleanup remained unknown; source admission is closed' ) admission_open = False if staged is not None: staged_count += 1 if staged.source_failure: source_failure_count += 1 if first_source_failure is None: first_source_failure = staged admission_open = False safe_print( f'Staged {staged_count} event(s): {staged.target} ' f'findings={staged.finding_count} errors={staged.error_count}', flush=True, ) if unresolved_handoff_error is not None: raise unresolved_handoff_error cycle_status = str(discovery_info.get('cycle_status') or 'completed') if backlog_only: cycle_status = 'backlog_only' if source_failure_count: cycle_status = 'source_failed' metrics = { 'authoritative_async': True, 'fetched_count': discovery_info.get('fetched_count', len(fetched_targets)), 'queued_new_count': discovery_info.get('queued_new_count', 0), 'queued_updated_count': discovery_info.get('queued_updated_count', 0), 'scan_requested_count': claimed_count, 'staged_count': staged_count, 'scanned_count': 0, 'clean_count': 0, 'found_count': 0, 'skipped_count': 0, 'error_count': 0, 'findings_count': 0, 'verified_findings_count': 0, 'unique_secrets_count': 0, 'unique_findings_count': 0, 'backlog_only': bool(backlog_only), 'docker_depth_collection_only': docker_depth_collection_only, 'source_failure_count': source_failure_count, 'source_failure_category': first_source_failure.source_failure_category if first_source_failure else '', 'source_failure_auth_related': first_source_failure.source_failure_auth_related if first_source_failure else False, 'source_failure_message': first_source_failure.first_error if first_source_failure else '', 'source_failure_queue_ids': [], 'cycle_status': cycle_status, 'deep_dispatch_durable': bool(discovery_info.get('deep_dispatch_durable', False)), } for key in ( 'discovery_pages_fetched', 'discovery_retry_enqueued_count', 'discovery_retry_inserted_count', 'discovery_retry_coalesced_count', 'discovery_preexisting_count', 'discovery_known_page_count', 'discovery_stopped_on_preexisting', 'discovery_deep', 'discovery_pass_kind', ): if key in discovery_info: metrics[key] = discovery_info[key] partial_metrics.update(metrics) if cycle_id: db.finish_source_cycle( cycle_id, cycle_status, metrics, {'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None}, metrics.get('source_failure_message') or None, ) safe_print( f'Cycle staged {staged_count} durable bundle(s); PostgreSQL ingestion is asynchronous.', flush=True, ) return metrics def _run_cycle_legacy_compat( args, db=None, run_id=None, cycle_id=None, source_name=None, partial_metrics=None, ): source_for_db = source_name or args.platform partial_metrics = partial_metrics if isinstance(partial_metrics, dict) else {} postgres_db = bool(db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False)) spool = None ingest_outcomes = {} if postgres_db: db.require_runtime_safety_schema() spool = result_spool_for_args(args) wait_for_result_spool_ready( spool, db, outcomes=ingest_outcomes, stop_event=getattr(args, 'result_spool_stop_event', None), wait_seconds=float(getattr(args, 'result_spool_wait_sec', 1.0) or 1.0), diagnostic_interval=float(getattr(args, 'result_spool_diagnostic_interval_sec', 30.0) or 30.0), ) drained = drain_scan_publication_outbox(db, 100) if drained: print(f'Published {drained} pending scanner outbox event(s).') configured_slots = max(1, int(getattr(scan_config, 'max_active_scans', 1) or 1)) dispatch_limit = min(max(1, int(getattr(args, 'workers', 1) or 1)), configured_slots) if args.max_targets: dispatch_limit = min(dispatch_limit, max(1, int(args.max_targets))) require_scan_publication_capacity(db, args, additional_items=dispatch_limit) else: dispatch_limit = 0 slot_first_dispatch = bool(postgres_db and scan_limiter_enabled()) def has_backlog(): present = db.has_claimable_targets( source_for_db, args.platform, max_attempts=int(getattr(args, 'target_retry_max_attempts', 3) or 3), ) if db.last_error: raise RuntimeError(f'Unable to inspect target backlog for {source_for_db}: {db.last_error}') return present def claim_with_dispatch_slots(): leases = acquire_scan_slot_leases( ['dispatch-target', args.platform], dispatch_limit, getattr(args, 'timeout', None), ) if not leases: raise RuntimeError('slot-first dispatch did not acquire a configured scan slot') try: return prepare_targets( args, [], db, run_id, cycle_id, source_for_db, spool=spool, claim_limit_override=len(leases), dispatch_leases=leases, ) except BaseException: for lease in leases: if lease.heartbeat_thread is None: lease.release() raise fetched_targets = [] backlog_only = False prepared = None initial_backlog = has_backlog() if slot_first_dispatch else False if initial_backlog and args.platform == 'docker': resolve_due_docker_experiment_targets(db, source_for_db, args) if slot_first_dispatch and initial_backlog: prepared = claim_with_dispatch_slots() backlog_only = bool(prepared[0]) if not backlog_only: fetched_targets = fetch_targets(args, db, run_id, cycle_id, source_for_db) if not fetched_targets: print('No targets fetched. This may be an empty result or an API limit/error.') if slot_first_dispatch: _, todo_file, checked_file, discovery_info = prepare_targets( args, fetched_targets or [], db, run_id, cycle_id, source_for_db, spool=spool, enqueue_only=True, partial_metrics=partial_metrics, ) partial_metrics.update({ 'fetched_count': discovery_info.get('fetched_count', len(fetched_targets)), 'queued_new_count': discovery_info.get('queued_new_count', 0), 'queued_updated_count': discovery_info.get('queued_updated_count', 0), }) if has_backlog(): prepared = claim_with_dispatch_slots() prepared[3]['fetched_count'] = len(fetched_targets or []) prepared[3]['queued_new_count'] = discovery_info.get('queued_new_count', 0) prepared[3]['queued_updated_count'] = discovery_info.get('queued_updated_count', 0) prepared[3]['projection_error'] = discovery_info.get('projection_error', '') else: prepared = ([], todo_file, checked_file, discovery_info) else: prepared = prepare_targets( args, fetched_targets or [], db, run_id, cycle_id, source_for_db, spool=spool, partial_metrics=partial_metrics, ) targets_to_scan, todo_file, checked_file, queue_info = prepared partial_metrics.update({ 'fetched_count': queue_info.get('fetched_count', len(fetched_targets)), 'queued_new_count': queue_info.get('queued_new_count', 0), 'queued_updated_count': queue_info.get('queued_updated_count', 0), }) if args.max_targets and len(targets_to_scan) > args.max_targets: targets_to_scan = targets_to_scan[:args.max_targets] queue_info['scan_requested_count'] = len(targets_to_scan) if not targets_to_scan: print('No new targets to scan.') metrics = { 'fetched_count': queue_info.get('fetched_count', len(fetched_targets)), 'queued_new_count': queue_info.get('queued_new_count', 0), 'queued_updated_count': queue_info.get('queued_updated_count', 0), 'scan_requested_count': 0, 'backlog_only': backlog_only, **summarize_results([]), } partial_metrics.update(metrics) if db and cycle_id: counts = queue_counts(todo_file, checked_file) if ( not postgres_db or bool(getattr(args, 'sync_file_queues', True)) ) else { 'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None, } db.finish_source_cycle(cycle_id, 'completed', metrics, counts, 'No new targets to scan') return metrics queue_event_meta = {} queue_claims = {key: dict(value) for key, value in (queue_info.get('queue_claims') or {}).items()} active_lease_tokens = {claim.get('lease_token') for claim in queue_claims.values() if claim.get('lease_token')} reserved_lease_tokens = set(active_lease_tokens) active_lease_lock = threading.Lock() reservation_id = queue_info.get('spool_reservation_id') scan_slot_leases = list(queue_info.get('scan_slot_leases') or []) claim_guard = PostClaimRefundGuard( db, spool, reservation_id, queue_claims.values(), active_lease_tokens, active_lease_lock, dispatch_leases=scan_slot_leases, ) if postgres_db else None if claim_guard: scan_kwargs = claim_guard.call( 'infrastructure scan option preparation failure', prepare_scan_options, args, len(targets_to_scan), ) else: scan_kwargs = prepare_scan_options(args, len(targets_to_scan)) def progress(completed, total, target): if target != 'Scan completed': safe_print(f'Progress: {completed}/{total} finished - {target}', flush=True) if claim_guard: event_scan_options = claim_guard.call( 'infrastructure scan event option preparation failure', lambda: {key: value for key, value in scan_kwargs.items() if key != 'token'}, ) else: event_scan_options = {key: value for key, value in scan_kwargs.items() if key != 'token'} def spool_result(result): assign_finding_uids(result) queue_claim = queue_claims.get(str(result.get('target', ''))) or {} queue_id = queue_claim.get('id') lease_token = queue_claim.get('lease_token') if queue_id is None or not lease_token: raise RuntimeError(f"completed target has no fenced queue claim: {result.get('target', '')}") skipped_reason = str(result.get('skipped') or '') attempts = int(queue_claim.get('attempts') or 0) max_attempts = max(1, int(getattr(args, 'target_retry_max_attempts', 3) or 3)) reset_attempts = False if skipped_reason in CI_SOFT_SKIP_REASONS.get(args.platform, set()): cooldown_days = max(1, int(getattr(args, 'ci_soft_cooldown_days', 7) or 7)) available_after = (datetime.now(timezone.utc) + timedelta(days=cooldown_days)).isoformat(timespec='seconds') queue_status = 'deferred' queue_error = skipped_reason reset_attempts = True elif result.get('errors'): layer_disposition = docker_layer_queue_disposition( result, args, attempts=attempts, ) if layer_disposition is not None: queue_status, available_after, reset_attempts = layer_disposition else: queue_status, available_after, attempts, max_attempts = queue_error_disposition( db, source_for_db, args.platform, result.get('target', ''), result, args, queue_claim ) reset_attempts = queue_result_resets_attempts(result) queue_error = first_error_line(result) else: queue_status = 'done' available_after = None queue_error = None derived_targets = collect_postman_targets([result]) strip_nearby_context_for_persistence(result) try: event = prepare_scan_event({ 'version': 1, 'scan_event_id': result.get('scan_event_id'), 'run_id': run_id, 'cycle_id': cycle_id, 'source': source_for_db, 'query': args.query, 'target': result.get('target', ''), 'result': result, 'scan_options': event_scan_options, 'queue_id': queue_id, 'claim_lease_token': lease_token, 'claim_lease_owner': queue_claim.get('lease_owner'), 'queue_status': queue_status, 'queue_error': queue_error, 'available_after': available_after, 'reset_attempts': reset_attempts, 'derived_postman_targets': derived_targets, }) record = write_reserved_result_spool_event( spool, event, reservation_id, queue_id, lease_token, reserved_lease_tokens, active_lease_tokens, active_lease_lock, ) except Exception as exc: result['persistence_failure'] = str(exc)[:500] with active_lease_lock: active_lease_tokens.discard(lease_token) reserved_lease_tokens.discard(lease_token) refunded = db.refund_target_claim( queue_id, lease_token, f'infrastructure persistence failure: {exc}', ) reservation_released = False if refunded: try: reservation_released = spool.release_reserved_claim( reservation_id, queue_id, lease_token, ) except Exception: logger.exception('Unable to release failed result-spool reservation') logger.critical( 'COMPLETED RESULT COULD NOT BE PERSISTED target=%s queue_id=%s claim_refunded=%s reservation_released=%s error=%s', result.get('target', ''), queue_id, refunded, reservation_released, exc, ) if not refunded: raise RuntimeError( f'infrastructure persistence failed and claim could not be refunded for queue row {queue_id}: {exc}' ) from exc raise RuntimeError(f'infrastructure result persistence failure: {exc}') from exc queue_event_meta[record.event_id] = { 'queue_status': queue_status, 'attempts': attempts, 'max_attempts': max_attempts, } wait_for_result_spool_ready( spool, db, outcomes=ingest_outcomes, stop_event=getattr(args, 'result_spool_stop_event', None), wait_seconds=float(getattr(args, 'result_spool_wait_sec', 1.0) or 1.0), diagnostic_interval=float(getattr(args, 'result_spool_diagnostic_interval_sec', 30.0) or 30.0), ) outcome = ingest_outcomes.get(record.event_id) if outcome is None: confirm = getattr(db, 'confirmed_scan_event', None) outcome = confirm(record.event_id, record.event_hash) if confirm else None if not scan_outcome_matches(record, outcome): raise RuntimeError(f'database did not confirm durable handoff for scan event {record.event_id}') ingest_outcomes[record.event_id] = outcome heartbeat_stop = claim_guard.call( 'infrastructure heartbeat stop-token setup failure', threading.Event, ) if claim_guard else threading.Event() heartbeat_thread = None heartbeat_db = None heartbeat_failed = claim_guard.call( 'infrastructure heartbeat failure-token setup failure', threading.Event, ) if claim_guard else threading.Event() lease_owner = queue_info.get('lease_owner') lease_seconds = int(queue_info.get('lease_seconds') or 0) if db and lease_owner and lease_seconds: try: heartbeat_db = ScannerDB(db_path=getattr(db, 'path', None), db_url=getattr(db, 'url', None), initialize=False) if not heartbeat_db.enabled: raise RuntimeError(f'Unable to open dedicated lease heartbeat DB for {source_for_db}') set_application_name = getattr(heartbeat_db, 'set_application_name', None) if set_application_name: set_application_name(f'truf-heartbeat:{source_for_db}') heartbeat_db.require_runtime_safety_schema() if getattr(heartbeat_db.conn, 'is_postgres', False): heartbeat_db.conn.execute("SELECT set_config('statement_timeout', '10000ms', false)") heartbeat_db.conn.execute("SELECT set_config('lock_timeout', '5000ms', false)") heartbeat_db.conn.commit() except Exception as exc: if heartbeat_db is not None: heartbeat_db.close() heartbeat_db = None claim_guard.refund_undurable(f'infrastructure heartbeat setup failure: {exc}') raise def renew_leases(): interval = max(10, min(60, lease_seconds // 3)) last_success = time.monotonic() failure_grace = max(60, lease_seconds // 2) while not heartbeat_stop.wait(interval): with active_lease_lock: current_tokens = list(active_lease_tokens) if not current_tokens: return renewed = heartbeat_db.renew_target_leases(lease_owner, lease_seconds, current_tokens) if renewed == len(current_tokens): if reservation_id: try: if not renew_active_result_spool_reservation( spool, reservation_id, lease_seconds, reserved_lease_tokens, active_lease_lock, ): raise RuntimeError('reservation is absent or expired') except Exception as exc: logger.critical('Result-spool reservation renewal failed: %s', exc) heartbeat_failed.set() return last_success = time.monotonic() continue if recheck_active_lease_ownership( heartbeat_db, lease_owner, current_tokens, active_lease_tokens, active_lease_lock, ): if reservation_id: try: if not renew_active_result_spool_reservation( spool, reservation_id, lease_seconds, reserved_lease_tokens, active_lease_lock, ): raise RuntimeError('reservation is absent or expired') except Exception as exc: logger.critical('Result-spool reservation renewal failed: %s', exc) heartbeat_failed.set() return last_success = time.monotonic() continue if heartbeat_db.last_error and time.monotonic() - last_success < failure_grace: continue heartbeat_failed.set() return try: candidate_thread = threading.Thread(target=renew_leases, name=f'lease-heartbeat-{source_for_db}', daemon=True) candidate_thread.start() heartbeat_thread = candidate_thread except Exception as exc: if heartbeat_db is not None: heartbeat_db.close() heartbeat_db = None claim_guard.refund_undurable(f'infrastructure heartbeat thread start failure: {exc}') raise results = [] batch_sink_error = None try: try: results = scan_targets_batch( targets_to_scan, args.platform, progress_callback=progress, max_workers=args.workers, persist_results=not bool(queue_info.get('lease_owner')), result_sink=spool_result if postgres_db else None, scan_slot_leases=scan_slot_leases, sink_within_scan_slot=bool(postgres_db and scan_slot_leases), **scan_kwargs, ) except ResultSinkError as exc: results = exc.results batch_sink_error = exc except Exception as exc: if claim_guard: claim_guard.refund_undurable(f'infrastructure scan submission failure: {exc}') raise finally: heartbeat_stop.set() heartbeat_stuck = False if heartbeat_thread: # The Postgres heartbeat has a 10s statement timeout. Give an in-flight # renewal time to finish instead of closing its connection underneath it. heartbeat_thread.join(timeout=15) heartbeat_stuck = heartbeat_thread.is_alive() if heartbeat_db and not heartbeat_stuck: heartbeat_db.close() if heartbeat_stuck: if claim_guard: claim_guard.refund_undurable('infrastructure lease heartbeat shutdown failure') raise RuntimeError(f'Lease heartbeat did not stop for {source_for_db}') if claim_guard: with active_lease_lock: undurable_claims = bool(active_lease_tokens) if undurable_claims: claim_guard.refund_undurable('scan batch returned without a durable result for every fenced claim') raise RuntimeError(f'Scan batch returned without durable results for every claim in {source_for_db}') if spool: wait_for_result_spool_ready( spool, db, outcomes=ingest_outcomes, stop_event=getattr(args, 'result_spool_stop_event', None), wait_seconds=float(getattr(args, 'result_spool_wait_sec', 1.0) or 1.0), diagnostic_interval=float(getattr(args, 'result_spool_diagnostic_interval_sec', 30.0) or 30.0), ) if reservation_id: spool.release_reservation(reservation_id) if heartbeat_failed.is_set(): raise RuntimeError(f'Lease heartbeat lost ownership for {source_for_db}') if batch_sink_error: raise batch_sink_error db_queue = postgres_db project_files = not db_queue or bool(getattr(args, 'sync_file_queues', True)) results_for_checked = [result for result in results if not result.get('errors')] if db_queue: successful_results = [] accepted_results = [] for result in results: event_id = str(result.get('scan_event_id') or '') persisted = ingest_outcomes.get(event_id) if not persisted or not persisted.get('ingested'): raise RuntimeError(f'No confirmed database ingestion for scan event {event_id or ""}') event_meta = queue_event_meta.get(event_id) or {} queue_status = event_meta.get('queue_status') accepted_results.append(result) if persisted.get('stale'): safe_print( f"Preserved stale scan result for {result.get('target', '')}; " 'newer queue ownership/status was left unchanged' ) elif queue_status in ('done', 'failed'): successful_results.append(result) if queue_status == 'failed' and persisted.get('queue_completion_applied'): safe_print( f"Target retry exhausted/non-retryable: {result.get('target', '')} " f"class={result.get('error_class', 'unknown')} " f"attempts={event_meta.get('attempts', 0)}/{event_meta.get('max_attempts', 0)}" ) results_for_checked = successful_results drain_scan_publication_outbox(db, 100) elif db and cycle_id: accepted_results = results recorded = db.record_target_results(run_id, cycle_id, source_for_db, args.query, results, scan_kwargs) successful_results = [] for result in results: normalized = normalize_target(result.get('target', ''), source_for_db) target_scan_id = (recorded or {}).get(normalized) if target_scan_id is None: target_scan_id = (recorded or {}).get(str(result.get('target', ''))) if target_scan_id is not None and not result.get('errors'): successful_results.append(result) results_for_checked = successful_results else: accepted_results = results harvested_postman_targets = collect_postman_targets(accepted_results) if harvested_postman_targets and project_files: try: queued_postman = enqueue_targets_for_platform(getattr(args, 'queue_dir', None) or args.save_dir, 'postman', harvested_postman_targets) action = 'projected' if db_queue else 'queued' print(f'Harvested {len(harvested_postman_targets)} Postman artifact(s); {action} {queued_postman} new Postman target(s).') except Exception as exc: if not db_queue: raise safe_print(f'Warning: harvested Postman targets were committed to Postgres but file projection failed: {exc}') if project_files: try: mark_checked(results_for_checked, todo_file, checked_file, args.platform) except Exception as exc: if not db_queue: raise safe_print(f'Warning: queue completion was committed to Postgres but checked/todo projection failed: {exc}') print_result_report(results, args.save_dir) metrics = summarize_results(results) source_failures = [result for result in results if result.get('source_failure')] source_failure_queue_ids = [] for result in source_failures: claim = (queue_info.get('queue_claims') or {}).get(str(result.get('target', ''))) or {} if claim.get('id') is not None: source_failure_queue_ids.append(claim['id']) metrics.update({ 'fetched_count': queue_info.get('fetched_count', len(fetched_targets)), 'queued_new_count': queue_info.get('queued_new_count', 0), 'queued_updated_count': queue_info.get('queued_updated_count', 0), 'scan_requested_count': queue_info.get('scan_requested_count', len(targets_to_scan)), 'backlog_only': backlog_only, 'source_failure_count': len(source_failures), 'source_failure_category': source_failures[0].get('source_failure_category') if source_failures else '', 'source_failure_auth_related': bool(source_failures and source_failures[0].get('source_failure_auth_related')), 'source_failure_message': first_error_line(source_failures[0]) if source_failures else '', 'source_failure_queue_ids': source_failure_queue_ids, }) partial_metrics.update(metrics) queue_after = queue_counts(todo_file, checked_file) if project_files else { 'todo_count': 0, 'checked_count': 0, 'todo_file': None, 'checked_file': None, } if db and cycle_id: db.finish_source_cycle( cycle_id, 'source_failed' if source_failures else 'completed', metrics, queue_after, metrics.get('source_failure_message') or None, ) print( f"Cycle complete: {metrics['scanned_count']} scanned, " f"{metrics['findings_count']} findings, {metrics['error_count']} errors, " f"{metrics.get('degraded_count', 0)} degraded." ) return metrics def run_cycle(args, db=None, run_id=None, cycle_id=None, source_name=None, partial_metrics=None): source = source_name or args.platform postgres_db = bool( db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False) ) if postgres_db and hasattr(db, 'reserve_and_claim_target'): return run_cycle_v2(args, db, run_id, cycle_id, source, partial_metrics) # SQLite and narrow test doubles retain the explicit v1 compatibility adapter. return _run_cycle_legacy_compat(args, db, run_id, cycle_id, source, partial_metrics) def load_config(path, *, managed_postgres=None, final_cutover=None): try: import yaml except ImportError as e: raise SystemExit('PyYAML is required for --config mode. Run: python -m pip install -r requirements.txt') from e if not os.path.exists(path): raise FileNotFoundError(path) with open(path, 'r', encoding='utf-8') as f: loaded = yaml.safe_load(f) if loaded is None: loaded = {} validated = validate_docker_depth_config( loaded, managed_postgres=managed_postgres, final_cutover=final_cutover, ) return apply_path_config(validated.config, path) def resolve_config_path(config_path, value): if not value: return value if os.path.isabs(value): return value return os.path.join(os.path.dirname(os.path.abspath(config_path)), value) def load_secrets(config, config_path): global_config = config.get('global', {}) secrets_file = global_config.get('secrets_file') if not secrets_file: return {} path = resolve_config_path(config_path, secrets_file) if not os.path.exists(path): print(f'Info: secrets file {path} not found. Falling back to env/source tokens where configured.') return {} try: import yaml except ImportError as e: raise SystemExit('PyYAML is required for secrets.yaml. Run: python -m pip install -r requirements.txt') from e with open(path, 'r', encoding='utf-8') as f: return yaml.safe_load(f) or {} def parse_state_time(value): if not value: return None try: parsed = datetime.fromisoformat(str(value).replace('Z', '+00:00')) if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=timezone.utc) return parsed.astimezone(timezone.utc) except ValueError: return None def utc_now(): return datetime.now(timezone.utc) def normalize_queries(value): if value is None: return [''] if isinstance(value, str): queries = [item.strip() for item in value.split(',')] else: queries = [str(item).strip() for item in value] queries = [query for query in queries if query] return list(dict.fromkeys(queries)) or [''] def _valid_dockerhub_policy_sha256(value): return bool( isinstance(value, str) and len(value) == 64 and all(character in '0123456789abcdef' for character in value) ) def _valid_dockerhub_deep_record(record): if not isinstance(record, dict): return None policy_sha256 = record.get('policy_sha256') timestamp = record.get('last_dispatched_at') if not _valid_dockerhub_policy_sha256(policy_sha256) or not isinstance(timestamp, str): return None try: parsed = datetime.fromisoformat(timestamp.replace('Z', '+00:00')) except (ValueError, OverflowError): return None if parsed.tzinfo is None or parsed.utcoffset() != timedelta(0): return None return { 'policy_sha256': policy_sha256, 'last_dispatched_at': timestamp, '_parsed': parsed.astimezone(timezone.utc), } def prepare_dockerhub_discovery_state( source_state, configured_policies, query, now=None, *, force_deep=False, ): if not isinstance(force_deep, bool): raise ValueError('DockerHub forced deep-pass authority is invalid') configured_queries = set(configured_policies) parent = source_state.get(DOCKERHUB_DISCOVERY_STATE_KEY) raw_records = parent.get('deep_by_query') if isinstance(parent, dict) else None raw_incomplete = ( parent.get('incomplete_by_query') if isinstance(parent, dict) else None ) if not ( isinstance(parent, dict) and parent.get('schema') == 1 and not isinstance(parent.get('schema'), bool) and isinstance(raw_records, dict) ): raw_records = {} raw_incomplete = {} records = {} parsed_records = {} incomplete = {} for configured_query in configured_queries: valid = _valid_dockerhub_deep_record(raw_records.get(configured_query)) if valid is not None: records[configured_query] = { 'policy_sha256': valid['policy_sha256'], 'last_dispatched_at': valid['last_dispatched_at'], } parsed_records[configured_query] = valid['_parsed'] incomplete_record = ( raw_incomplete.get(configured_query) if isinstance(raw_incomplete, dict) else None ) configured_policy = configured_policies.get(configured_query) if ( isinstance(incomplete_record, dict) and isinstance(configured_policy, dict) and incomplete_record.get('policy_sha256') == configured_policy.get('policy_sha256') ): incomplete[configured_query] = { 'policy_sha256': incomplete_record['policy_sha256'], } source_state[DOCKERHUB_DISCOVERY_STATE_KEY] = { 'schema': 1, 'deep_by_query': records, 'incomplete_by_query': incomplete, } policy = configured_policies.get(query) if not isinstance(policy, dict) or not _valid_dockerhub_policy_sha256( policy.get('policy_sha256') ): raise ValueError('DockerHub discovery policy state is unavailable') current = records.get(query) now = (now or utc_now()).astimezone(timezone.utc) last_dispatched = parsed_records.get(query) deep_due = bool( force_deep or current is None or current.get('policy_sha256') != policy['policy_sha256'] or last_dispatched is None or last_dispatched > now or now - last_dispatched >= DOCKERHUB_DEEP_INTERVAL or query in incomplete ) return { 'deep': deep_due, 'pass_kind': 'deep' if deep_due else 'ordinary', 'policy_sha256': policy['policy_sha256'], } def mark_dockerhub_deep_dispatched(source_state, query, policy_sha256, now=None): if not _valid_dockerhub_policy_sha256(policy_sha256): raise ValueError('DockerHub discovery policy hash is invalid') parent = source_state.get(DOCKERHUB_DISCOVERY_STATE_KEY) if not isinstance(parent, dict) or parent.get('schema') != 1: parent = {'schema': 1, 'deep_by_query': {}, 'incomplete_by_query': {}} source_state[DOCKERHUB_DISCOVERY_STATE_KEY] = parent records = parent.get('deep_by_query') if not isinstance(records, dict): records = {} parent['deep_by_query'] = records dispatched_at = (now or utc_now()).astimezone(timezone.utc).isoformat( timespec='seconds', ) records[query] = { 'policy_sha256': policy_sha256, 'last_dispatched_at': dispatched_at, } def mark_dockerhub_discovery_incomplete(source_state, query, policy_sha256): if not _valid_dockerhub_policy_sha256(policy_sha256): raise ValueError('DockerHub discovery policy hash is invalid') parent = source_state.get(DOCKERHUB_DISCOVERY_STATE_KEY) if not isinstance(parent, dict) or parent.get('schema') != 1: parent = {'schema': 1, 'deep_by_query': {}, 'incomplete_by_query': {}} source_state[DOCKERHUB_DISCOVERY_STATE_KEY] = parent incomplete = parent.get('incomplete_by_query') if not isinstance(incomplete, dict): incomplete = {} parent['incomplete_by_query'] = incomplete incomplete[query] = {'policy_sha256': policy_sha256} def clear_dockerhub_discovery_incomplete(source_state, query): parent = source_state.get(DOCKERHUB_DISCOVERY_STATE_KEY) incomplete = parent.get('incomplete_by_query') if isinstance(parent, dict) else None if isinstance(incomplete, dict): incomplete.pop(query, None) def annotate_dockerhub_discovery_args(args, annotation): args.dockerhub_discovery_deep = bool(annotation['deep']) args.dockerhub_discovery_pass_kind = str(annotation['pass_kind']) args.dockerhub_discovery_policy_sha256 = str(annotation['policy_sha256']) return args def get_state_path(config, config_path): env_state_file = os.getenv('RUNNER_STATE_FILE') or os.getenv('SCANNER_STATE_FILE') if env_state_file: return resolve_config_path(config_path, env_state_file) global_config = config.get('global', {}) state_file = global_config.get('state_file') if state_file: return state_file results_dir = global_config.get('results_dir') or scan_config.results_dir return os.path.join(results_dir, 'runner_state.json') def default_source_state(): return { 'query_index': 0, 'auth_index': 0, 'auth_status': {}, 'auth_endpoint_status': {}, 'last_auth': None, 'last_query': None, 'last_started_at': None, 'last_completed_at': None, 'last_status': None, 'cycles': 0, } def load_state(path, config): if os.path.exists(path): try: with open(path, 'r', encoding='utf-8') as f: state = json.load(f) except (OSError, json.JSONDecodeError): state = {} else: state = {} state.setdefault('version', 1) state.setdefault('sources', {}) for source_name in config.get('sources', {}): source_state = state['sources'].setdefault(source_name, default_source_state()) for key, value in default_source_state().items(): source_state.setdefault(key, value) return state def save_state(path, state): parent = os.path.dirname(path) if parent: ensure_private_directory(parent, reject_reparse=True) tmp_path = f'{path}.{os.getpid()}.{time.time_ns()}.tmp' with open(tmp_path, 'w', encoding='utf-8') as f: json.dump(state, f, indent=2, ensure_ascii=False) delays = (0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 3.0) last_error = None for attempt, delay in enumerate((*delays, None), 1): try: os.replace(tmp_path, path) harden_private_file(path) return except PermissionError as e: last_error = e if delay is None: break time.sleep(delay) except OSError as e: if getattr(e, 'winerror', None) != 5 or delay is None: raise last_error = e time.sleep(delay) try: with open(path, 'w', encoding='utf-8') as f: json.dump(state, f, indent=2, ensure_ascii=False) except OSError as e: raise e from last_error finally: try: if os.path.exists(tmp_path): os.remove(tmp_path) except OSError: pass def enabled_source_names(config, selected_source=None): sources = config.get('sources', {}) if selected_source: selected_source = 'dockerhub' if selected_source == 'docker' else selected_source if selected_source not in sources: raise SystemExit(f'Source {selected_source} is not present in config') return [selected_source] return [name for name, source in sources.items() if source.get('enabled', False)] def auth_pool_entries(source_config, secrets): pool_name = source_config.get('auth_pool') if not pool_name: return [], None entries = (secrets.get('auth_pools') or {}).get(pool_name, []) normalized = [] for index, entry in enumerate(entries): if not isinstance(entry, dict): continue entry = dict(entry) entry.setdefault('name', f'{pool_name}_{index + 1}') normalized.append(entry) return normalized, pool_name def auth_status_kind(status): if not isinstance(status, dict): return 'ok' if status.get('disabled_until') == 'manual' or status.get('status') == 'dead' or status.get('disabled_reason') == 'auth_invalid': return 'dead' disabled_until = parse_state_time(status.get('disabled_until')) if disabled_until and disabled_until > utc_now(): return 'limited' return 'ok' def auth_entry_is_available(source_name, entry, state): source_state = state['sources'].setdefault(source_name, default_source_state()) auth_status = source_state.setdefault('auth_status', {}) status = auth_status.get(entry.get('name'), {}) if auth_status_kind(status) == 'dead': return False disabled_until = parse_state_time(status.get('disabled_until')) return not disabled_until or disabled_until <= utc_now() def refresh_auth_summary(source_name, source_config, state, secrets, current_auth_name=None): entries, pool_name = auth_pool_entries(source_config, secrets) source_state = state['sources'].setdefault(source_name, default_source_state()) auth_status = source_state.setdefault('auth_status', {}) counts = {'ok': 0, 'dead': 0, 'limited': 0} rate_limit_errors = 0 auth_invalid_errors = 0 for entry in entries: name = entry.get('name') item = auth_status.setdefault(name, {}) kind = auth_status_kind(item) counts[kind] = counts.get(kind, 0) + 1 rate_limit_errors += int(item.get('rate_limit_count', 0) or 0) rate_limit_errors += int(item.get('secondary_rate_limit_count', 0) or 0) auth_invalid_errors += int(item.get('auth_invalid_count', 0) or 0) if current_auth_name is None: current_auth_name = source_state.get('last_auth') source_state['auth_summary'] = { 'pool': pool_name or '', 'current': current_auth_name or source_state.get('last_auth') or 'none', 'total': len(entries), 'ok': counts.get('ok', 0), 'dead': counts.get('dead', 0), 'limited': counts.get('limited', 0), 'rate_limit_errors': rate_limit_errors, 'auth_invalid_errors': auth_invalid_errors, } return source_state['auth_summary'] def available_auth_entries(source_name, source_config, state, secrets): entries, _ = auth_pool_entries(source_config, secrets) return [entry for entry in entries if auth_entry_is_available(source_name, entry, state)] def select_auth_entry(source_name, source_config, state, secrets): entries, pool_name = auth_pool_entries(source_config, secrets) if not entries: return None available = [entry for entry in entries if auth_entry_is_available(source_name, entry, state)] if not available: return None source_state = state['sources'].setdefault(source_name, default_source_state()) rotation = source_config.get('auth_rotation', 'per_cycle') if rotation == 'none': entry = available[0] else: start = int(source_state.get('auth_index', 0)) entry = None for offset in range(len(entries)): candidate = entries[(start + offset) % len(entries)] if candidate in available: entry = candidate source_state['auth_index'] = (start + offset + 1) % len(entries) break if entry is None: entry = available[0] source_state['last_auth'] = entry.get('name') return entry def mark_auth_rate_limited(source_name, auth_entry, state, reset_at=None, cooldown=3600, category='rate_limit', message=None): if not auth_entry: return source_state = state['sources'].setdefault(source_name, default_source_state()) auth_status = source_state.setdefault('auth_status', {}) name = auth_entry.get('name') now = utc_now() disabled_until = reset_at if category == 'auth_invalid': disabled_until = 'manual' elif not disabled_until: disabled_until = (now + timedelta(seconds=int(cooldown))).isoformat(timespec='seconds') item = auth_status.setdefault(name, {}) if ( category != 'auth_invalid' and ( item.get('disabled_until') == 'manual' or item.get('status') == 'dead' or item.get('disabled_reason') == 'auth_invalid' ) ): return item['disabled_until'] = disabled_until item['disabled_reason'] = category item['status'] = 'dead' if category == 'auth_invalid' else 'limited' item['last_error'] = str(message or '')[:500] item['last_error_at'] = now.isoformat(timespec='seconds') item['last_rate_limited_at'] = now.isoformat(timespec='seconds') item['failures'] = int(item.get('failures', 0)) + 1 if category == 'auth_invalid': item['auth_invalid_count'] = int(item.get('auth_invalid_count', 0)) + 1 item['dead_at'] = now.isoformat(timespec='seconds') elif category == 'secondary_rate_limit': item['secondary_rate_limit_count'] = int(item.get('secondary_rate_limit_count', 0)) + 1 else: item['rate_limit_count'] = int(item.get('rate_limit_count', 0)) + 1 def clear_auth_rate_limit(source_name, auth_entry, state): if not auth_entry: return source_state = state['sources'].setdefault(source_name, default_source_state()) auth_status = source_state.setdefault('auth_status', {}) item = auth_status.setdefault(auth_entry.get('name'), {}) if item.get('disabled_until') == 'manual' or item.get('status') == 'dead': return item['disabled_until'] = None item['disabled_reason'] = None item['status'] = 'ok' item['last_success_at'] = utc_now().isoformat(timespec='seconds') item['success_count'] = int(item.get('success_count', 0)) + 1 def auth_token_from_entry(entry): return entry.get('token') if entry else None def current_query_for_source(source_name, source_config, state): queries = normalize_queries(source_config.get('queries')) source_state = state['sources'].setdefault(source_name, default_source_state()) index = int(source_state.get('query_index', 0)) % len(queries) source_state['query_index'] = index return queries[index], index, len(queries) def advance_query_for_source(source_name, source_config, state): queries = normalize_queries(source_config.get('queries')) source_state = state['sources'].setdefault(source_name, default_source_state()) source_state['query_index'] = (int(source_state.get('query_index', 0)) + 1) % len(queries) def source_to_platform(source_name): return 'docker' if source_name in ('docker', 'dockerhub') else source_name def normalize_archive_event_types(value): if value is None: return ['PushEvent', 'CreateEvent', 'PublicEvent'] if isinstance(value, str): return [item.strip() for item in value.split(',') if item.strip()] return [str(item).strip() for item in value if str(item).strip()] def normalize_package_sources(value): if value is None: return 'npm,pypi' if isinstance(value, str): return ','.join(item.strip() for item in value.split(',') if item.strip()) return ','.join(str(item).strip() for item in value if str(item).strip()) def normalize_search_kinds(value): if value is None: return 'collection,environment' if isinstance(value, str): return ','.join(item.strip() for item in value.split(',') if item.strip()) return ','.join(str(item).strip() for item in value if str(item).strip()) def github_token_entries_for_postman(source_config, secrets): entries, _ = auth_pool_entries(source_config, secrets) token_entries = [] for entry in entries: if entry.get('token'): token_entries.append({'name': entry.get('name'), 'token': entry.get('token')}) if source_config.get('token'): token_entries.append({'name': 'source_token', 'token': source_config.get('token')}) return token_entries def build_args_from_source_config(source_name, source_config, global_config, query, auth_entry=None): query_overrides = source_config.get('query_overrides') or {} if not isinstance(query_overrides, dict): raise ValueError(f'{source_name} query_overrides must be a mapping') query_override = query_overrides.get(query) if query_override is not None: if not isinstance(query_override, dict): raise ValueError(f'{source_name} query override for {query!r} must be a mapping') allowed = {'pages', 'per_page', 'max_targets'} unknown = sorted(set(query_override) - allowed) if unknown: raise ValueError( f'{source_name} query override for {query!r} has unsupported keys: ' + ', '.join(unknown) ) effective_source_config = dict(source_config) for key, value in query_override.items(): try: value = int(value) except (TypeError, ValueError) as e: raise ValueError( f'{source_name} query override {key} for {query!r} must be an integer' ) from e minimum = 0 if key == 'max_targets' else 1 if value < minimum: raise ValueError( f'{source_name} query override {key} for {query!r} must be >= {minimum}' ) effective_source_config[key] = value source_config = effective_source_config platform = source_to_platform(source_name) token = auth_token_from_entry(auth_entry) or source_config.get('token') package_sources = normalize_package_sources(source_config.get('package_sources', global_config.get('package_sources', ['npm', 'pypi']))) path_context = global_config return SimpleNamespace( platform=platform, configured_queries=tuple(normalize_queries(source_config.get('queries'))), mode=source_config.get('mode', 'recent'), query=query, query_file=None, pages=int(source_config.get('pages', 1)), per_page=int(source_config.get('per_page', 50)), target_file=resolve_optional_path(source_config.get('target_file'), path_context), token=token, docker_username=(auth_entry or {}).get('username') or source_config.get('docker_username'), docker_token=(auth_entry or {}).get('token') or source_config.get('docker_token'), workers=int(source_config.get('workers', 6)), timeout=int(source_config.get('timeout', scan_config.docker_timeout if platform == 'docker' else scan_config.git_timeout)), detectors=str(source_config.get('detectors', global_config.get('detectors', scan_config.detectors))), exclude_detectors=str(source_config.get('exclude_detectors', global_config.get('exclude_detectors', scan_config.exclude_detectors))), drop_detectors=source_config.get('drop_detectors', global_config.get('drop_detectors', getattr(scan_config, 'drop_detectors', []))), no_verification=bool(source_config.get('no_verification', global_config.get('no_verification', scan_config.no_verification))), save_dir=global_config.get('results_dir', scan_config.results_dir), runtime_dir=global_config.get('runtime_dir'), result_bundle_dir=global_config.get('result_bundle_dir', scan_config.result_bundle_dir), result_bundle_max_event_bytes=int(global_config.get('result_bundle_max_event_bytes', scan_config.result_bundle_max_event_bytes)), remote_assignment_reserve_bytes=int(global_config.get('remote_assignment_reserve_bytes', 2 * 1024 * 1024)), remote_assignment_max_active=int(global_config.get('remote_assignment_max_active', 50)), result_bundle_max_items=int(global_config.get('result_bundle_max_items', 10000)), result_bundle_max_total_bytes=int(global_config.get('result_bundle_max_total_bytes', 3 * 1024 * 1024 * 1024)), result_bundle_min_free_bytes=int(global_config.get('result_bundle_min_free_bytes', 20 * 1024 * 1024 * 1024)), projection_backlog_max_items=int(global_config.get('projection_backlog_max_items', 10000)), projection_backlog_max_bytes=int(global_config.get('projection_backlog_max_bytes', 2 * 1024 * 1024 * 1024)), projection_backlog_headroom_bytes=int(global_config.get( 'projection_backlog_headroom_bytes', 2 * int(global_config.get('result_bundle_max_event_bytes', scan_config.result_bundle_max_event_bytes)), )), keycheck_queue_max_items=int(global_config.get('keycheck_queue_max_items', 100000)), keycheck_queue_max_bytes=int(global_config.get('keycheck_queue_max_bytes', 512 * 1024 * 1024)), pipeline_quarantine_max_items=int(global_config.get('pipeline_quarantine_max_items', 10000)), pipeline_quarantine_max_bytes=int(global_config.get('pipeline_quarantine_max_bytes', 1024 * 1024 * 1024)), keycheck_candidates_per_event=int(global_config.get('keycheck_candidates_per_event', 2000)), keycheck_candidate_bytes_per_event=int(global_config.get('keycheck_candidate_bytes_per_event', 2 * 1024 * 1024)), result_spool_dir=global_config.get('result_spool_dir', scan_config.result_spool_dir), result_spool_max_event_bytes=int(global_config.get('result_spool_max_event_bytes', scan_config.result_spool_max_event_bytes)), result_spool_max_events=int(global_config.get('result_spool_max_events', scan_config.result_spool_max_events)), result_spool_max_total_bytes=int(global_config.get('result_spool_max_total_bytes', scan_config.result_spool_max_total_bytes)), result_spool_min_free_bytes=int(global_config.get('result_spool_min_free_bytes', scan_config.result_spool_min_free_bytes)), result_spool_wait_sec=float(global_config.get('result_spool_wait_sec', 1.0)), result_spool_diagnostic_interval_sec=float(global_config.get('result_spool_diagnostic_interval_sec', 30.0)), scan_outbox_max_pending_items=int(global_config.get('scan_outbox_max_pending_items', scan_config.scan_outbox_max_pending_items)), scan_outbox_max_pending_bytes=int(global_config.get('scan_outbox_max_pending_bytes', scan_config.scan_outbox_max_pending_bytes)), scan_outbox_max_pending_age_sec=int(global_config.get('scan_outbox_max_pending_age_sec', scan_config.scan_outbox_max_pending_age_sec)), database_path=global_config.get('database_path'), database_url=global_config.get('database_url'), queue_dir=global_config.get('queue_dir', getattr(scan_config, 'queue_dir', scan_config.results_dir)), work_dir=global_config.get('work_dir', scan_config.work_dir), trufflehog_path=global_config.get('trufflehog_path', scan_config.trufflehog_path), trufflehog_config=resolve_optional_path( source_config.get('trufflehog_config', global_config.get('trufflehog_config', getattr(scan_config, 'trufflehog_config', ''))), path_context, ), trufflehog_job_memory_limit_bytes=int(source_config.get( 'trufflehog_job_memory_limit_bytes', global_config.get('trufflehog_job_memory_limit_bytes', scan_config.trufflehog_job_memory_limit_bytes), )), trufflehog_concurrency=int(source_config.get( 'trufflehog_concurrency', global_config.get('trufflehog_concurrency', 0), )), loop=False, cooldown=int(source_config.get('cooldown', global_config.get('cooldown', 300))), max_cycles=0, max_targets=int(source_config.get('max_targets', 0)), target_retry_max_attempts=int(source_config.get('target_retry_max_attempts', global_config.get('target_retry_max_attempts', 3))), target_retry_base_delay_sec=int(source_config.get('target_retry_base_delay_sec', global_config.get('target_retry_base_delay_sec', 3600))), target_retry_max_delay_sec=int(source_config.get('target_retry_max_delay_sec', global_config.get('target_retry_max_delay_sec', 86400))), target_timeout_retry_delay_sec=int(source_config.get('target_timeout_retry_delay_sec', global_config.get('target_timeout_retry_delay_sec', 21600))), target_claim_batch_size=int(source_config.get('target_claim_batch_size', global_config.get('target_claim_batch_size', 0))), target_claim_order=str(source_config.get('target_claim_order', global_config.get('target_claim_order', 'oldest'))).strip().lower(), sync_file_queues=bool_config(source_config.get('sync_file_queues', global_config.get('sync_file_queues', True)), True), recent_hours=int(source_config.get('recent_hours', 24)), recent_days=int(source_config.get('recent_days', 7)), max_version_age_days=int(source_config.get('max_version_age_days', 0)), versions_per_package=int(source_config.get('versions_per_package', global_config.get('versions_per_package', 1))), package_sources=package_sources, refresh_registry=bool(source_config.get('refresh_registry', global_config.get('refresh_registry', False))), updated_target_rescan_enabled=bool_config( source_config.get('updated_target_rescan_enabled', False), False, ), updated_target_rescan_max_per_cycle=int( source_config.get('updated_target_rescan_max_per_cycle', 0) or 0 ), updated_target_rescan_cooldown_hours=int( source_config.get('updated_target_rescan_cooldown_hours', 0) or 0 ), search_kinds=normalize_search_kinds(source_config.get('search_kinds', global_config.get('search_kinds', ['collection', 'environment']))), gist_since=source_config.get('gist_since', global_config.get('gist_since', '')), postman_cache_dir=resolve_optional_path( source_config.get('postman_cache_dir') or global_config.get('postman_cache_dir') or getattr(scan_config, 'postman_cache_dir', None), path_context, ), postman_discovery_max_artifacts_per_cycle=int(source_config.get('postman_discovery_max_artifacts_per_cycle', global_config.get('postman_discovery_max_artifacts_per_cycle', getattr(scan_config, 'postman_discovery_max_artifacts_per_cycle', 1000)))), postman_discovery_max_artifacts_per_page=int(source_config.get('postman_discovery_max_artifacts_per_page', global_config.get('postman_discovery_max_artifacts_per_page', getattr(scan_config, 'postman_discovery_max_artifacts_per_page', 100)))), postman_discovery_max_bytes_per_cycle=int(source_config.get('postman_discovery_max_bytes_per_cycle', global_config.get('postman_discovery_max_bytes_per_cycle', getattr(scan_config, 'postman_discovery_max_bytes_per_cycle', 1024 * 1024 * 1024)))), postman_discovery_max_elapsed_sec=float(source_config.get('postman_discovery_max_elapsed_sec', global_config.get('postman_discovery_max_elapsed_sec', getattr(scan_config, 'postman_discovery_max_elapsed_sec', 300.0)))), gharchive_cache_dir=resolve_optional_path( source_config.get('gharchive_cache_dir') or global_config.get('gharchive_cache_dir'), path_context, ), max_artifact_size_mb=int(source_config.get('max_artifact_size_mb', global_config.get('max_artifact_size_mb', 50))), max_file_age_days=int(source_config.get('max_file_age_days', global_config.get('max_file_age_days', 365))), github_code_search_rpm=int(source_config.get('github_code_search_rpm', global_config.get('github_code_search_rpm', 8))), all_tokens_cooldown=int(source_config.get('all_tokens_cooldown', global_config.get('all_tokens_cooldown', 1800))), max_repo_age_days=int(source_config.get('max_repo_age_days', 0)), repo_age_field=source_config.get('repo_age_field', 'created_at'), max_commit_age_days=int(source_config.get('max_commit_age_days', 0)), commit_lookup_pages=int(source_config.get('commit_lookup_pages', 3)), skip_if_commit_lookup_fails=bool(source_config.get('skip_if_commit_lookup_fails', True)), raise_rate_limit=bool(source_config.get('retry_with_next_auth_on_rate_limit', True)), auth_name=(auth_entry or {}).get('name'), max_depth=int(source_config.get('max_depth', 0)), exact_git_planning_enabled=bool_config( source_config.get( 'exact_git_planning_enabled', global_config.get('exact_git_planning_enabled', False), ), False, ), git_baseline_depth=int(source_config.get( 'git_baseline_depth', source_config.get('max_depth', 100) or 100, )), git_ref_resolution_timeout_sec=float(source_config.get( 'git_ref_resolution_timeout_sec', global_config.get('git_ref_resolution_timeout_sec', 10), )), git_ref_resolution_attempts=int(source_config.get( 'git_ref_resolution_attempts', global_config.get('git_ref_resolution_attempts', 2), )), git_ref_resolution_max_bytes=int(source_config.get( 'git_ref_resolution_max_bytes', global_config.get('git_ref_resolution_max_bytes', 1 << 20), )), admission_resolution_attempts=int(source_config.get( 'admission_resolution_attempts', global_config.get('admission_resolution_attempts', 90), )), admission_resolution_seconds=float(source_config.get( 'admission_resolution_seconds', global_config.get('admission_resolution_seconds', 90), )), admission_resolution_retry_delay_sec=float(source_config.get( 'admission_resolution_retry_delay_sec', global_config.get('admission_resolution_retry_delay_sec', 1), )), scan_full_history=bool(source_config.get('scan_full_history', False)), sort_by=source_config.get('sort_by', 'updated'), sort_order=source_config.get('sort_order', 'desc'), created_filter=source_config.get('created_filter', 'any'), gitlab_sort_by=source_config.get('gitlab_sort_by', source_config.get('sort_by', 'last_activity_at')), gitlab_visibility=source_config.get('gitlab_visibility', source_config.get('visibility', 'public')), gitlab_discovery_request_attempts=int(source_config.get('discovery_request_attempts', 1)), gitlab_discovery_retry_delay=int(source_config.get('discovery_retry_delay', 0)), huggingface_discovery_request_attempts=int(source_config.get('discovery_request_attempts', 1)), huggingface_discovery_retry_delay=int(source_config.get('discovery_retry_delay', 0)), external_trufflehog_lifecycle=bool(source_config.get('external_trufflehog_lifecycle', False)), docker_sort_by=source_config.get('docker_sort_by', source_config.get('sort_by', 'updated_at')), fetch_workers=int(source_config.get('fetch_workers', global_config.get('fetch_workers', 8))), fetch_timeout=int(source_config.get('fetch_timeout', global_config.get('fetch_timeout', 15))), tag_fetch_workers=int(source_config.get('tag_fetch_workers', global_config.get('tag_fetch_workers', 4))), tag_retry_count=int(source_config.get('tag_retry_count', global_config.get('tag_retry_count', 2))), tag_retry_delay=int(source_config.get('tag_retry_delay', global_config.get('tag_retry_delay', 5))), tag_resolve_limit=int(source_config.get('tag_resolve_limit', global_config.get('tag_resolve_limit', 100))), docker_platform_filter_enabled=bool(source_config.get('docker_platform_filter_enabled', global_config.get('docker_platform_filter_enabled', True))), docker_platform_os=str(source_config.get('docker_platform_os', global_config.get('docker_platform_os', 'linux'))), docker_platform_arch=str(source_config.get('docker_platform_arch', global_config.get('docker_platform_arch', 'amd64'))), docker_platform_candidate_tags=int(source_config.get('docker_platform_candidate_tags', global_config.get('docker_platform_candidate_tags', 20))), docker_images_per_repository=docker_images_per_repository_limit( source_config.get('docker_images_per_repository', global_config.get('docker_images_per_repository', 1)) ), docker_content_scan_mode=str(source_config.get( 'docker_content_scan_mode', global_config.get('docker_content_scan_mode', 'full'), )).strip().lower(), docker_layer_canary_basis_points=int(source_config.get( 'docker_layer_canary_basis_points', global_config.get('docker_layer_canary_basis_points', 0), )), docker_adaptive_canary_basis_points=int(source_config.get( 'docker_adaptive_canary_basis_points', global_config.get('docker_adaptive_canary_basis_points', 0), )), docker_adaptive_gate_max_age_sec=int(source_config.get( 'docker_adaptive_gate_max_age_sec', global_config.get('docker_adaptive_gate_max_age_sec', 604800), )), docker_layer_config_max_bytes=int(source_config.get( 'docker_layer_config_max_bytes', global_config.get('docker_layer_config_max_bytes', 1 << 20), )), docker_layer_max_bytes=int(source_config.get( 'docker_layer_max_bytes', global_config.get('docker_layer_max_bytes', 256 << 20), )), docker_layer_image_max_bytes=int(source_config.get( 'docker_layer_image_max_bytes', global_config.get('docker_layer_image_max_bytes', 1 << 30), )), docker_layer_max_layers=int(source_config.get( 'docker_layer_max_layers', global_config.get('docker_layer_max_layers', 8), )), docker_layer_archive_max_size_bytes=int(source_config.get( 'docker_layer_archive_max_size_bytes', global_config.get('docker_layer_archive_max_size_bytes', 256 << 20), )), docker_layer_archive_max_depth=int(source_config.get( 'docker_layer_archive_max_depth', global_config.get('docker_layer_archive_max_depth', 4), )), docker_layer_archive_timeout_sec=int(source_config.get( 'docker_layer_archive_timeout_sec', global_config.get('docker_layer_archive_timeout_sec', 30), )), docker_layer_blob_timeout_sec=int(source_config.get( 'docker_layer_blob_timeout_sec', global_config.get('docker_layer_blob_timeout_sec', 600), )), docker_layer_filesystem_concurrency=int(source_config.get( 'docker_layer_filesystem_concurrency', global_config.get('docker_layer_filesystem_concurrency', 2), )), docker_layer_blob_max_attempts=int(source_config.get( 'docker_layer_blob_max_attempts', global_config.get('docker_layer_blob_max_attempts', 3), )), docker_layer_blob_lease_sec=int(source_config.get( 'docker_layer_blob_lease_sec', global_config.get('docker_layer_blob_lease_sec', 1800), )), docker_layer_min_free_bytes=int(source_config.get( 'docker_layer_min_free_bytes', global_config.get('docker_layer_min_free_bytes', 20 << 30), )), docker_layer_checkpoint_delay_sec=int(source_config.get( 'docker_layer_checkpoint_delay_sec', global_config.get('docker_layer_checkpoint_delay_sec', 60), )), docker_adaptive_checkpoint_max_blobs=int(source_config.get( 'docker_adaptive_checkpoint_max_blobs', global_config.get('docker_adaptive_checkpoint_max_blobs', 4), )), docker_adaptive_checkpoint_max_bytes=int(source_config.get( 'docker_adaptive_checkpoint_max_bytes', global_config.get('docker_adaptive_checkpoint_max_bytes', 512 << 20), )), docker_repository_refresh_interval_sec=max(0, int( source_config.get( 'docker_repository_refresh_interval_sec', global_config.get('docker_repository_refresh_interval_sec', 86400), ) or 0 )), docker_repository_refresh_max_per_cycle=max(0, min(1, int( source_config.get( 'docker_repository_refresh_max_per_cycle', global_config.get('docker_repository_refresh_max_per_cycle', 0), ) or 0 ))), archive_hours_back=int(source_config.get('archive_hours_back', global_config.get('archive_hours_back', 6))), archive_max_repos_per_cycle=int(source_config.get('archive_max_repos_per_cycle', global_config.get('archive_max_repos_per_cycle', 200))), archive_max_files_per_cycle=int(source_config.get('archive_max_files_per_cycle', global_config.get('archive_max_files_per_cycle', 300))), archive_max_commit_lookups=int(source_config.get('archive_max_commit_lookups', global_config.get('archive_max_commit_lookups', 200))), archive_rescan_cooldown_hours=int(source_config.get('archive_rescan_cooldown_hours', global_config.get('archive_rescan_cooldown_hours', 48))), archive_event_types=normalize_archive_event_types(source_config.get('archive_event_types', global_config.get('archive_event_types'))), ci_seed_sources=source_config.get('ci_seed_sources', global_config.get('ci_seed_sources', 'github,gitlab,package_git')), ci_use_finding_seeds=bool(source_config.get('ci_use_finding_seeds', global_config.get('ci_use_finding_seeds', False))), ci_max_repos_per_cycle=int(source_config.get('ci_max_repos_per_cycle', global_config.get('ci_max_repos_per_cycle', 50))), ci_seed_scan_limit=int(source_config.get('ci_seed_scan_limit', global_config.get('ci_seed_scan_limit', 5000))), ci_seed_query_batch_size=int(source_config.get('ci_seed_query_batch_size', global_config.get('ci_seed_query_batch_size', 250))), ci_soft_cooldown_days=int(source_config.get('ci_soft_cooldown_days', global_config.get('ci_soft_cooldown_days', 7))), ci_runs_per_repo=int(source_config.get('ci_runs_per_repo', global_config.get('ci_runs_per_repo', 5))), ci_pipelines_per_project=int(source_config.get('ci_pipelines_per_project', global_config.get('ci_pipelines_per_project', 5))), ci_jobs_per_pipeline=int(source_config.get('ci_jobs_per_pipeline', global_config.get('ci_jobs_per_pipeline', 20))), ci_lookback_days=int(source_config.get('ci_lookback_days', global_config.get('ci_lookback_days', 30))), ci_max_log_archive_mb=int(source_config.get('ci_max_log_archive_mb', global_config.get('ci_max_log_archive_mb', 50))), ci_max_log_file_mb=int(source_config.get('ci_max_log_file_mb', global_config.get('ci_max_log_file_mb', 20))), ci_max_trace_mb=int(source_config.get('ci_max_trace_mb', global_config.get('ci_max_trace_mb', 20))), ci_failed_first=bool(source_config.get('ci_failed_first', global_config.get('ci_failed_first', True))), ci_scan_artifacts=bool(source_config.get('ci_scan_artifacts', global_config.get('ci_scan_artifacts', False))), ci_max_artifacts_per_run=int(source_config.get('ci_max_artifacts_per_run', global_config.get('ci_max_artifacts_per_run', 3))), ci_max_artifacts_per_pipeline=int(source_config.get('ci_max_artifacts_per_pipeline', global_config.get('ci_max_artifacts_per_pipeline', 5))), ci_max_artifact_archive_mb=int(source_config.get('ci_max_artifact_archive_mb', global_config.get('ci_max_artifact_archive_mb', 50))), ci_max_artifact_file_mb=int(source_config.get('ci_max_artifact_file_mb', global_config.get('ci_max_artifact_file_mb', 10))), ci_max_artifact_files=int(source_config.get('ci_max_artifact_files', global_config.get('ci_max_artifact_files', 1000))), ci_target_max_download_mb=int(source_config.get('ci_target_max_download_mb', global_config.get('ci_target_max_download_mb', 500))), stop_on_seen_pages=bool(source_config.get('stop_on_seen_pages', global_config.get('stop_on_seen_pages', False))), seen_page_threshold=int(source_config.get('seen_page_threshold', global_config.get('seen_page_threshold', 2))), min_pages_before_stop=int(source_config.get('min_pages_before_stop', global_config.get('min_pages_before_stop', 1))), cleanup_temp_age_min=int(global_config.get('cleanup_temp_age_min', 120)), cleanup_only=False, ) def apply_global_config(global_config): scan_config.runtime_dir = global_config.get('runtime_dir', getattr(scan_config, 'runtime_dir', None)) scan_config.work_dir = global_config.get('work_dir', scan_config.work_dir) scan_config.results_dir = global_config.get('results_dir', scan_config.results_dir) scan_config.result_bundle_dir = global_config.get('result_bundle_dir', scan_config.result_bundle_dir) scan_config.result_bundle_max_event_bytes = int(global_config.get( 'result_bundle_max_event_bytes', scan_config.result_bundle_max_event_bytes, )) scan_config.result_spool_dir = global_config.get('result_spool_dir', scan_config.result_spool_dir) scan_config.result_spool_max_event_bytes = int(global_config.get('result_spool_max_event_bytes', scan_config.result_spool_max_event_bytes)) scan_config.result_spool_max_events = int(global_config.get('result_spool_max_events', scan_config.result_spool_max_events)) scan_config.result_spool_max_total_bytes = int(global_config.get('result_spool_max_total_bytes', scan_config.result_spool_max_total_bytes)) scan_config.result_spool_min_free_bytes = int(global_config.get('result_spool_min_free_bytes', scan_config.result_spool_min_free_bytes)) scan_config.scan_outbox_max_pending_items = int(global_config.get('scan_outbox_max_pending_items', scan_config.scan_outbox_max_pending_items)) scan_config.scan_outbox_max_pending_bytes = int(global_config.get('scan_outbox_max_pending_bytes', scan_config.scan_outbox_max_pending_bytes)) scan_config.scan_outbox_max_pending_age_sec = int(global_config.get('scan_outbox_max_pending_age_sec', scan_config.scan_outbox_max_pending_age_sec)) scan_config.queue_dir = global_config.get('queue_dir', getattr(scan_config, 'queue_dir', scan_config.results_dir)) scan_config.keycheck_dir = global_config.get('keycheck_dir', getattr(scan_config, 'keycheck_dir', None)) scan_config.postman_cache_dir = global_config.get('postman_cache_dir', getattr(scan_config, 'postman_cache_dir', None)) scan_config.postman_cache_max_items = int(global_config.get('postman_cache_max_items', getattr(scan_config, 'postman_cache_max_items', 100000))) scan_config.postman_cache_max_bytes = int(global_config.get('postman_cache_max_bytes', getattr(scan_config, 'postman_cache_max_bytes', 20 * 1024 * 1024 * 1024))) scan_config.postman_cache_min_free_bytes = int(global_config.get('postman_cache_min_free_bytes', getattr(scan_config, 'postman_cache_min_free_bytes', 5 * 1024 * 1024 * 1024))) scan_config.postman_cache_lock_timeout_sec = int(global_config.get('postman_cache_lock_timeout_sec', getattr(scan_config, 'postman_cache_lock_timeout_sec', 30))) scan_config.postman_discovery_max_artifacts_per_cycle = int(global_config.get('postman_discovery_max_artifacts_per_cycle', getattr(scan_config, 'postman_discovery_max_artifacts_per_cycle', 1000))) scan_config.postman_discovery_max_artifacts_per_page = int(global_config.get('postman_discovery_max_artifacts_per_page', getattr(scan_config, 'postman_discovery_max_artifacts_per_page', 100))) scan_config.postman_discovery_max_bytes_per_cycle = int(global_config.get('postman_discovery_max_bytes_per_cycle', getattr(scan_config, 'postman_discovery_max_bytes_per_cycle', 1024 * 1024 * 1024))) scan_config.postman_discovery_max_elapsed_sec = float(global_config.get('postman_discovery_max_elapsed_sec', getattr(scan_config, 'postman_discovery_max_elapsed_sec', 300.0))) scan_config.postman_package_harvest_max_artifacts = int(global_config.get('postman_package_harvest_max_artifacts', getattr(scan_config, 'postman_package_harvest_max_artifacts', 100))) scan_config.postman_package_harvest_max_bytes = int(global_config.get('postman_package_harvest_max_bytes', getattr(scan_config, 'postman_package_harvest_max_bytes', 128 * 1024 * 1024))) scan_config.postman_package_harvest_max_elapsed_sec = float(global_config.get('postman_package_harvest_max_elapsed_sec', getattr(scan_config, 'postman_package_harvest_max_elapsed_sec', 30.0))) scan_config.postman_context_max_input_bytes = int(global_config.get('postman_context_max_input_bytes', getattr(scan_config, 'postman_context_max_input_bytes', 16 * 1024 * 1024))) scan_config.postman_context_max_nodes = int(global_config.get('postman_context_max_nodes', getattr(scan_config, 'postman_context_max_nodes', 100000))) scan_config.postman_context_max_depth = int(global_config.get('postman_context_max_depth', getattr(scan_config, 'postman_context_max_depth', 64))) scan_config.postman_context_max_scalar_bytes = int(global_config.get('postman_context_max_scalar_bytes', getattr(scan_config, 'postman_context_max_scalar_bytes', 16 * 1024 * 1024))) scan_config.postman_context_max_items = int(global_config.get('postman_context_max_items', getattr(scan_config, 'postman_context_max_items', 50000))) scan_config.context_enrichment_max_source_bytes = int(global_config.get('context_enrichment_max_source_bytes', getattr(scan_config, 'context_enrichment_max_source_bytes', 16 * 1024 * 1024))) scan_config.context_enrichment_max_findings = int(global_config.get('context_enrichment_max_findings', getattr(scan_config, 'context_enrichment_max_findings', 2000))) scan_config.context_enrichment_max_postman_comparisons = int(global_config.get('context_enrichment_max_postman_comparisons', getattr(scan_config, 'context_enrichment_max_postman_comparisons', 200000))) scan_config.context_enrichment_max_elapsed_sec = float(global_config.get('context_enrichment_max_elapsed_sec', getattr(scan_config, 'context_enrichment_max_elapsed_sec', 5.0))) scan_config.trufflehog_diagnostic_max_lines = int(global_config.get('trufflehog_diagnostic_max_lines', getattr(scan_config, 'trufflehog_diagnostic_max_lines', 2000))) scan_config.trufflehog_diagnostic_max_line_chars = int(global_config.get('trufflehog_diagnostic_max_line_chars', getattr(scan_config, 'trufflehog_diagnostic_max_line_chars', 8192))) scan_config.trufflehog_diagnostic_max_line_bytes = int(global_config.get('trufflehog_diagnostic_max_line_bytes', getattr(scan_config, 'trufflehog_diagnostic_max_line_bytes', 8192))) scan_config.trufflehog_diagnostic_max_errors = int(global_config.get('trufflehog_diagnostic_max_errors', getattr(scan_config, 'trufflehog_diagnostic_max_errors', 200))) scan_config.trufflehog_diagnostic_max_warnings = int(global_config.get('trufflehog_diagnostic_max_warnings', getattr(scan_config, 'trufflehog_diagnostic_max_warnings', 200))) scan_config.trufflehog_diagnostic_max_unclassified = int(global_config.get('trufflehog_diagnostic_max_unclassified', getattr(scan_config, 'trufflehog_diagnostic_max_unclassified', 20))) scan_config.keycheck_input_max_line_bytes = int(global_config.get('keycheck_input_max_line_bytes', getattr(scan_config, 'keycheck_input_max_line_bytes', 16 * 1024 * 1024))) scan_config.keycheck_candidate_artifact_max_items = int(global_config.get('keycheck_candidate_artifact_max_items', getattr(scan_config, 'keycheck_candidate_artifact_max_items', 2000))) scan_config.keycheck_candidate_artifact_max_bytes = int(global_config.get('keycheck_candidate_artifact_max_bytes', getattr(scan_config, 'keycheck_candidate_artifact_max_bytes', 2 * 1024 * 1024))) scan_config.keycheck_candidate_file_max_items = int(global_config.get('keycheck_candidate_file_max_items', getattr(scan_config, 'keycheck_candidate_file_max_items', 100000))) scan_config.keycheck_candidate_file_max_bytes = int(global_config.get('keycheck_candidate_file_max_bytes', getattr(scan_config, 'keycheck_candidate_file_max_bytes', 32 * 1024 * 1024))) scan_config.keycheck_candidate_line_max_bytes = int(global_config.get('keycheck_candidate_line_max_bytes', getattr(scan_config, 'keycheck_candidate_line_max_bytes', 8192))) scan_config.gharchive_cache_dir = global_config.get('gharchive_cache_dir', getattr(scan_config, 'gharchive_cache_dir', None)) scan_config.gharchive_cache_max_items = int(global_config.get('gharchive_cache_max_items', getattr(scan_config, 'gharchive_cache_max_items', 48))) scan_config.gharchive_cache_max_bytes = int(global_config.get('gharchive_cache_max_bytes', getattr(scan_config, 'gharchive_cache_max_bytes', 8 * 1024 * 1024 * 1024))) scan_config.gharchive_cache_min_free_bytes = int(global_config.get('gharchive_cache_min_free_bytes', getattr(scan_config, 'gharchive_cache_min_free_bytes', 5 * 1024 * 1024 * 1024))) scan_config.gharchive_download_max_bytes = int(global_config.get('gharchive_download_max_bytes', getattr(scan_config, 'gharchive_download_max_bytes', 512 * 1024 * 1024))) scan_config.gharchive_decompressed_max_bytes = int(global_config.get('gharchive_decompressed_max_bytes', getattr(scan_config, 'gharchive_decompressed_max_bytes', 8 * 1024 * 1024 * 1024))) scan_config.gharchive_max_events = int(global_config.get('gharchive_max_events', getattr(scan_config, 'gharchive_max_events', 5000000))) scan_config.gharchive_max_line_bytes = int(global_config.get('gharchive_max_line_bytes', getattr(scan_config, 'gharchive_max_line_bytes', 8 * 1024 * 1024))) scan_config.gharchive_cache_lock_timeout_sec = int(global_config.get('gharchive_cache_lock_timeout_sec', getattr(scan_config, 'gharchive_cache_lock_timeout_sec', 600))) scan_config.proxy_file = global_config.get('proxy_file', getattr(scan_config, 'proxy_file', None)) scan_config.api_proxy_enabled = bool_config(global_config.get('api_proxy_enabled'), getattr(scan_config, 'api_proxy_enabled', False)) scan_config.api_proxy_file = resolve_optional_path( global_config.get('api_proxy_file') or scan_config.proxy_file, global_config, ) scan_config.api_proxy_timeout = int(global_config.get('api_proxy_timeout', getattr(scan_config, 'api_proxy_timeout', 5))) scan_config.api_proxy_max_retries = int(global_config.get('api_proxy_max_retries', getattr(scan_config, 'api_proxy_max_retries', 100))) scan_config.api_proxy_retry_delay = int(global_config.get('api_proxy_retry_delay', getattr(scan_config, 'api_proxy_retry_delay', 5))) scan_config.download_proxy_enabled = bool_config(global_config.get('download_proxy_enabled'), getattr(scan_config, 'download_proxy_enabled', False)) scan_config.download_proxy_file = resolve_optional_path( global_config.get('download_proxy_file') or getattr(scan_config, 'download_proxy_file', ''), global_config, ) if (global_config.get('download_proxy_file') or getattr(scan_config, 'download_proxy_file', '')) else '' scan_config.max_active_scans = int(global_config.get('max_active_scans', getattr(scan_config, 'max_active_scans', 0))) scan_config.opportunistic_scan_slots = max(0, min(1, int(global_config.get( 'opportunistic_scan_slots', getattr(scan_config, 'opportunistic_scan_slots', 0), )))) scan_config.opportunistic_scan_sources = csv_items(global_config.get( 'opportunistic_scan_sources', getattr(scan_config, 'opportunistic_scan_sources', []), )) scan_config.opportunistic_scan_reserve_overhead_bytes = max(0, int(global_config.get( 'opportunistic_scan_reserve_overhead_bytes', getattr(scan_config, 'opportunistic_scan_reserve_overhead_bytes', 1024 * 1024 * 1024), ))) scan_config.opportunistic_scan_min_available_after_reserve_bytes = max(0, int(global_config.get( 'opportunistic_scan_min_available_after_reserve_bytes', getattr(scan_config, 'opportunistic_scan_min_available_after_reserve_bytes', 4 * 1024 * 1024 * 1024), ))) scan_config.opportunistic_scan_min_commit_after_reserve_bytes = max(0, int(global_config.get( 'opportunistic_scan_min_commit_after_reserve_bytes', getattr(scan_config, 'opportunistic_scan_min_commit_after_reserve_bytes', 6 * 1024 * 1024 * 1024), ))) scan_config.scan_limiter_db = resolve_optional_path( global_config.get('scan_limiter_db') or getattr(scan_config, 'scan_limiter_db', ''), global_config, ) scan_config.scan_slot_wait_sec = float(global_config.get('scan_slot_wait_sec', getattr(scan_config, 'scan_slot_wait_sec', 0.5))) scan_config.scan_slot_wait_log_sec = int(global_config.get('scan_slot_wait_log_sec', getattr(scan_config, 'scan_slot_wait_log_sec', 30))) scan_config.scan_slot_stale_sec = int(global_config.get('scan_slot_stale_sec', getattr(scan_config, 'scan_slot_stale_sec', 7200))) scan_config.low_space_cleanup_max_items = int(global_config.get('low_space_cleanup_max_items', getattr(scan_config, 'low_space_cleanup_max_items', 50))) scan_config.drop_detectors = csv_items(global_config.get('drop_detectors', getattr(scan_config, 'drop_detectors', []))) scan_config.jsonl_rotation_enabled = bool_config(global_config.get('jsonl_rotation_enabled'), getattr(scan_config, 'jsonl_rotation_enabled', False)) scan_config.found_secrets_max_mb = int(global_config.get('found_secrets_max_mb', getattr(scan_config, 'found_secrets_max_mb', 512))) scan_config.scan_results_max_mb = int(global_config.get('scan_results_max_mb', getattr(scan_config, 'scan_results_max_mb', 1024))) scan_config.scan_errors_max_mb = int(global_config.get('scan_errors_max_mb', getattr(scan_config, 'scan_errors_max_mb', 64))) scan_config.scan_errors_keep = int(global_config.get('scan_errors_keep', getattr(scan_config, 'scan_errors_keep', 5))) scan_config.jsonl_lock_stale_sec = int(global_config.get('jsonl_lock_stale_sec', getattr(scan_config, 'jsonl_lock_stale_sec', 300))) scan_config.jsonl_max_segments = int(global_config.get('jsonl_max_segments', getattr(scan_config, 'jsonl_max_segments', 16))) scan_config.jsonl_ledger_max_rows = int(global_config.get('jsonl_ledger_max_rows', getattr(scan_config, 'jsonl_ledger_max_rows', 1000000))) scan_config.jsonl_ledger_max_bytes = int(global_config.get('jsonl_ledger_max_bytes', getattr(scan_config, 'jsonl_ledger_max_bytes', 512 * 1024 * 1024))) scan_config.jsonl_legacy_index_max_bytes = int(global_config.get('jsonl_legacy_index_max_bytes', getattr(scan_config, 'jsonl_legacy_index_max_bytes', 16 * 1024 * 1024))) scan_config.jsonl_tail_scan_max_bytes = int(global_config.get('jsonl_tail_scan_max_bytes', getattr(scan_config, 'jsonl_tail_scan_max_bytes', 8 * 1024 * 1024))) scan_config.jsonl_torn_quarantine_max_bytes = int(global_config.get('jsonl_torn_quarantine_max_bytes', getattr(scan_config, 'jsonl_torn_quarantine_max_bytes', 64 * 1024))) scan_config.dockerhub_tag_cache_path = resolve_optional_path( global_config.get('dockerhub_tag_cache_path') or getattr(scan_config, 'dockerhub_tag_cache_path', ''), global_config, ) if (global_config.get('dockerhub_tag_cache_path') or getattr(scan_config, 'dockerhub_tag_cache_path', '')) else '' scan_config.dockerhub_tag_cache_ttl_sec = int(global_config.get('dockerhub_tag_cache_ttl_sec', getattr(scan_config, 'dockerhub_tag_cache_ttl_sec', 21600))) scan_config.dockerhub_tag_negative_cache_ttl_sec = int(global_config.get('dockerhub_tag_negative_cache_ttl_sec', getattr(scan_config, 'dockerhub_tag_negative_cache_ttl_sec', 3600))) scan_config.dockerhub_tag_rate_limit_cache_ttl_sec = int(global_config.get('dockerhub_tag_rate_limit_cache_ttl_sec', getattr(scan_config, 'dockerhub_tag_rate_limit_cache_ttl_sec', 1800))) scan_config.dockerhub_tag_cache_max_rows = int(global_config.get('dockerhub_tag_cache_max_rows', getattr(scan_config, 'dockerhub_tag_cache_max_rows', 50000))) scan_config.dockerhub_tag_cache_max_age_sec = int(global_config.get('dockerhub_tag_cache_max_age_sec', getattr(scan_config, 'dockerhub_tag_cache_max_age_sec', 7 * 86400))) scan_config.dockerhub_tag_cache_max_bytes = int(global_config.get('dockerhub_tag_cache_max_bytes', getattr(scan_config, 'dockerhub_tag_cache_max_bytes', 256 * 1024 * 1024))) scan_config.dockerhub_tag_cache_min_free_bytes = int(global_config.get('dockerhub_tag_cache_min_free_bytes', getattr(scan_config, 'dockerhub_tag_cache_min_free_bytes', 512 * 1024 * 1024))) scan_config.trufflehog_path = global_config.get('trufflehog_path', scan_config.trufflehog_path) scan_config.trufflehog_config = resolve_optional_path( global_config.get('trufflehog_config') or getattr(scan_config, 'trufflehog_config', ''), global_config, ) if (global_config.get('trufflehog_config') or getattr(scan_config, 'trufflehog_config', '')) else '' scan_config.trufflehog_job_memory_limit_bytes = int(global_config.get( 'trufflehog_job_memory_limit_bytes', getattr(scan_config, 'trufflehog_job_memory_limit_bytes', 4096 * 1024 * 1024), )) scan_config.trufflehog_windows_job_cpu_weight = int(global_config.get( 'trufflehog_windows_job_cpu_weight', getattr(scan_config, 'trufflehog_windows_job_cpu_weight', 0), )) scan_config.trufflehog_windows_memory_priority = int(global_config.get( 'trufflehog_windows_memory_priority', getattr(scan_config, 'trufflehog_windows_memory_priority', 0), )) scan_config.trufflehog_stdout_max_mb = int(global_config.get( 'trufflehog_stdout_max_mb', getattr(scan_config, 'trufflehog_stdout_max_mb', 32), )) scan_config.trufflehog_stderr_max_mb = int(global_config.get( 'trufflehog_stderr_max_mb', getattr(scan_config, 'trufflehog_stderr_max_mb', 8), )) if 'min_free_gb' in global_config: scan_config.min_free_gb = float(global_config['min_free_gb']) if 'no_verification' in global_config: scan_config.no_verification = bool(global_config['no_verification']) if 'strict_git_provider_token_filter' in global_config: scan_config.strict_git_provider_token_filter = bool(global_config['strict_git_provider_token_filter']) def configure_source_auth(source_name, source_config, state=None, secrets=None, auth_entry=None): if source_to_platform(source_name) != 'docker': return if secrets is not None and state is not None and source_config.get('auth_pool'): entries = available_auth_entries(source_name, source_config, state, secrets) accounts = [ { 'name': entry.get('name'), 'username': entry.get('username'), 'token': entry.get('token'), } for entry in entries if entry.get('username') and entry.get('token') ] configure_docker_discovery_accounts( accounts, cooldown_sec=int(source_config.get('rate_limit_cooldown', 1800)), ) endpoint_status = state['sources'].setdefault( source_name, default_source_state(), ).get('auth_endpoint_status') if endpoint_status: restore_docker_endpoint_cooldowns(endpoint_status) print(f"Docker auth pool: configured {len(accounts)} available account(s)") return docker_token = ( (auth_entry or {}).get('token') or source_config.get('docker_token') or source_config.get('token') or os.getenv('DOCKERHUB_TOKEN') or os.getenv('DOCKER_TOKEN') or os.getenv('DOCKER_TOKENS') ) docker_username = (auth_entry or {}).get('username') or source_config.get('docker_username') or os.getenv('DOCKERHUB_USERNAME') or os.getenv('DOCKER_USERNAME') configure_docker_discovery_tokens(docker_token, docker_username) def persist_docker_auth_events(source_name, source_config, state, secrets): if source_to_platform(source_name) != 'docker': return entries, _ = auth_pool_entries(source_config, secrets) entries_by_name = {str(entry.get('name')): entry for entry in entries} cooldown = int(source_config.get('rate_limit_cooldown', 1800)) source_state = state['sources'].setdefault(source_name, default_source_state()) endpoint_status = source_state.setdefault('auth_endpoint_status', {}) for event in drain_docker_auth_events(): entry = entries_by_name.get(str(event.get('name') or '')) if not entry: continue category = str(event.get('category') or '') endpoint = str(event.get('endpoint') or '') if endpoint == 'hub_search' and category != 'auth_invalid': account_status = endpoint_status.setdefault(endpoint, {}).setdefault( str(entry.get('name')), {}, ) now = utc_now() if category == 'ok': account_status['disabled_until'] = None account_status['disabled_reason'] = None account_status['status'] = 'ok' account_status['last_success_at'] = now.isoformat(timespec='seconds') account_status['success_count'] = int( account_status.get('success_count', 0) ) + 1 else: account_status['disabled_until'] = event.get('reset_at') or ( now + timedelta(seconds=cooldown) ).isoformat(timespec='seconds') account_status['disabled_reason'] = category or 'rate_limit' account_status['status'] = 'limited' account_status['last_error'] = str(event.get('message') or '')[:500] account_status['last_error_at'] = now.isoformat(timespec='seconds') account_status['failures'] = int(account_status.get('failures', 0)) + 1 continue if category == 'ok': clear_auth_rate_limit(source_name, entry, state) else: mark_auth_rate_limited( source_name, entry, state, reset_at=event.get('reset_at'), cooldown=cooldown, category=category or 'rate_limit', message=event.get('message') or 'Docker authentication unavailable', ) refresh_auth_summary(source_name, source_config, state, secrets) def run_configured_source( source_name, config, state, state_path, secrets, db=None, run_id=None, cycle_runner=None, ): selected_cycle_runner = run_cycle if cycle_runner is None else cycle_runner if not callable(selected_cycle_runner): raise TypeError('Configured source cycle runner must be callable') source_config = config['sources'][source_name] global_config = config.get('global', {}) query, query_index, query_count = current_query_for_source(source_name, source_config, state) source_state = state['sources'].setdefault(source_name, default_source_state()) print(f'\n=== {source_name} source cycle ===') print(f'Query: {query_index + 1}/{query_count} -> {query!r}') print(f"Mode: {source_config.get('mode', 'recent')}") source_state['last_query'] = query source_state['last_started_at'] = datetime.now().isoformat(timespec='seconds') source_state['last_status'] = 'running' save_state(state_path, state) auth_entry = select_auth_entry(source_name, source_config, state, secrets) auth_name = auth_entry.get('name') if auth_entry else 'none' refresh_auth_summary(source_name, source_config, state, secrets, auth_name) print(f"Auth: {auth_name} (rotation={source_config.get('auth_rotation', 'per_cycle')})") save_state(state_path, state) args = build_args_from_source_config(source_name, source_config, global_config, query, auth_entry) dockerhub_policies = {} dockerhub_annotation = None dockerhub_experiment = None if ( source_name == 'dockerhub' and args.platform == 'docker' and args.mode == 'search' ): validated_depth = validate_docker_depth_config( config, managed_postgres=bool( db and getattr(db, 'conn', None) and getattr(db.conn, 'is_postgres', False) ), final_cutover=( source_config.get( 'sync_file_queues', global_config.get('sync_file_queues', True), ) is False ), ) dockerhub_experiment = validated_depth.experiment dockerhub_policies = configured_dockerhub_discovery_policies(source_config) annotate_dockerhub_runtime_args( args, dockerhub_experiment, dockerhub_policies, ) if dockerhub_experiment is not None: policy_hashes = { policy['policy_sha256'] for policy in dockerhub_policies.values() } generation_check = getattr( db, 'dockerhub_discovery_generation_complete', None, ) if not callable(generation_check): raise RuntimeError( 'DockerHub collection generation authority is unavailable' ) generation_complete = generation_check( 'dockerhub', dockerhub_experiment.collection_generation, dockerhub_experiment.ordered_query_hash, len(dockerhub_experiment.queries), next(iter(policy_hashes)), ) else: generation_complete = True dockerhub_annotation = prepare_dockerhub_discovery_state( source_state, dockerhub_policies, query, force_deep=not generation_complete, ) annotate_dockerhub_discovery_args(args, dockerhub_annotation) mark_dockerhub_discovery_incomplete( source_state, query, dockerhub_annotation['policy_sha256'], ) save_state(state_path, state) scan_config.trufflehog_job_memory_limit_bytes = int(getattr( args, 'trufflehog_job_memory_limit_bytes', scan_config.trufflehog_job_memory_limit_bytes, )) if source_to_platform(source_name) == 'postman': source_state.setdefault('postman_auth_status', {}) args.github_tokens = github_token_entries_for_postman(source_config, secrets) args.github_auth_status = source_state['postman_auth_status'] if args.github_tokens: print(f"Postman GitHub auth pool: configured {len(args.github_tokens)} token(s)") configure_source_auth(source_name, source_config, state, secrets, auth_entry) cycle_id = None if db and run_id: cycle_id = db.start_source_cycle( run_id, source_name, args.platform, args.mode, query, query_index + 1, query_count, auth_name, source_config, queue_counts_for_args(args), ) if getattr(db.conn, 'is_postgres', False) and cycle_id is None: raise RuntimeError(f'Unable to create Postgres source cycle for {source_name}; refusing to claim targets') annotate_dockerhub_runtime_args( args, dockerhub_experiment, dockerhub_policies, cycle_id, ) all_auth_cooling_down = False final_status = 'completed' cycle_metrics = {} attempt_finished = False def failure_metrics(**extra): metrics = {'fetched_count': 0} metrics.update(cycle_metrics) metrics.update(extra) return metrics while True: try: attempt_finished = False cycle_metrics = selected_cycle_runner( args, db, run_id, cycle_id, source_name, partial_metrics=cycle_metrics, ) persist_docker_auth_events(source_name, source_config, state, secrets) attempt_finished = True returned_status = str(cycle_metrics.get('cycle_status') or '') if returned_status in ( 'completed', 'completed_with_retries', 'query_invalid', 'failed', 'source_failed', 'backlog_only', 'paused', ): final_status = returned_status elif cycle_metrics.get('backlog_only'): final_status = 'backlog_only' else: final_status = 'completed' scanned = int(cycle_metrics.get('staged_count') or cycle_metrics.get('scanned_count', 0)) if cycle_metrics.get('source_failure_count'): failure_category = cycle_metrics.get('source_failure_category') or 'source_failed' if failure_category == 'source_auth': failure_category = 'auth_invalid' raise RateLimitError( source_name, cycle_metrics.get('source_failure_message') or 'source-wide scanner failure', category=failure_category, auth_related=bool(cycle_metrics.get('source_failure_auth_related')), ) if source_to_platform(source_name) != 'docker' or not source_config.get('auth_pool'): clear_auth_rate_limit(source_name, auth_entry, state) refresh_auth_summary(source_name, source_config, state, secrets, auth_name) break except RateLimitError as e: persist_docker_auth_events(source_name, source_config, state, secrets) cooldown = int(source_config.get('rate_limit_cooldown', 3600)) category = getattr(e, 'category', 'rate_limit') print(f"{source_name} API {category}: {e}") if not getattr(e, 'auth_related', True): scanned = 0 if final_status != 'source_failed': final_status = 'query_invalid' if category == 'query_invalid' else 'failed' if db and cycle_id and not attempt_finished: db.finish_source_cycle(cycle_id, final_status, failure_metrics(), queue_counts_for_args(args), str(e)) break mark_auth_rate_limited(source_name, auth_entry, state, e.reset_at, cooldown, category, str(e)) refresh_auth_summary(source_name, source_config, state, secrets, auth_name) save_state(state_path, state) if not source_config.get('retry_with_next_auth_on_rate_limit', True): if final_status != 'source_failed': final_status = 'auth_failed' if category.startswith('auth_') else 'rate_limited' scanned = 0 if db and cycle_id and not attempt_finished: db.finish_source_cycle(cycle_id, final_status, failure_metrics(), queue_counts_for_args(args), str(e)) break auth_entry = select_auth_entry(source_name, source_config, state, secrets) if not auth_entry: print(f"All {source_name} auth tokens are unavailable; skipping this source cycle.") scanned = 0 all_auth_cooling_down = True if final_status != 'source_failed': final_status = 'auth_failed' if category.startswith('auth_') else 'rate_limited' refresh_auth_summary(source_name, source_config, state, secrets, 'none') if db and cycle_id and not attempt_finished: db.finish_source_cycle(cycle_id, final_status, failure_metrics(), queue_counts_for_args(args), str(e)) break auth_name = auth_entry.get('name') refresh_auth_summary(source_name, source_config, state, secrets, auth_name) print(f"Retrying {source_name} with auth {auth_name}") args = build_args_from_source_config(source_name, source_config, global_config, query, auth_entry) annotate_dockerhub_runtime_args( args, dockerhub_experiment, dockerhub_policies, ) if dockerhub_annotation is not None: annotate_dockerhub_discovery_args(args, dockerhub_annotation) scan_config.trufflehog_job_memory_limit_bytes = int(getattr( args, 'trufflehog_job_memory_limit_bytes', scan_config.trufflehog_job_memory_limit_bytes, )) if source_to_platform(source_name) == 'postman': args.github_tokens = github_token_entries_for_postman(source_config, secrets) args.github_auth_status = source_state.setdefault('postman_auth_status', {}) configure_source_auth(source_name, source_config, state, secrets, auth_entry) if db and run_id: cycle_id = db.start_source_cycle( run_id, source_name, args.platform, args.mode, query, query_index + 1, query_count, auth_name, source_config, queue_counts_for_args(args), ) if getattr(db.conn, 'is_postgres', False) and cycle_id is None: raise RuntimeError(f'Unable to create Postgres auth-retry cycle for {source_name}') annotate_dockerhub_runtime_args( args, dockerhub_experiment, dockerhub_policies, cycle_id, ) cycle_metrics = {} except Exception as e: persist_docker_auth_events(source_name, source_config, state, secrets) if db and getattr(db, 'conn', None): try: db.conn.rollback() except Exception: pass discovery_transport_failure = ( isinstance(e, GitLabDiscoveryTransportError) and source_name == 'gitlab' ) or ( isinstance(e, DockerHubDiscoveryTransportError) and source_to_platform(source_name) == 'docker' ) if ( isinstance(e, DockerHubDiscoveryTransportError) and not getattr(e, 'retryable', True) ): if db and cycle_id and not attempt_finished: db.finish_source_cycle( cycle_id, 'failed', failure_metrics(), queue_counts_for_args(args), str(e), ) raise if discovery_transport_failure: scanned = 0 final_status = 'failed' cycle_metrics = failure_metrics(discovery_transport_failed=True) print( f'{source_name} API network: discovery cycle failed ' f'after bounded retries: {e}' ) if db and cycle_id and not attempt_finished: db.finish_source_cycle( cycle_id, final_status, cycle_metrics, queue_counts_for_args(args), str(e), ) break if db and cycle_id and not attempt_finished: db.finish_source_cycle(cycle_id, 'failed', failure_metrics(), queue_counts_for_args(args), str(e)) raise source_state['last_completed_at'] = datetime.now().isoformat(timespec='seconds') source_state['last_status'] = final_status source_state['cycles'] = int(source_state.get('cycles', 0)) + 1 source_state['last_scanned'] = scanned if selected_cycle_runner is run_discovery_cycle: source_state['last_cycle_result'] = { 'status': final_status, 'fetched_count': max(0, int(cycle_metrics.get('fetched_count', 0) or 0)), 'queued_new_count': max(0, int(cycle_metrics.get('queued_new_count', 0) or 0)), 'queued_updated_count': max(0, int(cycle_metrics.get('queued_updated_count', 0) or 0)), } if final_status in ('completed', 'completed_with_retries'): source_state['last_discovery_success_at'] = source_state['last_completed_at'] source_state['last_error_category'] = '' else: source_state['last_error_category'] = ( final_status if final_status in ( 'auth_failed', 'failed', 'paused', 'query_invalid', 'rate_limited', 'source_failed', ) else 'runtime_error' ) if ( dockerhub_annotation is not None and dockerhub_annotation['deep'] and final_status in ('completed', 'completed_with_retries') and cycle_metrics.get('deep_dispatch_durable') is True ): mark_dockerhub_deep_dispatched( source_state, query, dockerhub_annotation['policy_sha256'], ) if ( dockerhub_annotation is not None and final_status in ('completed', 'completed_with_retries', 'query_invalid') ): clear_dockerhub_discovery_incomplete(source_state, query) if final_status in ('completed', 'completed_with_retries', 'query_invalid'): advance_query_for_source(source_name, source_config, state) cycle_metrics['cycle_status'] = final_status save_state(state_path, state) if dockerhub_policies: try: retry_metrics = process_dockerhub_discovery_retry( args, db, source_name, dockerhub_policies, ) cycle_metrics.update(retry_metrics) except Exception: logger.warning('DockerHub discovery retry lane failed') try: persist_docker_auth_events(source_name, source_config, state, secrets) save_state(state_path, state) except Exception: logger.warning('DockerHub discovery retry auth-state persistence failed') next_query, next_index, next_count = current_query_for_source(source_name, source_config, state) print(f'{source_name} complete. Next query: {next_index + 1}/{next_count} -> {next_query!r}') return cycle_metrics def print_state(state_path, state): print(f'State file: {state_path}') print(json.dumps(state, indent=2, ensure_ascii=False)) def current_process_private_bytes(): if os.name != 'nt': return 0 try: counters = _PROCESS_MEMORY_COUNTERS_EX() counters.cb = ctypes.sizeof(counters) if not _GET_PROCESS_MEMORY_INFO( _METRIC_GET_CURRENT_PROCESS(), ctypes.byref(counters), counters.cb, ): return 0 return int(counters.PrivateUsage) except Exception: return 0 def run_discovery_producer_mode(args, config=None): if config is None: config = load_config(args.config, managed_postgres=True, final_cutover=True) if args.show_state or args.cleanup_only or not args.once: raise SystemExit('discovery-producer requires an authenticated one-shot source invocation') source_name = str(args.source or '') if source_name not in DISCOVERY_PRODUCER_SOURCES: raise SystemExit('discovery-producer source is outside the canonical allowlist') sources = config.get('sources') or {} if source_name not in sources: raise SystemExit(f'Source {source_name} is not present in config') global_config = config.get('global', {}) preflight_lifecycle_paths( args.config, config, authority_profile=DISCOVERY_PRODUCER_ROLE, ) apply_global_config(global_config) require_sensitive_runtime_paths(global_config, create=False) state_path = get_state_path(config, args.config) state = load_state(state_path, config) secrets = load_secrets(config, args.config) refresh_auth_summary(source_name, sources[source_name], state, secrets) save_state(state_path, state) db = ScannerDB( scan_config.results_dir, global_config.get('database_path') or global_config.get('db_path'), db_url=global_config.get('database_url'), initialize=False, ) if not db.enabled or not getattr(getattr(db, 'conn', None), 'is_postgres', False): db.close() raise RuntimeError('discovery-producer requires the managed PostgreSQL authority') run_id = None run_status = 'completed' run_error = None try: set_application_name = getattr(db, 'set_application_name', None) if set_application_name: set_application_name(f'truf-discovery:{source_name}') db.require_runtime_safety_schema() db.require_final_cutover() run_id = db.start_run( DISCOVERY_PRODUCER_ROLE, sys.argv, selected_source=source_name, selected_platform=source_to_platform(source_name), config_path=args.config, config_hash=hash_file(args.config), enabled_sources=[source_name], global_config=global_config, ) if run_id is None: raise RuntimeError('Unable to create PostgreSQL discovery-producer run') run_configured_source( source_name, config, state, state_path, secrets, db, run_id, cycle_runner=run_discovery_cycle, ) except BaseException as exc: run_status = 'failed' run_error = type(exc).__name__ source_state = state['sources'].setdefault(source_name, default_source_state()) source_state['last_status'] = 'failed' source_state['last_error_category'] = 'runtime_error' source_state['last_cycle_result'] = {'status': 'failed'} save_state(state_path, state) raise finally: if run_id is not None: db.finish_run(run_id, run_status, run_error) db.close() def run_config_mode(args, config=None, runtime_role='scanner'): if runtime_role == DISCOVERY_PRODUCER_ROLE: return run_discovery_producer_mode(args, config=config) if runtime_role != 'scanner': raise SystemExit('unsupported console runtime role') if config is None: config = load_config( args.config, managed_postgres=bool(os.getenv('TRUF_MANAGED_POSTGRES_DSN')), ) global_config = config.get('global', {}) preflight_lifecycle_paths(args.config, config) apply_global_config(global_config) state_path = get_state_path(config, args.config) state = load_state(state_path, config) if args.show_state: print_state(state_path, state) return initialize_scanner_runtime(preflight_complete=True) require_sensitive_runtime_paths(global_config, create=False) secrets = load_secrets(config, args.config) for source_name, source_config in (config.get('sources') or {}).items(): refresh_auth_summary(source_name, source_config, state, secrets) save_state(state_path, state) if args.cleanup_only: raise SystemExit('Source-side cleanup is retired; use the authenticated janitor worker') if not check_dependencies(): raise SystemExit(1) source_names = enabled_source_names(config, args.source) if not source_names: raise SystemExit('No enabled sources found in config') loop_enabled = bool(global_config.get('loop', True)) and not args.once cooldown = int(global_config.get('cooldown', args.cooldown)) db = ScannerDB( scan_config.results_dir, global_config.get('database_path') or global_config.get('db_path'), db_url=global_config.get('database_url'), initialize=bool(global_config.get('database_initialize', False)), ) if db.enabled: set_application_name = getattr(db, 'set_application_name', None) if set_application_name: set_application_name( f'truf-source:{source_names[0]}' if len(source_names) == 1 else 'truf-scanner' ) if getattr(db, 'postgres_required', False) and not db.enabled: raise RuntimeError('Configured Postgres database is unavailable; refusing file-queue fallback') if db.enabled and getattr(db.conn, 'is_postgres', False): db.require_runtime_safety_schema() run_id = db.start_run( 'config', sys.argv, selected_source=args.source, config_path=args.config, config_hash=hash_file(args.config), enabled_sources=source_names, global_config=global_config, ) if db.enabled else None if db.enabled and getattr(db.conn, 'is_postgres', False) and run_id is None: raise RuntimeError('Unable to create Postgres scanner run; refusing target claims') run_status = 'completed' run_error = None try: while True: print(f'\n=== Config cycle started at {datetime.now().isoformat(timespec="seconds")} ===') backlog_progress = False for source_name in source_names: if not config['sources'][source_name].get('enabled', False) and not args.source: continue cycle_metrics = run_configured_source( source_name, config, state, state_path, secrets, db, run_id, ) or {} backlog_progress = backlog_progress or bool( cycle_metrics.get('backlog_only') and (cycle_metrics.get('staged_count') or cycle_metrics.get('scanned_count')) ) if not loop_enabled: break sleep_seconds = ( max(0.1, float(global_config.get('backlog_poll_sec', 0.5) or 0.5)) if backlog_progress else cooldown ) print(f'Sleeping {sleep_seconds:g} seconds...') time.sleep(sleep_seconds) except KeyboardInterrupt: run_status = 'stopped' save_state(state_path, state) print('Stopped. State preserved. Temp cleanup deferred to next run or --cleanup-only.') except UnresolvedHandoffInfrastructureError as e: run_status = 'failed' run_error = str(e) raise SystemExit(SOURCE_INFRASTRUCTURE_HOLD_EXIT) from e except Exception as e: run_status = 'failed' run_error = str(e) raise finally: if db.enabled and run_id: db.finish_run(run_id, run_status, run_error) db.close() def parse_args(): default_paths = default_project_paths() parser = argparse.ArgumentParser(description='Console runner for TruffleHog scans.') parser.add_argument('--config') parser.add_argument('--source', choices=['github', 'github_archive', 'github_archive_files', 'github_gists', 'gitlab', 'docker', 'dockerhub', 'npm', 'pypi', 'package_git', 'huggingface', 'postman', 'github_actions', 'gitlab_ci']) parser.add_argument('--once', action='store_true') parser.add_argument('--show-state', action='store_true') parser.add_argument('--platform', choices=['github', 'github_archive', 'github_archive_files', 'github_gists', 'gitlab', 'docker', 'dockerhub', 'npm', 'pypi', 'package_git', 'huggingface', 'postman', 'github_actions', 'gitlab_ci']) parser.add_argument('--mode', choices=['recent', 'search', 'custom', 'archive'], default='recent') parser.add_argument('--query', default='') parser.add_argument('--query-file') parser.add_argument('--pages', type=int, default=1) parser.add_argument('--per-page', type=int, default=50) parser.add_argument('--target-file') parser.add_argument('--token') parser.add_argument('--docker-username') parser.add_argument('--docker-token') parser.add_argument('--workers', type=int, default=6) parser.add_argument('--timeout', type=int, default=None) parser.add_argument('--detectors', default=scan_config.detectors) parser.add_argument('--exclude-detectors', default=scan_config.exclude_detectors) parser.add_argument('--drop-detectors', default=','.join(getattr(scan_config, 'drop_detectors', []))) parser.add_argument('--no-verification', action='store_true', default=scan_config.no_verification) parser.add_argument('--save-dir', default=scan_config.results_dir) parser.add_argument('--runtime-dir', default=default_paths['runtime_dir']) parser.add_argument('--result-bundle-dir', default=default_paths['result_bundle_dir']) parser.add_argument('--result-bundle-max-event-bytes', type=int, default=scan_config.result_bundle_max_event_bytes) parser.add_argument('--result-spool-dir', default=default_paths['result_spool_dir']) parser.add_argument('--result-spool-max-event-bytes', type=int, default=scan_config.result_spool_max_event_bytes) parser.add_argument('--result-spool-max-events', type=int, default=scan_config.result_spool_max_events) parser.add_argument('--result-spool-max-total-bytes', type=int, default=scan_config.result_spool_max_total_bytes) parser.add_argument('--result-spool-min-free-bytes', type=int, default=scan_config.result_spool_min_free_bytes) parser.add_argument('--result-spool-wait-sec', type=float, default=1.0) parser.add_argument('--result-spool-diagnostic-interval-sec', type=float, default=30.0) parser.add_argument('--scan-outbox-max-pending-items', type=int, default=scan_config.scan_outbox_max_pending_items) parser.add_argument('--scan-outbox-max-pending-bytes', type=int, default=scan_config.scan_outbox_max_pending_bytes) parser.add_argument('--scan-outbox-max-pending-age-sec', type=int, default=scan_config.scan_outbox_max_pending_age_sec) parser.add_argument('--queue-dir', default=default_paths['queue_dir']) parser.add_argument('--work-dir', default=scan_config.work_dir) parser.add_argument('--trufflehog-path', default=scan_config.trufflehog_path) parser.add_argument('--trufflehog-config', default=getattr(scan_config, 'trufflehog_config', '')) parser.add_argument('--trufflehog-concurrency', type=int, default=0) parser.add_argument('--loop', action='store_true') parser.add_argument('--cooldown', type=int, default=300) parser.add_argument('--max-cycles', type=int, default=0) parser.add_argument('--max-targets', type=int, default=0) parser.add_argument('--target-retry-max-attempts', type=int, default=3) parser.add_argument('--target-retry-base-delay-sec', type=int, default=3600) parser.add_argument('--target-retry-max-delay-sec', type=int, default=86400) parser.add_argument('--target-timeout-retry-delay-sec', type=int, default=21600) parser.add_argument('--target-claim-batch-size', type=int, default=0) parser.add_argument('--recent-hours', type=int, default=24) parser.add_argument('--recent-days', type=int, default=7) parser.add_argument('--max-version-age-days', type=int, default=0) parser.add_argument('--versions-per-package', type=int, default=1) parser.add_argument('--package-sources', default='npm,pypi') parser.add_argument('--package-git-refresh-registry', action='store_true') parser.add_argument('--search-kinds', default='collection,environment') parser.add_argument('--gist-since', default='') parser.add_argument('--postman-cache-dir', default=default_paths.get('postman_cache_dir')) parser.add_argument('--gharchive-cache-dir', default=default_paths.get('gharchive_cache_dir')) parser.add_argument('--max-artifact-size-mb', type=int, default=50) parser.add_argument('--max-file-age-days', type=int, default=365) parser.add_argument('--github-code-search-rpm', type=int, default=8) parser.add_argument('--all-tokens-cooldown', type=int, default=1800) parser.add_argument('--max-repo-age-days', type=int, default=0) parser.add_argument('--repo-age-field', default='created_at') parser.add_argument('--max-commit-age-days', type=int, default=0) parser.add_argument('--commit-lookup-pages', type=int, default=3) parser.add_argument('--no-skip-if-commit-lookup-fails', dest='skip_if_commit_lookup_fails', action='store_false') parser.set_defaults(skip_if_commit_lookup_fails=True) parser.add_argument('--max-depth', type=int, default=0) parser.add_argument('--exact-git-planning', dest='exact_git_planning_enabled', action='store_true') parser.add_argument('--git-baseline-depth', type=int, default=100) parser.add_argument('--git-ref-resolution-timeout-sec', type=float, default=10) parser.add_argument('--git-ref-resolution-attempts', type=int, default=2) parser.add_argument('--git-ref-resolution-max-bytes', type=int, default=1 << 20) parser.add_argument('--admission-resolution-attempts', type=int, default=90) parser.add_argument('--admission-resolution-seconds', type=float, default=90) parser.add_argument('--admission-resolution-retry-delay-sec', type=float, default=1) parser.add_argument('--scan-full-history', action='store_true') parser.add_argument('--sort-by', default='updated') parser.add_argument('--sort-order', choices=['asc', 'desc'], default='desc') parser.add_argument('--created-filter', choices=['any', 'today', 'week', 'month', 'year'], default='any') parser.add_argument('--gitlab-sort-by', default='last_activity_at') parser.add_argument('--gitlab-visibility', choices=['public', 'internal', 'private'], default='public') parser.add_argument('--docker-sort-by', default='updated_at') parser.add_argument('--fetch-workers', type=int, default=8) parser.add_argument('--fetch-timeout', type=int, default=15) parser.add_argument('--tag-fetch-workers', type=int, default=4) parser.add_argument('--tag-retry-count', type=int, default=2) parser.add_argument('--tag-retry-delay', type=int, default=5) parser.add_argument('--tag-resolve-limit', type=int, default=100) parser.add_argument('--no-docker-platform-filter', dest='docker_platform_filter_enabled', action='store_false') parser.add_argument('--docker-platform-os', default='linux') parser.add_argument('--docker-platform-arch', default='amd64') parser.add_argument('--docker-platform-candidate-tags', type=int, default=20) def docker_image_depth_argument(value): try: return docker_images_per_repository_limit(int(value, 10)) except (TypeError, ValueError) as exc: raise argparse.ArgumentTypeError( 'docker image depth must be an integer from 1 through 10' ) from exc parser.add_argument('--docker-images-per-repository', type=docker_image_depth_argument, default=1) parser.add_argument( '--docker-content-scan-mode', choices=['full', 'canary', 'layer', 'adaptive-canary', 'adaptive'], default='full', ) parser.add_argument('--docker-layer-canary-basis-points', type=int, default=0) parser.add_argument('--docker-adaptive-canary-basis-points', type=int, default=0) parser.add_argument('--docker-adaptive-gate-max-age-sec', type=int, default=604800) parser.add_argument('--docker-layer-config-max-bytes', type=int, default=1 << 20) parser.add_argument('--docker-layer-max-bytes', type=int, default=256 << 20) parser.add_argument('--docker-layer-image-max-bytes', type=int, default=1 << 30) parser.add_argument('--docker-layer-max-layers', type=int, default=8) parser.add_argument('--docker-layer-archive-max-size-bytes', type=int, default=256 << 20) parser.add_argument('--docker-layer-archive-max-depth', type=int, default=4) parser.add_argument('--docker-layer-archive-timeout-sec', type=int, default=30) parser.add_argument('--docker-layer-blob-timeout-sec', type=int, default=600) parser.add_argument('--docker-layer-filesystem-concurrency', type=int, default=2) parser.add_argument('--docker-layer-blob-max-attempts', type=int, default=3) parser.add_argument('--docker-layer-blob-lease-sec', type=int, default=1800) parser.add_argument('--docker-layer-min-free-bytes', type=int, default=20 << 30) parser.add_argument('--docker-layer-checkpoint-delay-sec', type=int, default=60) parser.add_argument('--docker-adaptive-checkpoint-max-blobs', type=int, default=4) parser.add_argument('--docker-adaptive-checkpoint-max-bytes', type=int, default=512 << 20) parser.add_argument('--docker-repository-refresh-interval-sec', type=int, default=86400) parser.add_argument('--docker-repository-refresh-max-per-cycle', type=int, default=0) parser.set_defaults(docker_platform_filter_enabled=True) parser.add_argument('--archive-hours-back', type=int, default=6) parser.add_argument('--archive-max-repos-per-cycle', type=int, default=200) parser.add_argument('--archive-max-files-per-cycle', type=int, default=300) parser.add_argument('--archive-max-commit-lookups', type=int, default=200) parser.add_argument('--archive-rescan-cooldown-hours', type=int, default=48) parser.add_argument('--archive-event-types', default='PushEvent,CreateEvent,PublicEvent') parser.add_argument('--ci-seed-sources', default='github,gitlab,package_git') parser.add_argument('--ci-use-finding-seeds', action='store_true') parser.add_argument('--ci-max-repos-per-cycle', type=int, default=50) parser.add_argument('--ci-seed-scan-limit', type=int, default=5000) parser.add_argument('--ci-soft-cooldown-days', type=int, default=7) parser.add_argument('--ci-runs-per-repo', type=int, default=5) parser.add_argument('--ci-pipelines-per-project', type=int, default=5) parser.add_argument('--ci-jobs-per-pipeline', type=int, default=20) parser.add_argument('--ci-lookback-days', type=int, default=30) parser.add_argument('--ci-max-log-archive-mb', type=int, default=50) parser.add_argument('--ci-max-log-file-mb', type=int, default=20) parser.add_argument('--ci-max-trace-mb', type=int, default=20) parser.add_argument('--ci-scan-artifacts', action='store_true') parser.add_argument('--ci-max-artifacts-per-run', type=int, default=3) parser.add_argument('--ci-max-artifacts-per-pipeline', type=int, default=5) parser.add_argument('--ci-max-artifact-archive-mb', type=int, default=50) parser.add_argument('--ci-max-artifact-file-mb', type=int, default=10) parser.add_argument('--ci-max-artifact-files', type=int, default=1000) parser.add_argument('--ci-target-max-download-mb', type=int, default=500) parser.add_argument('--no-ci-failed-first', dest='ci_failed_first', action='store_false') parser.set_defaults(ci_failed_first=True) parser.add_argument('--stop-on-seen-pages', action='store_true') parser.add_argument('--seen-page-threshold', type=int, default=2) parser.add_argument('--min-pages-before-stop', type=int, default=1) parser.add_argument('--cleanup-temp-age-min', type=int, default=120) parser.add_argument('--cleanup-only', action='store_true') return parser.parse_args() def main(): args = parse_args() if args.show_state and not args.config: raise SystemExit('--show-state requires an explicit --config path') config = None if args.config and os.path.exists(args.config): config = load_config( args.config, managed_postgres=bool(os.getenv('TRUF_MANAGED_POSTGRES_DSN')), ) runtime_role = str(os.getenv(CHILD_KIND_ENV) or 'scanner').strip().lower() if not args.show_state: if not args.config: raise SystemExit('Legacy direct scanner mode is retired; use canonical supervisor source commands.') try: require_active_supervisor_child( args.config, child_kind=runtime_role, require_dsn=True, ) except LifecycleAuthorityError as exc: raise SystemExit(str(exc)) from exc if args.config: run_config_mode(args, config=config, runtime_role=runtime_role) return scan_config.work_dir = args.work_dir scan_config.results_dir = args.save_dir scan_config.result_spool_dir = args.result_spool_dir scan_config.result_spool_max_event_bytes = args.result_spool_max_event_bytes scan_config.result_spool_max_events = args.result_spool_max_events scan_config.result_spool_max_total_bytes = args.result_spool_max_total_bytes scan_config.result_spool_min_free_bytes = args.result_spool_min_free_bytes scan_config.queue_dir = args.queue_dir scan_config.postman_cache_dir = args.postman_cache_dir scan_config.trufflehog_path = args.trufflehog_path scan_config.trufflehog_config = args.trufflehog_config security_paths = default_project_paths() security_paths.update({ 'runtime_dir': args.runtime_dir, 'results_dir': args.save_dir, 'result_spool_dir': args.result_spool_dir, 'queue_dir': args.queue_dir, 'work_dir': args.work_dir, }) require_sensitive_runtime_paths(security_paths, create=False) initialize_scanner_runtime(preflight_complete=True) if args.cleanup_only: cleanup_pending_command_work_dirs(attempts=3, delay=0.2, log_failures=True) cleanup_stale_temp_dirs(args.cleanup_temp_age_min) print('Cleanup complete.') return if skip_startup_cleanup(): print('Startup temp cleanup skipped (supervised child).') else: cleanup_pending_command_work_dirs(attempts=3, delay=0.2, log_failures=True) cleanup_stale_temp_dirs(args.cleanup_temp_age_min) if not args.platform: raise SystemExit('--platform is required unless --cleanup-only is used') if args.platform == 'dockerhub': args.platform = 'docker' if args.platform == 'docker': docker_token = args.docker_token or args.token or os.getenv('DOCKERHUB_TOKEN') or os.getenv('DOCKER_TOKEN') or os.getenv('DOCKER_TOKENS') docker_username = args.docker_username or os.getenv('DOCKERHUB_USERNAME') or os.getenv('DOCKER_USERNAME') configure_docker_tokens(docker_token, docker_username) if args.timeout is None: args.timeout = scan_config.docker_timeout if args.platform == 'docker' else scan_config.git_timeout if not check_dependencies(): raise SystemExit(1) db = ScannerDB(args.save_dir) if getattr(db, 'postgres_required', False) and not db.enabled: raise RuntimeError('Configured Postgres database is unavailable; refusing file-queue fallback') if db.enabled and getattr(db.conn, 'is_postgres', False): db.require_runtime_safety_schema() run_id = db.start_run( 'legacy', sys.argv, selected_platform=args.platform, enabled_sources=[args.platform], global_config=vars(args), ) if db.enabled else None if db.enabled and getattr(db.conn, 'is_postgres', False) and run_id is None: raise RuntimeError('Unable to create Postgres scanner run; refusing target claims') run_status = 'completed' run_error = None cycle = 0 try: while True: cycle += 1 print(f'\n=== Cycle {cycle} started at {datetime.now().isoformat(timespec="seconds")} ===') cycle_id = db.start_source_cycle( run_id, args.platform, args.platform, args.mode, args.query, cycle, args.max_cycles or None, None, vars(args), queue_counts_for_args(args), ) if db.enabled and run_id else None cycle_metrics = {} try: run_cycle( args, db, run_id, cycle_id, args.platform, partial_metrics=cycle_metrics, ) except RateLimitError as e: if db.enabled and cycle_id: metrics = {'fetched_count': 0} metrics.update(cycle_metrics) db.finish_source_cycle(cycle_id, 'rate_limited', metrics, queue_counts_for_args(args), str(e)) raise except Exception as e: if db.enabled and getattr(db, 'conn', None): try: db.conn.rollback() except Exception: pass if db.enabled and cycle_id: metrics = {'fetched_count': 0} metrics.update(cycle_metrics) db.finish_source_cycle(cycle_id, 'failed', metrics, queue_counts_for_args(args), str(e)) raise if not args.loop: break if args.max_cycles and cycle >= args.max_cycles: break print(f'Sleeping {args.cooldown} seconds...') time.sleep(args.cooldown) except KeyboardInterrupt: run_status = 'stopped' print('Stopped. Temp cleanup deferred to next run or --cleanup-only.') except Exception as e: run_status = 'failed' run_error = str(e) raise finally: if db.enabled and run_id: db.finish_run(run_id, run_status, run_error) db.close() if __name__ == '__main__': main()