Знакът на КАГАМИ КАГАМИ
kagami.bg/academy · lesson · machine-readable viewUPDATED 2026-10-03
IDENTITY
module
GX10-04-145 · A live multi-site dashboard with FastAPI, WebSocket and PostgreSQL NOTIFY
series
GX10 (local AI server class: NVIDIA GB10, e.g. ASUS Ascent GX10 / DGX Spark)
level
Advanced
duration
3–4 h
prerequisites
PostgreSQL 14+, Python 3.10+, FFmpeg, packages fastapi, uvicorn[standard], asyncpg, httpx, pydantic v2; optional Grafana; a central server (a GB10-class machine works) and one small computer per site
trust_label
UPDATED 2026-10-03 (rewritten against PostgreSQL NOTIFY, asyncpg, FastAPI and Grafana documentation as of 2026-10-03) · NOT TESTED (no GB10 machine and no running database; Python files syntax-checked, client.js checked with node --check; FFmpeg command and Grafana queries not run)
versions
PostgreSQL 14+ (current: 18) · Python 3.10+ · packages not pinned (use current releases, then pin)
language
human view: bg · english edition: /en/academy/gx10/ (same file name)
previous / next
04-144_Ajax_AI_Integration.html / 04-146_CCTV_DPIA_Auto_Generator.html
PURPOSE

Build a live operations dashboard for many sites: edge agents push events and a 30-second heartbeat over HTTPS (one API key per site, stored only as a SHA-256 hash); a central FastAPI service stores them in PostgreSQL; an AFTER INSERT trigger sends pg_notify with only the row id; a dedicated asyncpg listener reads the row and broadcasts it over WebSocket to authenticated operators; if the listener reconnects, clients receive a resync and reload a full snapshot. Camera frames are fetched on request with FFmpeg, kept in memory for 10 seconds and never written to disk. Grafana queries cover history and SLA. Personal data is minimised: counts per shift, no names or positions. Consult a lawyer before using real cameras or staff data.

KEY CONCEPTS
COMMANDS / PATHS
CHECKLIST
NEXT MODULE

04-146 · CCTV DPIA draft with a local AI (04-146_CCTV_DPIA_Auto_Generator.html) · series index: kagami.bg/academy/gx10/ · offer: Quick experiment (kagami.bg/stalbata/)

SOURCES
TAGS
gx10nvidia-gb10dashboardwebsocketpostgresqlnotifyfastapigrafanaffmpeggdpr
ОБНОВЕНО · 03.10.2026

Табло за много обекти на живо: WebSocket и PostgreSQL NOTIFY

Организация с много обекти иска един екран: кой обект е в ред, къде има сигнал и колко бързо е реагирано. Строим го с централна PostgreSQL база, малък агент на всеки обект, FastAPI с WebSocket и PostgreSQL NOTIFY за доставка на живо. Пазим минимум данни.

⏱ 3–4 ч Напреднало GX10 FastAPI · PostgreSQL · FFmpeg · Grafana
FastAPI · asyncpg · PostgreSQL🔒 локално FFmpeg🔒 локално Grafana🔒 локално
🔄
ОБНОВЕНО · 03.10.2026 — какво е променено
Урокът е написан наново. Махнахме примерите с конкретни видове обекти и клиенти, твърденията, че на всеки обект стои GX10, и „под секунда“ (не е мерено). Поправихме грешки в старата версия: SQL за таблица с месечно разделяне не беше валиден (разделена таблица без PARTITION BY и общ ключ), в таблицата с камери стоеше адрес с потребител и парола, WebSocket приемаше всеки, а името на оператора идваше от съобщението на браузъра, кадрите се взимаха с блокираща команда в асинхронна услуга и се пишеха във временна папка, а за часовете на кеширане се смесваха два различни часовника. Старият JSON за Grafana беше непълен (без източник на данни) и е заменен със заявки. Добавихме: ключ на обект (проверка с хеш), уникален uid срещу дубликати, удостоверяване на оператора през първо съобщение и проверка на произхода, lifespan вместо остарелия on_event, нов старт на слушателя с пълно презареждане, кадър в паметта без запис на диск и раздел за данните и правната рамка. Позициите на отделни охранители не се пазят — само брой на смяна.
⚠️
Какво не сме пускали сами
При проверката нямахме машина от класа GB10 и работеща PostgreSQL база. Python файловете са проверени за синтактични грешки, а client.js — със node --check; нищо не е пускано, затова няма етикет „ТЕСТВАНО“ нито „ПРОВЕРЕНО“. Командата за FFmpeg и заявките за Grafana не са пускани. Поведението на NOTIFY е по документацията на PostgreSQL, не по наш опит. Скоростта и броят обекти, които издържа, не са мерени.

