346 lines
20 KiB
Python
346 lines
20 KiB
Python
"""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 "<none>"), 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 "<none>")].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()
|