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

7157 lines
326 KiB
Python

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 '<unnamed>')[: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 "<missing>"}')
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()