The KAGAMI mark КАГАМИ
kagami.bg/academy · lesson · machine-readable viewUPDATED 2026-10-03
IDENTITY
module
GX10-04-149 · An incident classifier with Redis Streams and a local multimodal model
series
GX10 (local AI server class: NVIDIA GB10, e.g. ASUS Ascent GX10 / DGX Spark)
level
Advanced
duration
about 1 h 30 min
prerequisites
A GB10-class machine with DGX OS and Ollama (llama3.2-vision:11b), Redis 6.2+ on loopback, PostgreSQL 14+, Python 3.10+, packages redis, httpx, asyncpg
trust_label
UPDATED 2026-10-03 (rewritten against Redis Streams documentation, the Ollama chat API reference and the Ollama library page for llama3.2-vision as of 2026-10-03) · NOT TESTED (no GB10 machine and no running Redis, Ollama or PostgreSQL; the Python files were syntax-checked only)
versions
llama3.2-vision:11b (7.8 GB, 128K context; 90b is 55 GB, not tried) · Redis 6.2+ (XAUTOCLAIM since 6.2.0) · PostgreSQL 14+ · Python 3.10+ · packages not pinned (use current releases, then pin)
language
human view: en · bulgarian edition: /academy/gx10/ (same file name)
previous / next
04-148_Shift_Manager.html / 04-151_CCTV_Health_Monitor.html
PURPOSE

Build an incident classification pipeline. Producers write events (text and an optional frame file name) to the Redis stream incidents:raw. A worker in a consumer group reads them with XREADGROUP, loads the frame only from an allowed folder, calls Ollama /api/chat with the frame in the images list and a JSON schema in format, and writes label, confidence and reason to incidents:classified. The priority comes from a fixed table by label, never from the model; low confidence caps the priority at P3, and an unreadable answer becomes P2 unclassified. An escalator worker writes an audit row (idempotent) and posts a suggestion to a webhook for a human for P1 and P2; failures leave messages pending and XAUTOCLAIM re-delivers them. Nothing is dispatched automatically. Consult a lawyer before using real cameras; this is not a certified safety system.

KEY CONCEPTS
COMMANDS / PATHS
CHECKLIST
NEXT MODULE

04-151 · Camera health monitoring (04-151_CCTV_Health_Monitor.html) · series index: kagami.bg/en/academy/gx10/ · offer: Quick experiment (kagami.bg/stalbata/)

SOURCES
TAGS
gx10nvidia-gb10redis-streamsclassificationmultimodalollamaauditpostgresqlhuman-in-the-loop
UPDATED · 03.10.2026

An Incident Classifier with Redis Streams and a Local Multimodal Model

People and cameras send notes and frames about incidents. We build a pipeline: events enter a Redis Stream, a worker classifies them with a local multimodal model, and a second worker writes an audit trail and notifies a human about the more urgent ones. The priority comes from a table, not from the model, and nothing is sent automatically to authorities.

⏱ 1 h 30 min Advanced GX10 Redis · Ollama · PostgreSQL · Python
Redis Streams🔒 local Ollama (llama3.2-vision)🔒 local PostgreSQL🔒 local
🔄
UPDATED · 03.10.2026 — what changed
The lesson was rewritten. We removed the name "Llama 14B" (the code in the old version used a model with 11 billion parameters), the claim that the picture reduces errors "by about 40%" (there was no source), the specific response times of one real organisation (only a scale remains, and your deadlines are set by your contract), automatic routing to the police and an ambulance, the code with a chat bot key and the blocking loops in asynchronous code. We also fixed an error: the picture was sent in a format that is not that of Ollama's native API — there the frame goes in the images list of the message. We added: priority from a table rather than from the model; behaviour for low confidence and for an unreadable answer (towards more attention, never towards "nothing"); the frame as a file name rather than inside Redis; re-claiming forgotten messages (XAUTOCLAIM); an audit table with columns for a human's judgement and queries to measure the errors.
⚠️
What we have not run ourselves
We had no GB10-class machine at hand, nor running Redis, Ollama and PostgreSQL. The Python files were checked for syntax errors but were not run — which is why there is neither a "TESTED" nor a "VERIFIED" label. The quality of the model on your incidents, the speed and how many errors it makes were not measured. The llama3.2-vision model officially supports only English when it also receives a picture; we have not tried Bulgarian text with a frame.

01What you'll learn

02Before you start

