From c48ed8eb1192b47692bb604f7146292ef9b24ff4 Mon Sep 17 00:00:00 2001 From: thierrypdamiba Date: Sun, 15 Mar 2026 14:09:06 -0700 Subject: [PATCH] address pr review comments --- .../tutorials-build-essentials/_index.md | 6 +- .../video-anomaly-edge-part-1.md | 78 ++++------ .../video-anomaly-edge-part-2.md | 143 +++++++++++++----- .../video-anomaly-edge-part-3.md | 83 ++++++---- 4 files changed, 192 insertions(+), 118 deletions(-) diff --git a/qdrant-landing/content/documentation/tutorials-build-essentials/_index.md b/qdrant-landing/content/documentation/tutorials-build-essentials/_index.md index 7baadfe09..2622eabcd 100644 --- a/qdrant-landing/content/documentation/tutorials-build-essentials/_index.md +++ b/qdrant-landing/content/documentation/tutorials-build-essentials/_index.md @@ -13,9 +13,9 @@ partition: build | [Discord RAG Bot](/documentation/tutorials-build-essentials/agentic-rag-camelai-discord/) | Develop a functional bot with CAMEL-AI. | OpenAI | 45m | Intermediate | | [Agentic RAG with CrewAI](/documentation/tutorials-build-essentials/agentic-rag-crewai-zoom/) | Step-by-step multi-agent RAG system. | CrewAI | 45m | Beginner | | [n8n Workflow Automation](/documentation/tutorials-build-essentials/qdrant-n8n/) | Combine Qdrant with low-code n8n workflows. | n8n | 45m | Intermediate | -| [Video Anomaly Detection Part 1](/documentation/tutorials-build-essentials/video-anomaly-edge-part-1/) | Architecture, Twelve Labs, and NVIDIA VSS integration. | Python | 60m | Advanced | -| [Video Anomaly Detection Part 2](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/) | Two-shard Qdrant Edge architecture and escalation pipeline. | Python | 90m | Advanced | -| [Video Anomaly Detection Part 3](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) | Scoring, baseline governance, and deployment on Vultr. | Python | 60m | Advanced | +| [Video Anomaly Detection Part I](/documentation/tutorials-build-essentials/video-anomaly-edge-part-1/) | Architecture, Twelve Labs, and NVIDIA VSS integration. | Python | 60m | Advanced | +| [Video Anomaly Detection Part II](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/) | Two-shard Qdrant Edge architecture and escalation pipeline. | Python | 90m | Advanced | +| [Video Anomaly Detection Part III](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) | Scoring, baseline governance, and deployment on Vultr. | Python | 60m | Advanced | diff --git a/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-1.md b/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-1.md index f0c981dd7..eec8709d3 100644 --- a/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-1.md +++ b/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-1.md @@ -1,5 +1,5 @@ --- -title: "Video Anomaly Detection Part 1 | Architecture, Twelve Labs, and NVIDIA VSS" +title: "Video Anomaly Detection Part I | Architecture, Twelve Labs, and NVIDIA VSS" weight: 9 partition: build social_preview_image: /articles_data/video-anomaly-edge/preview/social_preview.jpg @@ -10,15 +10,15 @@ aliases: # Video Anomaly Detection: Architecture, Twelve Labs, and NVIDIA VSS -| Time: 60 min | Level: Advanced | Stack: Qdrant Edge, Twelve Labs Marengo 3.0, NVIDIA VSS, Vultr | Output: [GitHub](https://github.com/qdrant/examples/tree/master/video-anomaly-edge) | -| --- | ----------- | ----------- | ----------- | +| Time: 120 min | Level: Advanced | Output: [GitHub](https://github.com/qdrant/video-anomaly-edge) | +| --- | ----------- | ----------- | -*This is Part 1 of a 3-part series on building real-time video anomaly detection from edge to cloud. We'll go from architecture and integrations to a production-grade detection pipeline.* +*This is Part I of a 3-part series on building real-time video anomaly detection from edge to cloud. We'll go from architecture and integrations to a production-grade detection pipeline.* **Series:** -- Part 1 | Architecture, Twelve Labs, and NVIDIA VSS (here) -- [Part 2 | Edge-to-Cloud Pipeline](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/) -- [Part 3 | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) +- Part I | Architecture, Twelve Labs, and NVIDIA VSS (here) +- [Part II | Edge-to-Cloud Pipeline](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/) +- [Part III | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) --- @@ -48,15 +48,17 @@ Specifically, you will build a platform that transforms live surveillance stream +![Tech stack overview: NVIDIA Jetson, Qdrant Edge, Twelve Labs, and Qdrant Cloud connected in an edge-to-cloud pipeline](/articles_data/video-anomaly-edge/tech-stack-overview.png) + --- ## Application Demo Before we begin coding, check out the project repository and live demo to get familiarized with what we'll be building. -**GitHub**: [qdrant/examples/video-anomaly-edge](https://github.com/qdrant/examples/tree/master/video-anomaly-edge) +**GitHub**: [qdrant/video-anomaly-edge](https://github.com/qdrant/video-anomaly-edge) -**Live Demo**: [avenue-demo.vercel.app](https://avenue-demo.vercel.app/) +**Live Demo**: [qdrant-edge-video-anomaly.vercel.app](https://qdrant-edge-video-anomaly.vercel.app/) ![Sentinel dashboard screenshot](/articles_data/video-anomaly-edge/sentinel-screenshot.png) @@ -99,8 +101,8 @@ In this series you will: **1.** Clone the repository into your local environment. ```bash -git clone https://github.com/qdrant/examples.git -cd examples/video-anomaly-edge +git clone https://github.com/qdrant/video-anomaly-edge.git +cd video-anomaly-edge ``` **2.** Install dependencies with uv. @@ -126,7 +128,7 @@ TWELVE_LABS_API_KEY= TWELVE_LABS_API_URL=https://api.twelvelabs.io/v1.3 TWELVE_LABS_MARENGO_INDEX_NAME=anomaly-marengo-search TWELVE_LABS_PEGASUS_INDEX_NAME=anomaly-pegasus-summary -TWELVE_LABS_MARENGO_MODEL=marengo2.7 +TWELVE_LABS_MARENGO_MODEL=marengo3.0 TWELVE_LABS_PEGASUS_MODEL=pegasus1.2 # NVIDIA VSS @@ -142,8 +144,7 @@ ANOMALY_THRESHOLD=0.15 **4.** Clone the NVIDIA VSS framework (with Twelve Labs integration) for reference. ```bash -git clone https://github.com/james-le-twelve-labs/nvidia-vss -git clone https://github.com/nathanchess/twelvelabs-nvidia-vss-sample +git clone https://github.com/qdrant/qdrant-twelvelabs-nvidia-vss ``` **5.** Start the full stack with Docker Compose. @@ -172,7 +173,7 @@ Binary classifiers require labeled examples of every anomaly type you want to de **Concept drift.** What counts as "normal" changes over time. A school hallway looks different during class hours versus recess. kNN baselines can be updated continuously without retraining. -![Why classifiers fail: CLIP single-frame scores 0.23 AUC-ROC while Twelve Labs Marengo temporal embeddings score 0.97, a 4.2x improvement](/articles_data/video-anomaly-edge/why-classifiers-fail.png) +![Why classifiers fail: CLIP single-frame scores 0.23 AUC-ROC while Twelve Labs Marengo temporal embeddings score 0.9696, a 4.2x improvement](/articles_data/video-anomaly-edge/why-classifiers-fail.png) The kNN approach is simple and effective. Embed video clips into a vector space, build a baseline of normal embeddings in Qdrant, and flag clips whose nearest neighbors are far away: @@ -204,8 +205,6 @@ This is why we use Twelve Labs Marengo in the cloud. It's purpose-built for vide The system uses a three-tier architecture designed around a simple principle: **cheap, fast triage at the edge; accurate, rich analysis in the cloud.** -![Tech stack overview: NVIDIA Jetson, Qdrant Edge, Twelve Labs, and Qdrant Cloud connected in an edge-to-cloud pipeline](/articles_data/video-anomaly-edge/tech-stack-overview.png) - **Architecture:** - **Edge Tier (NVIDIA Jetson)** @@ -260,29 +259,7 @@ Twelve Labs provides two key models for our architecture. **Marengo** handles em Let's look at how we integrate them into our backend. -### Twelve Labs Client - -`/backend/twelvelabs_client.py` - -```python -from twelvelabs import TwelveLabs - -TWELVE_LABS_API_KEY = os.getenv("TWELVE_LABS_API_KEY", "") -MARENGO_MODEL = os.getenv("TWELVE_LABS_MARENGO_MODEL", "marengo2.7") -PEGASUS_MODEL = os.getenv("TWELVE_LABS_PEGASUS_MODEL", "pegasus1.2") - -_client: Optional[TwelveLabs] = None - -def get_client() -> TwelveLabs: - global _client - if _client is None: - if not TWELVE_LABS_API_KEY: - raise RuntimeError("TWELVE_LABS_API_KEY not set") - _client = TwelveLabs(api_key=TWELVE_LABS_API_KEY) - return _client -``` - -We use a singleton pattern for the client, initialized once and reused across all requests. The two models need separate indexes in Twelve Labs: +The client is a simple singleton initialized from `TWELVE_LABS_API_KEY` in your `.env`. The two models need separate indexes in Twelve Labs: ```python def _ensure_index(index_name: str, models: list[dict]) -> str: @@ -404,7 +381,7 @@ The true value is in **modularity**. Our architecture uses the Twelve Labs integ ### Video Chunking for VSS -The first step in the VSS pipeline is chunking. Following the pattern from the [Twelve Labs x NVIDIA VSS manufacturing sample](https://github.com/nathanchess/twelvelabs-nvidia-vss-sample), we split videos using FFmpeg's segment muxer: +The first step in the VSS pipeline is chunking. Following the pattern from the [Twelve Labs x NVIDIA VSS manufacturing sample](https://github.com/qdrant/qdrant-twelvelabs-nvidia-vss), we split videos using FFmpeg's segment muxer: `/backend/vss.py` @@ -447,7 +424,7 @@ def chunk_video( return sorted(output_dir.glob(f"{input_path.stem}_chunk_*.mp4")) ``` -**Why chunk at all?** The same cost issue from the [manufacturing automation tutorial](https://www.twelvelabs.io/blog/manufacturing-automation) applies here. Processing 24 hours of raw video is expensive. Our edge tier already filters ~85% of footage, and chunking the remaining escalated clips further optimizes the cloud pipeline. Only chunks of interest flow through the full VSS stack. +**Why chunk at all?** The same cost issue from the [manufacturing sample](https://github.com/qdrant/qdrant-twelvelabs-nvidia-vss) applies here. Processing 24 hours of raw video is expensive. Our edge tier already filters ~85% of footage, and chunking the remaining escalated clips further optimizes the cloud pipeline. Only chunks of interest flow through the full VSS stack. ### Async Upload to VSS @@ -458,10 +435,13 @@ async def upload_to_vss(file_path: str | Path) -> Optional[str]: """Upload a single video file to NVIDIA VSS.""" file_path = Path(file_path) + with open(file_path, "rb") as f: + content = f.read() + timeout = aiohttp.ClientTimeout(total=VSS_UPLOAD_TIMEOUT) async with aiohttp.ClientSession(timeout=timeout) as session: data = aiohttp.FormData() - data.add_field("file", open(file_path, "rb"), + data.add_field("file", content, filename=file_path.name, content_type="video/mp4") data.add_field("purpose", "vision") data.add_field("media_type", "video") @@ -525,6 +505,8 @@ services: - driver: nvidia count: 1 capabilities: [gpu] + networks: + - anomaly-net backend: environment: @@ -540,20 +522,20 @@ services: ## Recap -In Part 1, you set up the project, learned why kNN anomaly detection in Qdrant outperforms traditional classifiers for open-world surveillance, integrated Twelve Labs Marengo and Pegasus for video embeddings and Q&A, and connected NVIDIA VSS for GPU-accelerated ingestion. The architecture is in place. Now we need to build the edge. +In Part I, you set up the project, learned why kNN anomaly detection in Qdrant outperforms traditional classifiers for open-world surveillance, integrated Twelve Labs Marengo and Pegasus for video embeddings and Q&A, and connected NVIDIA VSS for GPU-accelerated ingestion. The architecture is in place. Now we need to build the edge. ## What's Next -In **[Part 2 | Edge-to-Cloud Pipeline](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/)**, we'll implement the two-shard Qdrant Edge architecture, edge triage scoring, escalation flow with ensemble scoring, and offline resilience. +In **[Part II | Edge-to-Cloud Pipeline](/documentation/tutorials-build-essentials/video-anomaly-edge-part-2/)**, we'll implement the two-shard Qdrant Edge architecture, edge triage scoring, escalation flow with ensemble scoring, and offline resilience. -In **[Part 3 | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/)**, we'll cover incident formation, baseline governance, unified retrieval, results on UCF-Crime, and deployment on Vultr Cloud GPUs. +In **[Part III | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/)**, we'll cover incident formation, baseline governance, unified retrieval, results on UCF-Crime, and deployment on Vultr Cloud GPUs. --- Additional Resources: -- **Project Repository**: [qdrant/examples/video-anomaly-edge](https://github.com/qdrant/examples/tree/master/video-anomaly-edge) -- **NVIDIA VSS Twelve Labs Integration**: [james-le-twelve-labs/nvidia-vss](https://github.com/james-le-twelve-labs/nvidia-vss) +- **Project Repository**: [qdrant/video-anomaly-edge](https://github.com/qdrant/video-anomaly-edge) +- **NVIDIA VSS Twelve Labs Integration**: [qdrant/qdrant-twelvelabs-nvidia-vss](https://github.com/qdrant/qdrant-twelvelabs-nvidia-vss) - **Twelve Labs Documentation**: [docs.twelvelabs.io](https://docs.twelvelabs.io/) - **Qdrant Documentation**: [qdrant.tech/documentation](https://qdrant.tech/documentation/) - **Vultr Cloud GPUs**: [vultr.com/products/cloud-gpu](https://www.vultr.com/products/cloud-gpu/) diff --git a/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-2.md b/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-2.md index 2b9aeaeb6..29e077495 100644 --- a/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-2.md +++ b/qdrant-landing/content/documentation/tutorials-build-essentials/video-anomaly-edge-part-2.md @@ -1,5 +1,5 @@ --- -title: "Video Anomaly Detection Part 2 | Edge-to-Cloud Pipeline" +title: "Video Anomaly Detection Part II | Edge-to-Cloud Pipeline" weight: 10 partition: build social_preview_image: /articles_data/video-anomaly-edge/preview/social_preview.jpg @@ -9,15 +9,15 @@ aliases: # Video Anomaly Detection: Edge-to-Cloud Pipeline -| Time: 90 min | Level: Advanced | Stack: Qdrant Edge, Twelve Labs Marengo 3.0, NVIDIA VSS, Vultr | Output: [GitHub](https://github.com/qdrant/examples/tree/master/video-anomaly-edge) | -| --- | ----------- | ----------- | ----------- | +| Time: 90 min | Level: Advanced | Output: [GitHub](https://github.com/qdrant/video-anomaly-edge) | +| --- | ----------- | ----------- | -*This is Part 2 of a 3-part series on building real-time video anomaly detection from edge to cloud.* +*This is Part II of a 3-part series on building real-time video anomaly detection from edge to cloud.* **Series:** -- [Part 1 | Architecture, Twelve Labs, and NVIDIA VSS](/documentation/tutorials-build-essentials/video-anomaly-edge-part-1/) -- Part 2 | Edge-to-Cloud Pipeline (here) -- [Part 3 | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) +- [Part I | Architecture, Twelve Labs, and NVIDIA VSS](/documentation/tutorials-build-essentials/video-anomaly-edge-part-1/) +- Part II | Edge-to-Cloud Pipeline (here) +- [Part III | Scoring, Governance, and Deployment](/documentation/tutorials-build-essentials/video-anomaly-edge-part-3/) --- @@ -64,8 +64,11 @@ from qdrant_edge import ( Distance as EdgeDistance, EdgeConfig, EdgeShard, + FieldCondition, + Filter, Point, Query, + RangeFloat, SearchRequest, UpdateOperation, VectorDataConfig, @@ -118,7 +121,8 @@ def score_local(self, embedding: np.ndarray) -> float: top_k = results[:K_NEIGHBORS] if not top_k: - return 0.0 + # No baseline yet: treat as maximally anomalous, not normal + return 1.0 sims = [r.score for r in top_k] return 1.0 - float(np.mean(sims)) @@ -148,6 +152,8 @@ for i in range(500): This preserves the baseline's coverage while fitting comfortably in edge device memory. +When `EDGE_SNAPSHOT_KMEANS=500` is set, the backend's `/api/snapshots/full` endpoint runs MiniBatchKMeans on the cloud baseline before creating the snapshot, producing a 500-vector representative subset for the edge immutable shard. + --- ## Edge Triage: Why Imperfect Is the Point @@ -186,8 +192,8 @@ When an edge device scores a clip above the escalation threshold, it sends the c 3. **Upload**: Clip and edge metadata sent to cloud 4. **Cloud re-analysis**: Twelve Labs Marengo produces embeddings, kNN score in Qdrant Cloud + semantic signals 5. **VSS enrichment**: VLM captioning, audio transcription, CV pipeline (if enabled) -6. **Ensemble scoring**: 70% cloud score + 30% edge score -7. **Confirmation**: Final score compared against cloud threshold (0.038) +6. **Ensemble scoring**: 70% cloud score + 30% edge score with temporal boost +7. **Confirmation**: Ensemble score compared against threshold (0.038) The escalation handler in our backend supports both Twelve Labs and local model server paths: @@ -236,12 +242,10 @@ async def handle_escalation(request: EscalationRequest) -> EscalationResult: k=CONFIRMATION_K, ) - # Ensemble scoring + # Return cloud score — confirmation happens after ensemble in the endpoint cloud_score = cloud_result.anomaly_score - is_confirmed = cloud_score > ESCALATION_THRESHOLD return EscalationResult( cloud_score=cloud_score, - is_confirmed_anomaly=is_confirmed, ... ) ``` @@ -253,13 +257,18 @@ async def handle_escalation(request: EscalationRequest) -> EscalationResult: The ensemble weighting reflects the accuracy differential between tiers: ```python -ensemble_score = ( - DEFAULT_CLOUD_WEIGHT * cloud_score + # 0.7 - DEFAULT_EDGE_WEIGHT * edge_score # 0.3 +ens = ensemble_scorer.score( + edge_score=edge_score, + cloud_score=cloud_score, + device_id=edge_device_id, + timestamp=time.time(), ) + +# Confirmation uses ensemble score, not raw cloud score +is_confirmed = ens.is_anomaly # ensemble_score > threshold (0.038) ``` -A low cloud score suppresses a high edge score (false positive). A high cloud score confirms a high edge score (true anomaly). This is why the cloud threshold (0.038) is lower than the edge threshold (0.06) because the cloud model is more accurate and can set a tighter decision boundary. +A low cloud score suppresses a high edge score (false positive). A high cloud score confirms a high edge score (true anomaly). The threshold (0.038) is lower than the edge threshold (0.06) because the cloud model is more accurate and can set a tighter decision boundary. ### Temporal Boosting @@ -299,6 +308,9 @@ The immutable shard stays current through snapshot syncing. A full sync download ```python def sync_from_server(self, full: bool = False) -> None: + # Flush pending uploads first; track whether they succeeded + upload_ok = self._upload_batch(pending_items) # returns True on success + if full or not self._immutable_shard: # Full sync: download complete snapshot resp = requests.post(f"{CLOUD_API_URL}/api/snapshots/full", stream=True) @@ -320,18 +332,20 @@ def sync_from_server(self, full: bool = False) -> None: f.write(chunk) self._immutable_shard.update_from_snapshot(str(snapshot_path)) - # Clean synced points from mutable shard - self._mutable_shard.update( - UpdateOperation.delete_points_by_filter( - Filter(must=[FieldCondition( - key="sync_timestamp", - range=RangeFloat(lte=sync_timestamp), - )]) + # Only purge mutable shard if upload succeeded. + # If upload failed, points are not in the cloud yet -- deleting them would cause data loss. + if upload_ok: + self._mutable_shard.update( + UpdateOperation.delete_points_by_filter( + Filter(must=[FieldCondition( + key="sync_timestamp", + range=RangeFloat(lte=sync_timestamp), + )]) + ) ) - ) ``` -After syncing, points that were already uploaded to the cloud are purged from the mutable shard to prevent double-counting during kNN queries. +After syncing, points that were already uploaded to the cloud are purged from the mutable shard to prevent double-counting during kNN queries. The cleanup is gated on `upload_ok` to avoid deleting local data that never reached the cloud. ### Offline Resilience @@ -340,14 +354,23 @@ After syncing, points that were already uploaded to the cloud are purged from th If the cloud is unreachable, escalation data is persisted to disk as JSON files: ```python -async def escalate_to_cloud(self, clip_path, edge_embedding, edge_score): +async def escalate_to_cloud(self, clip_path, edge_embedding, edge_score, timestamp_ms=0, scene_id=""): + payload = {"edge_score": edge_score, "embedding": edge_embedding.tolist()} try: async with httpx.AsyncClient(timeout=30.0) as client: - resp = await client.post( - f"{CLOUD_API_URL}/api/escalate", - data={"metadata": json.dumps(payload)}, - files={"clip": (Path(clip_path).name, f, "video/mp4")}, - ) + data = { + "edge_device_id": EDGE_DEVICE_ID, + "edge_score": str(edge_score), + "edge_embedding": json.dumps(edge_embedding.tolist()), + "timestamp_ms": str(timestamp_ms), + "scene_id": scene_id, + } + with open(clip_path, "rb") as f: + resp = await client.post( + f"{CLOUD_API_URL}/api/escalate", + data=data, + files={"file": (Path(clip_path).name, f, "video/mp4")}, + ) resp.raise_for_status() except Exception: # Persist to offline queue for later flush @@ -445,20 +468,68 @@ def _evict_by_score_priority(self) -> None: --- +## Fleet Management + +A real deployment has dozens of edge devices, not one. The cloud backend tracks each Jetson independently through a device registry. + +Each device registers itself on first boot: + +```python +POST /api/edges/register +{ + "name": "north-entrance-cam", + "location": "Building A - Zone 1", + "device_id": "edge-a1b2c3d4" # optional, auto-generated if omitted +} +``` + +From that point, escalations are tagged with `edge_device_id` so the cloud can track per-device performance: + +```python +POST /api/escalate + edge_device_id = "edge-a1b2c3d4" + edge_score = 0.12 + edge_embedding = [...] + clip =