import sys sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('keycheck runner could not disable bytecode writes') import argparse import copy import json import math import os import re import stat import subprocess import threading import time from collections import deque from concurrent.futures import FIRST_COMPLETED, ThreadPoolExecutor, wait from datetime import datetime, timedelta, timezone from paths import apply_path_config, default_project_paths from owned_process import OwnedProcess from scanner_db import ScannerDB, extract_finding_location, json_dumps, extract_raw_secret, redact_argv, sha256_text from keycheckers.keycheck_common import ( STATUS_TRANSACTION_JOURNAL_FILENAME, checkpoint_matches_signature, keycheck_status_group, mask_secret, private_atomic_writer, ) from postgres_runtime import load_postgres_environment from lifecycle_authority import ( LifecycleAuthorityError, require_active_supervisor_child, supervised_child_environment, strip_supervisor_credentials, ) from runtime_security import ( ensure_private_directory, preflight_lifecycle_paths, require_private_file, require_sensitive_runtime_paths, ) SERVICES = { "anthropic": os.path.join("keycheckers", "anthropic", "anthropicKeycheck.py"), "aws": os.path.join("keycheckers", "aws", "awsKeycheck.py"), "azure": os.path.join("keycheckers", "azure", "azureKeycheck.py"), "deepseek": os.path.join("keycheckers", "deepseek", "deepseekKeycheck.py"), "dockerhub": os.path.join("keycheckers", "dockerhub", "dockerhubKeycheck.py"), "gcp": os.path.join("keycheckers", "gcp", "gcpKeycheck.py"), "gemini": os.path.join("keycheckers", "gemini", "geminiKeycheck.py"), "groq": os.path.join("keycheckers", "groq", "groqKeycheck.py"), "github": os.path.join("keycheckers", "github", "githubKeycheck.py"), "gitlab": os.path.join("keycheckers", "gitlab", "gitlabKeycheck.py"), "kimi": os.path.join("keycheckers", "kimi", "kimiKeycheck.py"), "openai": os.path.join("keycheckers", "openai", "Keycheck.py"), "openrouter": os.path.join("keycheckers", "openrouter", "OpenrouterKeycheck.py"), "provider_resolver": os.path.join("keycheckers", "provider_resolver", "providerResolverKeycheck.py"), "qwen": os.path.join("keycheckers", "qwen", "qwenKeycheck.py"), "replicate": os.path.join("keycheckers", "replicate", "replicateKeycheck.py"), "xai": os.path.join("keycheckers", "xai", "xaiKeycheck.py"), "huggingface": os.path.join("keycheckers", "huggingface", "huggingfaceKeycheck.py"), "zai": os.path.join("keycheckers", "zai", "zaiKeycheck.py"), } KEYCHECK_CAPACITY_BLOCKED_EXIT = 76 _unconfirmed_provider_processes = [] SERVICE_CAPABILITIES = { "anthropic": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_no_balance", "recheck_all"}}, "aws": {"proxy": True, "flags": {"retry_network", "retry_unknown", "retry_valid", "recheck_all"}}, "azure": {"proxy": True, "flags": {"retry_network", "retry_unknown", "retry_valid", "recheck_all"}}, "deepseek": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_no_balance", "retry_valid", "recheck_all"}}, "dockerhub": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "recheck_all"}}, "gcp": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_valid", "recheck_all"}}, "gemini": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_valid", "recheck_all"}}, "groq": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, "github": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "recheck_all"}}, "gitlab": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "recheck_all"}}, "kimi": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, "openai": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "recheck_all"}}, "openrouter": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_no_balance", "retry_valid", "recheck_all"}}, "provider_resolver": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, "qwen": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, "replicate": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "no_resource_probe", "recheck_all"}}, "xai": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, "huggingface": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "recheck_all"}}, "zai": {"proxy": True, "flags": {"retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all"}}, } RETRY_STATUS_DEFAULT_SUFFIXES = { "retry_network": ("Network",), "retry_limited": ("Limited",), "retry_unknown": ("Unknown",), "retry_restricted": ("Restricted",), "retry_no_balance": ("NoBalance",), "retry_valid": ("Alive",), } RETRY_STATUS_FILE_OVERRIDES = { ("anthropic", "retry_no_balance"): ("anthropicNoQuota.txt",), ("aws", "retry_valid"): ("awsAlive.txt", "awsBedrock.txt", "awsAdmin.txt"), ("azure", "retry_network"): ( "azureNetwork.txt", "azureOpenAIBadEndpoint.txt", "azureFoundryBadEndpoint.txt", ), ("azure", "retry_unknown"): ( "azureUnknown.txt", "azureOpenAIUnresolved.txt", "azureFoundryUnresolved.txt", ), ("azure", "retry_valid"): ("azureAlive.txt", "azureFoundryLLM.txt"), ("dockerhub", "retry_limited"): ("dockerhubRateLimited.txt",), ("deepseek", "retry_unknown"): ("deepseekUnknown.txt", "deepseekNoContext.txt"), ("gcp", "retry_limited"): ("gcpUnknown.txt",), ("gcp", "retry_valid"): ("gcpAlive.txt", "gcpVertex.txt"), ("gemini", "retry_limited"): ("geminiRateLimited.txt", "geminiAliveRateLimited.txt"), ("gemini", "retry_valid"): ("geminiAlive.txt", "geminiAliveRateLimited.txt"), ("huggingface", "retry_unknown"): ("huggingfaceUnknown.txt", "huggingfaceNoContext.txt"), ("github", "retry_limited"): ("githubRateLimited.txt",), ("gitlab", "retry_limited"): ("gitlabRateLimited.txt",), ("openai", "retry_no_balance"): ("openaiLimited.txt",), ("qwen", "retry_unknown"): ("qwenUnknown.txt", "qwenNoContext.txt"), ("replicate", "retry_unknown"): ("replicateUnknown.txt", "replicateNoContext.txt"), ("xai", "retry_unknown"): ("xaiUnknown.txt", "xaiNoContext.txt"), } POSTGRES_STATUS_FILE_NAMES = { "anthropic": { "VALID": "anthropicAlive.txt", "NO_QUOTA": "anthropicNoQuota.txt", "DEAD": "anthropicDead.txt", "LIMITED": "anthropicLimited.txt", "RESTRICTED": "anthropicRestricted.txt", "NETWORK": "anthropicNetwork.txt", "UNKNOWN": "anthropicUnknown.txt", }, "aws": { "VALID": "awsAlive.txt", "BEDROCK": "awsBedrock.txt", "ADMIN": "awsAdmin.txt", "CANARY": "awsCanary.txt", "QUARANTINED": "awsQuarantined.txt", "ACCESS_DENIED": "awsAccessDenied.txt", "DEAD": "awsDead.txt", "NETWORK": "awsNetwork.txt", "UNKNOWN": "awsUnknown.txt", }, "azure": { "VALID": "azureAlive.txt", "DEAD": "azureDead.txt", "RESTRICTED": "azureRestricted.txt", "NETWORK": "azureNetwork.txt", "UNKNOWN": "azureUnknown.txt", "OPENAI_UNRESOLVED": "azureOpenAIUnresolved.txt", "OPENAI_BAD_ENDPOINT": "azureOpenAIBadEndpoint.txt", "FOUNDRY": "azureFoundryLLM.txt", "FOUNDRY_UNRESOLVED": "azureFoundryUnresolved.txt", "FOUNDRY_BAD_ENDPOINT": "azureFoundryBadEndpoint.txt", }, "deepseek": { "VALID": "deepseekAlive.txt", "NO_BALANCE": "deepseekNoBalance.txt", "DEAD": "deepseekDead.txt", "LIMITED": "deepseekLimited.txt", "NETWORK": "deepseekNetwork.txt", "NO_CONTEXT": "deepseekNoContext.txt", "UNKNOWN": "deepseekUnknown.txt", }, "dockerhub": { "VALID": "dockerhubAlive.txt", "VALID_2FA": "dockerhubAlive.txt", "DEAD": "dockerhubDead.txt", "NO_USERNAME": "dockerhubNoUsername.txt", "RATE_LIMITED": "dockerhubRateLimited.txt", "NETWORK": "dockerhubNetwork.txt", "UNKNOWN": "dockerhubUnknown.txt", }, "gcp": { "VALID": "gcpAlive.txt", "VERTEX": "gcpVertex.txt", "DEAD": "gcpDead.txt", "RESTRICTED": "gcpRestricted.txt", "NETWORK": "gcpNetwork.txt", "RATE_LIMITED": "gcpUnknown.txt", "UNKNOWN": "gcpUnknown.txt", "NO_CONTEXT": "gcpNoContext.txt", }, "gemini": { "VALID": "geminiAlive.txt", "VALID_RATE_LIMITED": "geminiAliveRateLimited.txt", "INVALID": "geminiDead.txt", "EXPIRED": "geminiExpired.txt", "LEAKED_REVOKED": "geminiLeaked.txt", "API_DISABLED": "geminiDisabled.txt", "RESTRICTED": "geminiRestricted.txt", "RATE_LIMITED": "geminiRateLimited.txt", "NETWORK_ERROR": "geminiNetwork.txt", "UNKNOWN": "geminiUnknown.txt", }, "groq": { "VALID": "groqAlive.txt", "NO_BALANCE": "groqNoBalance.txt", "DEAD": "groqDead.txt", "RESTRICTED": "groqRestricted.txt", "LIMITED": "groqLimited.txt", "NO_CONTEXT": "groqNoContext.txt", "NETWORK": "groqNetwork.txt", "UNKNOWN": "groqUnknown.txt", }, "github": { "VALID": "githubAlive.txt", "DEAD": "githubDead.txt", "RESTRICTED": "githubRestricted.txt", "RATE_LIMITED": "githubRateLimited.txt", "NETWORK": "githubNetwork.txt", "UNKNOWN": "githubUnknown.txt", "REFRESH_TOKEN": "githubRefreshToken.txt", }, "gitlab": { "VALID": "gitlabAlive.txt", "DEAD": "gitlabDead.txt", "RESTRICTED": "gitlabRestricted.txt", "RATE_LIMITED": "gitlabRateLimited.txt", "NETWORK": "gitlabNetwork.txt", "UNKNOWN": "gitlabUnknown.txt", }, "kimi": { "VALID": "kimiAlive.txt", "NO_BALANCE": "kimiNoBalance.txt", "DEAD": "kimiDead.txt", "RESTRICTED": "kimiRestricted.txt", "LIMITED": "kimiLimited.txt", "NETWORK": "kimiNetwork.txt", "UNKNOWN": "kimiUnknown.txt", }, "openai": { "ALIVE": "openaiAlive.txt", "DEAD": "openaiDead.txt", "INVALID_OR_REVOKED": "openaiDead.txt", "NETWORK_ERROR": "openaiNetwork.txt", "LIMITED": "openaiLimited.txt", "LIMITED_OR_NO_BALANCE": "openaiLimited.txt", "RESTRICTED": "openaiRestricted.txt", "UNKNOWN": "openaiUnknown.txt", "NO_TARGET_MODELS": "openaiNoTarget.txt", }, "openrouter": { "VALID": "openrouterAlive.txt", "NO_BALANCE": "openrouterNoBalance.txt", "DEAD": "openrouterDead.txt", "LIMITED": "openrouterLimited.txt", "NETWORK": "openrouterNetwork.txt", "UNKNOWN": "openrouterUnknown.txt", }, "provider_resolver": { "VALID": "providerResolverAlive.txt", "NO_BALANCE": "providerResolverNoBalance.txt", "DEAD": "providerResolverDead.txt", "RESTRICTED": "providerResolverRestricted.txt", "LIMITED": "providerResolverLimited.txt", "NETWORK": "providerResolverNetwork.txt", "NO_CONTEXT": "providerResolverNoContext.txt", "UNKNOWN": "providerResolverUnknown.txt", }, "qwen": { "VALID": "qwenAlive.txt", "DEAD": "qwenDead.txt", "RESTRICTED": "qwenRestricted.txt", "LIMITED": "qwenLimited.txt", "NO_BALANCE": "qwenNoBalance.txt", "NO_CONTEXT": "qwenNoContext.txt", "NETWORK": "qwenNetwork.txt", "UNKNOWN": "qwenUnknown.txt", }, "replicate": { "VALID": "replicateAlive.txt", "DEAD": "replicateDead.txt", "RESTRICTED": "replicateRestricted.txt", "LIMITED": "replicateLimited.txt", "NO_BALANCE": "replicateNoBalance.txt", "NETWORK": "replicateNetwork.txt", "NO_CONTEXT": "replicateNoContext.txt", "UNKNOWN": "replicateUnknown.txt", }, "xai": { "VALID": "xaiAlive.txt", "NO_BALANCE": "xaiNoBalance.txt", "DEAD": "xaiDead.txt", "RESTRICTED": "xaiRestricted.txt", "LIMITED": "xaiLimited.txt", "NO_CONTEXT": "xaiNoContext.txt", "NETWORK": "xaiNetwork.txt", "UNKNOWN": "xaiUnknown.txt", }, "huggingface": { "VALID": "huggingfaceAlive.txt", "DEAD": "huggingfaceDead.txt", "RESTRICTED": "huggingfaceRestricted.txt", "LIMITED": "huggingfaceLimited.txt", "NETWORK": "huggingfaceNetwork.txt", "NO_CONTEXT": "huggingfaceNoContext.txt", "UNKNOWN": "huggingfaceUnknown.txt", }, "zai": { "VALID": "zaiAlive.txt", "NO_BALANCE": "zaiNoBalance.txt", "DEAD": "zaiDead.txt", "RESTRICTED": "zaiRestricted.txt", "LIMITED": "zaiLimited.txt", "NETWORK": "zaiNetwork.txt", "UNKNOWN": "zaiUnknown.txt", }, } POSTGRES_AUXILIARY_STATUS_FILE_NAMES = { "gcp": ("gcpVertexGemini.txt", "gcpVertexAnthropic.txt"), "azure": ("azureOpenAILLM.txt",), } LEGACY_GCP_VERTEX_STATUS_FILES = ("gcpVertexGemini.txt", "gcpVertexAnthropic.txt") def load_config(config_path): if not config_path: return {"global": default_project_paths(), "keychecks": {}} try: import yaml except ImportError as e: raise SystemExit("PyYAML is required for --config") from e with open(config_path, "r", encoding="utf-8") as f: return apply_path_config(yaml.safe_load(f) or {}, config_path) def load_layout(config_path): return (load_config(config_path).get("global") or {}) def load_postgres_env(config_path, layout): return load_postgres_environment(config_path, {'global': layout or {}}) def service_list(value): if isinstance(value, (list, tuple)): value = ",".join(str(item) for item in value) if value == "all": return list(SERVICES) return [item.strip().lower() for item in value.split(",") if item.strip()] def count_nonempty_lines(path): if not os.path.exists(path): return 0 with open(path, "r", encoding="utf-8", errors="replace") as f: return sum(1 for line in f if line.strip()) def file_line_stats(path, unique_keys=False): # Operational summary should be fast even with large historical files. if not unique_keys: count = 0 tail = b'' with open(path, "rb") as f: while True: chunk = f.read(1024 * 1024) if not chunk: break count += chunk.count(b"\n") tail = chunk try: if os.path.getsize(path) > 0 and not tail.endswith(b"\n"): count += 1 except OSError: pass return count, None # Avoid expensive unique parsing on very large checked files; line count is enough for dashboard health. if os.path.getsize(path) > 5 * 1024 * 1024: count, _ = file_line_stats(path, unique_keys=False) return count, count line_count = 0 keys = set() with open(path, "r", encoding="utf-8", errors="replace") as f: for line in f: if not line.strip(): continue line_count += 1 key = status_key(line) if key: keys.add(key) return line_count, len(keys) def status_key(line): line = line.strip() if not line: return None if "\t" in line: return line.split("\t", 1)[0].strip() return line.split(":", 1)[0].strip() def _row_value(row, name, default=None): try: value = row[name] except (KeyError, TypeError): value = getattr(row, name, default) return default if value is None else value def _projection_metadata(row): value = _row_value(row, "metadata_json", "") if isinstance(value, dict): return value try: parsed = json.loads(str(value or "{}")) except (TypeError, ValueError, json.JSONDecodeError): return {} return parsed if isinstance(parsed, dict) else {} def _bounded_projection_text(value, max_bytes): text = str(value or "").replace("\x00", "").replace("\t", " ").replace("\r", " ").replace("\n", " ") encoded = text.encode("utf-8", errors="strict") if len(encoded) <= max_bytes: return text return encoded[:max_bytes].decode("utf-8", errors="ignore") def postgres_status_projection_detail(metadata, max_bytes=16000): metadata = metadata if isinstance(metadata, dict) else {} probe = metadata.get("probe") if isinstance(metadata.get("probe"), dict) else {} probe_model = metadata.get("llm_probe_model") or probe.get("model") or "" probe_status = metadata.get("llm_probe_status") or probe.get("status") or "" inventory = metadata.get("model_inventory") if not isinstance(inventory, list): inventory = metadata.get("models") if isinstance(metadata.get("models"), list) else [] models = [] seen = set() for value in inventory[:500]: model = str(value or "").replace("\x00", "").replace("\t", " ").replace("\r", " ").replace("\n", " ")[:200] if model and model not in seen: seen.add(model) models.append(model) parts = [] message = str(metadata.get("message") or "").strip() if message: parts.append(message[:1000]) if probe_model: parts.append(f"probe_model={probe_model}") if probe_status: parts.append(f"probe_status={probe_status}") model_count = metadata.get("model_count") if model_count is not None: parts.append(f"model_count={model_count}") if models: parts.append(f"models={','.join(models)}") return _bounded_projection_text("; ".join(parts), max(0, int(max_bytes))) def _metadata_enabled(metadata, name): value = metadata.get(name) if isinstance(value, bool): return value return str(value or "").strip().lower() in ("1", "true", "yes", "on") def postgres_status_projection_targets(service, status, metadata=None): service = str(service or "").strip().lower() status = str(status or "UNKNOWN").strip().upper() filename = (POSTGRES_STATUS_FILE_NAMES.get(service) or {}).get(status) if not filename: return [] targets = [filename] metadata = metadata or {} if service == "gcp" and status == "VERTEX": if _metadata_enabled(metadata, "vertex_google_enabled"): targets.append("gcpVertexGemini.txt") if _metadata_enabled(metadata, "vertex_anthropic_enabled"): targets.append("gcpVertexAnthropic.txt") if ( service == "azure" and status == "VALID" and int(metadata.get("deployment_count") or 0) > 0 and (not metadata.get("route_probe") or metadata.get("route_probe") == "accepted_auth_route") ): targets.append("azureOpenAILLM.txt") return list(dict.fromkeys(targets)) def _read_status_projection(path, max_items, max_bytes, max_line_bytes): if not os.path.exists(path): return [] require_private_file(path) if os.path.getsize(path) > max_bytes: raise RuntimeError(f"keycheck status projection exceeds its byte bound: {path}") lines = [] total = 0 with open(path, "rb") as handle: for raw_line in handle: if len(raw_line) > max_line_bytes: raise RuntimeError(f"keycheck status projection line exceeds its byte bound: {path}") if not raw_line.strip(): continue total += len(raw_line) if total > max_bytes or len(lines) >= max_items: raise RuntimeError(f"keycheck status projection exceeds its aggregate bound: {path}") lines.append(raw_line.decode("utf-8", errors="replace").rstrip("\r\n") + "\n") return lines def _write_status_projection(path, lines, max_items, max_bytes): if len(lines) > max_items: raise RuntimeError(f"keycheck status projection exceeds its item bound: {path}") encoded_bytes = sum(len(line.encode("utf-8", errors="strict")) for line in lines) if encoded_bytes > max_bytes: raise RuntimeError(f"keycheck status projection exceeds its byte bound: {path}") with private_atomic_writer(path, suffix=".status-projection.tmp") as handle: handle.writelines(lines) def project_postgres_status_files(layout, services, rows, managed_rows=None): keycheck_dir = layout["keycheck_dir"] ensure_private_directory(keycheck_dir, reject_reparse=True) selected = {str(service or "").strip().lower() for service in services} max_items = max(1, env_int("KEYCHECK_STATUS_PROJECTION_MAX_ITEMS", 100000)) max_bytes = max(1024, env_int("KEYCHECK_STATUS_PROJECTION_MAX_BYTES", 32 * 1024 * 1024)) max_line_bytes = max(1024, env_int("KEYCHECK_INPUT_MAX_LINE_BYTES", 16 * 1024 * 1024)) rows = list(rows or []) if len(rows) > max_items: raise RuntimeError("PostgreSQL keycheck status projection exceeds its row bound") grouped = {service: [] for service in selected if service in POSTGRES_STATUS_FILE_NAMES} managed_by_service = {service: set() for service in grouped} for row in managed_rows or (): service = str(_row_value(row, "service", "") or "").strip().lower() if service not in managed_by_service: continue key = str( _row_value(row, "secret_text", "") or _row_value(row, "secret_json", "") or "" ) if ( key and not any(character in key for character in ("\x00", "\t", "\r", "\n")) and len(key.encode("utf-8", errors="strict")) <= max_line_bytes ): managed_by_service[service].add(key) skipped = 0 for row in rows: service = str(_row_value(row, "service", "") or "").strip().lower() if service not in grouped: continue key = str( _row_value(row, "secret_text", "") or _row_value(row, "secret_json", "") or "" ) status = str(_row_value(row, "status", "UNKNOWN") or "UNKNOWN").strip().upper() metadata = _projection_metadata(row) targets = postgres_status_projection_targets(service, status, metadata) if ( not key or not targets or any(character in key for character in ("\x00", "\t", "\r", "\n")) or len(key.encode("utf-8", errors="strict")) > max_line_bytes ): skipped += 1 continue grouped[service].append({ "key": key, "status": status, "checked_at": str(_row_value(row, "checked_at", "") or ""), "message": postgres_status_projection_detail( metadata, min(16000, max(0, max_line_bytes - len(key.encode("utf-8", errors="strict")) - 128)), ), "targets": targets, }) changed_files = 0 projected_rows = 0 for service, service_rows in grouped.items(): service_dir = os.path.join(keycheck_dir, service) ensure_private_directory(service_dir, reject_reparse=True) filenames = list(dict.fromkeys([ *(POSTGRES_STATUS_FILE_NAMES[service].values()), *POSTGRES_AUXILIARY_STATUS_FILE_NAMES.get(service, ()), ])) paths = {name: os.path.join(service_dir, name) for name in filenames} snapshots = { name: _read_status_projection(path, max_items, max_bytes, max_line_bytes) for name, path in paths.items() } managed_keys = managed_by_service[service] | {row["key"] for row in service_rows} rewritten = { name: [line for line in lines if status_key(line) not in managed_keys] for name, lines in snapshots.items() } for row in service_rows: status_line = f'{row["key"]}\t{row["status"]}\t{row["message"]}\tpostgres-authoritative\n' for filename in row["targets"]: rewritten[filename].append(status_line) checked_name = f"{service}Checked.txt" checked_path = os.path.join(service_dir, checked_name) checked_snapshot = _read_status_projection( checked_path, max_items, max_bytes, max_line_bytes, ) checked_rewritten = [ line for line in checked_snapshot if status_key(line) not in managed_keys ] checked_rewritten.extend( f'{row["key"]}\t{row["status"]}\t{row["checked_at"]}\n' for row in service_rows ) for name, lines in rewritten.items(): if lines == snapshots[name]: continue _write_status_projection(paths[name], lines, max_items, max_bytes) changed_files += 1 if checked_rewritten != checked_snapshot: _write_status_projection(checked_path, checked_rewritten, max_items, max_bytes) changed_files += 1 projected_rows += len(service_rows) return { "projected_rows": projected_rows, "changed_files": changed_files, "skipped_rows": skipped, } def count_unique_status_keys(path): if not os.path.exists(path): return 0 with open(path, "r", encoding="utf-8", errors="replace") as f: return len({key for key in (status_key(line) for line in f) if key}) def utc_now_iso(): return datetime.now(timezone.utc).isoformat(timespec="seconds") def env_int(name, default): try: return int(os.getenv(name, default)) except (TypeError, ValueError): return default def truthy_env(name, default=False): value = os.getenv(name) if value is None: return default return str(value).strip().lower() in ("1", "true", "yes", "on") def file_mtime_iso(path): try: return datetime.fromtimestamp(os.path.getmtime(path), timezone.utc).isoformat(timespec="seconds") except OSError: return utc_now_iso() def status_from_payload(service, payload): status = str(payload.get("status") or "UNKNOWN").upper() if service == "gemini" and status == "VALID": probe = payload.get("probe") if isinstance(payload.get("probe"), dict) else {} if str(probe.get("status") or "").upper() == "RATE_LIMITED": return "VALID_RATE_LIMITED" return status def parse_source_line(source_line): text = str(source_line or "") if ":plain:" in text: return None, None marker = ":byte:" if marker in text: path, byte_text = text.rsplit(marker, 1) if path and byte_text.isdigit(): return path, ("byte", int(byte_text)) path, sep, line_text = text.rpartition(":") if not sep or not path or not line_text.isdigit(): return None, None return path, int(line_text) def load_jsonl_line(path, line_number, cache): if not path or not line_number or not os.path.exists(path): return None path = os.path.normpath(path) max_line_bytes = max(1024, env_int("KEYCHECK_INPUT_MAX_LINE_BYTES", 16 * 1024 * 1024)) if isinstance(line_number, tuple) and line_number[0] == "byte": try: with open(path, "rb") as f: f.seek(int(line_number[1])) raw_line = f.readline(max_line_bytes + 1) if len(raw_line) > max_line_bytes or not raw_line.endswith(b"\n"): return None line = raw_line.decode("utf-8", errors="replace") return json.loads(line) except (OSError, ValueError, json.JSONDecodeError): return None try: wanted = int(line_number) except (TypeError, ValueError): return None max_line_number = max(1, env_int("KEYCHECK_FINDING_LINE_MAX_NUMBER", 10000000)) if wanted <= 0 or wanted > max_line_number: return None cache_key = (path, wanted) line = cache.get(cache_key) if line is None: try: with open(path, "rb") as handle: raw_line = b'' for _ in range(wanted): raw_line = handle.readline(max_line_bytes + 1) if not raw_line or len(raw_line) > max_line_bytes: return None if not raw_line.endswith(b"\n"): return None line = raw_line.decode("utf-8", errors="replace") except OSError: return None max_cache_rows = max(1, env_int("KEYCHECK_FINDING_LINE_CACHE_ITEMS", 128)) max_cache_bytes = max(1024, env_int("KEYCHECK_FINDING_LINE_CACHE_BYTES", 16 * 1024 * 1024)) line_bytes = len(line.encode("utf-8", errors="replace")) while cache and ( len(cache) >= max_cache_rows or sum(len(value.encode("utf-8", errors="replace")) for value in cache.values()) + line_bytes > max_cache_bytes ): cache.pop(next(iter(cache))) if line_bytes <= max_cache_bytes: cache[cache_key] = line try: return json.loads(line) except json.JSONDecodeError: return None def finding_from_payload(payload, line_cache): finding = payload.get("finding") if isinstance(finding, dict) and finding: return finding path, line_number = parse_source_line(payload.get("source")) return load_jsonl_line(path, line_number, line_cache) def safe_metadata(payload): out = {} for key, value in payload.items(): lowered = str(key).lower() if lowered in ("finding", "raw", "raw_v2", "rawsecret", "raw_secret"): continue out[key] = value return out def infer_finding_attribution(finding): if not isinstance(finding, dict): return {} metadata = finding.get("SourceMetadata") if isinstance(finding.get("SourceMetadata"), dict) else {} data = metadata.get("Data") if isinstance(metadata.get("Data"), dict) else {} detector = finding.get("DetectorName") or finding.get("DetectorType") or "" for source_type, details in data.items(): if not isinstance(details, dict): continue lowered = str(source_type or "").lower() repo = str(details.get("repository") or details.get("repo") or "") target = repo or details.get("image") or details.get("link") or "" source = "" query = None if lowered == "docker" or details.get("image"): source = "dockerhub" target = details.get("image") or target elif lowered == "huggingface" or "huggingface.co/" in repo: source = "huggingface" query = "spaces" if "/spaces/" in repo or details.get("resource_type") == "space" else None elif "gitlab.com" in repo: source = "gitlab" elif "github.com" in repo: source = "github" elif lowered == "git": source = "git" if source: return { "source": source, "query": query, "target": target, "detector_name": detector, "found_at": details.get("timestamp") or "", } return {"detector_name": detector} def sync_keycheck_results_to_db(layout, services, reset=False): if reset: raise SystemExit('--reset-keycheck-results is retired; online keycheck result deletion is not supported') db_path = layout.get("database_path") db_url = layout.get("database_url") or os.getenv("SCANNER_DB_URL") or os.getenv("DATABASE_URL") keycheck_dir = layout.get("keycheck_dir") if not db_path and not db_url: raise SystemExit("global.database_path or global.database_url is required for keycheck DB sync") if not keycheck_dir or not os.path.isdir(keycheck_dir): raise SystemExit(f"keycheck_dir not found: {keycheck_dir}") db = ScannerDB(db_path=db_path, db_url=db_url, initialize=False) if not db.enabled: raise SystemExit(f"Unable to open scanner DB: {db.db_display or db_path or 'configured database'}") try: db.require_runtime_safety_schema() except Exception as exc: db.close() raise SystemExit(f'Keycheck DB schema is incomplete; offline migration required: {exc}') from exc line_cache = {} inserted = 0 skipped = 0 missing_finding = 0 service_names = services if services and services != ["all"] else sorted(name for name in os.listdir(keycheck_dir) if os.path.isdir(os.path.join(keycheck_dir, name))) for service in service_names: service_dir = os.path.join(keycheck_dir, service) if not os.path.isdir(service_dir): continue result_files = [name for name in os.listdir(service_dir) if name.lower().endswith("results.jsonl")] for name in result_files: path = os.path.join(service_dir, name) checked_at = file_mtime_iso(path) with open(path, "r", encoding="utf-8", errors="replace") as f: for line in f: if not line.strip(): continue try: payload = json.loads(line) except json.JSONDecodeError: skipped += 1 continue status = status_from_payload(service, payload) finding = finding_from_payload(payload, line_cache) if not finding: missing_finding += 1 raw_secret = extract_raw_secret(finding or {}) finding_hash = sha256_text(raw_secret) if raw_secret else "" key_hash = payload.get("key_hash") or finding_hash secret_hash = finding_hash or payload.get("secret_hash") or key_hash key_masked = payload.get("key_masked") or mask_secret(raw_secret) if not key_hash and not key_masked: skipped += 1 continue detector = payload.get("detector") or (finding or {}).get("DetectorName") or "" detector_secret_hash = payload.get("detector_secret_hash") or (sha256_text("|".join([str(detector or ""), secret_hash])) if secret_hash or detector else "") finding_uid = payload.get("finding_uid") or (finding or {}).get("finding_uid") or "" event_id = payload.get("event_id") or "" message = payload.get("message") or "" error = payload.get("error") if not message and isinstance(error, dict): message = error.get("message") or json.dumps(error, ensure_ascii=False)[:1000] metadata = safe_metadata(payload) metadata["backfill_source_file"] = path metadata["finding_uid"] = finding_uid metadata["event_id"] = event_id matches = [None] now = utc_now_iso() if event_id: duplicate_event = db.conn.execute( "SELECT id FROM keycheck_results WHERE event_id = ? LIMIT 1", (event_id,), ).fetchone() if duplicate_event: skipped += 1 continue duplicate = db.conn.execute(''' SELECT id FROM keycheck_results WHERE service = ? AND status = ? AND checked_at = ? AND COALESCE(source_line, '') = ? AND ( (key_hash != '' AND key_hash = ?) OR (key_hash = '' AND key_masked = ?) ) LIMIT 1 ''', ( service, status, payload.get("checked_at") or checked_at, payload.get("source") or "", key_hash, key_masked, )).fetchone() if duplicate: skipped += 1 continue for row in matches: keycheck_result_id = db.conn.insert_returning_id( '''INSERT INTO keycheck_results ( service, status, status_group, checked_at, key_hash, secret_hash, key_masked, finding_id, target_scan_id, cycle_id, run_id, source, query, target, detector_name, found_at, message, metadata_json, source_line, detector_secret_hash, event_id, finding_uid, link_status, linked_at, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( service, status, keycheck_status_group(status), payload.get("checked_at") or checked_at, key_hash, secret_hash, key_masked, row["id"] if row else None, row["target_scan_id"] if row else None, row["cycle_id"] if row else None, row["run_id"] if row else None, row["source"] if row else None, row["query"] if row else None, row["target"] if row else None, row["detector_name"] if row else detector, row["created_at"] if row else None, str(message or "")[:1000].replace("\n", " "), json.dumps(metadata, ensure_ascii=False, default=str, sort_keys=True), payload.get("source") or "", detector_secret_hash, event_id, finding_uid, "linked" if row else "pending", now if row else None, now, ), ) if event_id and keycheck_result_id: db.conn.execute( '''INSERT INTO keycheck_event_map (event_id, keycheck_result_id, created_at) VALUES (?, ?, ?) ON CONFLICT(event_id) DO NOTHING''', (event_id, keycheck_result_id, now), ) inserted += 1 if inserted % 1000 == 0: db.conn.commit() db.conn.commit() db_display = db.db_display db.close() print(f"keycheck_results synced: inserted={inserted} skipped={skipped} missing_finding={missing_finding} db={db_display}") return inserted def set_linker_timeouts(db): if getattr(db.conn, "is_postgres", False): db.conn.execute("SELECT set_config('statement_timeout', ?, false)", (f"{int(os.getenv('KEYCHECK_LINK_STATEMENT_TIMEOUT_MS', '10000'))}ms",)) db.conn.execute("SELECT set_config('lock_timeout', ?, false)", (f"{int(os.getenv('KEYCHECK_LINK_LOCK_TIMEOUT_MS', '2000'))}ms",)) def update_link_status(db, row_id, status, error=""): db.conn.execute( '''UPDATE keycheck_results SET link_status = ?, link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = CASE WHEN ? IN ('linked', 'source_only') THEN ? ELSE linked_at END, link_error = ? WHERE id = ?''', (status, status, utc_now_iso(), str(error or '')[:1000], row_id), ) def backfill_keycheck_event_map(db, limit=5000): limit = max(0, int(limit or 0)) if limit <= 0: return 0 if not db.conn.table_exists("keycheck_event_map"): return 0 rows = db.conn.execute( '''SELECT id, event_id, created_at FROM ( SELECT id, event_id, created_at FROM keycheck_results ORDER BY id DESC LIMIT ? ) recent WHERE event_id IS NOT NULL AND event_id != '' ORDER BY id DESC''', (limit,), ).fetchall() inserted = 0 for row in rows: cur = db.conn.execute( '''INSERT INTO keycheck_event_map (event_id, keycheck_result_id, created_at) VALUES (?, ?, ?) ON CONFLICT(event_id) DO NOTHING''', (row["event_id"], row["id"], row["created_at"] or utc_now_iso()), ) if getattr(cur, "rowcount", 0) > 0: inserted += 1 db.conn.commit() return inserted def backfill_finding_uid_map(db, limit=5000): limit = max(0, int(limit or 0)) if limit <= 0: return 0 if not db.conn.table_exists("finding_uid_map"): return 0 rows = db.conn.execute( '''SELECT id, finding_uid, created_at FROM ( SELECT id, finding_uid, created_at FROM findings ORDER BY id DESC LIMIT ? ) recent WHERE finding_uid IS NOT NULL AND finding_uid != '' ORDER BY id DESC''', (limit,), ).fetchall() inserted = 0 for row in rows: cur = db.conn.execute( '''INSERT INTO finding_uid_map (finding_uid, finding_id, created_at) VALUES (?, ?, ?) ON CONFLICT(finding_uid) DO NOTHING''', (row["finding_uid"], row["id"], row["created_at"] or utc_now_iso()), ) if getattr(cur, "rowcount", 0) > 0: inserted += 1 db.conn.commit() return inserted def read_json_file(path, default=None): try: with open(path, "r", encoding="utf-8") as f: return json.load(f) except (OSError, ValueError): return default def write_json_file(path, data): with private_atomic_writer(path) as f: json.dump(data, f, ensure_ascii=False, indent=2, sort_keys=True) def jsonl_manifest_path(path): base, _ = os.path.splitext(path) return f"{base}.manifest.json" def keycheck_result_paths(current_path): manifest_path = jsonl_manifest_path(current_path) if os.path.exists(manifest_path): if os.path.getsize(manifest_path) > 1024 * 1024: raise RuntimeError(f"keycheck result manifest exceeds bounded size: {manifest_path}") with open(manifest_path, "r", encoding="utf-8") as handle: manifest = json.load(handle) if not isinstance(manifest, dict): raise RuntimeError(f"invalid keycheck result manifest: {manifest_path}") else: manifest = {} root = os.path.dirname(os.path.abspath(current_path)) base, extension = os.path.splitext(os.path.basename(current_path)) pattern = re.compile(rf"^{re.escape(base)}\.(\d{{6}}){re.escape(extension)}$") physical = [] with os.scandir(root) as entries: for entry in entries: match = pattern.fullmatch(entry.name) if not match: continue if entry.is_symlink() or not entry.is_file(follow_symlinks=False): raise RuntimeError(f"unsafe keycheck result segment: {entry.path}") physical.append((int(match.group(1)), os.path.abspath(entry.path))) physical.sort() paths = [path for _, path in physical] listed = { os.path.basename(str(segment.get("path") or segment.get("name") or "")) for segment in (manifest.get("segments") or []) if isinstance(segment, dict) } present = {os.path.basename(path) for path in paths} missing = sorted(name for name in listed if name and name not in present) if missing: raise RuntimeError(f"manifest-listed keycheck result segment is missing: {missing[0]}") declared_current = os.path.abspath(manifest.get("current_path") or current_path) if declared_current != os.path.abspath(current_path): raise RuntimeError("keycheck result manifest current path mismatch") if os.path.exists(current_path): paths.append(os.path.abspath(current_path)) if not paths and os.path.exists(current_path): paths.append(os.path.abspath(current_path)) seen = set() output = [] for path in paths: if path not in seen: seen.add(path) output.append(path) return output def keycheck_result_current_skip(current_path): # A manifest skip can become stale between validation and file read. Event # IDs make duplicate replay safe, while skipping a new row is not safe. return 0 def service_result_paths(service_dir): output = [] for name in sorted(os.listdir(service_dir)): if not name.lower().endswith("results.jsonl"): continue path = os.path.join(service_dir, name) if os.path.isfile(path): output.extend(keycheck_result_paths(path)) seen = set() deduped = [] for path in output: if path not in seen: seen.add(path) deduped.append(path) return deduped def keycheck_ingest_tail_bytes(): value = os.getenv("KEYCHECK_DB_INGEST_TAIL_MB", "0") try: return max(0, int(float(value) * 1024 * 1024)) except (TypeError, ValueError): return 64 * 1024 * 1024 def file_signature(path, handle=None): details = os.fstat(handle.fileno()) if handle is not None else os.stat(path) return { "path": os.path.abspath(path), "device": int(getattr(details, "st_dev", 0) or 0), "size": int(details.st_size), "mtime_ns": int(getattr(details, "st_mtime_ns", int(details.st_mtime * 1_000_000_000))), "inode": int(getattr(details, "st_ino", 0) or 0), } def ingest_start_offset(path, file_state, signature): file_state = file_state or {} if checkpoint_matches_signature(path, file_state, signature): offset = int(file_state.get("offset", 0) or 0) if 0 <= offset <= signature["size"]: return offset, "offset" return 0, "full_replay" if file_state else "full_initial" def finding_uid_lookup_available(db): if str(os.getenv("KEYCHECK_FINDING_UID_LOOKUP", "1")).strip().lower() in ("0", "false", "no", "off"): return False if str(os.getenv("KEYCHECK_ASSUME_FINDING_UID_INDEX", "")).strip().lower() in ("1", "true", "yes", "on"): return True try: if db.conn.table_exists("finding_uid_map"): return True if getattr(db.conn, "is_postgres", False): row = db.conn.execute( '''SELECT 1 AS ok FROM pg_catalog.pg_indexes WHERE schemaname = 'public' AND tablename = 'findings' AND indexname = 'idx_findings_finding_uid' LIMIT 1''' ).fetchone() return bool(row) rows = db.conn.execute("PRAGMA index_list(findings)").fetchall() return any(str(row["name"] or "") == "idx_findings_finding_uid" for row in rows) except Exception: return False def finding_uid_index_available(db): if str(os.getenv("KEYCHECK_ASSUME_FINDING_UID_INDEX", "")).strip().lower() in ("1", "true", "yes", "on"): return True try: if getattr(db.conn, "is_postgres", False): row = db.conn.execute( '''SELECT 1 AS ok FROM pg_catalog.pg_indexes WHERE schemaname = 'public' AND tablename = 'findings' AND indexname = 'idx_findings_finding_uid' LIMIT 1''' ).fetchone() return bool(row) rows = db.conn.execute("PRAGMA index_list(findings)").fetchall() return any(str(row["name"] or "") == "idx_findings_finding_uid" for row in rows) except Exception: return False def select_finding_by_uid(db, finding_uid, lookup_available=True): if not lookup_available: return None if not finding_uid: return None if db.conn.table_exists("finding_uid_map"): rows = db.conn.execute( '''SELECT f.id, f.run_id, f.cycle_id, f.target_scan_id, f.source, f.query, f.target, f.detector_name, f.created_at FROM finding_uid_map m JOIN findings f ON f.id = m.finding_id WHERE m.finding_uid = ? LIMIT 2''', (finding_uid,), ).fetchall() if len(rows) == 1: return rows[0] if len(rows) > 1 or not finding_uid_index_available(db): return None rows = db.conn.execute( '''SELECT id, run_id, cycle_id, target_scan_id, source, query, target, detector_name, created_at FROM findings WHERE finding_uid = ? ORDER BY id DESC LIMIT 2''', (finding_uid,), ).fetchall() else: rows = db.conn.execute( '''SELECT id, run_id, cycle_id, target_scan_id, source, query, target, detector_name, created_at FROM findings WHERE finding_uid = ? ORDER BY id DESC LIMIT 2''', (finding_uid,), ).fetchall() return rows[0] if len(rows) == 1 else None def insert_keycheck_payload(db, service, path, line_offset, payload, line_cache, uid_lookup_available=True): status = status_from_payload(service, payload) finding = finding_from_payload(payload, line_cache) raw_secret = extract_raw_secret(finding or {}) finding_hash = sha256_text(raw_secret) if raw_secret else "" key_hash = payload.get("key_hash") or finding_hash secret_hash = payload.get("secret_hash") or finding_hash or key_hash key_masked = payload.get("key_masked") or mask_secret(raw_secret) if not key_hash and not key_masked: return "skipped" detector = payload.get("detector") or (finding or {}).get("DetectorName") or "" detector_secret_hash = payload.get("detector_secret_hash") or (sha256_text("|".join([str(detector or ""), secret_hash])) if secret_hash or detector else "") checked_at = payload.get("checked_at") or file_mtime_iso(path) source_line = payload.get("source") or payload.get("source_line") or "" finding_uid = payload.get("finding_uid") or (finding or {}).get("finding_uid") or "" event_id = payload.get("event_id") or sha256_text("|".join([ str(service or ""), os.path.abspath(path), str(line_offset), str(key_hash or key_masked or ""), str(status or ""), str(source_line or ""), ])) duplicate = db.conn.execute( "SELECT event_id FROM keycheck_event_map WHERE event_id = ? LIMIT 1", (event_id,), ).fetchone() if duplicate: return "duplicate" reservation = db.conn.execute( '''INSERT INTO keycheck_event_map (event_id, keycheck_result_id, created_at) VALUES (?, NULL, ?) ON CONFLICT(event_id) DO NOTHING''', (event_id, utc_now_iso()), ) if getattr(reservation, "rowcount", -1) == 0: return "duplicate" metadata = safe_metadata(payload) metadata["event_id"] = event_id metadata["finding_uid"] = finding_uid metadata["ingest_source_file"] = path metadata["ingest_source_offset"] = int(line_offset) message = payload.get("message") or "" error = payload.get("error") if not message and isinstance(error, dict): message = error.get("message") or json.dumps(error, ensure_ascii=False)[:1000] now = utc_now_iso() match = select_finding_by_uid(db, finding_uid, uid_lookup_available) keycheck_result_id = db.conn.insert_returning_id( '''INSERT INTO keycheck_results ( service, status, status_group, checked_at, key_hash, secret_hash, key_masked, finding_id, target_scan_id, cycle_id, run_id, source, query, target, detector_name, found_at, message, metadata_json, source_line, detector_secret_hash, event_id, finding_uid, link_status, link_attempts, linked_at, link_error, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( service, status, keycheck_status_group(status), checked_at, key_hash, secret_hash, key_masked, match["id"] if match else None, match["target_scan_id"] if match else None, match["cycle_id"] if match else None, match["run_id"] if match else None, match["source"] if match else None, match["query"] if match else None, match["target"] if match else None, match["detector_name"] if match else detector, match["created_at"] if match else None, str(message or "")[:1000].replace("\n", " "), json.dumps(metadata, ensure_ascii=False, default=str, sort_keys=True), source_line, detector_secret_hash, event_id, finding_uid, "linked" if match else "pending", 0, now if match else None, "", now, ), ) if keycheck_result_id: db.conn.execute( "UPDATE keycheck_event_map SET keycheck_result_id = ? WHERE event_id = ?", (keycheck_result_id, event_id), ) return "inserted" def ingest_keycheck_results_to_db(layout, services=None, max_rows=1000): db_path = layout.get("database_path") db_url = layout.get("database_url") or os.getenv("SCANNER_DB_URL") or os.getenv("DATABASE_URL") keycheck_dir = layout.get("keycheck_dir") if not keycheck_dir or not os.path.isdir(keycheck_dir): return 0 if not db_path and not db_url: return 0 db = ScannerDB(db_path=db_path, db_url=db_url, initialize=False) if not db.enabled: print(f"keycheck_results ingest skipped: unable to open DB: {db.db_display or db_path or 'configured database'}") return 0 set_linker_timeouts(db) try: db.require_runtime_safety_schema() except Exception as exc: db.close() raise RuntimeError(f'Keycheck ingest schema is incomplete; offline migration required: {exc}') from exc backfill_keycheck_event_map(db, env_int("KEYCHECK_EVENT_MAP_BACKFILL_ROWS", 0)) backfill_finding_uid_map(db, env_int("KEYCHECK_UID_MAP_BACKFILL_ROWS", 0)) uid_lookup_available = finding_uid_lookup_available(db) requested = [service for service in (services or []) if service and service != "all"] if set(requested) == set(SERVICES): requested = [] service_names = requested or sorted(name for name in os.listdir(keycheck_dir) if os.path.isdir(os.path.join(keycheck_dir, name))) state_identity = db_url or os.path.abspath(db_path or "") state_suffix = sha256_text(state_identity)[:12] if state_identity else "default" state_path = os.path.join(keycheck_dir, f"db_ingest_state_{state_suffix}.json") state = read_json_file(state_path, {}) or {} files_state = state.get("files") if isinstance(state.get("files"), dict) else {} inserted = 0 duplicates = 0 skipped = 0 errors = 0 rows_seen = 0 line_cache = {} max_rows = max(0, int(max_rows or 0)) try: for service in service_names: if max_rows and rows_seen >= max_rows: break service_dir = os.path.join(keycheck_dir, service) if not os.path.isdir(service_dir): continue service_paths = service_result_paths(service_dir) current_skips = { os.path.abspath(candidate): keycheck_result_current_skip(candidate) for candidate in service_paths if os.path.basename(candidate).lower().endswith("results.jsonl") } for path in service_paths: if max_rows and rows_seen >= max_rows: break if not os.path.isfile(path): continue key = os.path.abspath(path) processed = 0 max_line_bytes = max(1024, env_int("KEYCHECK_DB_INGEST_MAX_LINE_BYTES", 16 * 1024 * 1024)) with open(path, "rb") as f: signature = file_signature(path, f) offset, mode = ingest_start_offset(path, files_state.get(key), signature) offset = max(offset, min(int(current_skips.get(key, 0)), signature["size"])) final_offset = offset f.seek(offset) if mode.startswith("tail"): skipped_line = f.readline(max_line_bytes + 1) if skipped_line and not skipped_line.endswith(b"\n"): raise RuntimeError(f"oversized keycheck ingest tail boundary in {path}") final_offset = f.tell() while True: if max_rows and rows_seen >= max_rows: break line_offset = f.tell() raw_line = f.readline(max_line_bytes + 1) if not raw_line: break final_offset = f.tell() if len(raw_line) > max_line_bytes: final_offset = line_offset raise RuntimeError(f"oversized committed keycheck result at {path}:{line_offset}") if not raw_line.endswith(b"\n"): final_offset = line_offset raise RuntimeError(f"torn committed keycheck result at {path}:{line_offset}") try: payload = json.loads(raw_line.decode("utf-8", errors="replace")) except ValueError as exc: final_offset = line_offset kind = "torn" if not raw_line.endswith(b"\n") else "invalid" raise RuntimeError( f"{kind} committed keycheck result at {path}:{line_offset}" ) from exc if not isinstance(payload, dict): skipped += 1 continue rows_seen += 1 try: outcome = insert_keycheck_payload(db, service, path, line_offset, payload, line_cache, uid_lookup_available) if outcome == "inserted": inserted += 1 elif outcome == "duplicate": duplicates += 1 else: skipped += 1 db.conn.commit() processed += 1 except Exception as exc: errors += 1 try: db.conn.rollback() set_linker_timeouts(db) except Exception: pass final_offset = line_offset print(f"keycheck_results ingest row failed: service={service} file={path} offset={line_offset} error={str(exc)[:300]}") break checkpoint_signature = file_signature(path, f) files_state[key] = { **checkpoint_signature, "offset": int(final_offset), "mode": mode, "processed": int(processed), "updated_at": utc_now_iso(), } db.conn.commit() finally: write_json_file(state_path, {"files": files_state, "updated_at": utc_now_iso()}) db_display = db.db_display db.close() print(f"keycheck_results ingest: scanned={rows_seen} inserted={inserted} duplicates={duplicates} skipped={skipped} errors={errors} db={db_display}") return inserted def repair_keycheck_links(layout, services=None, batch_size=500, max_rows=0, max_attempts=3, since_days=14): db_path = layout.get("database_path") db_url = layout.get("database_url") or os.getenv("SCANNER_DB_URL") or os.getenv("DATABASE_URL") if not db_path and not db_url: raise SystemExit("global.database_path or global.database_url is required for keycheck link repair") db = ScannerDB(db_path=db_path, db_url=db_url, initialize=False) if not db.enabled: raise SystemExit(f"Unable to open scanner DB: {db.db_display or db_path or 'configured database'}") service_detectors = { "anthropic": ("anthropic",), "aws": ("aws",), "azure": ( "azure", "azureopenai", "azurecontainerregistry", "azurefoundry", "azurefoundryendpointbeforekey", "azurefoundrykeybeforeendpoint", ), "deepseek": ("deepseek", "deepseekapikey", "deepseek_api_key"), "dockerhub": ("dockerhub",), "gcp": ("gcp", "gcpapplicationdefaultcredentials"), "gemini": ("googleai", "googleaistudio"), "groq": ("groq",), "github": ("github", "githuboauth2"), "gitlab": ("gitlab",), "kimi": ("kimimoonshot", "moonshotai", "moonshot", "kimi"), "openai": ("openai",), "openrouter": ("openrouter",), "provider_resolver": ( "qwendashscope", "qwen_dashscope", "qwen", "dashscope", "deepseek", "deepseekapikey", "deepseek_api_key", "kimimoonshot", "moonshotai", "moonshot", "kimi", "zaiglm", ), "qwen": ("qwendashscope", "qwen_dashscope"), "replicate": ("replicate",), "xai": ("xai",), "huggingface": ("huggingface",), "zai": ("zaiglm",), } set_linker_timeouts(db) try: db.require_runtime_safety_schema() except Exception as exc: db.close() raise RuntimeError(f'Keycheck repair schema is incomplete; offline migration required: {exc}') from exc backfill_finding_uid_map(db, env_int("KEYCHECK_UID_MAP_BACKFILL_ROWS", 5000)) requested = [service for service in (services or []) if service and service != "all"] if set(requested) == set(SERVICES): requested = [] repaired = 0 repaired_metadata = 0 not_found = 0 errors = 0 scanned = 0 line_cache = {} uid_lookup_available = finding_uid_lookup_available(db) window_size = max(batch_size * 5, 100) max_scanned_candidates = max(max_rows * 100, window_size) if max_rows else env_int("KEYCHECK_REPAIR_MAX_SCAN_CANDIDATES", 10000) cutoff = (datetime.now(timezone.utc) - timedelta(days=max(0, int(since_days or 0)))).isoformat(timespec="seconds") last_id = 9223372036854775807 scanned_candidates = 0 while True: if max_rows and scanned >= max_rows: break if max_scanned_candidates and scanned_candidates >= max_scanned_candidates: break params = [last_id, window_size] try: candidates = db.conn.execute(f''' SELECT id, service, checked_at, metadata_json, source_line, detector_name, detector_secret_hash, key_hash, secret_hash, COALESCE(NULLIF(secret_hash, ''), NULLIF(key_hash, '')) AS hash_value, COALESCE(link_status, 'pending') AS link_status, COALESCE(link_attempts, 0) AS link_attempts, finding_id, finding_uid FROM keycheck_results WHERE id < ? ORDER BY id DESC LIMIT ? ''', params).fetchall() except Exception as exc: try: db.conn.rollback() except Exception: pass print(f"link repair stopped: candidate fetch failed after scanned={scanned}: {str(exc)[:300]}") break if not candidates: break scanned_candidates += len(candidates) last_id = candidates[-1]['id'] rows = [] for row in candidates: if row['checked_at'] and row['checked_at'] < cutoff: continue if requested and row['service'] not in requested: continue if row['finding_id'] is not None: continue if row['link_status'] not in ('pending', 'error', 'not_found', 'source_only'): continue if int(row['link_attempts'] or 0) >= max_attempts: continue rows.append(row) if not rows: continue for row in rows: if max_rows and scanned >= max_rows: break scanned += 1 try: detectors = service_detectors.get(row["service"], ()) try: metadata = json.loads(row["metadata_json"] or "{}") except json.JSONDecodeError: metadata = {} row_finding_uid = row["finding_uid"] or metadata.get("finding_uid") or "" if row_finding_uid: uid_match = select_finding_by_uid(db, row_finding_uid, uid_lookup_available) if uid_match: db.conn.execute(''' UPDATE keycheck_results SET finding_id = ?, run_id = ?, cycle_id = ?, target_scan_id = ?, source = ?, query = ?, target = ?, detector_name = ?, found_at = ?, link_status = 'linked', finding_uid = ?, link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = ?, link_error = '' WHERE id = ? ''', ( uid_match["id"], uid_match["run_id"], uid_match["cycle_id"], uid_match["target_scan_id"], uid_match["source"], uid_match["query"], uid_match["target"], uid_match["detector_name"], uid_match["created_at"], row_finding_uid, utc_now_iso(), row["id"], )) db.conn.commit() repaired += 1 continue source_ref = row["source_line"] or metadata.get("source_line") or metadata.get("source") finding = None inferred = {} stored_hashes = [value for value in (row["secret_hash"], row["key_hash"]) if value] source_ref_is_finding = False source_hash_mismatch = False if source_ref: path, line_number = parse_source_line(source_ref) source_ref_is_finding = bool(path and line_number) finding = load_jsonl_line(path, line_number, line_cache) if isinstance(finding, dict) and isinstance(finding.get("finding"), dict): finding = finding["finding"] source_hash = row["secret_hash"] or row["key_hash"] or "" if finding and stored_hashes: raw_secret = extract_raw_secret(finding) raw_hash = sha256_text(raw_secret) if raw_secret else "" if raw_hash and raw_hash not in stored_hashes: finding = None source_hash_mismatch = True elif raw_hash: source_hash = raw_hash inferred = infer_finding_attribution(finding) source_match = None if finding and source_hash: location = extract_finding_location(finding) detector = str(finding.get("DetectorName") or finding.get("DetectorType") or row["detector_name"] or "") clauses = ["secret_hash = ?"] params = [source_hash] if detector: clauses.append("LOWER(detector_name) = LOWER(?)") params.append(detector) exact_clauses = list(clauses) exact_params = list(params) exact_clauses.append("raw_finding_json = ?") exact_params.append(json_dumps(finding)) source_matches = db.conn.execute(f''' SELECT id, run_id, cycle_id, target_scan_id, source, query, target, detector_name, created_at FROM findings WHERE {' AND '.join(exact_clauses)} ORDER BY id DESC LIMIT 2 ''', exact_params).fetchall() if len(source_matches) == 1: source_match = source_matches[0] if location.get("file_path"): clauses.append("file_path = ?") params.append(location.get("file_path")) if location.get("line_number"): clauses.append("line_number = ?") params.append(str(location.get("line_number"))) if location.get("commit_hash"): clauses.append("commit_hash = ?") params.append(location.get("commit_hash")) if not source_match and (location.get("commit_hash") or inferred.get("source") or inferred.get("target")) and len(clauses) > 2: if inferred.get("source"): clauses.append("source = ?") params.append(inferred.get("source")) if inferred.get("target"): clauses.append("target = ?") params.append(inferred.get("target")) source_matches = db.conn.execute(f''' SELECT id, run_id, cycle_id, target_scan_id, source, query, target, detector_name, created_at FROM findings WHERE {' AND '.join(clauses)} ORDER BY id DESC LIMIT 2 ''', params).fetchall() if len(source_matches) == 1: source_match = source_matches[0] if source_match: db.conn.execute(''' UPDATE keycheck_results SET finding_id = ?, run_id = ?, cycle_id = ?, target_scan_id = ?, source = ?, query = ?, target = ?, detector_name = ?, found_at = ?, link_status = 'linked', link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = ?, link_error = '' WHERE id = ? ''', ( source_match["id"], source_match["run_id"], source_match["cycle_id"], source_match["target_scan_id"], source_match["source"], source_match["query"], source_match["target"], source_match["detector_name"], source_match["created_at"], utc_now_iso(), row["id"], )) db.conn.commit() repaired += 1 continue if row_finding_uid or source_ref_is_finding: if inferred.get("source") or inferred.get("target"): db.conn.execute(''' UPDATE keycheck_results SET source = COALESCE(source, ?), query = COALESCE(query, ?), target = COALESCE(target, ?), detector_name = COALESCE(NULLIF(detector_name, ''), ?), found_at = COALESCE(NULLIF(found_at, ''), ?), link_status = 'source_only', link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = ?, link_error = '' WHERE id = ? ''', ( inferred.get("source"), inferred.get("query"), inferred.get("target"), inferred.get("detector_name"), inferred.get("found_at"), utc_now_iso(), row["id"], )) db.conn.commit() repaired_metadata += 1 else: update_link_status(db, row["id"], "not_found", "exact finding_uid/source finding not present in DB") db.conn.commit() not_found += 1 continue hashes = [] if row["secret_hash"]: hashes.append(("secret_hash", row["secret_hash"])) if not row["secret_hash"] and row["key_hash"]: hashes.append(("secret_hash", row["key_hash"])) if row["detector_secret_hash"] and row["service"] not in ("gemini", "qwen"): hashes.insert(0, ("detector_secret_hash", row["detector_secret_hash"])) match = None for column, hash_value in hashes: params = [hash_value] detector_clause = "" extra_name = db.conn.json_extract('raw_finding_json', '$.ExtraData.name') if column == "secret_hash" and row["service"] == "gemini": detector_clause = f""" AND ( LOWER(detector_name) IN (?, ?) OR ( LOWER(detector_name) = 'customregex' AND LOWER(COALESCE({extra_name}, '')) = 'googleaistudio' ) ) """ params.extend(detectors) elif column == "secret_hash" and row["service"] == "qwen": detector_clause = f""" AND ( LOWER(detector_name) IN ({','.join('?' for _ in detectors)}) OR ( LOWER(detector_name) = 'customregex' AND LOWER(COALESCE({extra_name}, '')) IN ({','.join('?' for _ in detectors)}) ) ) """ params.extend((*detectors, *detectors)) elif column == "secret_hash" and detectors: detector_clause = " AND LOWER(detector_name) IN ({})".format(','.join('?' for _ in detectors)) params.extend(detectors) match_rows = db.conn.execute(f''' SELECT id, run_id, cycle_id, target_scan_id, source, query, target, detector_name, created_at FROM findings WHERE {column} = ? {detector_clause} ORDER BY id DESC LIMIT 2 ''', params).fetchall() if len(match_rows) == 1: match = match_rows[0] break if match: db.conn.execute(''' UPDATE keycheck_results SET finding_id = ?, run_id = ?, cycle_id = ?, target_scan_id = ?, source = ?, query = ?, target = ?, detector_name = ?, found_at = ?, link_status = 'linked', link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = ?, link_error = '' WHERE id = ? ''', ( match["id"], match["run_id"], match["cycle_id"], match["target_scan_id"], match["source"], match["query"], match["target"], match["detector_name"], match["created_at"], utc_now_iso(), row["id"], )) db.conn.commit() repaired += 1 continue if inferred.get("source") or inferred.get("target"): db.conn.execute(''' UPDATE keycheck_results SET source = COALESCE(source, ?), query = COALESCE(query, ?), target = COALESCE(target, ?), detector_name = COALESCE(NULLIF(detector_name, ''), ?), found_at = COALESCE(NULLIF(found_at, ''), ?), link_status = 'source_only', link_attempts = COALESCE(link_attempts, 0) + 1, linked_at = ?, link_error = '' WHERE id = ? ''', ( inferred.get("source"), inferred.get("query"), inferred.get("target"), inferred.get("detector_name"), inferred.get("found_at"), utc_now_iso(), row["id"], )) db.conn.commit() repaired_metadata += 1 continue update_link_status(db, row["id"], "not_found", "no matching finding") db.conn.commit() not_found += 1 except Exception as exc: try: db.conn.rollback() set_linker_timeouts(db) update_link_status(db, row["id"], "error", str(exc)) db.conn.commit() errors += 1 except Exception as mark_exc: try: db.conn.rollback() except Exception: pass print(f"link repair error: row={row['id']} error={str(exc)[:200]} mark_error_failed={str(mark_exc)[:200]}") errors += 1 db.conn.commit() print(f"link repair progress: scanned={scanned} linked={repaired} source_only={repaired_metadata} not_found={not_found} errors={errors}") db.conn.commit() db_display = db.db_display db.close() print(f"keycheck_results link repair: scanned={scanned} linked={repaired} source_only={repaired_metadata} not_found={not_found} errors={errors} db={db_display}") return repaired + repaired_metadata def service_summary(service, layout): output_dir = os.path.join(layout["keycheck_dir"], service) counts = { "service": service, "alive": 0, "alive_rate_limited": 0, "checked": 0, "dead": 0, "limited": 0, "restricted": 0, "network": 0, "unknown": 0, "no_quota": 0, "no_balance": 0, "no_username": 0, "no_target": 0, "results": 0, "total_status": 0, "output_dir": output_dir, "updated_at": utc_now_iso(), } file_counts = {} if not os.path.isdir(output_dir): counts["file_counts"] = file_counts return counts for name in sorted(os.listdir(output_dir)): path = os.path.join(output_dir, name) if not os.path.isfile(path) or not name.lower().endswith((".txt", ".jsonl")): continue lower = name.lower() if lower.endswith("results.jsonl") and os.path.getsize(path) > 5 * 1024 * 1024: file_counts[name] = None continue is_checked = lower.endswith("checked.txt") line_count, unique_checked = file_line_stats(path, unique_keys=is_checked) file_counts[name] = line_count if lower.endswith("results.jsonl"): counts["results"] += line_count if is_checked: counts["checked"] += int(unique_checked or 0) status_bucket = None if "aliveratelimited" in lower: status_bucket = "alive_rate_limited" elif lower.endswith("alive.txt") or any(item in lower for item in ("bedrock", "admin", "canary", "vertex", "foundryllm")): status_bucket = "alive" elif "noquota" in lower: status_bucket = "no_quota" elif "nobalance" in lower: status_bucket = "no_balance" elif "nousername" in lower: status_bucket = "no_username" elif "notarget" in lower: status_bucket = "no_target" elif any(item in lower for item in ("restricted", "disabled", "accessdenied", "quarantined")): status_bucket = "restricted" elif any(item in lower for item in ("expired", "leaked", "revoked", "dead")): status_bucket = "dead" elif "ratelimited" in lower or "limited" in lower: status_bucket = "limited" elif any(item in lower for item in ("network", "badendpoint", "unresolved", "nocontext", "unknown")): status_bucket = "unknown" if "network" not in lower else "network" if status_bucket == "alive_rate_limited": counts["alive_rate_limited"] += line_count counts["total_status"] += line_count elif status_bucket == "alive": counts["alive"] += line_count counts["total_status"] += line_count elif status_bucket == "dead": counts["dead"] += line_count counts["total_status"] += line_count elif status_bucket == "limited": counts["limited"] += line_count counts["total_status"] += line_count elif status_bucket == "restricted": counts["restricted"] += line_count counts["total_status"] += line_count elif status_bucket == "network": counts["network"] += line_count counts["total_status"] += line_count elif status_bucket == "unknown": counts["unknown"] += line_count counts["total_status"] += line_count elif status_bucket == "no_quota": counts["no_quota"] += line_count counts["total_status"] += line_count elif status_bucket == "no_balance": counts["no_balance"] += line_count counts["total_status"] += line_count elif status_bucket == "no_username": counts["no_username"] += line_count counts["total_status"] += line_count elif status_bucket == "no_target": counts["no_target"] += line_count counts["total_status"] += line_count counts["file_counts"] = file_counts return counts def write_summary( layout, services, summary_tsv=None, summary_json=None, alive_summary_tsv=None, input_mode=None, ): ensure_private_directory(layout["keycheck_dir"], reject_reparse=True) summary_tsv = summary_tsv or os.path.join(layout["keycheck_dir"], "summary.tsv") summary_json = summary_json or os.path.join(layout["keycheck_dir"], "summary.json") alive_summary_tsv = alive_summary_tsv or os.path.join(layout["keycheck_dir"], "alive_summary.tsv") if str(input_mode or os.getenv('KEYCHECK_INPUT_MODE') or 'jsonl').lower() == 'postgres': db = ScannerDB(db_url=layout.get('database_url') or os.getenv('KEYCHECK_DB_URL'), initialize=False) if not db.enabled: raise RuntimeError('PostgreSQL keycheck summary database is unavailable') try: placeholders = ','.join('?' for _ in services) state_rows = db.conn.execute( f'''SELECT service, status, status_group, COUNT(*) AS count, MAX(updated_at) AS updated_at FROM keycheck_current_state WHERE service IN ({placeholders}) GROUP BY service, status, status_group''', tuple(services), ).fetchall() projection_rows = db.conn.execute( f'''SELECT c.service, c.secret_text, c.secret_json, s.status, s.status_group, s.checked_at, s.metadata_json FROM keycheck_current_state s JOIN keycheck_credentials c ON c.id = s.credential_id WHERE c.service IN ({placeholders}) ORDER BY c.service, c.id''', tuple(services), ).fetchall() managed_projection_rows = db.conn.execute( f'''SELECT service, secret_text, secret_json FROM keycheck_credentials WHERE service IN ({placeholders})''', tuple(services), ).fetchall() history_rows = db.conn.execute( f'''SELECT service, COUNT(*) AS count FROM keycheck_results WHERE service IN ({placeholders}) GROUP BY service''', tuple(services), ).fetchall() db.conn.commit() finally: db.close() grouped = {service: {} for service in services} rate_limited_alive = {service: 0 for service in services} updated = {service: '' for service in services} for row in state_rows: service_counts = grouped[row['service']] service_counts[row['status_group']] = ( int(service_counts.get(row['status_group'], 0)) + int(row['count'] or 0) ) if row['status'] == 'VALID_RATE_LIMITED': rate_limited_alive[row['service']] += int(row['count'] or 0) updated[row['service']] = max(updated[row['service']], str(row['updated_at'] or '')) history = {row['service']: int(row['count'] or 0) for row in history_rows} rows = [] for service in services: counts = grouped.get(service) or {} alive_rate_limited = int(rate_limited_alive.get(service, 0)) row = { 'service': service, 'alive': max(0, int(counts.get('alive', 0)) - alive_rate_limited), 'alive_rate_limited': alive_rate_limited, 'checked': sum(counts.values()), 'dead': int(counts.get('dead', 0)), 'limited': int(counts.get('limited', 0)), 'restricted': int(counts.get('restricted', 0)), 'network': int(counts.get('network', 0)), 'unknown': int(counts.get('unknown', 0)), 'no_quota': int(counts.get('no_quota', 0)), 'no_balance': int(counts.get('no_balance', 0)), 'no_username': int(counts.get('no_username', 0)), 'no_target': int(counts.get('no_context', 0)), 'results': history.get(service, 0), 'total_status': sum(counts.values()), 'updated_at': updated.get(service) or utc_now_iso(), 'output_dir': os.path.join(layout['keycheck_dir'], service), } rows.append(row) projection = project_postgres_status_files( layout, services, projection_rows, managed_rows=managed_projection_rows, ) print( 'status_projection: ' f"rows={projection['projected_rows']} files={projection['changed_files']} " f"skipped={projection['skipped_rows']}" ) else: rows = [service_summary(service, layout) for service in services] columns = [ "service", "alive", "alive_rate_limited", "checked", "dead", "limited", "restricted", "network", "unknown", "no_quota", "no_balance", "no_username", "no_target", "results", "total_status", "updated_at", "output_dir", ] with private_atomic_writer(summary_tsv) as f: f.write("\t".join(columns) + "\n") for row in rows: f.write("\t".join(str(row.get(column, "")) for column in columns) + "\n") with private_atomic_writer(summary_json) as f: json.dump({"updated_at": utc_now_iso(), "services": rows}, f, ensure_ascii=False, indent=2) with private_atomic_writer(alive_summary_tsv) as f: f.write("service\talive\talive_rate_limited\talive_total\tupdated_at\toutput_dir\n") for row in rows: alive_total = int(row.get("alive", 0)) + int(row.get("alive_rate_limited", 0)) f.write("\t".join([ str(row.get("service", "")), str(row.get("alive", 0)), str(row.get("alive_rate_limited", 0)), str(alive_total), str(row.get("updated_at", "")), str(row.get("output_dir", "")), ]) + "\n") print(f"summary_tsv: {summary_tsv}") print(f"summary_json: {summary_json}") print(f"alive_summary_tsv: {alive_summary_tsv}") return rows def maybe_add(command, flag, value): if value: command.extend([flag, value]) def list_value(value): if not value: return [] if isinstance(value, str): return [item for item in value.split() if item] return [str(item) for item in value] 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 apply_keycheck_config_defaults(args, keychecks_config): for flag in ( "retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "no_resource_probe", "recheck_all", "summary_only", "no_summary", ): if not getattr(args, flag, False) and bool_config(keychecks_config.get(flag), False): setattr(args, flag, True) if not args.max_keys and keychecks_config.get("max_keys"): args.max_keys = int(keychecks_config.get("max_keys") or 0) def service_extra_args(service, keychecks_config): output = [] output.extend(list_value(keychecks_config.get("common_args"))) per_service = keychecks_config.get("service_args") or keychecks_config.get("per_service_args") or {} if isinstance(per_service, dict): output.extend(list_value(per_service.get(service))) return output def service_supports_proxy(service): return bool((SERVICE_CAPABILITIES.get(service) or {}).get("proxy", True)) def service_supports_flag(service, flag): capabilities = SERVICE_CAPABILITIES.get(service) or {} return flag in capabilities.get("flags", set()) def nonempty_private_status_file(output_dir, filename): output_dir = os.path.abspath(output_dir) path = os.path.abspath(os.path.join(output_dir, filename)) if os.path.commonpath((output_dir, path)) != output_dir: raise RuntimeError(f"keycheck status path escapes provider output directory: {filename}") if not os.path.lexists(path): return False require_private_file(path) details = os.stat(path, follow_symlinks=False) if not stat.S_ISREG(details.st_mode): raise RuntimeError(f"keycheck status path is not a regular file: {path}") return details.st_size > 0 def retry_status_file_names(service, flag): override = RETRY_STATUS_FILE_OVERRIDES.get((service, flag)) if override is not None: return override return tuple( f"{service}{suffix}.txt" for suffix in RETRY_STATUS_DEFAULT_SUFFIXES.get(flag, ()) ) def provider_replay_required(service, args, output_dir, extra_args=None): extra_flags = {str(value) for value in (extra_args or []) if str(value).startswith("--")} if service == "qwen" and nonempty_private_status_file( output_dir, STATUS_TRANSACTION_JOURNAL_FILENAME, ): return True if service_supports_flag(service, "recheck_all") and ( getattr(args, "recheck_all", False) or "--recheck-all" in extra_flags ): return True for flag in RETRY_STATUS_DEFAULT_SUFFIXES: requested = getattr(args, flag, False) or "--" + flag.replace("_", "-") in extra_flags if not requested or not service_supports_flag(service, flag): continue filenames = retry_status_file_names(service, flag) if any(nonempty_private_status_file(output_dir, filename) for filename in filenames): return True return False def auto_repair_config(keychecks_config): value = keychecks_config.get("repair_links") or keychecks_config.get("auto_repair_links") if isinstance(value, bool): value = {"enabled": value} if not isinstance(value, dict): value = {} return { "enabled": bool_config(value.get("enabled"), False), "services": value.get("services"), "batch_size": max(1, int(value.get("batch_size", 25) or 25)), "max_rows": max(0, int(value.get("max_rows", 250) or 250)), "max_attempts": max(1, int(value.get("max_attempts", 3) or 3)), "since_days": max(0, int(value.get("since_days", 2) or 2)), } def db_ingest_config(keychecks_config): value = keychecks_config.get("db_ingest") or keychecks_config.get("ingest_results") if isinstance(value, bool): value = {"enabled": value} if not isinstance(value, dict): value = {} return { "enabled": bool_config(value.get("enabled"), True), "services": value.get("services"), "max_rows": max(0, int(value.get("max_rows", 1000) or 1000)), } def maybe_ingest_keycheck_results(layout, requested_services, keychecks_config): cfg = db_ingest_config(keychecks_config) if not cfg["enabled"]: return 0 services = service_list(str(cfg["services"])) if cfg.get("services") else requested_services print(f"db_ingest_keycheck_results: services={','.join(services)} max_rows={cfg['max_rows']}") return ingest_keycheck_results_to_db(layout, services, max_rows=cfg["max_rows"]) def maybe_auto_repair_links(layout, requested_services, keychecks_config): cfg = auto_repair_config(keychecks_config) if not cfg["enabled"]: return services = service_list(str(cfg["services"])) if cfg.get("services") else requested_services print( "auto_repair_links: " f"services={','.join(services)} batch_size={cfg['batch_size']} " f"max_rows={cfg['max_rows']} since_days={cfg['since_days']}" ) repair_keycheck_links( layout, services, batch_size=cfg["batch_size"], max_rows=cfg["max_rows"], max_attempts=cfg["max_attempts"], since_days=cfg["since_days"], ) def sync_service_outputs(script_dir, output_dir): ensure_private_directory(output_dir, reject_reparse=True) for name in os.listdir(script_dir): if not name.lower().endswith((".txt", ".jsonl")): continue src = os.path.join(script_dir, name) dst = os.path.join(output_dir, name) if os.path.isfile(src): with open(src, "rb") as fsrc, private_atomic_writer(dst, binary=True) as fdst: while True: chunk = fsrc.read(1024 * 1024) if not chunk: break fdst.write(chunk) def _provider_diagnostic_path(output_dir): return os.path.join(output_dir, '.provider-last-run.log') def _record_provider_startup_failure(service, output_dir, exc): diagnostic_path = _provider_diagnostic_path(output_dir) payload = f'provider startup failed: {type(exc).__name__}\n'.encode('ascii', errors='replace') with private_atomic_writer(diagnostic_path, binary=True, suffix='.provider.tmp') as diagnostic: diagnostic.write(payload) print( f'provider_exit: service={service} code=1 outcome=infrastructure_failure ' f'diagnostic={diagnostic_path}', flush=True, ) return 1 def _stop_failed_provider_process(process): if process is None: return True try: process.terminate() except Exception: pass try: if process.wait(timeout=5) is not None: return True except subprocess.TimeoutExpired: pass except Exception: pass try: process.kill() if process.wait(timeout=5) is not None: return True except Exception: pass # Keep the exact owner reachable and fail closed until this runner exits. _unconfirmed_provider_processes.append(process) return False def run_service(service, script_path, args, layout, extra_args=None): if _unconfirmed_provider_processes: raise RuntimeError('provider containment cleanup was not confirmed') deadline = float(getattr(args, '_provider_deadline', time.monotonic() + 1800)) if not math.isfinite(deadline): raise ValueError('provider deadline must be finite') project_dir = os.path.dirname(os.path.abspath(__file__)) script = os.path.join(project_dir, script_path) output_dir = os.path.join(layout["keycheck_dir"], service) ensure_private_directory(output_dir, reject_reparse=True) startup_errors = (LifecycleAuthorityError, OSError, subprocess.SubprocessError) try: metadata = require_active_supervisor_child(args.config, child_kind='keycheck', require_dsn=True) except startup_errors as exc: return _record_provider_startup_failure(service, output_dir, exc) canonical_dsn = os.getenv('TRUF_MANAGED_POSTGRES_DSN') or '' extra_args = list(extra_args or []) input_mode = str(getattr(args, 'input_mode', 'postgres') or 'postgres').lower() input_file = args.input or os.path.join(layout["results_dir"], "found_secrets.jsonl") proxy_file = args.proxy_file or layout["proxy_file"] if input_mode == 'postgres' and any( str(value).lower() in ('--input', '--plain') or str(value).lower().startswith(('--input=', '--plain=')) for value in extra_args ): raise ValueError('PostgreSQL provider configuration forbids compatibility input/plain arguments') if not os.path.exists(script): print(f"{service}: script not found yet: {script}") return 2 command = [ sys.executable, '-I', '-S', '-B', os.path.join(project_dir, 'child_bootstrap.py'), 'keycheck-provider', str(script_path).replace('\\', '/'), '--', ] if input_mode == 'jsonl': maybe_add(command, "--input", input_file) if service_supports_proxy(service): maybe_add(command, "--proxy-file", proxy_file) if args.max_keys and input_mode == 'jsonl': command.extend(["--max-keys", str(args.max_keys)]) for flag in ("retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "no_resource_probe", "recheck_all"): if ( getattr(args, flag, False) and service_supports_flag(service, flag) and (input_mode == 'jsonl' or flag == 'no_resource_probe') ): command.append("--" + flag.replace("_", "-")) command.extend(extra_args) env = os.environ.copy() strip_supervisor_credentials(env) env.update(supervised_child_environment(metadata, canonical_dsn, 'keycheck-provider')) env['SCANNER_DB_URL'] = canonical_dsn env['DATABASE_URL'] = canonical_dsn env["KEYCHECK_INPUT_FILE"] = input_file env["KEYCHECK_INPUT_MODE"] = input_mode env["KEYCHECK_PROXY_FILE"] = proxy_file env["KEYCHECK_OUTPUT_DIR"] = output_dir env["KEYCHECK_SERVICE"] = service env["KEYCHECK_STATE_DIR"] = output_dir env["KEYCHECK_DB_PATH"] = layout.get("database_path", "") env["KEYCHECK_DB_URL"] = canonical_dsn env["KEYCHECK_DB_INLINE"] = "0" if input_mode == 'postgres': env['KEYCHECK_PROVIDER_SLICE_KEYS'] = str(max(1, int(args.max_keys or 1))) env['KEYCHECK_RESULT_PROJECTION_RESERVE_BYTES'] = str(int( layout.get('keycheck_result_projection_reserve_bytes', 3 * 1024 * 1024) )) env['KEYCHECK_PROJECTION_MAX_ITEMS'] = str(int( layout.get('projection_backlog_max_items', 10000) )) env['KEYCHECK_PROJECTION_MAX_BYTES'] = str(int( layout.get('projection_backlog_max_bytes', 2 * 1024 * 1024 * 1024) )) env["PYTHONIOENCODING"] = "utf-8" env["PYTHONPATH"] = os.pathsep.join([project_dir, os.path.join(project_dir, "keycheckers"), env.get("PYTHONPATH", "")]) env.pop("KEYCHECK_DISABLE_HIGH_WATERMARK", None) if input_mode == 'jsonl' and provider_replay_required(service, args, output_dir, extra_args): env["KEYCHECK_DISABLE_HIGH_WATERMARK"] = "1" if time.monotonic() >= deadline: print(f'provider_exit: service={service} code=124 outcome=deadline_exceeded', flush=True) return 124 print(f"{service}: {' '.join(redact_argv(command))}") try: completed = OwnedProcess( command, cwd=os.path.dirname(script), env=env, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, creationflags=subprocess.CREATE_NO_WINDOW if os.name == "nt" else 0, ) except startup_errors as exc: return _record_provider_startup_failure(service, output_dir, exc) stream = getattr(completed, 'stdout', None) output_tail = bytearray() def copy_provider_output(): read = getattr(stream, 'read1', stream.read) while True: chunk = read(64 * 1024) if not chunk: return output_tail.extend(chunk) if len(output_tail) > 64 * 1024: del output_tail[:-64 * 1024] try: binary_output = getattr(sys.stdout, 'buffer', None) if binary_output is not None: binary_output.write(chunk) binary_output.flush() else: sys.stdout.write(chunk.decode('utf-8', errors='replace')) sys.stdout.flush() except (OSError, ValueError): pass output_thread = None exited = False cleanup_confirmed = True failure = None timed_out = False try: if stream is not None: output_thread = threading.Thread(target=copy_provider_output, name=f'{service}-provider-output', daemon=True) output_thread.start() code = completed.wait(timeout=max(0.0, deadline - time.monotonic())) exited = True except startup_errors as exc: timed_out = isinstance(exc, subprocess.TimeoutExpired) failure = type(exc).__name__ code = 124 if timed_out else 1 finally: if not exited: cleanup_confirmed = _stop_failed_provider_process(completed) if output_thread is not None: output_thread.join(timeout=5) if output_thread.is_alive(): # Closing a buffered pipe here can block on the reader's lock. failure = failure or 'OutputDrainTimeout' code = code or 1 if output_thread is not None or failure: diagnostic_path = _provider_diagnostic_path(output_dir) with private_atomic_writer(diagnostic_path, binary=True, suffix='.provider.tmp') as diagnostic: diagnostic.write(bytes(output_tail[-64 * 1024:])) if failure: diagnostic.write(f'\nprovider runtime failed: {failure}\n'.encode('ascii')) if not cleanup_confirmed: raise RuntimeError('provider containment cleanup was not confirmed') outcome = ( 'deadline_exceeded' if timed_out else 'ok' if code == 0 else 'capacity_blocked' if code == KEYCHECK_CAPACITY_BLOCKED_EXIT else 'infrastructure_failure' ) detail = f" diagnostic={diagnostic_path}" if output_thread is not None and code != 0 else "" print(f"provider_exit: service={service} code={code} outcome={outcome}{detail}", flush=True) return code def run_provider_services( requested, args, layout, keychecks_config, work_probe=None, monotonic=time.monotonic, ): exit_code = 0 provider_failures = [] if not requested: return exit_code, provider_failures configured_workers = min(4, max(1, int( keychecks_config.get('scheduler_workers', 4) or 4 ))) worker_count = min(configured_workers, len(requested)) batch_keys = max(1, int(keychecks_config.get('scheduler_batch_keys', 1000) or 1000)) explicit_limit = max(0, int(getattr(args, 'max_keys', 0) or 0)) deadline = monotonic() + max( 1.0, float(keychecks_config.get('scheduler_deadline_sec', 1800) or 1800), ) if not math.isfinite(deadline): raise ValueError('provider deadline must be finite') probe_db = None if work_probe is None and layout.get('database_url'): probe_db = ScannerDB(db_url=layout['database_url'], initialize=False) if not probe_db.enabled: raise RuntimeError('keycheck scheduler work database is unavailable') probe_db.set_application_name('truf-keycheck-scheduler') probe_db.require_runtime_safety_schema() probe_db.require_final_cutover() work_probe = probe_db.keycheck_service_has_work def run_bounded(service): service_args = copy.copy(args) service_args.max_keys = explicit_limit or batch_keys service_args._provider_deadline = deadline return run_service( service, SERVICES[service], service_args, layout, service_extra_args(service, keychecks_config), ) outcomes = {service: 0 for service in requested} queue = deque( service for service in requested if work_probe is None or work_probe(service) ) active = {} try: with ThreadPoolExecutor(max_workers=worker_count, thread_name_prefix='keycheck-provider') as executor: while queue or active: while queue and len(active) < worker_count and monotonic() < deadline: service = queue.popleft() active[executor.submit(run_bounded, service)] = service if not active: break done, _ = wait(tuple(active), return_when=FIRST_COMPLETED) for future in done: service = active.pop(future) try: code = future.result() except Exception as exc: print( f'provider_exit: service={service} code=1 ' f'outcome=infrastructure_failure error={type(exc).__name__}' ) code = 1 capacity_blocked = code == KEYCHECK_CAPACITY_BLOCKED_EXIT if capacity_blocked: print( f'provider_backpressure: service={service} ' 'reason=projection_capacity_saturated', flush=True, ) code = 0 if code and not outcomes[service]: outcomes[service] = code if ( not code and not capacity_blocked and not explicit_limit and monotonic() < deadline and work_probe is not None and work_probe(service) ): queue.append(service) finally: if probe_db is not None: probe_db.close() for service in requested: code = outcomes[service] if code: if not exit_code: exit_code = code provider_failures.append((service, code)) return exit_code, provider_failures def enqueue_requested_postgres_rechecks(layout, requested, args): if str(getattr(args, 'input_mode', 'postgres')).lower() != 'postgres': return 0 groups = set() mapping = { 'retry_network': {'network'}, 'retry_limited': {'limited', 'no_balance'}, 'retry_unknown': {'unknown', 'no_context'}, 'retry_restricted': {'restricted'}, 'retry_no_balance': {'no_balance'}, 'retry_valid': {'alive'}, } for flag, values in mapping.items(): if getattr(args, flag, False): groups.update(values) if not groups and not getattr(args, 'recheck_all', False): return 0 database_url = layout.get('database_url') if not database_url: return 0 db = ScannerDB(db_url=database_url, initialize=False) if not db.enabled: raise RuntimeError('PostgreSQL recheck candidate database is unavailable') total = 0 try: db.set_application_name('truf-keycheck-recheck-generator') db.require_runtime_safety_schema() db.require_final_cutover() for service in requested: total += db.enqueue_keycheck_rechecks( service, None if getattr(args, 'recheck_all', False) else groups, max_items=max(1, int( getattr(args, 'max_keys', 0) or ((layout.get('keycheck_recheck_batch_items') or 10000)) )), queue_max_items=int(layout.get('keycheck_queue_max_items', 100000)), queue_max_bytes=int(layout.get('keycheck_queue_max_bytes', 512 * 1024 * 1024)), ) finally: db.close() return total def collect_legacy_gcp_vertex_credentials(layout, max_items=100): service_dir = os.path.join(layout["keycheck_dir"], "gcp") max_items = max(1, int(max_items)) max_bytes = max(1024, env_int("KEYCHECK_STATUS_PROJECTION_MAX_BYTES", 32 * 1024 * 1024)) max_line_bytes = max(1024, env_int("KEYCHECK_INPUT_MAX_LINE_BYTES", 16 * 1024 * 1024)) entries = {} skipped = 0 for filename in LEGACY_GCP_VERTEX_STATUS_FILES: path = os.path.join(service_dir, filename) if not os.path.exists(path): continue for line in _read_status_projection(path, max_items, max_bytes, max_line_bytes): raw = status_key(line) try: parsed = json.loads(str(raw or "")) except (TypeError, ValueError, json.JSONDecodeError): skipped += 1 continue private_key = str(parsed.get("private_key") or "") if isinstance(parsed, dict) else "" if ( not isinstance(parsed, dict) or not parsed.get("client_email") or not parsed.get("private_key_id") or "-----BEGIN" not in private_key or "-----END" not in private_key ): skipped += 1 continue canonical = json.dumps(parsed, ensure_ascii=False, sort_keys=True, separators=(",", ":")) if len(canonical.encode("utf-8")) > max_line_bytes: skipped += 1 continue entries.setdefault(canonical, set()).add(filename) if len(entries) > max_items: raise RuntimeError("legacy GCP Vertex credential import exceeds its item bound") return entries, skipped def import_legacy_gcp_vertex_credentials(layout, max_items=100): entries, skipped = collect_legacy_gcp_vertex_credentials(layout, max_items=max_items) db = ScannerDB(db_url=layout.get("database_url") or os.getenv("SCANNER_DB_URL"), initialize=False) if not db.enabled: raise RuntimeError("PostgreSQL legacy GCP credential import database is unavailable") db.set_application_name("truf-keycheck-legacy-gcp-import") db.require_runtime_safety_schema() db.require_final_cutover() now = utc_now_iso() queue_max_items = int(layout.get("keycheck_queue_max_items", 100000)) queue_max_bytes = int(layout.get("keycheck_queue_max_bytes", 512 * 1024 * 1024)) imported = 0 existing = 0 queued = 0 linked_results = 0 missing_history = 0 queued_bytes = 0 try: capacity = db.conn.execute("SELECT * FROM pipeline_capacity WHERE id = 1 FOR UPDATE").fetchone() if not capacity: raise RuntimeError("pipeline capacity row is unavailable") for secret_json, source_files in entries.items(): provider_key_hash = sha256_text(secret_json) present = db.conn.execute( "SELECT id FROM keycheck_credentials WHERE service = ? AND provider_key_hash = ?", ("gcp", provider_key_hash), ).fetchone() if present: existing += 1 continue latest = db.conn.execute( '''SELECT * FROM keycheck_results WHERE service = ? AND (key_hash = ? OR secret_hash = ?) ORDER BY checked_at DESC, id DESC LIMIT 1''', ("gcp", provider_key_hash, provider_key_hash), ).fetchone() if not latest: missing_history += 1 continue conflicts = db.conn.execute( '''SELECT DISTINCT credential_id FROM keycheck_results WHERE service = ? AND (key_hash = ? OR secret_hash = ?) AND credential_id IS NOT NULL''', ("gcp", provider_key_hash, provider_key_hash), ).fetchall() if conflicts: raise RuntimeError("legacy GCP history already belongs to a different credential") credential_hash = sha256_text("|".join(("truf-credential-v2", "gcp", secret_json))) credential_id = db.conn.insert_returning_id( '''INSERT INTO keycheck_credentials( service, credential_hash, provider_key_hash, candidate_kind, secret_text, secret_json, key_masked, endpoint, principal, metadata_json, created_at, updated_at ) VALUES (?, ?, ?, 'gcp_json', NULL, ?, ?, '', '', ?, ?, ?) ON CONFLICT(service, provider_key_hash) DO NOTHING''', ( "gcp", credential_hash, provider_key_hash, secret_json, mask_secret(secret_json), json_dumps({ "legacy_status_import": True, "source_files": sorted(source_files), }), now, now, ), ) if credential_id is None: raise RuntimeError("legacy GCP credential insert lost its uniqueness race") updated = db.conn.execute( '''UPDATE keycheck_results SET credential_id = ? WHERE service = ? AND (key_hash = ? OR secret_hash = ?) AND credential_id IS NULL''', (credential_id, "gcp", provider_key_hash, provider_key_hash), ) linked_results += max(0, int(getattr(updated, "rowcount", 0) or 0)) db.conn.execute( '''INSERT INTO keycheck_current_state( credential_id, service, status, status_group, last_result_id, result_source, checked_at, recheck_after, state_version, metadata_json, updated_at ) VALUES (?, 'gcp', ?, ?, ?, ?, ?, NULL, 1, ?, ?)''', ( credential_id, latest["status"], latest["status_group"], latest["id"], latest["result_source"] or "legacy_status_import", latest["checked_at"], latest["metadata_json"] or "{}", now, ), ) capacity_bytes = len(secret_json.encode("utf-8")) + 512 if ( int(capacity["keycheck_items"]) + queued + 1 > queue_max_items or int(capacity["keycheck_bytes"]) + queued_bytes + capacity_bytes > queue_max_bytes ): raise RuntimeError("legacy GCP credential import would exceed keycheck queue capacity") candidate_uid = sha256_text("|".join(( "truf-keycheck-legacy-status-v1", "gcp", provider_key_hash, ))) candidate_id = db.conn.insert_returning_id( '''INSERT INTO keycheck_candidates( candidate_uid, credential_id, service, routed_service, secret_hash, source, query, target, detector_name, found_at, finding_uid, metadata_json, state, capacity_bytes, created_at, updated_at ) VALUES (?, ?, 'gcp', 'gcp', ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?, ?) ON CONFLICT(candidate_uid) DO NOTHING''', ( candidate_uid, credential_id, latest["secret_hash"] or provider_key_hash, latest["source"] or "legacy_keycheck_status", latest["query"] or "", latest["target"] or "", latest["detector_name"] or "GCP", latest["found_at"] or latest["checked_at"], latest["finding_uid"] or "", json_dumps({ "legacy_status_import": True, "recheck_of_status": latest["status"], "source_files": sorted(source_files), }), capacity_bytes, now, now, ), ) imported += 1 if candidate_id is not None: queued += 1 queued_bytes += capacity_bytes if queued: db.conn.execute( '''UPDATE pipeline_capacity SET keycheck_items = keycheck_items + ?, keycheck_bytes = keycheck_bytes + ?, updated_at = ? WHERE id = 1''', (queued, queued_bytes, now), ) db.conn.commit() except Exception: db.conn.rollback() raise finally: db.close() return { "file_credentials": len(entries), "imported": imported, "existing": existing, "queued": queued, "linked_results": linked_results, "missing_history": missing_history, "skipped": skipped, } def parse_args(): parser = argparse.ArgumentParser(description="Run normalized keycheckers from the unified project layout.") parser.add_argument("--config", default="config.yaml") parser.add_argument("--service", default="all", help="Service name, comma list, or all") parser.add_argument("--input") parser.add_argument("--input-mode", choices=('postgres', 'jsonl'), default='postgres') parser.add_argument("--proxy-file") parser.add_argument("--max-keys", type=int, default=0) parser.add_argument("--retry-network", action="store_true") parser.add_argument("--retry-limited", action="store_true") parser.add_argument("--retry-unknown", action="store_true") parser.add_argument("--retry-restricted", action="store_true") parser.add_argument("--retry-no-balance", action="store_true") parser.add_argument("--retry-valid", action="store_true") parser.add_argument("--no-resource-probe", action="store_true", help="Pass through to services that support read-only resource probes, currently Replicate.") parser.add_argument("--recheck-all", action="store_true") parser.add_argument("--import-legacy-gcp-vertex", action="store_true") parser.add_argument("--print-plan", action="store_true") parser.add_argument("--summary-only", action="store_true") parser.add_argument("--no-summary", action="store_true") parser.add_argument("--summary-tsv") parser.add_argument("--summary-json") parser.add_argument("--alive-summary-tsv") parser.add_argument("--no-db-ingest", action="store_true", help="Do not import new *Results.jsonl events into keycheck_results") parser.add_argument("--ingest-keychecks-to-db", action="store_true", help="Import durable *Results.jsonl events into keycheck_results without running provider checks") parser.add_argument("--ingest-max-rows", type=int, default=None, help="Max result events to ingest for --ingest-keychecks-to-db; defaults to keychecks.db_ingest.max_rows") parser.add_argument("--no-auto-repair-links", action="store_true") parser.add_argument("--sync-keychecks-to-db", action="store_true", help="Backfill keycheck_results from existing *Results.jsonl files") parser.add_argument("--reset-keycheck-results", action="store_true", help="Retired; online keycheck result deletion is not supported") parser.add_argument("--repair-keycheck-links", action="store_true", help="Link unattributed keycheck_results to matching findings now present in scanner DB") parser.add_argument("--repair-batch-size", type=int, default=500) parser.add_argument("--repair-max-rows", type=int, default=0) parser.add_argument("--repair-max-attempts", type=int, default=3) parser.add_argument("--repair-since-days", type=int, default=14) return parser.parse_args() PLAN_MUTATING_MODES = ( "sync_keychecks_to_db", "reset_keycheck_results", "repair_keycheck_links", "ingest_keychecks_to_db", "import_legacy_gcp_vertex", "summary_only", ) def plan_requires_authority(args): return any(bool(getattr(args, name, False)) for name in PLAN_MUTATING_MODES) def print_keycheck_plan(args, layout, requested, keychecks_config): """Describe requested work without opening databases or changing canonical files.""" if args.sync_keychecks_to_db or args.reset_keycheck_results: print("mode: reset-keycheck-results-retired" if args.reset_keycheck_results else "mode: sync-keychecks-to-db") print(f"reset_keycheck_results: {bool(args.reset_keycheck_results)}") print("services: " + ",".join(requested)) return if args.repair_keycheck_links: print("mode: repair-keycheck-links") print("services: " + ",".join(requested)) print( "repair_limits: " f"batch_size={max(1, args.repair_batch_size)} " f"max_rows={max(0, args.repair_max_rows)} " f"max_attempts={max(1, args.repair_max_attempts)} " f"since_days={max(0, args.repair_since_days)}" ) return if args.ingest_keychecks_to_db: cfg = db_ingest_config(keychecks_config) max_rows = cfg["max_rows"] if args.ingest_max_rows is None else args.ingest_max_rows print("mode: ingest-keychecks-to-db") print("services: " + ",".join(requested)) print(f"max_rows: {max_rows}") return if args.summary_only: print("mode: summary-only") else: print("mode: run-keychecks") active_flags = [ flag for flag in ( "retry_network", "retry_limited", "retry_unknown", "retry_restricted", "retry_no_balance", "retry_valid", "recheck_all", ) if getattr(args, flag, False) ] if active_flags: print("active_flags: " + " ".join("--" + flag.replace("_", "-") for flag in active_flags)) for service in requested: extra = service_extra_args(service, keychecks_config) suffix = f" args={' '.join(extra)}" if extra else "" print(f"{service}: {os.path.join(layout['project_dir'], SERVICES[service])} -> {os.path.join(layout['keycheck_dir'], service)}{suffix}") if not args.no_summary: print(f"summary_tsv: {args.summary_tsv or os.path.join(layout['keycheck_dir'], 'summary.tsv')}") print(f"summary_json: {args.summary_json or os.path.join(layout['keycheck_dir'], 'summary.json')}") print(f"alive_summary_tsv: {args.alive_summary_tsv or os.path.join(layout['keycheck_dir'], 'alive_summary.tsv')}") def main(): args = parse_args() if getattr(args, 'reset_keycheck_results', False) and not args.print_plan: raise SystemExit('--reset-keycheck-results is retired; online keycheck result deletion is not supported') if not args.print_plan or plan_requires_authority(args): try: require_active_supervisor_child(args.config, child_kind='keycheck', require_dsn=True) except LifecycleAuthorityError as exc: raise SystemExit(str(exc)) from exc config = load_config(args.config) layout = config.get("global") or {} if not args.print_plan: preflight_lifecycle_paths(args.config, config, authority_profile='server') require_sensitive_runtime_paths(layout, create=False) load_postgres_env(args.config, layout) keychecks_config = config.get("keychecks") or {} apply_keycheck_config_defaults(args, keychecks_config) configured_input_mode = str(keychecks_config.get('input_mode') or '').strip().lower() args.input_mode = str(getattr(args, 'input_mode', 'postgres') or 'postgres').lower() if args.input_mode == 'postgres' and configured_input_mode in ('postgres', 'jsonl'): args.input_mode = configured_input_mode if args.input and args.input_mode != 'jsonl': raise SystemExit('--input is valid only with explicit --input-mode jsonl') requested = service_list(args.service) unknown = [service for service in requested if service not in SERVICES] if unknown: raise SystemExit(f"Unknown service(s): {', '.join(unknown)}") summary_services = list(SERVICES) if args.input_mode == 'postgres' else requested input_file = args.input or os.path.join(layout["results_dir"], "found_secrets.jsonl") proxy_file = args.proxy_file or layout["proxy_file"] print(f"input_mode: {args.input_mode}") print(f"input: {input_file if args.input_mode == 'jsonl' else 'postgres:keycheck_candidates'}") print(f"proxy: {proxy_file}") print(f"keycheck_dir: {layout['keycheck_dir']}") if args.print_plan: print_keycheck_plan(args, layout, requested, keychecks_config) return if args.sync_keychecks_to_db: if args.input_mode != 'jsonl': raise SystemExit('--sync-keychecks-to-db requires explicit --input-mode jsonl') sync_keycheck_results_to_db(layout, requested) return if args.repair_keycheck_links: repair_keycheck_links( layout, requested, batch_size=max(1, args.repair_batch_size), max_rows=max(0, args.repair_max_rows), max_attempts=max(1, args.repair_max_attempts), since_days=max(0, args.repair_since_days), ) return if args.ingest_keychecks_to_db: if args.input_mode != 'jsonl': raise SystemExit('--ingest-keychecks-to-db requires explicit --input-mode jsonl') cfg = db_ingest_config(keychecks_config) if args.ingest_max_rows is not None and args.ingest_max_rows <= 0: raise SystemExit("--ingest-max-rows must be > 0; omit it to use keychecks.db_ingest.max_rows") max_rows = cfg["max_rows"] if args.ingest_max_rows is None else args.ingest_max_rows ingest_keycheck_results_to_db(layout, requested, max_rows=max_rows) return if args.summary_only: write_summary( layout, summary_services, args.summary_tsv, args.summary_json, args.alive_summary_tsv, input_mode=args.input_mode, ) return if args.import_legacy_gcp_vertex: if requested != ["gcp"]: raise SystemExit("--import-legacy-gcp-vertex requires --service gcp") imported = import_legacy_gcp_vertex_credentials(layout) print( "legacy_gcp_vertex_import: " + " ".join(f"{name}={value}" for name, value in imported.items()), flush=True, ) if args.input_mode != 'postgres': raise SystemExit( 'Normal provider execution requires PostgreSQL candidate leases; ' 'legacy JSONL compatibility work is offline-only.' ) queued_rechecks = enqueue_requested_postgres_rechecks(layout, requested, args) if queued_rechecks: print(f'postgres_recheck_candidates: {queued_rechecks}', flush=True) exit_code, provider_failures = run_provider_services( requested, args, layout, keychecks_config, ) if not args.no_summary: write_summary( layout, summary_services, args.summary_tsv, args.summary_json, args.alive_summary_tsv, input_mode=args.input_mode, ) if args.input_mode == 'jsonl' and not args.no_db_ingest: maybe_ingest_keycheck_results(layout, requested, keychecks_config) if args.input_mode == 'jsonl' and not args.no_auto_repair_links: maybe_auto_repair_links(layout, requested, keychecks_config) if provider_failures: detail = ",".join(f"{service}={code}" for service, code in provider_failures) print(f"provider_failures: {detail}", flush=True) else: print("provider_failures: none", flush=True) raise SystemExit(exit_code) if __name__ == "__main__": main()