Files
llm/scripts/audit_1c_adapter_coverage.py
T
2026-08-14 09:40:51 +03:00

563 lines
21 KiB
Python

#!/usr/bin/env python3
from __future__ import annotations
import argparse
import concurrent.futures
import copy
import json
import os
import re
import sys
import time
import threading
import urllib.error
import urllib.request
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable
DEFAULT_BASE_URL = "http://docker.cin.su:8011"
DEFAULT_BASE_ID = "upo_test"
# Metadata kinds that either own application data or expose values through the
# public data facade. Kinds absent from a concrete base remain in the audit so
# that coverage cannot be declared only from a convenient test configuration.
DATA_KINDS = {
"AccountingRegister",
"AccumulationRegister",
"BusinessProcess",
"CalculationRegister",
"Catalog",
"ChartOfAccounts",
"ChartOfCalculationTypes",
"ChartOfCharacteristicTypes",
"Constant",
"Document",
"Enum",
"ExchangePlan",
"InformationRegister",
"Sequence",
"Task",
}
Rpc = Callable[[str, str, str, dict[str, Any], float], dict[str, Any]]
def utc_now() -> str:
return datetime.now(timezone.utc).isoformat()
def write_json_atomic(path: Path, value: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_name(path.name + ".tmp")
temporary.write_text(json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
temporary.replace(path)
def rpc(base_url: str, token: str, method: str, payload: dict[str, Any], timeout: float) -> dict[str, Any]:
request = urllib.request.Request(
base_url.rstrip("/") + "/rpc",
data=json.dumps({"method": method, "payload": payload}, ensure_ascii=False).encode("utf-8"),
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json; charset=utf-8"},
method="POST",
)
with urllib.request.urlopen(request, timeout=timeout) as response:
value = json.loads(response.read().decode("utf-8"))
return value if isinstance(value, dict) else {"status": "error", "error": "response_not_object"}
def safe_rpc(
rpc_call: Rpc,
base_url: str,
token: str,
method: str,
payload: dict[str, Any],
timeout: float,
) -> tuple[dict[str, Any], int]:
started = time.monotonic()
try:
result = rpc_call(base_url, token, method, payload, timeout)
except (TimeoutError, urllib.error.URLError, OSError, ValueError) as exc:
result = {
"status": "transport_error",
"error": type(exc).__name__,
"diagnostics": {"message": str(exc)[:500]},
}
return result, round((time.monotonic() - started) * 1000)
def first_object(response: dict[str, Any]) -> dict[str, Any] | None:
for key in ("objects", "items"):
values = response.get(key)
if isinstance(values, list) and values and isinstance(values[0], dict):
return values[0]
return None
def public_selector(sample: dict[str, Any]) -> dict[str, Any]:
return {
key: sample[key]
for key in ("ref", "kind", "name", "guid")
if sample.get(key) not in {None, ""}
}
def operation_summary(result: dict[str, Any], duration_ms: int, **values: Any) -> dict[str, Any]:
summary = {"status": result.get("status") or "unknown", "duration_ms": duration_ms, **values}
diagnostics = result.get("diagnostics")
if summary["status"] != "ok" and isinstance(diagnostics, dict) and diagnostics.get("message"):
summary["message"] = str(diagnostics["message"])[:500]
if result.get("error"):
summary["error"] = str(result["error"])[:200]
return summary
def record_ref_from_row(row: Any) -> str | None:
if not isinstance(row, dict):
return None
value = row.get("ref")
if isinstance(value, dict):
value = value.get("hex")
compact = str(value or "").replace("-", "").strip()
return compact if re.fullmatch(r"[0-9a-fA-F]{32}", compact) else None
def audit_data_kind(
base_url: str,
base_id: str,
token: str,
kind: str,
timeout: float,
include_reads: bool,
rpc_call: Rpc = rpc,
existing: dict[str, Any] | None = None,
progress: Callable[[dict[str, Any]], None] | None = None,
) -> dict[str, Any]:
common = {"base_id": base_id, "timeout_seconds": max(1, int(timeout))}
result = copy.deepcopy(existing) if isinstance(existing, dict) else {}
result.update({"kind": kind, "status": "degraded"})
result.setdefault("operations", {})
def save_progress() -> None:
if progress is not None:
progress(copy.deepcopy(result))
prior_list = result["operations"].get("metadata.objects.list") or {}
if prior_list.get("status") == "ok" and isinstance(result.get("sample"), dict):
listed = {"status": "ok"}
sample = {**(result["sample"].get("selector") or {}), "name": result["sample"].get("name")}
else:
listed, list_ms = safe_rpc(
rpc_call,
base_url,
token,
"metadata.objects.list",
{**common, "kind": kind, "limit": 1},
timeout,
)
sample = first_object(listed)
result["operations"]["metadata.objects.list"] = operation_summary(
listed,
list_ms,
objects=len(listed.get("objects") or listed.get("items") or []),
)
if sample:
result["sample"] = {"selector": public_selector(sample), "name": sample.get("name")}
save_progress()
if not sample:
result["status"] = "absent" if listed.get("status") == "ok" else "degraded"
result["reason"] = "no_sample_object"
save_progress()
return result
selector = public_selector(sample)
result["sample"] = {"selector": selector, "name": sample.get("name")}
prior_schema = result["operations"].get("data.schema") or {}
if prior_schema.get("status") == "ok":
schema_status = "ok"
else:
schema, schema_ms = safe_rpc(
rpc_call,
base_url,
token,
"data.schema",
{**common, **selector},
timeout,
)
table = schema.get("table") if isinstance(schema.get("table"), dict) else {}
result["operations"]["data.schema"] = operation_summary(
schema,
schema_ms,
fields=len(schema.get("fields") or []),
table=table.get("name"),
cache=(schema.get("cache") or {}).get("status") if isinstance(schema.get("cache"), dict) else None,
)
schema_status = str(schema.get("status") or "unknown")
save_progress()
if schema_status != "ok" or not include_reads:
result["status"] = "ok" if schema_status == "ok" else "degraded"
save_progress()
return result
prior_data_list = result["operations"].get("data.list") or {}
must_repeat_list = prior_data_list.get("status") != "ok" or (
"sample_record_ref" not in result and "sample_record_ref_status" not in result
)
if must_repeat_list:
data_list, data_list_ms = safe_rpc(
rpc_call,
base_url,
token,
"data.list",
{**common, **selector, "limit": 1},
timeout,
)
rows = data_list.get("rows") if isinstance(data_list.get("rows"), list) else []
list_summary = operation_summary(data_list, data_list_ms, rows=len(rows))
if data_list.get("status") == "ok":
result["operations"]["data.list"] = list_summary
result["operations"].pop("data.list_retry", None)
ref = record_ref_from_row(rows[0]) if rows else None
if ref:
result["sample_record_ref"] = ref
result.pop("sample_record_ref_status", None)
else:
result.pop("sample_record_ref", None)
result["sample_record_ref_status"] = "empty_object" if not rows else "object_has_no_reference_key"
elif prior_data_list.get("status") == "ok":
# A successful operation is evidence. Do not downgrade it only
# because a later attempt to recover the sample ref timed out.
result["operations"]["data.list_retry"] = list_summary
ref = result.get("sample_record_ref")
else:
result["operations"]["data.list"] = list_summary
result.pop("sample_record_ref", None)
result["sample_record_ref_status"] = "data_list_failed"
ref = None
save_progress()
else:
ref = result.get("sample_record_ref")
if (result["operations"].get("data.count") or {}).get("status") != "ok":
counted, count_ms = safe_rpc(
rpc_call,
base_url,
token,
"data.count",
{**common, **selector},
timeout,
)
result["operations"]["data.count"] = operation_summary(counted, count_ms, count=counted.get("count"))
save_progress()
prior_get_status = (result["operations"].get("data.get") or {}).get("status")
if prior_get_status not in {"ok", "not_applicable"}:
if ref:
fetched, get_ms = safe_rpc(
rpc_call,
base_url,
token,
"data.get",
{**common, **selector, "record_ref": ref},
timeout,
)
fetched_rows = fetched.get("rows") if isinstance(fetched.get("rows"), list) else []
result["operations"]["data.get"] = operation_summary(fetched, get_ms, rows=len(fetched_rows))
elif (result["operations"].get("data.list") or {}).get("status") == "ok":
reason = str(result.get("sample_record_ref_status") or "object_has_no_reference_key")
result["operations"]["data.get"] = {"status": "not_applicable", "reason": reason, "duration_ms": 0}
else:
result["operations"]["data.get"] = {"status": "blocked", "reason": "data_list_failed", "duration_ms": 0}
save_progress()
required = ("data.schema", "data.list", "data.count")
failures = [name for name in required if result["operations"].get(name, {}).get("status") != "ok"]
get_status = result["operations"]["data.get"]["status"]
if get_status not in {"ok", "not_applicable"}:
failures.append("data.get")
result["status"] = "ok" if not failures else "degraded"
if failures:
result["failed_operations"] = failures
else:
result.pop("failed_operations", None)
save_progress()
return result
def load_checkpoint(
path: Path | None,
base_url: str,
base_id: str,
resume: bool,
include_reads: bool,
) -> dict[str, Any]:
fresh = {
"schema": "onec_adapter_data_audit_checkpoint.v1",
"base_url": base_url,
"base_id": base_id,
"include_reads": include_reads,
"started_at": utc_now(),
"updated_at": utc_now(),
"checks": {},
}
if not resume or path is None or not path.exists():
return fresh
value = json.loads(path.read_text(encoding="utf-8-sig"))
if value.get("schema") != fresh["schema"]:
raise ValueError(f"unsupported checkpoint schema in {path}")
if value.get("base_url") != base_url or value.get("base_id") != base_id:
raise ValueError(f"checkpoint {path} belongs to another adapter or base")
if bool(value.get("include_reads")) != include_reads:
raise ValueError(f"checkpoint {path} was created for another data audit mode")
if not isinstance(value.get("checks"), dict):
raise ValueError(f"checkpoint {path} has no checks object")
return value
def run_data_checks(
base_url: str,
base_id: str,
token: str,
kinds: list[str],
timeout: float,
include_reads: bool,
workers: int,
checkpoint_path: Path | None,
resume: bool,
retry_degraded: bool = False,
rpc_call: Rpc = rpc,
) -> tuple[dict[str, dict[str, Any]], int]:
checkpoint = load_checkpoint(checkpoint_path, base_url, base_id, resume, include_reads)
checks = checkpoint["checks"]
checkpoint_lock = threading.Lock()
reusable = {
kind
for kind in kinds
if kind in checks and (not retry_degraded or checks[kind].get("status") == "ok")
}
resumed = len(reusable)
pending = [kind for kind in kinds if kind not in reusable]
def execute(kind: str) -> dict[str, Any]:
def save_partial(value: dict[str, Any]) -> None:
with checkpoint_lock:
checks[kind] = value
checkpoint["updated_at"] = utc_now()
if checkpoint_path is not None:
write_json_atomic(checkpoint_path, checkpoint)
return audit_data_kind(
base_url,
base_id,
token,
kind,
timeout,
include_reads,
rpc_call,
existing=checks.get(kind),
progress=save_partial,
)
with concurrent.futures.ThreadPoolExecutor(max_workers=max(1, workers)) as executor:
futures = {executor.submit(execute, kind): kind for kind in pending}
for future in concurrent.futures.as_completed(futures):
kind = futures[future]
try:
checks[kind] = future.result()
except Exception as exc: # a single kind must not discard completed evidence
checks[kind] = {
"kind": kind,
"status": "degraded",
"error": type(exc).__name__,
"message": str(exc)[:500],
}
with checkpoint_lock:
checkpoint["updated_at"] = utc_now()
if checkpoint_path is not None:
write_json_atomic(checkpoint_path, checkpoint)
return {kind: checks[kind] for kind in kinds if kind in checks}, resumed
def build_report(
base_url: str,
base_id: str,
token: str,
timeout: float,
sample_objects: bool,
sample_schemas: bool,
*,
sample_reads: bool = False,
workers: int = 1,
checkpoint_path: Path | None = None,
resume: bool = False,
retry_degraded: bool = False,
rpc_call: Rpc = rpc,
) -> dict[str, Any]:
audit, _ = safe_rpc(
rpc_call,
base_url,
token,
"metadata.adapter.audit",
{"base_id": base_id, "include_missing": True, "include_unmapped": True, "timeout_seconds": int(timeout)},
timeout,
)
if audit.get("status") != "ok":
return {"schema": "onec_adapter_coverage_audit.v1", "status": "error", "audit": audit}
matrix: list[dict[str, Any]] = []
for support in audit.get("metadata_kinds") or []:
if not isinstance(support, dict):
continue
kind = str(support.get("kind") or "")
count = int(support.get("count") or 0)
matrix.append({
"kind": kind,
"kind_ru": support.get("kind_ru"),
"objects": count,
"capabilities": support.get("capabilities") or [],
"discovery": "present" if count else "absent_in_base",
})
data_checks: dict[str, dict[str, Any]] = {}
resumed_checks = 0
if sample_schemas or sample_reads:
present_data_kinds = sorted(row["kind"] for row in matrix if row["objects"] and row["kind"] in DATA_KINDS)
data_checks, resumed_checks = run_data_checks(
base_url,
base_id,
token,
present_data_kinds,
timeout,
sample_reads,
workers,
checkpoint_path,
resume,
retry_degraded,
rpc_call,
)
for row in matrix:
check = data_checks.get(row["kind"])
if check:
row["data_check"] = check
row["list_status"] = check.get("operations", {}).get("metadata.objects.list", {}).get("status")
schema = check.get("operations", {}).get("data.schema")
if schema:
row["data_schema"] = schema
elif sample_objects:
for row in matrix:
if not row["objects"]:
continue
listed, _ = safe_rpc(
rpc_call,
base_url,
token,
"metadata.objects.list",
{"base_id": base_id, "kind": row["kind"], "limit": 1, "timeout_seconds": int(timeout)},
timeout,
)
row["list_status"] = listed.get("status")
sample = first_object(listed)
if sample:
row["sample_selector"] = public_selector(sample)
missing = [row["kind"] for row in matrix if row["discovery"] == "absent_in_base"]
failures = [
row["kind"]
for row in matrix
if sample_objects and row.get("objects") and "list_status" in row and row.get("list_status") != "ok"
]
data_failures = sorted(kind for kind, check in data_checks.items() if check.get("status") != "ok")
status = "ok" if not failures and not data_failures else "degraded"
return {
"schema": "onec_adapter_coverage_audit.v1",
"status": status,
"generated_at": utc_now(),
"base_url": base_url,
"base_id": base_id,
"sampling": {
"objects": sample_objects,
"data_schemas": sample_schemas or sample_reads,
"data_reads": sample_reads,
"workers": workers,
"resumed_checks": resumed_checks,
},
"policy": {
"application_data": "read_only",
"metadata_structure": "read_only",
"sql_identity": "configured_base_credentials_only",
"writes": ["ConfigSave", "ConfigCASSave"],
},
"counts": {
"kinds": len(matrix),
"present_kinds": sum(1 for row in matrix if row["objects"]),
"absent_kinds": len(missing),
"list_failures": len(failures),
"data_kinds_declared": len(DATA_KINDS),
"data_kinds_checked": len(data_checks),
"data_check_failures": len(data_failures),
},
"absent_in_base": missing,
"list_failures": failures,
"data_check_failures": data_failures,
"matrix": matrix,
"data_checks": data_checks,
"child_objects": audit.get("child_objects") or {},
"not_yet_decoded": audit.get("not_yet_decoded") or [],
"optional_deep_reads": audit.get("optional_deep_reads") or [],
"unmapped_source_roles": audit.get("unmapped_source_roles") or audit.get("unknown_source_roles") or {},
"write_capabilities": audit.get("write_capabilities") or {},
}
def main() -> int:
parser = argparse.ArgumentParser(description="Audit live 1C adapter coverage without exposing SQL credentials.")
parser.add_argument("--base-url", default=DEFAULT_BASE_URL)
parser.add_argument("--base-id", default=DEFAULT_BASE_ID)
parser.add_argument("--token-env", default="ONEC_ADAPTER_TOKEN")
parser.add_argument("--timeout", type=float, default=120.0, help="Timeout for each adapter call, in seconds.")
parser.add_argument("--workers", type=int, default=1, help="Concurrent data-kind checks (default: 1).")
parser.add_argument("--sample-objects", action="store_true", help="List one object for each present metadata kind.")
parser.add_argument("--sample-data-schemas", action="store_true", help="Decode one logical data schema for every present data kind.")
parser.add_argument("--sample-data-reads", action="store_true", help="Run schema, list, get (when applicable), and count for every present data kind.")
parser.add_argument("--checkpoint", type=Path, help="Atomically save progress after every completed data kind.")
parser.add_argument("--resume", action="store_true", help="Reuse completed kinds from --checkpoint.")
parser.add_argument("--retry-degraded", action="store_true", help="With --resume, rerun checkpoint entries whose status is not ok.")
parser.add_argument("--output", type=Path)
args = parser.parse_args()
if args.workers < 1:
parser.error("--workers must be >= 1")
if args.resume and args.checkpoint is None:
parser.error("--resume requires --checkpoint")
if args.retry_degraded and not args.resume:
parser.error("--retry-degraded requires --resume")
token = os.environ.get(args.token_env, "").strip()
if not token:
parser.error(f"adapter token is required in environment variable {args.token_env}")
sample_objects = args.sample_objects or args.sample_data_schemas or args.sample_data_reads
try:
report = build_report(
args.base_url,
args.base_id,
token,
args.timeout,
sample_objects,
args.sample_data_schemas,
sample_reads=args.sample_data_reads,
workers=args.workers,
checkpoint_path=args.checkpoint,
resume=args.resume,
retry_degraded=args.retry_degraded,
)
except ValueError as exc:
parser.error(str(exc))
rendered = json.dumps(report, ensure_ascii=False, indent=2)
if args.output:
write_json_atomic(args.output, report)
print(rendered)
return 0 if report.get("status") == "ok" else 1
if __name__ == "__main__":
sys.exit(main())