2843 lines
127 KiB
Python
2843 lines
127 KiB
Python
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()
|