from __future__ import annotations import hashlib import json from datetime import datetime, timezone from threading import Lock from typing import Any from backend import config _REGISTRY_LOCK = Lock() _AUDIT_LOCK = Lock() def _utc_now() -> str: return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") def _load_registry() -> dict[str, str]: with _REGISTRY_LOCK: if not config.USER_REGISTRY_PATH.exists(): return {} raw = config.USER_REGISTRY_PATH.read_text(encoding="utf-8").strip() if not raw: return {} try: data = json.loads(raw) except json.JSONDecodeError: data = {} if not isinstance(data, dict): return {} return {str(k): str(v) for k, v in data.items()} def _save_registry(registry: dict[str, str]) -> None: with _REGISTRY_LOCK: config.USER_REGISTRY_PATH.write_text( json.dumps(registry, indent=2, ensure_ascii=True), encoding="utf-8", ) def resolve_user(ip: str) -> str: safe_ip = ip or "unknown" registry = _load_registry() existing = registry.get(safe_ip) if existing: return existing assigned = f"user_{len(registry) + 1}" registry[safe_ip] = assigned _save_registry(registry) return assigned def audit_event( event_type: str, payload: dict[str, Any], *, user_n: str, ip: str, job_id: str | None = None, ) -> None: item = { "ts": _utc_now(), "event_type": event_type, "user_n": user_n, "ip": ip or "unknown", "job_id": job_id, "payload": payload, } with _AUDIT_LOCK: with config.AUDIT_LOG_PATH.open("a", encoding="utf-8") as handle: handle.write(json.dumps(item, ensure_ascii=True) + "\n") def summarize_text(value: str, *, max_chars: int = 2000) -> dict[str, Any]: safe = value or "" digest = hashlib.sha256(safe.encode("utf-8")).hexdigest() if len(safe) <= max_chars: return { "text_truncated": safe, "length": len(safe), "sha256": digest, "truncated": False, } return { "text_truncated": safe[:max_chars], "length": len(safe), "sha256": digest, "truncated": True, }