Fix Dockhand auth for topology flow and accurate CDC change stats.

Pass DOCKHAND_API_TOKEN to container inventory calls so pipeline_active
and topology animation work after Authentik. Replace ring-buffer-only CDC
stats with minute rollups (no 1000 cap), add 15m/1h/6h/24h window selector
on the Live Changes tab, and poll recent events on an interval.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
root
2026-07-01 15:14:48 +02:00
parent cdc8aa4bf7
commit b46e7f01dd
10 changed files with 537 additions and 106 deletions
+112 -43
View File
@@ -18,7 +18,7 @@ import json
import os
import re
import uuid
from collections import deque
from collections import defaultdict, deque
from datetime import datetime, timezone
from typing import Any
@@ -28,7 +28,8 @@ from fastapi.responses import JSONResponse
router = APIRouter(prefix="/api/changes", tags=["changes"])
KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "10.0.21.36:9092")
RING_SIZE = int(os.getenv("CDC_RING_SIZE", "1000"))
RING_SIZE = int(os.getenv("CDC_RING_SIZE", "5000"))
METRICS_RETENTION_MIN = int(os.getenv("CDC_METRICS_RETENTION_MIN", str(24 * 60)))
# Active CDC topic prefixes -> logical source. Matches Debezium topic.prefix.
PREFIX_SOURCE = {
@@ -45,6 +46,8 @@ TOPIC_PATTERN = re.compile(
OP_MAP = {"c": "insert", "u": "update", "d": "delete", "r": "snapshot"}
_ring: deque[dict[str, Any]] = deque(maxlen=RING_SIZE)
# Minute-level rollups — accurate counts beyond the display ring buffer cap.
_minute_rollups: dict[int, dict[str, Any]] = {}
_state: dict[str, Any] = {
"started": False,
"connected": False,
@@ -127,6 +130,83 @@ def _parse(topic: str, value: bytes | None) -> dict[str, Any] | None:
}
def _record_metrics(entry: dict[str, Any]) -> None:
"""Track per-minute aggregates for accurate window stats (not ring-limited)."""
try:
ts = datetime.fromisoformat(entry["ts"].replace("Z", "+00:00"))
except Exception:
ts = datetime.now(timezone.utc)
minute_key = int(ts.timestamp()) // 60
bucket = _minute_rollups.setdefault(
minute_key,
{"total": 0, "by_source": defaultdict(int), "by_op": defaultdict(int), "by_table": defaultdict(int)},
)
bucket["total"] += 1
bucket["by_source"][entry["source"]] += 1
bucket["by_op"][entry["op"]] += 1
bucket["by_table"][entry["table"]] += 1
cutoff = minute_key - METRICS_RETENTION_MIN
for old_key in list(_minute_rollups.keys()):
if old_key < cutoff:
del _minute_rollups[old_key]
def _aggregate_window(minutes: int) -> dict[str, Any]:
"""Aggregate rollup buckets for the requested time window."""
now = datetime.now(timezone.utc)
cutoff_minute = int((now.timestamp() - minutes * 60) // 60)
use_hour_buckets = minutes >= 120
by_source: dict[str, int] = defaultdict(int)
by_op: dict[str, int] = defaultdict(int)
by_table: dict[str, int] = defaultdict(int)
buckets: dict[str, int] = defaultdict(int)
total = 0
for minute_key, bucket in _minute_rollups.items():
if minute_key < cutoff_minute:
continue
total += bucket["total"]
for k, v in bucket["by_source"].items():
by_source[k] += v
for k, v in bucket["by_op"].items():
by_op[k] += v
for k, v in bucket["by_table"].items():
by_table[k] += v
ts = datetime.fromtimestamp(minute_key * 60, timezone.utc)
label = ts.strftime("%H:00") if use_hour_buckets else ts.strftime("%H:%M")
buckets[label] += bucket["total"]
# Fill chart timeline with zero buckets so the window span is visible.
if use_hour_buckets:
span = max(1, (minutes + 59) // 60)
filled: dict[str, int] = {}
for i in range(span):
t = datetime.fromtimestamp(now.timestamp() - (span - 1 - i) * 3600, timezone.utc)
filled[t.strftime("%H:00")] = buckets.get(t.strftime("%H:00"), 0)
buckets = filled
else:
span = max(1, minutes)
filled = {}
for i in range(span):
t = datetime.fromtimestamp(now.timestamp() - (span - 1 - i) * 60, timezone.utc)
filled[t.strftime("%H:%M")] = buckets.get(t.strftime("%H:%M"), 0)
buckets = filled
return {
"total": total,
"by_source": dict(by_source),
"by_op": dict(by_op),
"by_table": dict(by_table),
"buckets": [{"t": k, "n": v} for k, v in sorted(buckets.items())],
"inserts": by_op.get("insert", 0),
"updates": by_op.get("update", 0),
"deletes": by_op.get("delete", 0),
"rate_per_min": round(total / max(1, minutes), 1),
}
async def cdc_consumer_loop() -> None:
"""Background loop: consume CDC topics and publish each change live."""
_state["started"] = True
@@ -158,6 +238,7 @@ async def cdc_consumer_loop() -> None:
if not entry:
continue
_ring.append(entry)
_record_metrics(entry)
_state["consumed"] += 1
_state["last_ts"] = entry["ts"]
try:
@@ -182,30 +263,36 @@ async def cdc_consumer_loop() -> None:
def snapshot(minutes: int = 15) -> dict[str, Any]:
"""Lightweight CDC snapshot for other modules (Data Flow graph)."""
cutoff = datetime.now(timezone.utc).timestamp() - minutes * 60
by_source: dict[str, int] = {}
total = 0
for c in _ring:
try:
if datetime.fromisoformat(c["ts"]).timestamp() < cutoff:
continue
except Exception:
continue
total += 1
by_source[c["source"]] = by_source.get(c["source"], 0) + 1
return {"connected": _state["connected"], "consumed": _state["consumed"],
"buffered": len(_ring), "window_total": total, "by_source": by_source}
agg = _aggregate_window(minutes)
return {
"connected": _state["connected"],
"consumed": _state["consumed"],
"buffered": len(_ring),
"window_total": agg["total"],
"by_source": agg["by_source"],
}
# ── Endpoints ────────────────────────────────────────────────────────────────
@router.get("")
async def list_changes(
limit: int = Query(100, le=500),
limit: int = Query(200, le=1000),
source: str | None = None,
op: str | None = None,
table: str | None = None,
minutes: int | None = Query(None, le=1440),
) -> JSONResponse:
items = list(_ring)
if minutes is not None:
cutoff = datetime.now(timezone.utc).timestamp() - minutes * 60
filtered_by_time: list[dict[str, Any]] = []
for c in items:
try:
if datetime.fromisoformat(c["ts"].replace("Z", "+00:00")).timestamp() >= cutoff:
filtered_by_time.append(c)
except Exception:
continue
items = filtered_by_time
if source:
items = [c for c in items if c["source"] == source]
if op:
@@ -217,47 +304,29 @@ async def list_changes(
"ok": True,
"changes": items,
"buffered": len(_ring),
"buffer_cap": RING_SIZE,
"connected": _state["connected"],
"consumed": _state["consumed"],
"last_error": _state["last_error"],
"last_ts": _state["last_ts"],
})
@router.get("/stats")
async def change_stats(minutes: int = Query(15, le=240)) -> JSONResponse:
now = datetime.now(timezone.utc)
cutoff = now.timestamp() - minutes * 60
by_source: dict[str, int] = {}
by_op: dict[str, int] = {}
by_table: dict[str, int] = {}
buckets: dict[str, int] = {}
total = 0
for c in _ring:
try:
t = datetime.fromisoformat(c["ts"]).timestamp()
except Exception:
continue
if t < cutoff:
continue
total += 1
by_source[c["source"]] = by_source.get(c["source"], 0) + 1
by_op[c["op"]] = by_op.get(c["op"], 0) + 1
by_table[c["table"]] = by_table.get(c["table"], 0) + 1
bucket = datetime.fromtimestamp(t, timezone.utc).strftime("%H:%M")
buckets[bucket] = buckets.get(bucket, 0) + 1
async def change_stats(minutes: int = Query(15, le=1440)) -> JSONResponse:
agg = _aggregate_window(minutes)
return JSONResponse({
"ok": True,
"window_minutes": minutes,
"total": total,
"by_source": by_source,
"by_op": by_op,
"by_table": by_table,
"buckets": [{"t": k, "n": v} for k, v in sorted(buckets.items())],
**agg,
"connected": _state["connected"],
"consumed": _state["consumed"],
"buffered": len(_ring),
"buffer_cap": RING_SIZE,
"last_ts": _state["last_ts"],
})
@router.get("/status")
async def changes_status() -> JSONResponse:
return JSONResponse({"ok": True, "buffered": len(_ring), **_state})
return JSONResponse({"ok": True, "buffered": len(_ring), "buffer_cap": RING_SIZE, **_state})