#!/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-gpu.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())