Files
2026-07-17 12:57:18 +03:00

142 lines
5.3 KiB
Python

import logging
import time
from collections import defaultdict, deque
from threading import Lock, Semaphore
from typing import Any, Deque, Dict, Optional, Tuple
from config import Config
logger = logging.getLogger(__name__)
class SyncRuntime:
"""Потокобезопасные runtime-примитивы для sync API."""
def __init__(self) -> None:
self._semaphore = Semaphore(Config.SYNC_MAX_CONCURRENT)
self._rate_limits: Dict[str, Deque[float]] = defaultdict(deque)
self._pull_client_mutex = Lock()
self._pull_client_locks: Dict[str, Lock] = {}
# Дескрипторы flock (POSIX) — освобождать в release_pull_client.
self._ipc_pull_fps: Dict[str, Any] = {}
self._pull_durations_sec: Deque[float] = deque(maxlen=256)
def acquire_slot(self) -> bool:
return self._semaphore.acquire(blocking=False)
def release_slot(self) -> None:
self._semaphore.release()
def acquire_pull_client(self, client_id: str, *, wait_sec: Optional[float] = None) -> bool:
"""
Не более одного параллельного /pull на один node_id.
На POSIX — flock (видно всем процессам: reloader / несколько воркеров).
На Windows — только threading.Lock внутри процесса.
"""
cid = (client_id or "").strip()
if not cid:
return False
if wait_sec is None:
wait_sec = float(getattr(Config, "SYNC_PULL_CLIENT_LOCK_WAIT_SEC", 180.0))
from app.services.sync_pull_ipc import ipc_pull_acquire, ipc_pull_supported
if ipc_pull_supported():
fp = ipc_pull_acquire(cid, wait_sec)
if fp is None:
return False
self._ipc_pull_fps[cid] = fp
logger.debug("[SYNC-IPC] pull lock node %s…", cid[:12])
return True
with self._pull_client_mutex:
lk = self._pull_client_locks.get(cid)
if lk is None:
lk = Lock()
self._pull_client_locks[cid] = lk
if wait_sec <= 0:
return lk.acquire(blocking=False)
return lk.acquire(blocking=True, timeout=wait_sec)
def release_pull_client(self, client_id: str) -> None:
cid = (client_id or "").strip()
if not cid:
return
from app.services.sync_pull_ipc import ipc_pull_release, ipc_pull_supported
if ipc_pull_supported():
fp = self._ipc_pull_fps.pop(cid, None)
ipc_pull_release(fp)
return
lk = self._pull_client_locks.get(cid)
if not lk:
return
try:
lk.release()
except RuntimeError:
pass
def concurrency_status(self) -> Tuple[int, int]:
available = self._semaphore._value # type: ignore[attr-defined]
used = Config.SYNC_MAX_CONCURRENT - available
logger.debug(
"[SYNC-CONCURRENCY] Используется: %s/%s слотов синхронизации",
used,
Config.SYNC_MAX_CONCURRENT,
)
return used, available
def check_rate_limit(self, client_id: str) -> bool:
now = time.time()
window = 60
dq = self._rate_limits[client_id]
while dq and dq[0] <= now - window:
dq.popleft()
if len(dq) >= Config.SYNC_MAX_CONCURRENT * 10:
logger.warning("[SYNC-RATE-LIMIT] Клиент %s превысил лимит", client_id[:8])
return False
dq.append(now)
return True
def record_pull_duration(self, seconds: float) -> None:
self._pull_durations_sec.append(max(0.0, float(seconds)))
def pull_duration_p95(self) -> float:
items = sorted(self._pull_durations_sec)
if not items:
return 0.0
idx = int(0.95 * (len(items) - 1))
return float(items[idx])
def diagnostics_snapshot(self) -> Dict[str, Any]:
"""Метрики только текущего процесса (воркера), не кластера целиком."""
from app.services.sync_pull_ipc import ipc_pull_supported
used, available = self.concurrency_status()
tracked = sum(1 for dq in self._rate_limits.values() if dq)
ipc_clients = len(self._ipc_pull_fps)
thread_lock_clients = len(self._pull_client_locks)
return {
"scope": "process_local",
"note": (
"Счётчики слотов, p95 длительности pull и окна rate-limit относятся к этому "
"процессу Python. При нескольких воркерах uWSGI/gunicorn значения различаются."
),
"concurrency_used": used,
"concurrency_available": available,
"concurrency_max": int(Config.SYNC_MAX_CONCURRENT),
"pull_duration_p95_sec": self.pull_duration_p95(),
"pull_duration_samples": len(self._pull_durations_sec),
"pull_lock_ipc_supported": bool(ipc_pull_supported()),
"pull_lock_ipc_held_clients": ipc_clients,
"pull_lock_thread_clients": thread_lock_clients,
"rate_limit_windows_with_activity": tracked,
}
sync_runtime = SyncRuntime()