⚖️
Data, frames and the legal side
Camera frames show people, and incident descriptions often mention individuals. This is personal data; collecting, analysing and storing it falls under the GDPR and national law, and the legal basis, informing people, retention, access and the rights of those concerned depend on the site and the purpose. Recommendation: before you run a system with real cameras, check with a lawyer and write down the decisions. This lesson is technical and is not legal advice. Good practice that helps the lawyer: a frame only on an event, with a short retention period, no face recognition, and neither a frame nor a name in Redis or in the audit table. The classifier is a hint to an operator, not a certified safety system.

03Steps

  1. Redis Streams in plain words

    A stream is a log to which new entries are added at the end. A consumer group shares the entries between several workers, each entry going to one of them. The worker reads with XREADGROUP and, when done, says XACK. Until there is an acknowledgement the entry is "pending" and can be claimed again with XAUTOCLAIM — so nothing is lost if a worker crashes.

    text
    incidents:raw  --(worker: model)-->  incidents:classified  --(escalator)-->  PostgreSQL (audit)
       ^                                                              |
       |                                                              +--> webhook: a suggestion for a human
    event producers (text + a file name for the frame)

    Events are short: text and, if there is one, a file name for the frame. The picture itself is not put into Redis — it is heavy, and this way we do not copy it to one more place.

  2. The priority is a table, not the model's opinion

    The model returns only a label, a confidence and a short reason. The priority is taken from the table — it is predictable, easy to audit and does not change when you change the model. The scale below is an example; your contract sets the response times, so there are no numbers here.

    LabelPriorityHow we treat it
    intrusion · theft · fightP1straight to a human
    vandalism · trespassP2to a human
    suspiciousP3only recorded; a human sees it in the review
    false_alarmP4recorded; never deleted on its own
    low confidence (below 0.6)at most P3a human should look
    an unreadable answerP2 · unclassifiedmore attention, not less
  3. The common core: reading, acknowledging, re-claiming

    One small helper is used by both workers. The group is created with XGROUP CREATE ... MKSTREAM (if it already exists Redis answers BUSYGROUP, and that is normal). The worker takes new entries with > and then also the forgotten ones that have been idle for more than 60 seconds. An entry is acknowledged only if the processing returned True; an error or False leaves it pending.

    python · common.py
    import redis.asyncio as redis
    from redis.exceptions import ResponseError
    
    
    async def consume(r: redis.Redis, stream: str, group: str, consumer: str, handler,
                      idle_ms: int = 60000) -> None:
        """Read a stream through a consumer group.
        handler(msg_id, fields) -> True means 'done, acknowledge'. False or an error leaves the
        message pending; it is picked up again after idle_ms (XAUTOCLAIM)."""
        try:
            await r.xgroup_create(stream, group, id="0", mkstream=True)
        except ResponseError as e:
            if "BUSYGROUP" not in str(e):          # the group already exists: fine
                raise
        while True:
            batch = []
            res = await r.xreadgroup(group, consumer, {stream: ">"}, count=5, block=2000)
            for _stream, msgs in res or []:
                batch.extend(msgs)
            claimed = await r.xautoclaim(stream, group, consumer, min_idle_time=idle_ms,
                                         start_id="0-0", count=5)
            batch.extend(claimed[1])               # [next_id, messages, ...] depending on the version
            for msg_id, fields in batch:
                if fields is None:                 # the entry was deleted from the stream
                    await r.xack(stream, group, msg_id)
                    continue
                try:
                    ok = await handler(msg_id, fields)
                except Exception as e:
                    print("handler error:", type(e).__name__)
                    ok = False
                if ok:
                    await r.xack(stream, group, msg_id)
  4. The worker with the model

    The worker reads an event, loads the frame by name only from the designated folder (protection against paths such as ../), sends the text and the frame to Ollama and asks for an answer that follows a schema. The frame goes into the images list of the message, base64-encoded — that is how Ollama's API takes it. Temperature 0 makes the answers more repeatable.

    bash
    ollama pull llama3.2-vision:11b
    python · worker.py
    import asyncio
    import base64
    import json
    import os
    from pathlib import Path
    
    import httpx
    import redis.asyncio as redis
    
    from common import consume
    
    REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379")
    OLLAMA_URL = os.environ.get("OLLAMA_URL", "http://localhost:11434")
    MODEL = os.environ.get("INCIDENT_MODEL", "llama3.2-vision:11b")
    FRAME_DIR = Path(os.environ.get("FRAME_DIR", "/var/lib/incidents/frames")).resolve()
    CONSUMER = os.environ.get("CONSUMER", "worker-1")
    
    STREAM_IN, STREAM_OUT, GROUP = "incidents:raw", "incidents:classified", "classifiers"
    
    # Priority comes from this table, not from the model: it is predictable and easy to audit.
    PRIORITY = {"intrusion": 1, "theft": 1, "fight": 1, "vandalism": 2, "trespass": 2,
                "suspicious": 3, "false_alarm": 4}
    
    SCHEMA = {
        "type": "object",
        "properties": {
            "label": {"type": "string", "enum": list(PRIORITY)},
            "confidence": {"type": "number"},
            "reason": {"type": "string"},
        },
        "required": ["label", "confidence", "reason"],
    }
    SYSTEM = (
        "You classify incident reports for a human security operator. You may get a short "
        "description and one camera frame. Choose exactly one label: "
        + ", ".join(PRIORITY) + ". "
        "Use false_alarm only when the evidence clearly shows nothing is wrong. "
        "Give a confidence from 0 to 1 and a reason of at most 20 words. Reply with JSON only."
    )
    
    
    def load_frame(name: str):
        """Read a frame by file name, and only from inside FRAME_DIR."""
        path = (FRAME_DIR / name).resolve()
        if not path.is_relative_to(FRAME_DIR) or not path.is_file():
            return None
        return base64.b64encode(path.read_bytes()).decode()
    
    
    async def classify(text: str, frame_b64):
        msg = {"role": "user", "content": text or "(no description)"}
        if frame_b64:
            msg["images"] = [frame_b64]
        payload = {"model": MODEL, "stream": False, "format": SCHEMA,
                   "options": {"temperature": 0, "num_predict": 200},
                   "messages": [{"role": "system", "content": SYSTEM}, msg]}
        try:
            async with httpx.AsyncClient(timeout=120) as client:
                r = await client.post(f"{OLLAMA_URL}/api/chat", json=payload)
                r.raise_for_status()
                out = json.loads(r.json()["message"]["content"])
            label = out["label"]
            conf = float(out["confidence"])
            priority = PRIORITY[label]
            if conf < 0.6:                       # not sure: a human should look, never ignore
                priority = min(priority, 3)
            return {"label": label, "priority": priority, "confidence": conf,
                    "reason": str(out.get("reason", ""))[:200]}
        except Exception:
            # Fail toward attention: an unreadable answer is never treated as 'nothing happened'.
            return {"label": "unclassified", "priority": 2, "confidence": 0.0,
                    "reason": "the model did not return a usable answer"}
    
    
    async def main() -> None:
        r = redis.from_url(REDIS_URL, decode_responses=True)
    
        async def handle(msg_id: str, f: dict) -> bool:
            frame = load_frame(f["frame"]) if f.get("frame") else None
            res = await classify(f.get("text", ""), frame)
            await r.xadd(STREAM_OUT, {
                "event_id": msg_id, "zone_id": f.get("zone_id", ""),
                "label": res["label"], "priority": str(res["priority"]),
                "confidence": f"{res['confidence']:.2f}", "reason": res["reason"], "model": MODEL,
            }, maxlen=10000, approximate=True)
            print(f"{msg_id} -> P{res['priority']} {res['label']}")
            return True
    
        await consume(r, STREAM_IN, GROUP, CONSUMER, handle)
    
    
    if __name__ == "__main__":
        asyncio.run(main())

    If the confidence is below 0.6, the priority is at most P3 so that a human sees it. If the answer is unreadable or the model does not answer, the event gets P2 and the label unclassified: when in doubt we fall towards more attention. The stream incidents:classified has a length limit (maxlen); if a worker is stopped for long, the oldest entries can be trimmed, so watch the length (step 7).

  5. Audit and notifying a human

    The second worker writes every classified event into incident_audit (an event received again is not written twice: ON CONFLICT DO NOTHING) and, if the priority is P1 or P2, sends a message to an address that you set — for example a webhook in n8n that writes to the duty operator. The message is expressly a suggestion. If sending fails, the entry stays pending and is retried.

    sql · schema.sql
    CREATE TABLE incident_audit (
        event_id       TEXT PRIMARY KEY,            -- the Redis stream id of the raw event
        zone_id        INTEGER,
        label          TEXT NOT NULL,
        priority       INTEGER NOT NULL,            -- 1 (most urgent) .. 4
        confidence     REAL,
        reason         TEXT,
        model          TEXT,
        classified_at  TIMESTAMPTZ NOT NULL DEFAULT now(),
        -- filled in later by a person: this is how you measure the classifier
        human_label    TEXT,
        human_priority INTEGER,
        reviewed_by    TEXT,
        reviewed_at    TIMESTAMPTZ
    );
    python · escalator.py
    import asyncio
    import os
    
    import asyncpg
    import httpx
    import redis.asyncio as redis
    
    from common import consume
    
    REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379")
    DB_DSN = os.environ["DB_DSN"]
    NOTIFY_URL = os.environ.get("NOTIFY_WEBHOOK_URL")      # e.g. an n8n webhook that messages the duty operator
    CONSUMER = os.environ.get("CONSUMER", "escalator-1")
    
    STREAM_OUT, GROUP = "incidents:classified", "escalators"
    
    
    async def main() -> None:
        r = redis.from_url(REDIS_URL, decode_responses=True)
        pool = await asyncpg.create_pool(DB_DSN, min_size=1, max_size=3)
        http = httpx.AsyncClient(timeout=10)
    
        async def handle(msg_id: str, f: dict) -> bool:
            priority = int(f["priority"])
            async with pool.acquire() as con:
                await con.execute(
                    "INSERT INTO incident_audit (event_id, zone_id, label, priority, confidence, reason, model) "
                    "VALUES ($1,$2,$3,$4,$5,$6,$7) ON CONFLICT (event_id) DO NOTHING",
                    f["event_id"], int(f["zone_id"]) if f.get("zone_id") else None,
                    f["label"], priority, float(f["confidence"]), f["reason"], f.get("model"))
            if priority <= 2 and NOTIFY_URL:
                # A suggestion for a person. Nothing is dispatched automatically.
                resp = await http.post(NOTIFY_URL, json={
                    "event_id": f["event_id"], "suggested_priority": priority,
                    "label": f["label"], "confidence": f["confidence"], "reason": f["reason"],
                    "note": "AI suggestion - a human operator decides"})
                resp.raise_for_status()      # if the notification fails, the message stays pending and is retried
            return True
    
        await consume(r, STREAM_OUT, GROUP, CONSUMER, handle)
    
    
    if __name__ == "__main__":
        asyncio.run(main())
  6. Running it with invented events

    Start the three processes and feed in a few invented events:

    bash
    # 1) Redis (6.2 or newer) on 127.0.0.1, your own way
    # 2) the model
    ollama pull llama3.2-vision:11b
    
    # 3) the audit table
    psql "postgresql://<user>:<password>@localhost/incidents" -f schema.sql
    
    # 4) the environment
    export REDIS_URL="redis://127.0.0.1:6379"
    export DB_DSN="postgresql://<user>:<password>@localhost/incidents"
    export FRAME_DIR="/var/lib/incidents/frames"
    export NOTIFY_WEBHOOK_URL="<webhook-address>"      # e.g. an n8n webhook that messages the duty operator
    
    # 5) three terminals
    python3 worker.py
    python3 escalator.py
    python3 publish_test.py

    publish_test.py puts four invented descriptions into incidents:raw. You expect lines for the classified events in the two workers and rows in incident_audit. What exactly the model will decide we do not promise — we have not run it. Also try a failure: stop Ollama and feed an event — it should get P2 and unclassified.

  7. Monitoring and measuring the errors

    Watch the queue and measure. The columns human_label and human_priority in the audit table are filled in by a person who reviews events. Use them to count the most important thing — how often the classifier was less urgent than the human (the misses).

    WhatCommand
    How many wait for processing (the group) and who holds themredis-cli XPENDING incidents:raw classifiers
    How many events are in the streamredis-cli XLEN incidents:raw
    Information about the groupsredis-cli XINFO GROUPS incidents:raw
    sql
    -- how often the label differs from the human's
    SELECT label, human_label, count(*)
    FROM incident_audit
    WHERE human_label IS NOT NULL
    GROUP BY label, human_label
    ORDER BY label, human_label;
    
    -- the important number: how often the classifier was LESS urgent than the human
    SELECT count(*) AS misses
    FROM incident_audit
    WHERE human_priority IS NOT NULL AND priority > human_priority;
  8. Security — the minimum

    • Redis listens only on 127.0.0.1; if you need access from outside, add a password and TLS according to its documentation.
    • The frame folder has permissions only for the service and a deletion period agreed with the lawyer.
    • The notification address is a secret: keep it in an environment variable, not in code.
    • Do not put names and frames in the streams and the audit table.
    • Periodically review what the classifier decided against the human (step 7) and after every model change.

04Check

Quiz

1. What does XACK do?

2. Where does the priority come from in this lesson?

3. What do we do if the model's answer is unreadable?

4. Why do we not put the frame into Redis?

05What's next

06Sources

  1. Redis: Streams — consumer groups, XREADGROUP, XACK, XPENDING, XAUTOCLAIM (from 6.2.0).
  2. Ollama: llama3.2-vision — sizes, context, supported languages 🔒 local.
  3. Ollama: chat API — images, format, options 🔒 local.
  4. PostgreSQL: documentation — ON CONFLICT.