01Какво ще научиш

02Преди да започнеш

⚖️
Данни, видео и правна рамка
Кадрите от камери показват хора, а данните за смени и присъствие на служители са лични данни. Събирането, показването и съхранението им попадат под GDPR и българското право; основание, информиране на хората, срок на съхранение, достъп и права на засегнатите зависят от обекта и целта. Препоръка: преди да пуснеш система с реални обекти и камери, провери с юрист; за всичко, което засяга работното време и наблюдението на служители, също се допитай до юрист по трудово право. Урокът е технически и не е правна консултация. Технически добра практика, която улеснява юриста: кадър само при поискване и само в паметта, нищо на диска; в базата — само брой на смяна, не имена и позиции; запис кой кога е гледал.

03Стъпки

  1. Как е подредено всичко

    Агентът на обекта изпраща събития и сърдечен удар на всеки 30 секунди към централната услуга. Услугата записва в PostgreSQL. База с тригер праща NOTIFY, слушателят в услугата получава съобщението, чете реда и го разпраща по WebSocket до отворените браузъри.

    text
    обект (агент)  --HTTPS-->  централна услуга (FastAPI)  -->  PostgreSQL
                                                                  |
                                                           NOTIFY (id на реда)
                                                                  v
    браузър на оператора  <--WebSocket--  слушател в услугата

    Източникът на събития на обекта е твой: агентът чете файл (events.jsonl), в който записва каквото имаш — например изнесено от софтуера за мониторинг. Как се получават събитията (стъпка 1 от урок 04-144) зависи от твоето оборудване.

  2. Таблиците

    Всеки обект има собствен ключ: в базата стои само хешът му, така че изтичане на базата не дава ключовете. Полето uid е уникално — ако агентът изпрати същото събитие два пъти, втория път не се записва. Частичният индекс events_open_idx ускорява търсенето на непотвърдените. В камерите няма адрес с парола — само credential_ref.

    sql · schema.sql
    -- Needs PostgreSQL 14 or newer
    CREATE TABLE sites (
        id           SERIAL PRIMARY KEY,
        name         TEXT NOT NULL,                -- a neutral label, e.g. 'Site 001'
        sla_minutes  INTEGER NOT NULL DEFAULT 5,   -- agreed response time
        api_key_hash TEXT NOT NULL,                -- sha256 (hex) of this site's edge key
        active       BOOLEAN NOT NULL DEFAULT TRUE
    );
    
    CREATE TABLE site_health (
        site_id         INTEGER PRIMARY KEY REFERENCES sites(id),
        last_heartbeat  TIMESTAMPTZ,
        cameras_online  INTEGER,
        cameras_total   INTEGER,
        min_battery_pct INTEGER,
        guards_on_shift INTEGER
    );
    
    CREATE TABLE events (
        id          BIGSERIAL PRIMARY KEY,
        uid         TEXT NOT NULL UNIQUE,          -- set by the edge agent; blocks duplicates
        site_id     INTEGER NOT NULL REFERENCES sites(id),
        event_type  TEXT NOT NULL,
        severity    TEXT NOT NULL CHECK (severity IN ('critical','high','medium','low','info')),
        device_ref  TEXT,
        description TEXT,
        created_at  TIMESTAMPTZ NOT NULL DEFAULT now(),
        ack_by      TEXT,
        ack_at      TIMESTAMPTZ
    );
    CREATE INDEX ON events (created_at DESC);
    CREATE INDEX ON events (site_id, created_at DESC);
    CREATE INDEX events_open_idx ON events (created_at DESC) WHERE ack_at IS NULL;
    
    CREATE TABLE cameras (
        id             SERIAL PRIMARY KEY,
        site_id        INTEGER NOT NULL REFERENCES sites(id),
        name           TEXT NOT NULL,
        credential_ref TEXT NOT NULL    -- a reference; the stream address with its password is NOT stored here
    );
    
    -- Send only the row id: the listener reads the row itself.
    CREATE FUNCTION notify_new_event() RETURNS trigger
    LANGUAGE plpgsql AS $$
    BEGIN
        PERFORM pg_notify('new_event', NEW.id::text);
        RETURN NEW;
    END;
    $$;
    
    CREATE TRIGGER events_notify
    AFTER INSERT ON events
    FOR EACH ROW EXECUTE FUNCTION notify_new_event();

    За много големи обеми таблицата може да се разделя по време. Това изисква друг дизайн на ключовете и не е в този урок — виж раздела за разделяне на таблици в документацията на PostgreSQL.

  3. Агентът на обекта

    Агентът чете само нови редове от файла, праща ги един по един и мести отметката едва след успешно изпращане. Ако централата не отговаря, опитва пак след 30 секунди; нищо не се губи, а дубликатите ги спира uid. Не печата адреси и ключове в съобщенията за грешки.

    python · edge_agent.py
    import asyncio
    import json
    import os
    from pathlib import Path
    
    import httpx
    
    CENTRAL = os.environ["CENTRAL_API_URL"]        # https://<central-address>/api
    SITE_ID = int(os.environ["SITE_ID"])
    HEADERS = {"X-API-Key": os.environ["EDGE_API_KEY"]}
    SPOOL = Path(os.environ.get("EDGE_SPOOL", "/var/spool/edge/events.jsonl"))
    HEALTH_FILE = Path(os.environ.get("EDGE_HEALTH", "/var/spool/edge/health.json"))
    HEARTBEAT_EVERY = 30   # seconds
    
    
    def collect_health() -> dict:
        """Whatever local check you run writes this small JSON file; the agent only forwards it."""
        try:
            return json.loads(HEALTH_FILE.read_text())
        except (OSError, ValueError):
            return {}
    
    
    async def push_new_events(client: httpx.AsyncClient, offset: int) -> int:
        """Send every new line of the spool file. Move the offset only after a successful send."""
        if not SPOOL.exists():
            return offset
        with SPOOL.open("rb") as f:
            f.seek(offset)
            while True:
                line = f.readline()
                if not line or not line.endswith(b"\n"):
                    break                                  # no more complete lines
                try:
                    event = json.loads(line)               # needs: uid, event_type, severity
                except ValueError:
                    offset += len(line)                    # skip a broken line
                    continue
                r = await client.post(f"{CENTRAL}/sites/{SITE_ID}/events", json=event)
                if r.status_code not in (200, 201):
                    break                                  # try again on the next round
                offset += len(line)
        return offset
    
    
    async def main() -> None:
        offset = SPOOL.stat().st_size if SPOOL.exists() else 0   # start with new events only
        async with httpx.AsyncClient(headers=HEADERS, timeout=10) as client:
            while True:
                try:
                    offset = await push_new_events(client, offset)
                    await client.put(f"{CENTRAL}/sites/{SITE_ID}/health", json=collect_health())
                except httpx.HTTPError as e:
                    print("central unreachable:", type(e).__name__)   # do not print URLs or keys
                await asyncio.sleep(HEARTBEAT_EVERY)
    
    
    if __name__ == "__main__":
        asyncio.run(main())
    bash
    SITE_KEY="$(openssl rand -hex 32)"
    HASH="$(printf %s "$SITE_KEY" | sha256sum | cut -d' ' -f1)"
    psql "$DB_DSN" -c "INSERT INTO sites (name, api_key_hash) VALUES ('Site 001', '$HASH') RETURNING id;"
    echo "$SITE_KEY"    # put this key on the site, then clear your terminal history
    
    export CENTRAL_API_URL="https://<central-address>/api"
    export SITE_ID=1
    export EDGE_API_KEY="<site-key>"
    python3 edge_agent.py
  4. Централната услуга

    Един файл server.py прави четири неща: приема данни от обектите (по ключ на обект), отговаря с начален снимък, разпраща нови сигнали и дава кадри при поискване. Няколко решения, които си струва да обясним:

    • Оператор: първото съобщение по WebSocket е auth с лична лексема; без него връзката се затваря за 10 секунди. Името на оператора при потвърждение идва от лексемата, а не от съобщението.
    • Произход: ако е зададен ALLOWED_ORIGINS, връзка от друг произход се отказва.
    • Състоянието на обекта се смята в заявката от пулса и непотвърдените сигнали (таблицата по-долу).
    • Старт и спиране са в lifespan — начинът, препоръчан в документацията на FastAPI; on_event е обявен за остарял.
    СъстояниеКога
    offlineняма сърдечен удар от 90 секунди (три пропуснати по 30 секунди)
    redима непотвърден сигнал HIGH или CRITICAL
    amberнепотвърден сигнал MEDIUM или камера извън линия
    greenвсичко друго
    python · server.py
    import asyncio
    import hashlib
    import hmac
    import json
    import os
    import time
    from contextlib import asynccontextmanager
    from typing import Optional
    
    import asyncpg
    from fastapi import Depends, FastAPI, Header, HTTPException, Response, WebSocket, WebSocketDisconnect
    from pydantic import BaseModel, Field
    
    DB_DSN = os.environ["DB_DSN"]
    # "token1:Name One,token2:Name Two" - one long random token per operator.
    OPERATORS = dict(p.split(":", 1) for p in os.environ["OPERATOR_TOKENS"].split(",") if ":" in p)
    ALLOWED_ORIGINS = set(filter(None, os.environ.get("ALLOWED_ORIGINS", "").split(",")))
    
    
    class Hub:
        """Keeps the open browser connections and sends one message to all of them."""
    
        def __init__(self) -> None:
            self.clients: dict[WebSocket, str] = {}
    
        async def broadcast(self, message: dict) -> None:
            text = json.dumps(message, default=str)
            for ws in list(self.clients):
                try:
                    await ws.send_text(text)
                except Exception:
                    self.clients.pop(ws, None)
    
    
    hub = Hub()
    
    SITES_SQL = """
    SELECT * FROM (
      SELECT s.id, s.name,
        CASE
          WHEN h.last_heartbeat IS NULL OR h.last_heartbeat < now() - interval '90 seconds' THEN 'offline'
          WHEN EXISTS (SELECT 1 FROM events e WHERE e.site_id = s.id AND e.ack_at IS NULL
                       AND e.severity IN ('critical','high')) THEN 'red'
          WHEN EXISTS (SELECT 1 FROM events e WHERE e.site_id = s.id AND e.ack_at IS NULL
                       AND e.severity = 'medium')
               OR COALESCE(h.cameras_online < h.cameras_total, false) THEN 'amber'
          ELSE 'green'
        END AS status,
        h.cameras_online, h.cameras_total, h.min_battery_pct, h.guards_on_shift, h.last_heartbeat
      FROM sites s LEFT JOIN site_health h ON h.site_id = s.id
      WHERE s.active
    ) t
    ORDER BY CASE status WHEN 'red' THEN 1 WHEN 'offline' THEN 2 WHEN 'amber' THEN 3 ELSE 4 END, name
    """
    
    ALERTS_SQL = """
    SELECT id, site_id, event_type, severity, description, created_at
    FROM events
    WHERE ack_at IS NULL AND created_at > now() - interval '24 hours'
    ORDER BY CASE severity WHEN 'critical' THEN 1 WHEN 'high' THEN 2 WHEN 'medium' THEN 3 ELSE 4 END,
             created_at DESC
    LIMIT 100
    """
    
    
    async def snapshot(pool: asyncpg.Pool) -> dict:
        async with pool.acquire() as con:
            sites = [dict(r) for r in await con.fetch(SITES_SQL)]
            alerts = [dict(r) for r in await con.fetch(ALERTS_SQL)]
        return {"type": "snapshot", "sites": sites, "alerts": alerts}
    
    
    async def listener(app: FastAPI) -> None:
        """One dedicated connection listens for NOTIFY. If it drops, reconnect and tell clients to reload."""
        while True:
            try:
                conn = await asyncpg.connect(DB_DSN)
    
                def on_notify(_conn, _pid, _channel, payload):
                    asyncio.create_task(push_event(app, int(payload)))
    
                await conn.add_listener("new_event", on_notify)
                await hub.broadcast({"type": "resync"})        # notifications may have been missed
                while True:
                    await asyncio.sleep(30)
                    await conn.execute("SELECT 1")             # raises if the connection is gone
            except Exception as e:
                print("listener restart:", type(e).__name__)
                await asyncio.sleep(3)
    
    
    async def push_event(app: FastAPI, event_id: int) -> None:
        async with app.state.pool.acquire() as con:
            row = await con.fetchrow(
                "SELECT id, site_id, event_type, severity, description, created_at "
                "FROM events WHERE id = $1", event_id)
        if row:
            await hub.broadcast({"type": "new_alert", "data": dict(row)})
    
    
    @asynccontextmanager
    async def lifespan(app: FastAPI):
        app.state.pool = await asyncpg.create_pool(DB_DSN, min_size=1, max_size=8)
        app.state.thumbs: dict[int, tuple[float, bytes]] = {}
        app.state.locks: dict[int, asyncio.Lock] = {}
        task = asyncio.create_task(listener(app))
        yield
        task.cancel()
        await app.state.pool.close()
    
    
    app = FastAPI(title="Multi-site dashboard", lifespan=lifespan)
    
    
    # ---------- edge agents (one key per site) ----------
    async def site_auth(site_id: int, x_api_key: Optional[str] = Header(default=None)) -> int:
        if not x_api_key:
            raise HTTPException(status_code=401)
        async with app.state.pool.acquire() as con:
            stored = await con.fetchval("SELECT api_key_hash FROM sites WHERE id = $1 AND active", site_id)
        given = hashlib.sha256(x_api_key.encode()).hexdigest()
        if not stored or not hmac.compare_digest(given, stored):
            raise HTTPException(status_code=401)
        return site_id
    
    
    class EventIn(BaseModel):
        uid: str = Field(min_length=1, max_length=64)
        event_type: str = Field(min_length=1, max_length=40)
        severity: str = Field(pattern="^(critical|high|medium|low|info)$")
        device_ref: Optional[str] = Field(default=None, max_length=64)
        description: Optional[str] = Field(default=None, max_length=300)
    
    
    class HealthIn(BaseModel):
        cameras_online: Optional[int] = None
        cameras_total: Optional[int] = None
        min_battery_pct: Optional[int] = None
        guards_on_shift: Optional[int] = None
    
    
    @app.post("/api/sites/{site_id}/events", status_code=201)
    async def add_event(ev: EventIn, site_id: int = Depends(site_auth)):
        async with app.state.pool.acquire() as con:
            await con.execute(
                "INSERT INTO events (uid, site_id, event_type, severity, device_ref, description) "
                "VALUES ($1,$2,$3,$4,$5,$6) ON CONFLICT (uid) DO NOTHING",
                ev.uid, site_id, ev.event_type, ev.severity, ev.device_ref, ev.description)
        return {"ok": True}
    
    
    @app.put("/api/sites/{site_id}/health")
    async def put_health(h: HealthIn, site_id: int = Depends(site_auth)):
        async with app.state.pool.acquire() as con:
            await con.execute(
                "INSERT INTO site_health (site_id, last_heartbeat, cameras_online, cameras_total, "
                "min_battery_pct, guards_on_shift) VALUES ($1, now(), $2, $3, $4, $5) "
                "ON CONFLICT (site_id) DO UPDATE SET last_heartbeat = now(), "
                "cameras_online = $2, cameras_total = $3, min_battery_pct = $4, guards_on_shift = $5",
                site_id, h.cameras_online, h.cameras_total, h.min_battery_pct, h.guards_on_shift)
        return {"ok": True}
    
    
    # ---------- operators (browser) ----------
    @app.websocket("/ws/dashboard")
    async def dashboard(ws: WebSocket):
        if ALLOWED_ORIGINS and ws.headers.get("origin") not in ALLOWED_ORIGINS:
            await ws.close(code=1008)
            return
        await ws.accept()
        try:
            first = await asyncio.wait_for(ws.receive_json(), timeout=10)
            name = OPERATORS.get(first.get("token", "")) if first.get("type") == "auth" else None
            if not name:
                await ws.close(code=1008)
                return
            hub.clients[ws] = name
            await ws.send_text(json.dumps(await snapshot(app.state.pool), default=str))
            while True:
                msg = await ws.receive_json()
                if msg.get("type") == "get_snapshot":
                    await ws.send_text(json.dumps(await snapshot(app.state.pool), default=str))
                elif msg.get("type") == "ack_alert":
                    # the operator is the authenticated one, never a name sent by the browser
                    async with app.state.pool.acquire() as con:
                        done = await con.fetchval(
                            "UPDATE events SET ack_by = $2, ack_at = now() "
                            "WHERE id = $1 AND ack_at IS NULL RETURNING id", int(msg["event_id"]), name)
                    if done:
                        await hub.broadcast({"type": "ack", "event_id": done, "by": name})
        except (WebSocketDisconnect, asyncio.TimeoutError, ValueError):
            pass
        finally:
            hub.clients.pop(ws, None)
    
    
    async def grab_frame(url: str) -> Optional[bytes]:
        """One JPEG frame, straight into memory: nothing is written to disk."""
        proc = await asyncio.create_subprocess_exec(
            "ffmpeg", "-loglevel", "error", "-rtsp_transport", "tcp", "-i", url,
            "-frames:v", "1", "-q:v", "3", "-vf", "scale=640:-1", "-f", "image2pipe", "-vcodec", "mjpeg", "pipe:1",
            stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL)
        try:
            out, _ = await asyncio.wait_for(proc.communicate(), timeout=10)
        except asyncio.TimeoutError:
            proc.kill()
            await proc.wait()
            return None
        return out if proc.returncode == 0 and out else None
    
    
    def stream_url_for(credential_ref: str) -> Optional[str]:
        """Resolve the address from a secret store; here, from an environment variable CAM_<ref>."""
        return os.environ.get(f"CAM_{credential_ref}")
    
    
    @app.get("/api/cameras/{camera_id}/thumb.jpg")
    async def thumb(camera_id: int, authorization: Optional[str] = Header(default=None)):
        token = (authorization or "").removeprefix("Bearer ").strip()
        if token not in OPERATORS:
            raise HTTPException(status_code=401)
        cached = app.state.thumbs.get(camera_id)
        if cached and time.time() - cached[0] < 10:                  # 10-second cache
            return Response(cached[1], media_type="image/jpeg")
        lock = app.state.locks.setdefault(camera_id, asyncio.Lock())
        async with lock:
            cached = app.state.thumbs.get(camera_id)
            if cached and time.time() - cached[0] < 10:
                return Response(cached[1], media_type="image/jpeg")
            async with app.state.pool.acquire() as con:
                ref = await con.fetchval("SELECT credential_ref FROM cameras WHERE id = $1", camera_id)
            url = stream_url_for(ref) if ref else None
            data = await grab_frame(url) if url else None
            if not data:
                raise HTTPException(status_code=503, detail="camera offline")
            app.state.thumbs[camera_id] = (time.time(), data)
            return Response(data, media_type="image/jpeg")
    bash
    python3 -m venv .venv && . .venv/bin/activate
    pip install fastapi "uvicorn[standard]" asyncpg httpx pydantic
    psql "postgresql://<потребител>:<парола>@localhost/dashboard" -f schema.sql
    export DB_DSN="postgresql://<потребител>:<парола>@localhost/dashboard"
    export OPERATOR_TOKENS="$(openssl rand -hex 24):Operator One"
    export ALLOWED_ORIGINS="https://<central-address>"
    uvicorn server:app --host 127.0.0.1 --port 8002
  5. PostgreSQL NOTIFY: какво казва документацията

    Тригерът в schema.sql праща pg_notify('new_event', id). Три неща от документацията на PostgreSQL (03.10.2026) определят дизайна:

    • Съобщението носи най-много 8000 байта (при настройките по подразбиране); за повече се слага ключът на запис в таблица. Затова пращаме само id, а слушателят чете реда.
    • Съобщението се доставя след завършване на транзакцията; пази транзакциите кратки.
    • Който не слуша в момента, не получава нищо. Затова при прекъсване на връзката на слушателя услугата се свързва наново и праща на браузърите „resync“, а те зареждат пълен снимък.
    💡
    Слушател на отделна връзка
    asyncpg поддържа add_listener(канал, обратно_извикване); обратното извикване получава връзката, PID на сървъра, канала и съобщението. Слушателят държи собствена връзка, различна от пула.
  6. Браузърът

    Минимален клиент: свързва се, праща auth, получава снимък, добавя нови сигнали, маха потвърдените и при „resync“ иска нов снимък. Получените текстове се слагат с textContent, никога с innerHTML — иначе един сигнал с злонамерен текст може да изпълни код в страницата на оператора. Адресът е wss:// (зад обратен прокси с HTTPS).

    javascript · client.js
    // Minimal browser client. Replace <central-address> and keep the token out of the page source.
    const token = prompt("Operator token");           // for the lesson only; use a real sign-in
    const ws = new WebSocket("wss://<central-address>/ws/dashboard");
    const state = { sites: [], alerts: [] };
    
    ws.onopen = () => ws.send(JSON.stringify({ type: "auth", token }));
    
    ws.onmessage = (e) => {
      const m = JSON.parse(e.data);
      if (m.type === "snapshot") { state.sites = m.sites; state.alerts = m.alerts; }
      else if (m.type === "new_alert") { state.alerts.unshift(m.data); }
      else if (m.type === "ack") { state.alerts = state.alerts.filter(a => a.id !== m.event_id); }
      else if (m.type === "resync") { ws.send(JSON.stringify({ type: "get_snapshot" })); }
      render();
    };
    
    function ack(id) { ws.send(JSON.stringify({ type: "ack_alert", event_id: id })); }
    
    function render() {
      const el = document.getElementById("app");
      el.textContent = "";                              // textContent, never innerHTML, for received text
      for (const s of state.sites) {
        const d = document.createElement("div");
        d.className = "tile " + s.status;
        d.textContent = s.name + " · " + s.status;
        el.appendChild(d);
      }
      for (const a of state.alerts) {
        const b = document.createElement("button");
        b.textContent = a.severity + " · site " + a.site_id + " · " + a.event_type + " (acknowledge)";
        b.onclick = () => ack(a.id);
        el.appendChild(b);
      }
    }
    
    // A camera picture needs the header, so fetch it and show it as a blob.
    async function showThumb(img, cameraId) {
      const r = await fetch("/api/cameras/" + cameraId + "/thumb.jpg", { headers: { Authorization: "Bearer " + token } });
      if (r.ok) img.src = URL.createObjectURL(await r.blob());
    }
  7. Кадри от камери при поискване

    Когато операторът избере обект, браузърът иска кадър. Услугата пуска FFmpeg с -rtsp_transport tcp, взема един кадър (-frames:v 1), смалява го и го връща направо от паметта. Кадърът стои в кеш 10 секунди и не се записва на диск. Заявката иска Authorization: Bearer, а картинката се зарежда с fetch и blob, защото таг <img> не може да праща заглавки. Кодът е в server.py (grab_frame и thumb).

    • Адресът на потока (с парола) не се пази в базата: credential_ref сочи към тайник; в примера — променлива на средата CAM_<ref>.
    • ⚠️ Адресът е аргумент на FFmpeg и се вижда в списъка с процеси на машината. На машина с много потребители това е риск; ограничи достъпа до нея.
    • Не пишем адреса в журнала; грешките на FFmpeg са изключени.
  8. Справки в Grafana

    Живият екран е за деня. За история и договорени срокове добави в Grafana източник на данни PostgreSQL (с потребител само за четене) и четири панела със заявките по-долу. Макросът $__timeFilter(колона) е на източника и се замества с избрания времеви диапазон. Подробно за Grafana — в урок 04-122.

    sql
    -- 1) open alerts by severity (a stat panel)
    SELECT severity, count(*) AS alerts
    FROM events
    WHERE ack_at IS NULL AND created_at > now() - interval '24 hours'
    GROUP BY severity;
    
    -- 2) alerts per site in the chosen time range (a bar chart)
    SELECT s.name, count(*) AS alerts
    FROM events e JOIN sites s ON s.id = e.site_id
    WHERE $__timeFilter(e.created_at)
    GROUP BY s.name ORDER BY alerts DESC LIMIT 10;
    
    -- 3) response time against the agreed one (a table)
    SELECT s.name,
           round((avg(extract(epoch FROM (e.ack_at - e.created_at)) / 60))::numeric, 1) AS avg_minutes,
           s.sla_minutes,
           count(*) FILTER (WHERE e.ack_at - e.created_at > make_interval(mins => s.sla_minutes)) AS late
    FROM events e JOIN sites s ON s.id = e.site_id
    WHERE e.ack_at IS NOT NULL AND $__timeFilter(e.created_at)
    GROUP BY s.name, s.sla_minutes
    ORDER BY late DESC;
    
    -- 4) sites that have gone quiet (a table)
    SELECT s.name, h.last_heartbeat
    FROM sites s LEFT JOIN site_health h ON h.site_id = s.id
    WHERE s.active AND (h.last_heartbeat IS NULL OR h.last_heartbeat < now() - interval '90 seconds');

    Даваме заявките, а не готов JSON за импорт: версиите на Grafana се различават, а JSON без източник на данни не върши работа.

  9. Опитай с измислени данни

    Създай обект с ключ (виж стъпка 3), пусни услугата и прати тестово събитие:

    bash
    curl -s -X POST "http://127.0.0.1:8002/api/sites/1/events" \
      -H "X-API-Key: $SITE_KEY" -H "Content-Type: application/json" \
      -d '{"uid":"test-0001","event_type":"alarm","severity":"high","description":"test event"}'
    1. Отвори страница с client.js: сигналът трябва да се появи без презареждане, а потвърждаването да го махне на всички отворени екрани.
    2. Пусни агента, после го спри: след около 90 секунди обектът става offline.
    3. Спри и пусни наново PostgreSQL: екраните получават „resync“ и зареждат нов снимък. Какво точно ще видиш, не обещаваме — не сме го пускали.
  10. Сигурност — минимумът

    • Всеки обект има собствен дълъг ключ; в базата стои само хешът.
    • Всеки оператор има собствена лексема; за истинска употреба сложи нормален вход, не споделена лексема.
    • Услугата слуша на 127.0.0.1; отвън — само през обратен прокси с HTTPS и wss://.
    • Потребителят на Grafana е само за четене.
    • Копия на базата и срок за съхранение, уговорен с юриста.

04Проверка

Тест

1. Какво е правилно да носи съобщението на NOTIFY?

2. Кога се доставя NOTIFY, пратен в транзакция?

3. Защо след нов старт на слушателя браузърите зареждат пълен снимък?

4. Откъде идва името на оператора при потвърждение на сигнал?

05Какво следва

06Източници

  1. PostgreSQL: NOTIFY — 8000 байта, доставка след транзакцията, ключ на запис.
  2. asyncpg: add_listener — подпис на обратното извикване.
  3. FastAPI: WebSockets — WebSocket в FastAPI.
  4. FastAPI: lifespan — lifespan и остарелият on_event.
  5. Grafana: заявки към PostgreSQL — макроси като $__timeFilter.
  6. FFmpeg: документация — параметри за вход и изход (командата не е пускана).