Initial project import
This commit is contained in:
@@ -0,0 +1,345 @@
|
||||
"""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()
|
||||
Reference in New Issue
Block a user