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.
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.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
- What Redis Streams, consumer groups and message acknowledgement are.
- How to pick up forgotten messages if a worker stopped halfway.
- How to give text and a frame to a multimodal model through Ollama and get an answer that follows a schema.
- Why the priority is taken from a table and how to "fall towards more attention" when in doubt.
- How to keep an audit trail and notify a human without anything being sent automatically.
- How to measure how often the classifier is wrong.
02Before you start
- A machine of the NVIDIA GB10 class with DGX OS and Ollama (see lesson 04-01).
- The
llama3.2-vision:11bmodel: about 7.8 GB, 128K context, text and picture (per ollama.com as of 03.10.2026):ollama pull llama3.2-vision:11b. The large90bvariant is 55 GB — we have not tried it. - Redis 6.2 or newer (the
XAUTOCLAIMcommand is from 6.2.0), listening only on127.0.0.1. - PostgreSQL 14 or newer, Python 3.10 or newer and the packages
redis,httpx,asyncpg(not pinned — write the versions down after the first successful run). - A folder for frames that only the service can access.
03Steps
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
XREADGROUPand, when done, saysXACK. Until there is an acknowledgement the entry is "pending" and can be claimed again withXAUTOCLAIM— so nothing is lost if a worker crashes.textincidents: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.
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.
Label Priority How we treat it intrusion · theft · fight P1 straight to a human vandalism · trespass P2 to a human suspicious P3 only recorded; a human sees it in the review false_alarm P4 recorded; never deleted on its own low confidence (below 0.6) at most P3 a human should look an unreadable answer P2 · unclassified more attention, not less 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 answersBUSYGROUP, 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 returnedTrue; an error orFalseleaves it pending.python · common.pyimport 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)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 theimageslist of the message, base64-encoded — that is how Ollama's API takes it. Temperature 0 makes the answers more repeatable.bashollama pull llama3.2-vision:11bpython · worker.pyimport 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 streamincidents:classifiedhas a length limit (maxlen); if a worker is stopped for long, the oldest entries can be trimmed, so watch the length (step 7).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.sqlCREATE 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.pyimport 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())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.pypublish_test.pyputs four invented descriptions intoincidents:raw. You expect lines for the classified events in the two workers and rows inincident_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 andunclassified.Monitoring and measuring the errors
Watch the queue and measure. The columns
human_labelandhuman_priorityin 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).What Command How many wait for processing (the group) and who holds them redis-cli XPENDING incidents:raw classifiersHow many events are in the stream redis-cli XLEN incidents:rawInformation about the groups redis-cli XINFO GROUPS incidents:rawsql-- 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;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.
- Redis listens only on
04Check
- The group
classifiersexists andXPENDINGdoes not grow without limit. - An invented event passes through both streams and reaches
incident_audit. - With Ollama stopped the event gets P2 and
unclassified. - The message to the human is a suggestion, and nothing is sent automatically to authorities.
- Frames are only in the folder, not in Redis or in the table.
- You have at least 30 examples of your own with a human judgement and you count the misses.
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
- Redis: Streams — consumer groups,
XREADGROUP,XACK,XPENDING,XAUTOCLAIM(from 6.2.0). - Ollama: llama3.2-vision — sizes, context, supported languages 🔒 local.
- Ollama: chat API —
images,format,options🔒 local. - PostgreSQL: documentation —
ON CONFLICT.