"""Read-only operational observer for adapter-1c audit telemetry. This service never connects to 1C SQL storage and never mutates adapter data. It reads the adapter's privacy-safe rotated JSONL files from a read-only mount. """ from __future__ import annotations import json import math import mimetypes import os import re import time import threading from collections import Counter, defaultdict from http import HTTPStatus from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from urllib.parse import parse_qs, urlparse from urllib.request import Request, urlopen from urllib.error import URLError, HTTPError ROOT = Path(__file__).resolve().parent WEB_ROOT = ROOT / "web" AUDIT_DIR = Path(os.environ.get("ONEC_OBSERVER_AUDIT_DIR", "/audit")) MCP_AUDIT_DIR = Path(os.environ.get("ONEC_OBSERVER_MCP_AUDIT_DIR", "/mcp-audit")) STATE_DIR = Path(os.environ.get("ONEC_OBSERVER_STATE_DIR", "/state")) ADAPTER_URL = os.environ.get("ONEC_OBSERVER_ADAPTER_URL", "").rstrip("/") HOST = os.environ.get("ONEC_OBSERVER_HOST", "0.0.0.0") PORT = int(os.environ.get("ONEC_OBSERVER_PORT", "8031")) MAX_ROWS = 10000 def number(value: object) -> int: try: return int(value or 0) except (TypeError, ValueError): return 0 AUTO_COVERAGE_BASE = os.environ.get("ONEC_OBSERVER_COVERAGE_BASE_ID", "upo_test") AUTO_COVERAGE_INTERVAL = max(300, number(os.environ.get("ONEC_OBSERVER_COVERAGE_INTERVAL_SECONDS", "900"))) LAST_COVERAGE: dict[str, object] = {"status": "not_started"} def percentile(values: list[int], q: float) -> int: if not values: return 0 ordered = sorted(values) index = max(0, min(len(ordered) - 1, math.ceil(len(ordered) * q) - 1)) return ordered[index] def audit_files(directory: Path, prefix: str) -> list[Path]: if not directory.exists(): return [] paths = [p for p in directory.glob(f"{prefix}*") if p.is_file()] return sorted(paths, key=lambda p: p.stat().st_mtime) def read_events(directory: Path = AUDIT_DIR, prefix: str = "adapter-audit.jsonl", event_name: str = "adapter_rpc") -> tuple[list[dict], int]: events: list[dict] = [] malformed = 0 for path in audit_files(directory, prefix): try: with path.open("r", encoding="utf-8", errors="replace") as handle: for line in handle: try: row = json.loads(line) except json.JSONDecodeError: malformed += 1 continue if isinstance(row, dict) and row.get("event") == event_name: events.append(row) except OSError: continue return events[-MAX_ROWS:], malformed def event_view(row: dict, source: str = "rest") -> dict: request = row.get("request") if isinstance(row.get("request"), dict) else {} return { "source": source, "time": row.get("time"), "request_id": row.get("request_id"), "method": row.get("method"), "base_id": request.get("base_id"), "selector": {key: request.get(key) for key in ("ref", "kind", "name", "object_type", "object_name", "extension", "mode", "execution_mode") if request.get(key) not in (None, "")}, "status": row.get("status") or "unknown", "error": row.get("error") or "", "exception_type": row.get("exception_type") or "", "duration_ms": number(row.get("duration_ms")), "result_duration_ms": row.get("result_duration_ms"), "result_summary": row.get("result_summary") if isinstance(row.get("result_summary"), dict) else {}, } def correlations(rest: list[dict], mcp: list[dict]) -> list[dict]: rest_by_id = {str(row.get("request_id")): row for row in rest if row.get("request_id")} rows = [] for row in reversed(mcp): request_id = str(row.get("request_id") or "") if not request_id: continue rest_row = rest_by_id.get(request_id) mcp_view = event_view(row, "mcp") rows.append({"request_id": request_id, "mcp": mcp_view, "rest": event_view(rest_row, "rest") if rest_row else None, "correlation_status": "matched" if rest_row else "not_reached_rest"}) return rows[:1000] def adapter_rpc(method: str, payload: dict) -> dict: if not ADAPTER_URL: raise RuntimeError("adapter_url_not_configured") request = Request(f"{ADAPTER_URL}/rpc", data=json.dumps({"method": method, "payload": payload}).encode("utf-8"), headers={"Content-Type": "application/json; charset=utf-8"}, method="POST") try: with urlopen(request, timeout=45) as response: value = json.loads(response.read().decode("utf-8")) except (HTTPError, URLError, TimeoutError) as exc: raise RuntimeError(f"adapter_read_failed:{type(exc).__name__}") from exc if not isinstance(value, dict): raise RuntimeError("adapter_response_not_object") return value def coverage_snapshot(base_id: str) -> dict: methods = adapter_rpc("help.methods", {}) audit = adapter_rpc("metadata.adapter.audit", {"base_id": base_id, "include_missing": True, "include_unmapped": True, "timeout_seconds": 45}) snapshot = {"schema": "onec_adapter_observer_coverage.v1", "captured_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), "base_id": base_id, "methods": methods.get("methods") or [], "audit": audit} try: STATE_DIR.mkdir(parents=True, exist_ok=True) latest = STATE_DIR / f"coverage-{base_id}.json" previous = json.loads(latest.read_text(encoding="utf-8")) if latest.exists() else None if isinstance(previous, dict): old_methods = {str(item.get("name")) for item in previous.get("methods") or [] if isinstance(item, dict)} new_methods = {str(item.get("name")) for item in snapshot["methods"] if isinstance(item, dict)} def kind_counts(value: dict) -> dict[str, int]: audit_value = value.get("audit") if isinstance(value.get("audit"), dict) else {} return {str(item.get("kind")): number(item.get("count")) for item in audit_value.get("metadata_kinds") or [] if isinstance(item, dict)} old_kinds, new_kinds = kind_counts(previous), kind_counts(snapshot) changed_kinds = [{"kind": kind, "before": old_kinds.get(kind, 0), "after": new_kinds.get(kind, 0)} for kind in sorted(set(old_kinds) | set(new_kinds)) if old_kinds.get(kind, 0) != new_kinds.get(kind, 0)] def unresolved(value: dict) -> set[str]: audit_value = value.get("audit") if isinstance(value.get("audit"), dict) else {} return {json.dumps(item, ensure_ascii=False, sort_keys=True) if isinstance(item, dict) else str(item) for item in audit_value.get("not_yet_decoded") or []} old_unresolved, new_unresolved = unresolved(previous), unresolved(snapshot) snapshot["comparison"] = {"previous_captured_at": previous.get("captured_at"), "methods_added": sorted(new_methods - old_methods), "methods_removed": sorted(old_methods - new_methods), "kind_count_changes": changed_kinds, "undecoded_added": sorted(new_unresolved - old_unresolved), "undecoded_removed": sorted(old_unresolved - new_unresolved)} temporary = STATE_DIR / "coverage-latest.json.tmp" temporary.write_text(json.dumps(snapshot, ensure_ascii=False), encoding="utf-8") temporary.replace(latest) history_path = STATE_DIR / f"coverage-{base_id}.history.jsonl" history = history_path.read_text(encoding="utf-8", errors="replace").splitlines()[-49:] if history_path.exists() else [] history.append(json.dumps(snapshot, ensure_ascii=False)) history_path.write_text("\n".join(history) + "\n", encoding="utf-8") except OSError: snapshot["persistence_status"] = "unavailable" return snapshot def coverage_worker() -> None: """Best-effort periodic read-only snapshot; failure must not stop the UI.""" while True: try: snapshot = coverage_snapshot(AUTO_COVERAGE_BASE) LAST_COVERAGE.update({"status": "ok", "captured_at": snapshot.get("captured_at"), "base_id": AUTO_COVERAGE_BASE}) except RuntimeError as exc: LAST_COVERAGE.update({"status": "error", "base_id": AUTO_COVERAGE_BASE, "error": str(exc), "checked_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())}) time.sleep(AUTO_COVERAGE_INTERVAL) def recommendation(event: dict) -> str: error = str(event.get("error") or "") status = str(event.get("status") or "") if error == "time_budget_exhausted": return "Сузить публичный selector (ref, форма или routine) либо выполнить тяжёлую операцию как job." if error == "public_write_route_unresolved": return "Передать request_id и resolver summary разработчикам адаптера; не подбирать storage coordinates вручную." if error == "ambiguous_fragment": return "Уточнить routine_name или заменить модуль целиком; фрагмент не должен подбираться по совпадению." if error == "base_id_required": return "Передать base_id из списка сконфигурированных баз; не пытаться подставлять SQL-параметры." if status in {"unsupported", "blocked", "invalid_argument"}: return "Это безопасная остановка. Проверить публичный контракт метода и next_action в результате." if status == "exception": return "Найти совпадающий request_id в REST/MCP telemetry и воспроизвести только на upo_test." return "Повторить read-операцию с тем же публичным selector-ом и сравнить длительность/статус." def build_summary(events: list[dict], malformed: int) -> dict: durations = [number(row.get("duration_ms")) for row in events] failures = [row for row in events if str(row.get("status")) == "exception"] groups: dict[tuple[str, str, str], list[dict]] = defaultdict(list) for row in events: groups[(str(row.get("method") or ""), str(row.get("status") or "unknown"), str(row.get("error") or ""))].append(row) findings = [] normal_lifecycle = {"ok", "accepted", "running", "done", "cancelled", "not_found", "unknown"} for (method, status, error), rows in sorted(groups.items(), key=lambda item: len(item[1]), reverse=True): if status in normal_lifecycle and not error: continue finding = event_view(rows[-1]) findings.append({"method": method, "status": status, "error": error, "count": len(rows), "last_request_id": finding["request_id"], "recommendation": recommendation(finding)}) if len(findings) >= 20: break per_method: dict[str, list[int]] = defaultdict(list) for row in events: per_method[str(row.get("method") or "")].append(number(row.get("duration_ms"))) methods = [{"method": name, "calls": len(values), "p50_ms": percentile(values, .5), "p95_ms": percentile(values, .95), "max_ms": max(values)} for name, values in per_method.items()] slow = sorted((event_view(row) for row in events if number(row.get("duration_ms")) >= 5000), key=lambda row: row["duration_ms"], reverse=True)[:20] return {"schema": "onec_adapter_observer_summary.v1", "generated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), "events": len(events), "malformed_rows": malformed, "exceptions": len(failures), "p50_ms": percentile(durations, .5), "p95_ms": percentile(durations, .95), "max_ms": max(durations, default=0), "methods": sorted(methods, key=lambda row: row["p95_ms"], reverse=True), "findings": findings, "slow_events": slow} class Handler(BaseHTTPRequestHandler): server_version = "AdapterObserver/1.0" def log_message(self, _format: str, *_args: object) -> None: return def send_json(self, status: int, value: object) -> None: body = json.dumps(value, ensure_ascii=False).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.send_header("Cache-Control", "no-store") self.end_headers() self.wfile.write(body) def serve_file(self, relative: str) -> None: target = (WEB_ROOT / relative).resolve() if WEB_ROOT not in target.parents and target != WEB_ROOT or not target.is_file(): self.send_error(HTTPStatus.NOT_FOUND) return body = target.read_bytes() self.send_response(HTTPStatus.OK) self.send_header("Content-Type", mimetypes.guess_type(str(target))[0] or "application/octet-stream") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def do_GET(self) -> None: # noqa: N802 parsed = urlparse(self.path) if parsed.path == "/health": self.send_json(200, {"status": "ok", "service": "adapter-observer", "audit_dir": str(AUDIT_DIR), "files": [path.name for path in audit_files(AUDIT_DIR, "adapter-audit.jsonl")], "mcp_files": [path.name for path in audit_files(MCP_AUDIT_DIR, "mcp-audit.jsonl")], "coverage": LAST_COVERAGE}) return events, malformed = read_events() mcp_events, mcp_malformed = read_events(MCP_AUDIT_DIR, "mcp-audit.jsonl", "mcp_adapter_call") if parsed.path == "/api/summary": self.send_json(200, build_summary(events, malformed)) return if parsed.path == "/api/events": query = parse_qs(parsed.query) method, status, base_id = query.get("method", [""])[0], query.get("status", [""])[0], query.get("base_id", [""])[0] minimum_duration = number(query.get("min_duration_ms", [0])[0]) since, until = query.get("since", [""])[0], query.get("until", [""])[0] rows = [event_view(row) for row in reversed(events)] if method: rows = [row for row in rows if row["method"] == method] if status: rows = [row for row in rows if row["status"] == status] if base_id: rows = [row for row in rows if row["base_id"] == base_id] if minimum_duration: rows = [row for row in rows if row["duration_ms"] >= minimum_duration] if since: rows = [row for row in rows if str(row.get("time") or "") >= since] if until: rows = [row for row in rows if str(row.get("time") or "") <= until] limit = min(max(number(query.get("limit", [200])[0]), 1), 1000) self.send_json(200, {"schema": "onec_adapter_observer_events.v1", "events": rows[:limit], "total": len(rows), "malformed_rows": malformed}) return if parsed.path == "/api/mcp-events": self.send_json(200, {"schema": "onec_adapter_observer_mcp_events.v1", "events": [event_view(row, "mcp") for row in reversed(mcp_events[-1000:])], "total": len(mcp_events), "malformed_rows": mcp_malformed}) return if parsed.path == "/api/correlations": self.send_json(200, {"schema": "onec_adapter_observer_correlations.v1", "correlations": correlations(events, mcp_events), "mcp_events": len(mcp_events), "mcp_malformed_rows": mcp_malformed}) return if parsed.path == "/api/coverage": base_id = parse_qs(parsed.query).get("base_id", ["upo_test"])[0] if not re.fullmatch(r"[A-Za-zА-Яа-яЁё0-9_.-]{1,80}", base_id): self.send_json(400, {"error": "invalid_base_id"}) return try: self.send_json(200, coverage_snapshot(base_id)) except RuntimeError as exc: self.send_json(502, {"error": str(exc)}) return if parsed.path == "/api/objects": query = parse_qs(parsed.query) base_id, kind = query.get("base_id", ["upo"])[0], query.get("kind", [""])[0] if not kind: self.send_json(400, {"error": "kind_required"}) return try: started = time.monotonic() result = adapter_rpc("metadata.objects.list", {"base_id": base_id, "kind": kind, "limit": min(max(number(query.get("limit", [200])[0]), 1), 1000), "offset": max(number(query.get("offset", [0])[0]), 0), "refresh_cache": True, "exact_counts": True}) self.send_json(200, {**result, "observer": {"duration_ms": int((time.monotonic() - started) * 1000), "method": "metadata.objects.list"}}) except RuntimeError as exc: self.send_json(502, {"error": str(exc)}) return if parsed.path == "/api/object": query = parse_qs(parsed.query) base_id, ref = query.get("base_id", ["upo"])[0], query.get("ref", [""])[0] if not ref: self.send_json(400, {"error": "ref_required"}) return started = time.monotonic() sections = {} for name, method in (("attributes", "metadata.object.attributes"), ("forms", "metadata.object.forms"), ("modules", "metadata.object.modules"), ("templates", "metadata.object.templates")): section_started = time.monotonic() try: result = adapter_rpc(method, {"base_id": base_id, "ref": ref}) sections[name] = {"status": result.get("status", "unknown"), "data": result, "duration_ms": int((time.monotonic() - section_started) * 1000), "method": method} except RuntimeError as exc: sections[name] = {"status": "error", "error": str(exc), "duration_ms": int((time.monotonic() - section_started) * 1000), "method": method} self.send_json(200, {"schema": "onec_adapter_observer_object_node.v1", "base_id": base_id, "ref": ref, "status": "ok", "duration_ms": int((time.monotonic() - started) * 1000), "sections": sections}) return if parsed.path == "/api/object/action": query = parse_qs(parsed.query) base_id, ref = query.get("base_id", ["upo"])[0], query.get("ref", [""])[0] action = query.get("action", [""])[0] actions = { "card": "metadata.object.get", "properties": "metadata.object.properties", "attributes": "metadata.object.attributes", "forms": "metadata.object.forms", "commands": "metadata.object.commands", "modules": "metadata.object.modules", "templates": "metadata.object.templates", "related": "metadata.object.related", } if not ref: self.send_json(400, {"error": "ref_required"}) return if action not in actions: self.send_json(400, {"error": "unsupported_action", "supported_actions": list(actions)}) return started = time.monotonic() try: result = adapter_rpc(actions[action], {"base_id": base_id, "ref": ref}) self.send_json(200, {**result, "observer": {"duration_ms": int((time.monotonic() - started) * 1000), "method": actions[action]}}) except RuntimeError as exc: self.send_json(502, {"error": str(exc)}) return if parsed.path in {"/", "/index.html"}: self.serve_file("index.html") return if parsed.path.startswith("/assets/"): self.serve_file(parsed.path.lstrip("/")) return self.send_error(HTTPStatus.NOT_FOUND) def do_HEAD(self) -> None: # noqa: N802 parsed = urlparse(self.path) known = parsed.path in {"/", "/index.html", "/health", "/api/summary", "/api/events", "/api/mcp-events", "/api/correlations", "/api/coverage", "/api/objects", "/api/object", "/api/object/action"} or parsed.path.startswith("/assets/") self.send_response(HTTPStatus.OK if known else HTTPStatus.NOT_FOUND) self.send_header("Cache-Control", "no-store") self.end_headers() if __name__ == "__main__": if ADAPTER_URL: threading.Thread(target=coverage_worker, name="coverage-snapshot", daemon=True).start() ThreadingHTTPServer((HOST, PORT), Handler).serve_forever()