Cogniverse Ingestion Module Study Guide¶
Package: cogniverse_runtime (Application Layer) Module Location: - libs/runtime/cogniverse_runtime/ingestion/ — pipeline library (Whisper, ColPali, VLM, embeddings) - libs/runtime/cogniverse_runtime/ingestion_worker/ — Redis-queue consumer + submit/status REST APIs
Purpose: Configurable multi-modal content processing pipeline for video, document, audio, and image extraction and indexing.
The runtime exposes two ingest entrypoints with different execution models:
| Endpoint | Path | Pipeline call site |
|---|---|---|
| Synchronous | POST /ingestion/start | routers/ingestion.run_ingestion background task imports ingestion/ directly |
| Async (queued) | POST /ingestion/upload | Runtime resolves and validates the canonical tenant's requested or configured video profile, uploads to MinIO, then enqueues that exact profile on Redis Streams. The cogniverse-ingestor pod runs python -m cogniverse_runtime.ingestion_worker.worker, which consumes jobs and calls into ingestion/ |
Both paths invoke the same VideoIngestionPipeline followed by routers/ingestion._extract_graph_per_segment for KG provenance + face pipeline + back-ref PATCH.
Package Structure¶
libs/runtime/cogniverse_runtime/ingestion/
├── __init__.py # Package initialization
├── pipeline.py # Main VideoIngestionPipeline orchestrator
├── pipeline_builder.py # Pipeline construction utilities
├── strategy_factory.py # Creates strategy sets from config
├── strategy.py # Strategy/StrategyConfig — unified processing/ranking/storage
│ # strategy model used by cogniverse_core.registries.registry
│ # (distinct from the BaseStrategy processors below)
├── strategies.py # 18 BaseStrategy implementations (Frame, Chunk, SingleVector,
│ # Image, AudioFile, Document, DocumentVisual, Code, etc.)
├── processing_strategy_set.py # Strategy container with execution flow
├── processor_manager.py # Manages processor instances
├── processor_base.py # Base classes for processors and strategies
├── exceptions.py # Pipeline-specific exceptions
└── processors/
├── keyframe_processor.py # Histogram-based keyframe extraction
├── chunk_processor.py # FFmpeg-based chunk extraction
├── audio_processor.py # Whisper transcription
├── audio_transcriber.py # Audio transcription core logic
├── audio_embedding_generator.py # Audio embedding generation
├── vlm_processor.py # VLM description generation
├── vlm_descriptor.py # VLM description core logic
├── served_model.py # Served model-id discovery for /v1 endpoints
├── single_vector_processor.py # Sliding window segment processing
└── embedding_generator/
├── embedding_generator.py # Base classes and interfaces
├── embedding_generator_impl.py # Backend-agnostic embedding implementation
├── embedding_generator_factory.py # Factory for creating generators
├── token_pooling.py # Hierarchical token pooling (numpy/scipy), model_config.token_pool_factor
└── backend_factory.py # Backend client creation
ingestion_worker/ — service layer wrapping the pipeline¶
libs/runtime/cogniverse_runtime/ingestion_worker/
├── __init__.py
├── worker.py # Long-lived Redis-Streams consumer; pod entrypoint
├── queue.py # XADD / XREADGROUP / XACK helpers + IngestJob model
├── idempotency.py # SHA-keyed dedup: same file SHA → same ingest_id
├── redis_client.py # aioredis pool factory
├── minio_client.py # boto3 wrapper for the cogniverse-ingest bucket
├── submit_api.py # enqueue_ingestion — queue submission called by the POST /ingestion/upload handler
├── status_api.py # GET /ingestion/{id}/status + /events route handlers
├── backpressure.py # Per-tenant in-flight cap + circuit breaker
└── reaper.py # XAUTOCLAIM crash recovery for orphaned PEL entries
| Module | Role |
|---|---|
worker.py | python -m cogniverse_runtime.ingestion_worker.worker is the cogniverse-ingestor pod's CMD. Long-lived process that joins the configured Redis consumer group and processes one job at a time via _default_processor, which localizes the source, runs the cold builds (config manager + graph-factory install) off-loop via asyncio.to_thread, calls VideoIngestionPipeline.process_video_async from ingestion/, then runs the per-segment KG extraction + face pipeline + back-ref PATCH. At startup it also calls the shared cogniverse_runtime.entrypoint_env.configure_runtime_library_defaults() helper, which mirrors MinIO creds onto the AWS names used by fsspec and configures the S3 cache backend defaults before the first Redis read. After content feed and before graph extraction it persists ingest:graph-pending:<message_id>. A graph exception, partial write, cancellation, or positive finite INGEST_GRAPH_DEADLINE_SECONDS timeout raises the retryable, nonterminal GraphStageIncomplete condition and publishes retrying with that error_type, without clearing inflight state, releasing the tenant slot, or acknowledging the queue entry. When pipeline_cache enables s3 without MinIO settings, startup fails fast with the same cache-backend error used by the pipeline. Before it claims a job it subscribes to the config events channel (config_event_subscriber: CONFIG_EVENT_CHANNEL on REDIS_URL) as worker ingestion:<consumer id>:<pid>:<suffix>, handling CONFIG_EVENT_HANDLERS (configs_changed, backend_profiles_changed), so a config the runtime saves, restores or imports, and a backend profile it creates, updates or deletes, is dropped from every config manager the worker holds before that write answers; startup fails when the channel cannot be subscribed. Terminal cleanup (mark done, clear inflight, clear graph marker, decrement active, ack) runs only after the graph completes or before durable content exists. While a job runs the worker holds its ingestion task's lease on the shared task event store (TaskEventStore.attach, renewed by the worker's poller): a job cancelled while queued (POST /events/ingestion/{job_id}/cancel) settles without running — inflight cleared, graph marker cleared, tenant slot released, entry acked — and publishes {"state": "cancelled", "reason"}; a cancellation that arrives once the job runs lets its one video finish. The error and cleanup_error of the events it publishes carry every absolute path reduced to its file name (_without_local_paths), so a client never sees where the worker localised the upload. |
queue.py | Redis Streams primitives (submit, claim, ack, autoclaim, times_delivered, active counters, status streams) + the IngestJob dataclass that flows through the stream. publish_status_if_absent seeds a status stream only when it does not exist, in one server-side step, so concurrent callers restoring a reclaimed stream write exactly one event between them and a surviving stream keeps its history. |
idempotency.py | Computes a per-file SHA at upload time; subsequent uploads of the same file return the existing ingest_id instead of re-running the pipeline. force=true query-param bypasses this. get_done_ingest_id distinguishes a completed run from an in-flight one for the reaper. |
submit_api.py | No router defined here — enqueue_ingestion is the queue-submission helper called by the POST /ingestion/upload handler in routers/ingestion.py after the route has resolved the canonical tenant's profile and uploaded the file to MinIO. Takes the upload's original filename and writes it on the job's queued status event (and on a re-seeded snapshot). Computes the idempotency SHA, marks the inflight key and increments the tenant's active counter BEFORE the job becomes claimable, enqueues an IngestJob on Redis; once the job is queued it is listed as an active ingestion task in the task event store the caller passes as task_events (the runtime passes its per-process store, on its shared-state client; TaskEventStore.register_queued; a worker that claims it first records it itself, so a refused registration is logged); the router then returns 202 + ingest_id (or blocks for the terminal event when wait=true). The XADD and the committed marker (ingest:submitted:<sha>) commit in one MULTI/EXEC, so a submit failure verifies that marker on a fresh read before compensating: a genuine failure clears the inflight marker and restores the counter, while a reply lost after the transaction committed (the job is durably enqueued) is treated as success — no clear, no decrement — so a resubmit cannot duplicate the job. An idempotency hit returns the existing run's ingest_id carrying that run's state (complete once the done marker is set, otherwise in_flight) and re-seeds its status stream when Redis has already reclaimed it — the done marker outlives the stream by days — so every ingest_id this helper returns is readable through GET /ingestion/{id}/status. |
status_api.py | GET /ingestion/{ingest_id}/status and /events — read the job's status stream; complete, failed and cancelled are terminal (queue.TERMINAL_STATUS_STATES). The same stream serves /events/ingestion/{job_id}. |
backpressure.py | Per-tenant counter for in-flight jobs + automatic 429 when the cap is hit. |
reaper.py | A worker SIGKILLed between claim and ack strands its entry in a dead consumer's PEL — claim() reads only new entries, so the upload would be silently lost and XLEN would inflate until queue-depth backpressure 429s every submit. run_reaper_once XAUTOCLAIMs entries idle past INGEST_REAPER_MIN_IDLE_MS (default 5 min) and re-drives them idempotently: a sha already marked done is settled (ack + clear stale inflight and graph markers, no reprocess, no double decrement); a graph-pending entry is never dead-lettered, whatever its delivery count, and is re-driven once the hold since its last recorded failure (ingest:graph-redrive:<message_id>: re-drive count, time, cause) has elapsed — INGEST_REAPER_MIN_IDLE_MS doubled per re-drive so far, clamped at GRAPH_REDRIVE_HOLD_CAP_MS (6h), each re-drive logged with its number, the last cause and the next hold; anything else re-runs through _process_job like a fresh claim. After recovery the sweep drops consumer names idle past the same threshold that own no pending entry (queue.prune_consumers, one atomic script: snapshot, idle and pending checks, XGROUP DELCONSUMER) — every pod incarnation leaves one behind and Redis keeps them forever; the caller's own name is never dropped and a name that owns an entry is never touched. Live workers heartbeat their claim every INGEST_HEARTBEAT_INTERVAL_SECONDS (default 60s, XCLAIM-to-self with justid so the delivery counter is untouched), resetting the PEL idle clock — the min-idle threshold therefore only ever fires on entries whose owner stopped heartbeating (crashed), never on a live job whose pipeline outlives the threshold. An unmarked job redelivered more than INGEST_REAPER_MAX_DELIVERIES times (default 5) without completing — a pod-killing poison message — is abandoned to the ingest:queue:dead stream with a failed terminal event instead of crash-looping the ingestor; the settle (dead-stream entry, inflight clear, slot decrement) runs as one atomic exactly-once server-side step gated on an ingest:dead:<message_id> marker, so a crash-redelivery re-publishes the terminal and acks but can never double-free a tenant slot or duplicate the dead entry. worker.run() starts reaper_loop (interval INGEST_REAPER_INTERVAL_SECONDS, default 60s, first sweep after one full interval) unless INGEST_REAPER_ENABLED=false. |
Both packages are required: ingestion/ is "do the work" (synchronous, library-shaped); ingestion_worker/ is "consume jobs and call ingestion/" (async, service-shaped). The synchronous /ingestion/start path bypasses the worker and runs the pipeline in a FastAPI background task; the async /ingestion/upload path goes through MinIO + Redis and the worker.
Table of Contents¶
- Module Overview
- Architecture
- Core Components
- Processing Strategies
- Processors
- Data Flow
- Usage Examples
- Production Considerations
- Testing
Module Overview¶
Purpose and Responsibilities¶
The Ingestion Module transforms raw content files into searchable, multi-modal representations through a strategy-based processing pipeline. It orchestrates:
- Video Segmentation - Extract frames, chunks, or sliding windows
- Audio Transcription - Whisper-based speech-to-text
- Visual Description - VLM-based frame descriptions
- Embedding Generation - ColPali, X-CLIP, ColQwen, ColBERT embeddings
- Image Processing - Directory of images presented as keyframes (
ImageSegmentationStrategy) for ColPali embedding - Document Processing - Text extraction and ColBERT semantic embeddings for PDFs/documents
- Document Visual Processing - PDF pages rendered to images (pdf2image) with ColPali multi-vector page embeddings (
document_visual_colpaliprofile →document_visualschema) - Audio Processing - CLAP acoustic + ColBERT semantic dual embeddings for audio content, plus standalone audio-file discovery (
AudioFileSegmentationStrategy) - Code Processing - tree-sitter AST-aware chunking of source files (
CodeSegmentationStrategy) with ColBERT multi-vector embeddings (CodeTextEmbeddingStrategy) - Backend Ingestion - Feed documents to Vespa search
Required transcription failures stop the pipeline before embedding or feed and leave the Redis job failed and eligible for resubmission. A video without an audio stream produces an empty transcript successfully.
Text, source code and audio transcripts are embedded in model-sized windows: the served ColBERT model reports the character spans it encodes whole, and each span becomes its own document carrying chunk_index, chunk_count, chunk_start and chunk_end alongside its slice of the text. All windows of one source share that source's identity field, and the document_text_semantic, lateon_mv, code_lateon_mv and audio_clap_semantic profiles resolve results at source granularity, so a search returns one hit per source with its matched windows. Each span is re-tokenized on its own and shortened until the encoder takes it whole, and a reply whose spans leave a gap, overlap or stop short of the text fails the document.
Each run writes its keyframes — extracted or rehydrated from the cache — chunks, rendered pages, transcripts and their metadata under its own scratch directory beneath the profile output directory, and the pipeline removes that directory when the run ends, whether it completed, failed or was cancelled. A cancelled run lets the in-flight decoding stage settle before releasing the directory, and a directory that survives its release is logged as a warning. The ingestion worker builds its pipeline with retain_job_scratch=True, so the keyframes stay on disk through the graph stage, whose face pipeline reads them; it calls release_retained_scratch() once the graph stage ends, whether it completed or failed.
Key Features¶
- Strategy Pattern: Pluggable processors configured via YAML profiles
- Async Processing: Concurrent video processing with configurable parallelism
- Multi-Modal Support: Video, document, audio, and image embeddings
- Caching: Per-profile artifact caching for keyframes, transcripts, descriptions
- Profile-Based: Different strategies for different embedding models
- Format Conversion: Binary (int8) vs Float (bfloat16) embeddings
Architecture¶
1. Ingestion Pipeline Architecture¶
flowchart TB
Start["<span style='color:#000'>Video Input</span>"] --> Entry["<span style='color:#000'>VideoIngestionPipeline</span>"]
Entry --> ConfigRes["<span style='color:#000'>Configuration Resolution</span>"]
ConfigRes --> LoadProfile["<span style='color:#000'>Load Profile Config</span>"]
LoadProfile --> CreateStrategySet["<span style='color:#000'>Strategy Factory Creates ProcessingStrategySet</span>"]
CreateStrategySet --> InitProc["<span style='color:#000'>Processor Manager Initializes Required Processors</span>"]
InitProc --> AsyncProc["<span style='color:#000'>Async Video Processing</span>"]
AsyncProc --> StrategyExec["<span style='color:#000'>Strategy Orchestration</span>"]
StrategyExec --> Sequential["<span style='color:#000'>Sequential Strategy Execution</span>"]
Sequential --> Segment["<span style='color:#000'>1. Segmentation Strategy</span>"]
Segment --> SegmentType{"<span style='color:#000'>Segmentation Type</span>"}
SegmentType -->|Frame-Based| KFCache{"<span style='color:#000'>Keyframe Cache? shared tier</span>"}
KFCache -->|Hit| RehydrateKF["<span style='color:#000'>Load cached keyframes + rehydrate frame files to disk</span>"]
KFCache -->|Miss| ExtractFrames["<span style='color:#000'>Extract Keyframes + store in cache</span>"]
SegmentType -->|Chunk-Based| ExtractChunks["<span style='color:#000'>Extract Video Chunks</span>"]
SegmentType -->|Single-Vector| ExtractSegments["<span style='color:#000'>Extract Sliding Window Segments</span>"]
RehydrateKF --> Transcribe["<span style='color:#000'>2. Transcription Strategy</span>"]
ExtractFrames --> Transcribe
ExtractChunks --> Transcribe
ExtractSegments --> Transcribe
Transcribe --> TransCache{"<span style='color:#000'>Transcript Cache? shared tier</span>"}
TransCache -->|Hit| LoadCachedTrans["<span style='color:#000'>Load cached transcript</span>"]
TransCache -->|Miss| TranscribeAudio["<span style='color:#000'>Whisper Audio Transcription + store in cache</span>"]
LoadCachedTrans --> Describe["<span style='color:#000'>3. Description Strategy</span>"]
TranscribeAudio --> Describe
Describe --> DescType{"<span style='color:#000'>Description Type</span>"}
DescType -->|VLM Description| DescCache{"<span style='color:#000'>Descriptions Cache? shared tier</span>"}
DescCache -->|Hit| LoadCachedDesc["<span style='color:#000'>Load cached descriptions</span>"]
DescCache -->|Miss| GenerateDesc["<span style='color:#000'>Generate VLM Descriptions + store in cache</span>"]
DescType -->|No Description| SkipDesc["<span style='color:#000'>Skip Descriptions</span>"]
GenerateDesc --> Embed["<span style='color:#000'>4. Embedding Strategy</span>"]
LoadCachedDesc --> Embed
SkipDesc --> Embed
Embed --> EmbedGen["<span style='color:#000'>Embedding Generation & Backend Ingestion</span>"]
EmbedGen --> ModelInference["<span style='color:#000'>Model Inference</span>"]
ModelInference --> ModelType{"<span style='color:#000'>Embedding Model</span>"}
ModelType -->|ColPali| ColPaliInfer["<span style='color:#000'>ColPali Multi-Vector</span>"]
ModelType -->|X-CLIP| X-CLIPInfer["<span style='color:#000'>X-CLIP Single-Vector</span>"]
ModelType -->|ColQwen| ColQwenInfer["<span style='color:#000'>ColQwen Multi-Vector</span>"]
ColPaliInfer --> DocBuild["<span style='color:#000'>Document Building</span>"]
X-CLIPInfer --> DocBuild
ColQwenInfer --> DocBuild
DocBuild --> StrategyAware["<span style='color:#000'>Strategy-Aware Document Construction</span>"]
StrategyAware --> FormatConvert["<span style='color:#000'>Format Conversion Binary/Float</span>"]
FormatConvert --> BackendFeed["<span style='color:#000'>Backend Feeding</span>"]
BackendFeed --> FeedType{"<span style='color:#000'>Feed Type</span>"}
FeedType -->|Per-Document| SingleFeed["<span style='color:#000'>Single Document Upload</span>"]
FeedType -->|Batch| BatchFeed["<span style='color:#000'>Batch Upload feed_iterable</span>"]
SingleFeed --> Results["<span style='color:#000'>Pipeline Results</span>"]
BatchFeed --> Results
Results --> Metadata["<span style='color:#000'>Video Metadata ID, duration, path</span>"]
Results --> ProcessingRes["<span style='color:#000'>Processing Results keyframes, chunks, segments, transcript</span>"]
Results --> EmbedStats["<span style='color:#000'>Embedding Stats documents processed, fed</span>"]
Results --> Timing["<span style='color:#000'>Timing Metrics per-stage timing</span>"]
style Start fill:#90caf9,stroke:#1565c0,color:#000
style Entry fill:#ffcc80,stroke:#ef6c00,color:#000
style ConfigRes fill:#ffcc80,stroke:#ef6c00,color:#000
style LoadProfile fill:#b0bec5,stroke:#546e7a,color:#000
style CreateStrategySet fill:#ffcc80,stroke:#ef6c00,color:#000
style InitProc fill:#ffcc80,stroke:#ef6c00,color:#000
style AsyncProc fill:#ffcc80,stroke:#ef6c00,color:#000
style KFCache fill:#b0bec5,stroke:#546e7a,color:#000
style RehydrateKF fill:#a5d6a7,stroke:#388e3c,color:#000
style TransCache fill:#b0bec5,stroke:#546e7a,color:#000
style LoadCachedTrans fill:#a5d6a7,stroke:#388e3c,color:#000
style DescCache fill:#b0bec5,stroke:#546e7a,color:#000
style LoadCachedDesc fill:#a5d6a7,stroke:#388e3c,color:#000
style StrategyExec fill:#ffcc80,stroke:#ef6c00,color:#000
style Sequential fill:#ffcc80,stroke:#ef6c00,color:#000
style Segment fill:#ce93d8,stroke:#7b1fa2,color:#000
style SegmentType fill:#b0bec5,stroke:#546e7a,color:#000
style ExtractFrames fill:#ce93d8,stroke:#7b1fa2,color:#000
style ExtractChunks fill:#ce93d8,stroke:#7b1fa2,color:#000
style ExtractSegments fill:#ce93d8,stroke:#7b1fa2,color:#000
style Transcribe fill:#ce93d8,stroke:#7b1fa2,color:#000
style TranscribeAudio fill:#ce93d8,stroke:#7b1fa2,color:#000
style Describe fill:#ce93d8,stroke:#7b1fa2,color:#000
style DescType fill:#b0bec5,stroke:#546e7a,color:#000
style GenerateDesc fill:#ce93d8,stroke:#7b1fa2,color:#000
style SkipDesc fill:#b0bec5,stroke:#546e7a,color:#000
style Embed fill:#ce93d8,stroke:#7b1fa2,color:#000
style EmbedGen fill:#ffcc80,stroke:#ef6c00,color:#000
style ModelInference fill:#ffcc80,stroke:#ef6c00,color:#000
style ModelType fill:#b0bec5,stroke:#546e7a,color:#000
style ColPaliInfer fill:#ce93d8,stroke:#7b1fa2,color:#000
style X-CLIPInfer fill:#ce93d8,stroke:#7b1fa2,color:#000
style ColQwenInfer fill:#ce93d8,stroke:#7b1fa2,color:#000
style DocBuild fill:#ffcc80,stroke:#ef6c00,color:#000
style StrategyAware fill:#ffcc80,stroke:#ef6c00,color:#000
style FormatConvert fill:#ffcc80,stroke:#ef6c00,color:#000
style BackendFeed fill:#ffcc80,stroke:#ef6c00,color:#000
style FeedType fill:#b0bec5,stroke:#546e7a,color:#000
style SingleFeed fill:#a5d6a7,stroke:#388e3c,color:#000
style BatchFeed fill:#a5d6a7,stroke:#388e3c,color:#000
style Results fill:#a5d6a7,stroke:#388e3c,color:#000
style Metadata fill:#b0bec5,stroke:#546e7a,color:#000
style ProcessingRes fill:#b0bec5,stroke:#546e7a,color:#000
style EmbedStats fill:#b0bec5,stroke:#546e7a,color:#000
style Timing fill:#b0bec5,stroke:#546e7a,color:#000 2. Strategy Resolution Flow¶
sequenceDiagram
participant Pipeline as VideoIngestionPipeline
participant Factory as Strategy Factory
participant Config as Profile Config
participant Registry as Strategy Registry
participant StrategySet as ProcessingStrategySet
participant ProcMgr as Processor Manager
Pipeline->>Config: Load profile config
activate Config
Config-->>Pipeline: profile_config{strategies: {...}}
deactivate Config
Pipeline->>Factory: create_from_profile_config(profile_config)
activate Factory
Note over Factory: Parse strategies section
Factory->>Registry: Import strategy classes
activate Registry
loop For each strategy type
Factory->>Config: Get strategy config
Config-->>Factory: {class: "FrameSegmentationStrategy", params: {...}}
Factory->>Registry: Import class dynamically
Note over Registry: importlib.import_module<br/>("cogniverse_runtime.ingestion.strategies")
Registry-->>Factory: StrategyClass
Factory->>Factory: Instantiate strategy
Note over Factory: strategy = StrategyClass(**params)
Factory->>Factory: Add to strategy set
end
deactivate Registry
Factory->>StrategySet: Create ProcessingStrategySet
activate StrategySet
StrategySet->>StrategySet: Store strategies
Note over StrategySet: segmentation: FrameSegmentationStrategy<br/>transcription: AudioTranscriptionStrategy<br/>description: VLMDescriptionStrategy<br/>embedding: MultiVectorEmbeddingStrategy
deactivate StrategySet
Factory-->>Pipeline: strategy_set
deactivate Factory
Pipeline->>ProcMgr: initialize_from_strategies(strategy_set)
activate ProcMgr
loop For each strategy
ProcMgr->>StrategySet: strategy.get_required_processors()
StrategySet-->>ProcMgr: {processor_name: params}
ProcMgr->>ProcMgr: Create processor instance
Note over ProcMgr: processor = ProcessorClass(**params)
ProcMgr->>ProcMgr: Cache processor
end
ProcMgr-->>Pipeline: Processors initialized
deactivate ProcMgr
Note over Pipeline: Pipeline ready for video processing
Pipeline->>StrategySet: process(video_path, processor_manager, pipeline_context)
activate StrategySet
StrategySet->>StrategySet: Execute segmentation strategy
StrategySet->>StrategySet: Execute transcription strategy
StrategySet->>StrategySet: Execute description strategy
StrategySet->>StrategySet: Execute embedding strategy
StrategySet-->>Pipeline: combined_results
deactivate StrategySet Strategy Types Available (all 18 BaseStrategy subclasses in strategies.py):
flowchart LR
subgraph Segmentation["<span style='color:#000'>Segmentation Strategies</span>"]
Frame["<span style='color:#000'>FrameSegmentationStrategy<br/>Keyframe extraction</span>"]
Chunk["<span style='color:#000'>ChunkSegmentationStrategy<br/>Video chunks</span>"]
Single["<span style='color:#000'>SingleVectorSegmentationStrategy<br/>Sliding windows</span>"]
Image["<span style='color:#000'>ImageSegmentationStrategy<br/>Image directory as keyframes</span>"]
AudioFile["<span style='color:#000'>AudioFileSegmentationStrategy<br/>Discover audio files</span>"]
Doc["<span style='color:#000'>DocumentSegmentationStrategy<br/>Discover document files</span>"]
DocVis["<span style='color:#000'>DocumentVisualSegmentationStrategy<br/>PDF pages → images</span>"]
Code["<span style='color:#000'>CodeSegmentationStrategy<br/>tree-sitter AST chunks</span>"]
end
subgraph Transcription["<span style='color:#000'>Transcription Strategies</span>"]
Audio["<span style='color:#000'>AudioTranscriptionStrategy<br/>Whisper transcription</span>"]
NoTrans["<span style='color:#000'>NoTranscriptionStrategy<br/>Skip transcription</span>"]
end
subgraph Description["<span style='color:#000'>Description Strategies</span>"]
VLM["<span style='color:#000'>VLMDescriptionStrategy<br/>VLM frame descriptions</span>"]
NoDesc["<span style='color:#000'>NoDescriptionStrategy<br/>Skip descriptions</span>"]
end
subgraph Embedding["<span style='color:#000'>Embedding Strategies</span>"]
MultiVec["<span style='color:#000'>MultiVectorEmbeddingStrategy<br/>ColPali, ColQwen</span>"]
SingleVec["<span style='color:#000'>SingleVectorEmbeddingStrategy<br/>X-CLIP</span>"]
AudioEmb["<span style='color:#000'>AudioEmbeddingStrategy<br/>CLAP + ColBERT dual</span>"]
DocText["<span style='color:#000'>DocumentTextEmbeddingStrategy<br/>ColBERT for text</span>"]
DocVisEmb["<span style='color:#000'>DocumentVisualEmbeddingStrategy<br/>ColPali for PDF pages</span>"]
CodeText["<span style='color:#000'>CodeTextEmbeddingStrategy<br/>ColBERT for source code</span>"]
end
style Segmentation fill:#90caf9,stroke:#1565c0,color:#000
style Frame fill:#64b5f6,stroke:#1565c0,color:#000
style Chunk fill:#64b5f6,stroke:#1565c0,color:#000
style Single fill:#64b5f6,stroke:#1565c0,color:#000
style Image fill:#64b5f6,stroke:#1565c0,color:#000
style AudioFile fill:#64b5f6,stroke:#1565c0,color:#000
style Doc fill:#64b5f6,stroke:#1565c0,color:#000
style DocVis fill:#64b5f6,stroke:#1565c0,color:#000
style Code fill:#64b5f6,stroke:#1565c0,color:#000
style Transcription fill:#ffcc80,stroke:#ef6c00,color:#000
style Audio fill:#ffb74d,stroke:#ef6c00,color:#000
style NoTrans fill:#ffb74d,stroke:#ef6c00,color:#000
style Description fill:#ce93d8,stroke:#7b1fa2,color:#000
style VLM fill:#ba68c8,stroke:#7b1fa2,color:#000
style NoDesc fill:#ba68c8,stroke:#7b1fa2,color:#000
style Embedding fill:#a5d6a7,stroke:#388e3c,color:#000
style MultiVec fill:#81c784,stroke:#388e3c,color:#000
style SingleVec fill:#81c784,stroke:#388e3c,color:#000
style AudioEmb fill:#81c784,stroke:#388e3c,color:#000
style DocText fill:#81c784,stroke:#388e3c,color:#000
style DocVisEmb fill:#81c784,stroke:#388e3c,color:#000
style CodeText fill:#81c784,stroke:#388e3c,color:#000 3. Embedding Generation Flow¶
sequenceDiagram
participant Strategy as Embedding Strategy
participant Pipeline as VideoIngestionPipeline
participant Generator as EmbeddingGeneratorImpl
participant Model as Embedding Model / RemoteInferenceClient
participant Backend as IngestionBackend (Vespa)
Strategy->>Pipeline: generate_embeddings_with_processor(results, pipeline_context)
Pipeline->>Generator: generate_embeddings(video_data, output_dir)
activate Generator
Generator->>Generator: _extract_segments(video_data)
Note over Generator: checks document_pages/document_files/code_files/<br/>audio_files keys, then segments/keyframes/frames/<br/>chunks/video_chunks/single_vector_processing
alt "document_pages" in video_data
Generator->>Generator: _process_document_visual_segments()
else "document_files" in video_data
Generator->>Generator: _process_document_segments()
else "code_files" in video_data
Generator->>Generator: _process_code_segments()
else "audio_files" in video_data
Generator->>Generator: _process_audio_segments()
else storage_mode == "single_doc"
Generator->>Generator: _process_single_document()
else storage_mode == "multi_doc" (default)
Generator->>Generator: _process_multi_documents()
end
Generator->>Generator: _iter_segment_embeddings(segments, video_path, video_data)
activate Generator
loop for each segment (or remote frame batch)
alt frame segment + RemoteInferenceClient processor
Generator->>Model: _generate_frame_embeddings_batch(frame_paths)
Note over Model: batches up to _REMOTE_FRAME_BATCH frames per<br/>request; falls back to one-at-a-time on batch failure
else local model / non-frame segment
Generator->>Model: _generate_segment_embeddings(segment, video_path, video_data)
end
Model-->>Generator: embeddings (np.ndarray) or Exception
end
deactivate Generator
Generator->>Generator: _create_segment_document(video_id, segment, embeddings, transcript, description, source_url) -> Document
Generator->>Generator: batch.append(doc)
opt batch reaches _FEED_BATCH_SIZE (50) or loop ends
Generator->>Backend: _feed_documents(batch, errors) -> backend_client.ingest_documents(documents, schema_name)
Backend-->>Generator: success_count + failed_documents (rejections recorded into errors)
end
Generator-->>Pipeline: EmbeddingResult(video_id, total_documents,<br/>documents_processed, documents_fed,<br/>processing_time, errors, metadata)
deactivate Generator Embedding Processing Types (selected by which key video_data carries, checked in generate_embeddings()):
- Frame/Chunk/Window (
_process_multi_documents/_process_single_document, default path): ColPali multi-vector per frame, ColQwen multi-vector per chunk, or X-CLIP single-vector per segment/global — dispatch is bystorage_mode(multi_docvssingle_doc), not aprocessing_typefield - Document ColBERT (
document_filespresent →_process_document_segments): ColBERT 128-dim per-token multi-vector for text documents - Document Visual ColPali (
document_pagespresent →_process_document_visual_segments): ColPali (Tomoro ColQwen3) 320-dim per-patch multi-vector for PDF pages rendered to images (DocumentVisualSegmentationStrategy→DocumentVisualEmbeddingStrategy) - Code ColBERT (
code_filespresent →_process_code_segments): LateOn-Code-edge 48-dim per-token multi-vector for source-code chunks (CodeSegmentationStrategy→CodeTextEmbeddingStrategy) - Audio Dual (
audio_filespresent →_process_audio_segments): CLAP 512-dim acoustic single-vector + ColBERT 128-dim semantic multi-vector for audio content. The transcription result also suppliesaudio_languageandaudio_duration, which land in the schema fields of the same names. A profile that bindsinference_services.acoustic_embeddingfails the document when that service errors; without a bound service the acoustic vector is absent and the semantic document still feeds.
4. Vespa Upload Pipeline¶
flowchart TB
Start["<span style='color:#000'>Documents Ready for Upload</span>"] --> Builder["<span style='color:#000'>Document Builder</span>"]
Builder --> DocType{"<span style='color:#000'>Document Type</span>"}
DocType -->|Frame Documents| FrameDoc["<span style='color:#000'>Build Frame Documents</span>"]
DocType -->|Chunk Documents| ChunkDoc["<span style='color:#000'>Build Chunk Documents</span>"]
DocType -->|Video Documents| VideoDoc["<span style='color:#000'>Build Video Documents</span>"]
DocType -->|Segment Documents| SegmentDoc["<span style='color:#000'>Build Segment Documents</span>"]
FrameDoc --> FrameFields["<span style='color:#000'>Frame Document Fields</span>"]
FrameFields --> FVideoID["<span style='color:#000'>video_id: str</span>"]
FrameFields --> FFrameNum["<span style='color:#000'>frame_number: int</span>"]
FrameFields --> FTimestamp["<span style='color:#000'>timestamp: float</span>"]
FrameFields --> FPath["<span style='color:#000'>frame_path: str</span>"]
FrameFields --> FEmbed["<span style='color:#000'>embeddings: binary/float</span>"]
ChunkDoc --> ChunkFields["<span style='color:#000'>Chunk Document Fields</span>"]
ChunkFields --> CVideoID["<span style='color:#000'>video_id: str</span>"]
ChunkFields --> CChunkNum["<span style='color:#000'>chunk_number: int</span>"]
ChunkFields --> CStart["<span style='color:#000'>start_time: float</span>"]
ChunkFields --> CEnd["<span style='color:#000'>end_time: float</span>"]
ChunkFields --> CPath["<span style='color:#000'>chunk_path: str</span>"]
ChunkFields --> CEmbed["<span style='color:#000'>embeddings: binary/float</span>"]
VideoDoc --> VideoFields["<span style='color:#000'>Video Document Fields</span>"]
VideoFields --> VVideoID["<span style='color:#000'>video_id: str</span>"]
VideoFields --> VDuration["<span style='color:#000'>duration: float</span>"]
VideoFields --> VPath["<span style='color:#000'>video_path: str</span>"]
VideoFields --> VEmbed["<span style='color:#000'>embedding: float array</span>"]
SegmentDoc --> SegmentFields["<span style='color:#000'>Segment Document Fields</span>"]
SegmentFields --> SVideoID["<span style='color:#000'>video_id: str</span>"]
SegmentFields --> SSegmentID["<span style='color:#000'>segment_id: int</span>"]
SegmentFields --> SStart["<span style='color:#000'>start_time: float</span>"]
SegmentFields --> SEnd["<span style='color:#000'>end_time: float</span>"]
SegmentFields --> SText["<span style='color:#000'>text: str transcript</span>"]
SegmentFields --> SEmbed["<span style='color:#000'>embedding: float array</span>"]
FVideoID --> FormatConv["<span style='color:#000'>Format Conversion</span>"]
FFrameNum --> FormatConv
FTimestamp --> FormatConv
FPath --> FormatConv
FEmbed --> FormatConv
CVideoID --> FormatConv
CChunkNum --> FormatConv
CStart --> FormatConv
CEnd --> FormatConv
CPath --> FormatConv
CEmbed --> FormatConv
VVideoID --> FormatConv
VDuration --> FormatConv
VPath --> FormatConv
VEmbed --> FormatConv
SVideoID --> FormatConv
SSegmentID --> FormatConv
SStart --> FormatConv
SEnd --> FormatConv
SText --> FormatConv
SEmbed --> FormatConv
FormatConv --> ConvType{"<span style='color:#000'>Embedding Format</span>"}
ConvType -->|Binary| BinaryConv["<span style='color:#000'>Binary Conversion int8</span>"]
ConvType -->|Float| FloatConv["<span style='color:#000'>Float Conversion bfloat16</span>"]
BinaryConv --> HexEncode["<span style='color:#000'>Hex Encoding for Binary</span>"]
FloatConv --> FloatArray["<span style='color:#000'>Float Array for Float</span>"]
HexEncode --> Validate["<span style='color:#000'>Validate Documents</span>"]
FloatArray --> Validate
Validate --> CheckSchema["<span style='color:#000'>Check Vespa Schema Match</span>"]
CheckSchema --> CheckDims{"<span style='color:#000'>Embedding Dimensions Match?</span>"}
CheckDims -->|Yes| BatchPrep["<span style='color:#000'>Batch Preparation</span>"]
CheckDims -->|No| Error["<span style='color:#000'>Throw Dimension Mismatch Error</span>"]
BatchPrep --> BatchSize["<span style='color:#000'>Determine Batch Size</span>"]
BatchSize --> CreateBatches["<span style='color:#000'>Create Document Batches batch_size=50</span>"]
CreateBatches --> Upload["<span style='color:#000'>Bulk Upload to Vespa</span>"]
Upload --> VespaClient["<span style='color:#000'>Vespa PyClient</span>"]
VespaClient --> FeedIterable["<span style='color:#000'>feed_iterable documents, batch_size</span>"]
FeedIterable --> UploadLoop["<span style='color:#000'>Upload Loop</span>"]
UploadLoop --> Batch1{"<span style='color:#000'>Batch 1</span>"}
Batch1 -->|POST| VespaAPI1["<span style='color:#000'>Vespa HTTP API</span>"]
VespaAPI1 --> Result1{"<span style='color:#000'>Success?</span>"}
Result1 -->|Yes| TrackSuccess1["<span style='color:#000'>Track Success Count</span>"]
Result1 -->|No| TrackError1["<span style='color:#000'>Track Error Details</span>"]
TrackSuccess1 --> Batch2{"<span style='color:#000'>Batch 2</span>"}
TrackError1 --> Batch2
Batch2 -->|POST| VespaAPI2["<span style='color:#000'>Vespa HTTP API</span>"]
VespaAPI2 --> Result2{"<span style='color:#000'>Success?</span>"}
Result2 -->|Yes| TrackSuccess2["<span style='color:#000'>Track Success Count</span>"]
Result2 -->|No| TrackError2["<span style='color:#000'>Track Error Details</span>"]
TrackSuccess2 --> MoreBatches{"<span style='color:#000'>More Batches?</span>"}
TrackError2 --> MoreBatches
MoreBatches -->|Yes| UploadLoop
MoreBatches -->|No| Verify["<span style='color:#000'>Verify Upload</span>"]
Verify --> CheckCounts["<span style='color:#000'>Check Counts and Errors</span>"]
CheckCounts --> CountMatch{"<span style='color:#000'>All documents fed with no errors?</span>"}
CountMatch -->|Yes| Success["<span style='color:#000'>Upload Success</span>"]
CountMatch -->|No| Failed["<span style='color:#000'>Upload Failed with Exact Counts</span>"]
Success --> Complete["<span style='color:#000'>Upload Complete</span>"]
style Start fill:#90caf9,stroke:#1565c0,color:#000
style Builder fill:#ffcc80,stroke:#ef6c00,color:#000
style DocType fill:#b0bec5,stroke:#546e7a,color:#000
style FrameDoc fill:#ce93d8,stroke:#7b1fa2,color:#000
style ChunkDoc fill:#ce93d8,stroke:#7b1fa2,color:#000
style VideoDoc fill:#ce93d8,stroke:#7b1fa2,color:#000
style SegmentDoc fill:#ce93d8,stroke:#7b1fa2,color:#000
style FrameFields fill:#b0bec5,stroke:#546e7a,color:#000
style ChunkFields fill:#b0bec5,stroke:#546e7a,color:#000
style VideoFields fill:#b0bec5,stroke:#546e7a,color:#000
style SegmentFields fill:#b0bec5,stroke:#546e7a,color:#000
style FVideoID fill:#b0bec5,stroke:#546e7a,color:#000
style FFrameNum fill:#b0bec5,stroke:#546e7a,color:#000
style FTimestamp fill:#b0bec5,stroke:#546e7a,color:#000
style FPath fill:#b0bec5,stroke:#546e7a,color:#000
style FEmbed fill:#b0bec5,stroke:#546e7a,color:#000
style CVideoID fill:#b0bec5,stroke:#546e7a,color:#000
style CChunkNum fill:#b0bec5,stroke:#546e7a,color:#000
style CStart fill:#b0bec5,stroke:#546e7a,color:#000
style CEnd fill:#b0bec5,stroke:#546e7a,color:#000
style CPath fill:#b0bec5,stroke:#546e7a,color:#000
style CEmbed fill:#b0bec5,stroke:#546e7a,color:#000
style VVideoID fill:#b0bec5,stroke:#546e7a,color:#000
style VDuration fill:#b0bec5,stroke:#546e7a,color:#000
style VPath fill:#b0bec5,stroke:#546e7a,color:#000
style VEmbed fill:#b0bec5,stroke:#546e7a,color:#000
style SVideoID fill:#b0bec5,stroke:#546e7a,color:#000
style SSegmentID fill:#b0bec5,stroke:#546e7a,color:#000
style SStart fill:#b0bec5,stroke:#546e7a,color:#000
style SEnd fill:#b0bec5,stroke:#546e7a,color:#000
style SText fill:#b0bec5,stroke:#546e7a,color:#000
style SEmbed fill:#b0bec5,stroke:#546e7a,color:#000
style FormatConv fill:#ffcc80,stroke:#ef6c00,color:#000
style ConvType fill:#b0bec5,stroke:#546e7a,color:#000
style BinaryConv fill:#ffcc80,stroke:#ef6c00,color:#000
style FloatConv fill:#ffcc80,stroke:#ef6c00,color:#000
style HexEncode fill:#ffcc80,stroke:#ef6c00,color:#000
style FloatArray fill:#ffcc80,stroke:#ef6c00,color:#000
style Validate fill:#ffcc80,stroke:#ef6c00,color:#000
style CheckSchema fill:#ffcc80,stroke:#ef6c00,color:#000
style CheckDims fill:#b0bec5,stroke:#546e7a,color:#000
style BatchPrep fill:#ffcc80,stroke:#ef6c00,color:#000
style BatchSize fill:#b0bec5,stroke:#546e7a,color:#000
style CreateBatches fill:#ffcc80,stroke:#ef6c00,color:#000
style Upload fill:#ce93d8,stroke:#7b1fa2,color:#000
style VespaClient fill:#ce93d8,stroke:#7b1fa2,color:#000
style FeedIterable fill:#ce93d8,stroke:#7b1fa2,color:#000
style UploadLoop fill:#ffcc80,stroke:#ef6c00,color:#000
style Batch1 fill:#b0bec5,stroke:#546e7a,color:#000
style VespaAPI1 fill:#ce93d8,stroke:#7b1fa2,color:#000
style Result1 fill:#b0bec5,stroke:#546e7a,color:#000
style TrackSuccess1 fill:#a5d6a7,stroke:#388e3c,color:#000
style TrackError1 fill:#e53935,stroke:#c62828,color:#000
style Batch2 fill:#b0bec5,stroke:#546e7a,color:#000
style VespaAPI2 fill:#ce93d8,stroke:#7b1fa2,color:#000
style Result2 fill:#b0bec5,stroke:#546e7a,color:#000
style TrackSuccess2 fill:#a5d6a7,stroke:#388e3c,color:#000
style TrackError2 fill:#e53935,stroke:#c62828,color:#000
style MoreBatches fill:#b0bec5,stroke:#546e7a,color:#000
style Verify fill:#ffcc80,stroke:#ef6c00,color:#000
style CheckCounts fill:#ffcc80,stroke:#ef6c00,color:#000
style CountMatch fill:#b0bec5,stroke:#546e7a,color:#000
style Success fill:#a5d6a7,stroke:#388e3c,color:#000
style PartialSuccess fill:#ffcc80,stroke:#ef6c00,color:#000
style Complete fill:#a5d6a7,stroke:#388e3c,color:#000
style Error fill:#e53935,stroke:#c62828,color:#000 Vespa Upload Key Features:
- Batch Upload: feed_iterable for efficient bulk ingestion
- Format Conversion: Binary (hex-encoded int8) vs Float (bfloat16)
- Schema Validation: Dimension checking before upload
- Error Tracking: Per-batch success/failure monitoring
- Verification: Document count validation post-upload
5. Knowledge Graph Extraction (post-pipeline)¶
After VideoIngestionPipeline finishes feeding content to Vespa, the ingestion router (libs/runtime/cogniverse_runtime/routers/ingestion.py) runs a per-segment KG extraction pass against every text-emitting result it produced. Two functions drive it:
_iter_segments_for_graph(processing_results, source_doc_id) -> Iterator[SegmentRecord]— yields one record per Whisper transcript segment, VLM keyframe description, OCR/caption block, and document file. EachSegmentRecordcarriestextplus asegment_anchor: Mentionwith the timestamps the upstream processor produced._extract_graph_per_segment(processing_results, source_doc_id, tenant_id, config_manager) -> dict— extracts in two passes so the per-segment GLiNER and DSPy claim calls run concurrently (bounded by_KG_EXTRACT_CONCURRENCY) instead of one segment at a time. Pass 1 callsDocExtractor.extract_entities_from_text(...)for every segment (no cross-segment dependency); the coreference prior pool is then reconstructed in segment order, so segment N sees the entity names from segments 0..N-1 exactly as the serial path did; pass 2 callsDocExtractor.extract_claims_from_text(..., segment_entities=<pass-1 result>, prior_entities=<0..N-1 pool>)only for transcript, document, and code segments, while VLM and OCR segments still contribute entity anchors and back-refs. It accumulatesExtractionResult(anchoredNode.mentions: List[Mention]+ SPOEdgerows fromClaimExtractor) in segment order, runsCrossModalLinker.link()to addsame_asedges across modalities,GraphManager.upsert()s the merged result, then PATCHes per-segment back-refs (entity_ids/relation_ids/claim_idsarrays) onto the corresponding content documents.
Key product knobs:
| Symbol | Location | Purpose |
|---|---|---|
_GATE_PROMPT_SCAFFOLDING_CHARS = 2400 | orchestrator_agent.py | Folded into _evidence_token_estimate so the iterative-retrieval-loop token budget guard accounts for DSPy ChatAdapter wrapping (otherwise the estimate undercounts by ~30×). |
_ITER_GATE_RLM_PROMOTION_CHARS = 6000 | orchestrator_agent.py | Sufficient-context gate promotes from CoT to InstrumentedRLM above this evidence size. |
RLM_PROMOTION_TOKENS = 3000 | graph/claim_extractor.py | Claim extraction promotes to RLM above this segment-text size. |
PREDICATE_VOCABULARY (16 items) | graph/claim_extractor.py | Locked SPO predicate set; edges with relations outside the set are dropped after normalization. |
RLM_TRANSCRIPT_TURNS = 4 | graph/claim_extractor.py | REPL turns the promoted module may take, and the number of tool outputs the transcript budget is divided across. |
RLM_HARNESS_TOKENS = 1400 | graph/claim_extractor.py | Tokens the RLM harness occupies before any transcript (measured at 1258 for ClaimExtractionSignature). |
RLM_TRANSCRIPT_CHARS_PER_TOKEN = 2 | graph/claim_extractor.py | Characters per token for REPL output when sizing max_output_chars from a token allowance. |
PROMPT_TOKENIZER_MARGIN_SHARE = 0.10 | graph/claim_extractor.py | Share of the input allowance held back for the server tokenizer counting more than the client's estimate. |
ClaimExtractor sizes the promoted InstrumentedRLM from the window its endpoint serves: _serving_token_budget() reads the TokenBudget off the serving BudgetedLM and _rlm_output_chars() turns the input allowance, less the harness and the tokenizer margin, into per-turn max_output_chars. One module is cached per distinct cap. A window that leaves no transcript allowance raises RecursiveClaimBudgetError naming the window and the reservation; a prompt that cannot fit raises PromptBudgetExceededError before any request is sent.
Every promoted call runs under rlm_run_span, so it appears in the tenant's Phoenix project as an InstrumentedRLM.run span nested under KG_EXTRACT_SPAN_NAME (pipeline.kg.extract_per_segment). Alongside the seam's max_iterations / rlm_iterations, the claim path adds max_output_chars, context_window, input_chars, segment_id and source_doc_id. The single-prompt ChainOfThought path emits no run span.
User-facing detail with diagrams: docs/user/knowledge-graph.md.
Operational scripts: - scripts/seed_bright_corpus.py — one-shot ingest of the 30-row BRIGHT probe corpus into the test tenant's content schema. - scripts/setup_local_tests.sh — bring up the four kubectl port-forwards the live integration tests depend on (vLLM, Phoenix HTTP + gRPC).
Core Components¶
1. VideoIngestionPipeline¶
Purpose: Main orchestrator for video processing with async optimizations
Constructor:
def __init__(
self,
tenant_id: str, # REQUIRED - no default
config: PipelineConfig | None = None,
app_config: dict[str, Any] | None = None,
config_manager=None,
schema_loader=None,
schema_name: str | None = None,
debug_mode: bool = False,
event_queue: Optional[EventQueue] = None,
max_concurrent: int = 3,
)
Parameters:
tenant_id: Tenant identifier (REQUIRED - raises ValueError if not provided)config: Pipeline configuration (steps, thresholds, paths)app_config: Global application configconfig_manager: ConfigManager instance (required if app_config not provided)schema_loader: SchemaLoader instance (optional, for backend operations)schema_name: Profile name (e.g., "video_colpali_smol500_mv_frame")debug_mode: Enable detailed loggingevent_queue: Optional EventQueue for real-time progress notifications; a batch'sjob_idis its task idmax_concurrent: Maximum concurrent video processing tasks (default: 3)
Key Methods:
async process_video_async(video_path: Path | str, source_uri: str | None = None) -> dict[str, Any]¶
Process a single video through the entire pipeline. Pass a URI string (file://, s3://, pvc://, http(s)://) to have the pipeline localize it itself, or a local Path. When the caller has already localized the media (e.g. the ingestion worker downloads an s3:// object with its own object-store-configured MediaLocator), pass the local Path plus source_uri so every indexed document records the canonical source_url (the s3:// URL) rather than the temporary local path — answer-time keyframe resolution derives the object-store bucket from that recorded source_url.
result = await pipeline.process_video_async(Path("video.mp4"))
# Returns:
# {
# "video_id": "video",
# "video_path": "/path/to/video.mp4",
# "duration": 120.5,
# "status": "completed",
# "results": {
# "keyframes": {...}, # or "chunks" or "single_vector_processing"
# "transcript": {...},
# "descriptions": {...},
# "embeddings": {...}
# },
# "total_processing_time": 45.2
# }
async process_videos_concurrent(video_files: list[Path] | list[str] | list[Path | str], max_concurrent: int | None = None) -> dict[str, Any]¶
Process multiple videos concurrently with resource control. When max_concurrent is omitted it falls back to the pipeline's configured value (set via VideoIngestionPipelineBuilder.with_concurrency(...), default 3); a per-call value still overrides it.
video_files = [Path("v1.mp4"), Path("v2.mp4"), Path("v3.mp4")]
results = await pipeline.process_videos_concurrent(video_files, max_concurrent=2)
# Process 2 videos at once, queue remaining
# Returns: {"job_id": str, "results": [list of per-video results]}
Features:
- AsyncIO-based concurrent processing
- Semaphore-controlled resource limits
- Progress tracking per video
- Graceful error handling per video
def process_directory(video_dir: Path | None = None, max_concurrent: int = 3) -> dict[str, Any]¶
Synchronous entry point for batch processing.
results = pipeline.process_directory(
video_dir=Path("videos/"),
max_concurrent=3
)
# {
# "total_videos": 10,
# "processed_videos": [...], # Successful
# "failed_videos": [...], # Failed
# "total_processing_time": 300.5
# }
2. StrategyFactory¶
Purpose: Create strategy sets from explicit YAML configuration
Key Method:
@classmethod create_from_profile_config(profile_config: dict[str, Any]) -> ProcessingStrategySet¶
profile_config = {
"strategies": {
"segmentation": {
"class": "FrameSegmentationStrategy",
"params": {"fps": 0.5, "threshold": 0.999}
},
"transcription": {
"class": "AudioTranscriptionStrategy",
"params": {"model": "whisper-large-v3"}
},
"description": {
"class": "VLMDescriptionStrategy",
"params": {
"vlm_endpoint": "http://cogniverse-vllm-llm-student:8000/v1",
"batch_size": 500
}
},
"embedding": {
"class": "MultiVectorEmbeddingStrategy",
"params": {"model_name": "TomoroAI/tomoro-colqwen3-embed-4b"}
}
}
}
strategy_set = StrategyFactory.create_from_profile_config(profile_config)
Design:
- Uses dynamic imports (
importlib) to instantiate strategy classes - No hardcoded if/elif logic - fully config-driven
- All strategy classes must be in
cogniverse_runtime.ingestion.strategies
Profile-level inference_services injection (opt-in)¶
A profile can declare which inference service backs each strategy_type at the profile level:
inference_services:
embedding: vllm_colpali
transcription: vllm_asr
acoustic_embedding: clap_embed
strategies:
transcription:
class: AudioTranscriptionStrategy
params: { model: "openai/whisper-large-v3-turbo" }
embedding:
class: MultiVectorEmbeddingStrategy
params: {}
The factory looks up inference_services[<strategy_type>] and, if the strategy's __init__ explicitly declares an inference_service parameter, passes it as inference_service=<name>. Audio profiles also use inference_services.acoustic_embedding to resolve the CLAP sidecar URL into clap_endpoint_url. Strategies whose constructor does NOT declare the parameter are not injected — this is an opt-in contract, and **kwargs deliberately does not count as opting in (a **kwargs constructor absorbs the kwarg and silently drops it, so the factory would believe it was delivered).
# Strategy opts in by declaring the parameter:
class AudioTranscriptionStrategy(BaseStrategy):
def __init__(
self,
model: str = "base",
language: str = "auto",
inference_service: str | None = None,
):
...
Strict params: any other key in params that isn't a constructor parameter (and the constructor doesn't have **kwargs) raises TypeError at construction. Pre-fix the factory silently dropped unknown kwargs via a whole-signature filter, masking typos in profile JSON. Now misspelled params fail loudly so misconfiguration is observable.
3. ProcessingStrategySet¶
Purpose: Container for processing strategies with execution orchestration
Constructor:
Accepts any number of named strategies:
strategy_set = ProcessingStrategySet(
segmentation=FrameSegmentationStrategy(fps=0.5),
transcription=AudioTranscriptionStrategy(),
embedding=MultiVectorEmbeddingStrategy()
)
Key Method:
async process(video_path: Path, processor_manager, pipeline_context) -> dict[str, Any]¶
Execute all strategies, respecting their data dependencies.
results = await strategy_set.process(
video_path=Path("video.mp4"),
processor_manager=proc_manager,
pipeline_context=pipeline
)
# Returns combined results from all strategies
Execution Order: 1. Segmentation (keyframes/chunks) and Transcription (audio-to-text) run concurrently — they are independent, so the cv2 decode and Whisper transcription overlap. Each stage runs in its own telemetry span; the two are siblings under the pipeline span (asyncio.gather isolates their contexts). 2. Description → VLM descriptions (needs the keyframes from segmentation) 3. Embedding → Generate and feed embeddings (needs everything upstream)
Steps 2 and 3 run serially after the concurrent pair because each depends on the prior stage's output.
4. ProcessorManager (processor_manager.py)¶
Purpose: Manages processor lifecycle and provides instances to strategies
Key Methods:
def initialize_from_strategies(strategy_set, service_urls: dict[str, str])¶
Scan strategy requirements and create processors. service_urls is the {service_name: url} map from SystemConfig.inference_service_urls; any processor requirement carrying an inference_service key has it resolved to a concrete endpoint URL before construction. Pass {} when no remote inference services are deployed — strategies that request one then raise at init instead of silently falling back to a local model.
manager = ProcessorManager(logger)
manager.initialize_from_strategies(strategy_set, service_urls={})
# Internally calls strategy.get_required_processors() for each strategy
def get_processor(processor_name: str) -> BaseProcessor¶
Retrieve processor instance by name.
Supported Processors:
keyframe: KeyframeProcessorchunk: ChunkProcessoraudio: AudioProcessorvlm: VLMProcessorsingle_vector: SingleVectorVideoProcessor
5. BaseEmbeddingGenerator / EmbeddingResult¶
Purpose: Abstract base and result type for backend-agnostic embedding generation.
BaseEmbeddingGenerator defines the generate_embeddings(video_data, output_dir) -> EmbeddingResult contract. EmbeddingResult is the dataclass returned by every generator:
@dataclass
class EmbeddingResult:
video_id: str
total_documents: int
documents_processed: int
documents_fed: int
processing_time: float
errors: list[str]
metadata: dict
When an embedding stage runs, its result must contain the non-negative integer fields total_documents, documents_processed, and documents_fed, plus an errors list containing only non-empty strings. Processed documents cannot exceed the total, and fed documents cannot exceed the processed count. An explicit stage error, a malformed result, a non-empty errors list, or any partial feed produces a terminal failure with the embedding-stage context.
The concrete implementation is EmbeddingGeneratorImpl (below), constructed via EmbeddingGeneratorFactory / create_embedding_generator.
6. EmbeddingGeneratorImpl¶
Purpose: Concrete embedding generator supporting frame, chunk, video, document and audio content types via model_loader dispatch.
Frame and chunk embedding failures now raise EmbeddingGenerationError with the source path and failure detail, so unreadable media and remote inference errors fail the stage instead of disappearing as None. ffprobe failures in chunk extraction raise RuntimeError with the video path and probe output instead of returning an empty chunk list.
Extends: BaseEmbeddingGenerator
Model Loading: Uses ModelLoaderFactory with the model_loader config key to select the loader class directly:
# model_loader key maps directly to a loader class:
# "colpali" → ColPaliModelLoader
# "colqwen" → ColQwenModelLoader
# "xclip" → RemoteXClipLoader (remote only)
# "colbert" → ColBERTModelLoader
Additional Processing Methods (dispatched from generate_embeddings() by which key is present in video_data, checked in this order):
_process_document_visual_segments()-document_pagespresent: ColPali multi-vector page-image embeddings for PDF pages_process_document_segments()-document_filespresent: ColBERT 128-dim per-token multi-vector embeddings for text documents_process_code_segments()-code_filespresent: LateOn-Code-edge 48-dim per-token multi-vector embeddings for source-code chunks_process_audio_segments()-audio_filespresent: CLAP 512-dim acoustic + ColBERT 128-dim semantic dual embeddings for audio_process_single_document()/_process_multi_documents()- fallback for frame/chunk/window segments, chosen bystorage_mode
Required Profile Config Keys (enforced by EmbeddingGeneratorImpl.__init__):
embedding_model- Model identifier for the loader (required, raisesValueErrorif missing)model_loader- Selects the loader class inModelLoaderFactory(required, raisesValueErrorif missing)semantic_model- Secondary model for thecolbertmodel_loader (e.g.,lightonai/LateOn)
embedding_type is validated separately, at profile-registration time, by ProfileValidator._validate_embedding_type (libs/core/cogniverse_core/validation/profile_validator.py) — EmbeddingGeneratorImpl itself does not read or require it.
7. VideoIngestionPipelineBuilder / PipelineConfigBuilder (pipeline_builder.py)¶
Purpose: Fluent builder for VideoIngestionPipeline, avoiding a long positional/keyword constructor call at each call site.
Key Methods (VideoIngestionPipelineBuilder, all return self):
with_tenant_id(tenant_id),with_config_manager(config_manager),with_schema_loader(schema_loader)— required dependencieswith_config(config),with_app_config(app_config),with_schema(schema_name),with_debug(debug_mode=True)with_video_dir(video_dir),with_media_root_uri(media_root_uri),with_output_dir(output_dir),with_backend(backend),with_max_frames(max_frames)with_concurrency(max_concurrent)— sets the pipeline's defaultmax_concurrent(this is the "configured value"process_videos_concurrentfalls back to when its ownmax_concurrentargument is omitted)build() -> VideoIngestionPipeline— raisesValueErroriftenant_idorconfig_managerwas never set
pipeline = (
VideoIngestionPipelineBuilder()
.with_tenant_id("your_org:production")
.with_config_manager(config_manager)
.with_schema("video_colpali_smol500_mv_frame")
.with_concurrency(2)
.build()
)
PipelineConfigBuilder is a companion fluent builder for PipelineConfig itself, in the same module.
Processing Strategies¶
1. FrameSegmentationStrategy¶
Purpose: Extract individual frames from video
Parameters:
fps: Frames per second (default 1.0)threshold: Histogram similarity threshold (default 0.999)max_frames: Maximum frames to extract (default 3000)
Usage:
strategy = FrameSegmentationStrategy(fps=0.5, threshold=0.999, max_frames=3000)
requirements = strategy.get_required_processors()
# Returns: {"keyframe": {"fps": 0.5, "threshold": 0.999, "max_frames": 3000}}
Best For: ColPali multi-vector frame embeddings
2. ChunkSegmentationStrategy¶
Purpose: Extract video chunks for processing
Parameters:
chunk_duration: Duration of each chunk in seconds (default 30.0)chunk_overlap: Overlap between chunks in seconds (default 0.0)cache_chunks: Cache extracted chunks (default True)
Usage:
strategy = ChunkSegmentationStrategy(chunk_duration=30.0, chunk_overlap=0.0)
requirements = strategy.get_required_processors()
# Returns: {"chunk": {"chunk_duration": 30.0, "chunk_overlap": 0.0, "cache_chunks": True}}
Best For: ColQwen chunk-based video processing
Incompatible with VLMDescriptionStrategy: chunks are video-file segments, not keyframe images, so a chunk+VLM pairing cannot produce descriptions. ProcessingStrategySet raises ValueError at construction if the two are combined (chunk profiles pair with NoDescriptionStrategy).
3. SingleVectorSegmentationStrategy¶
Purpose: Process video with sliding windows for single-vector embeddings
Parameters:
strategy: Segmentation strategy ("sliding_window", "uniform")segment_duration: Segment duration in seconds (default 6.0)segment_overlap: Overlap between segments in seconds (default 1.0)sampling_fps: FPS for frame sampling within segments (default 2.0)max_frames_per_segment: Max frames per segment (default 12)store_as_single_doc: Store all segments in one document (default False)
Usage:
strategy = SingleVectorSegmentationStrategy(
strategy="sliding_window",
segment_duration=6.0,
segment_overlap=1.0,
sampling_fps=2.0,
max_frames_per_segment=12
)
Best For: X-CLIP single-vector embeddings
Custom Method:
async segment(video_path: Path, pipeline_context: Any, transcript_data: dict | None) -> dict[str, Any]¶
Directly processes video and returns segmented data:
result = await strategy.segment(
video_path=Path("video.mp4"),
pipeline_context=pipeline,
transcript_data=transcript
)
# Returns: {"single_vector_processing": {"segments": [...], "metadata": {...}}}
4. AudioTranscriptionStrategy¶
Purpose: Transcribe audio using Whisper
Parameters:
model: Whisper model (default"base", ~150MB; set"whisper-large-v3"etc. explicitly for higher accuracy)language: Language code or "auto" for detection (default"auto")inference_service: Optional named remote service (e.g."vllm_asr"); when set,AudioProcessorPOSTs to the vLLM Whisper pod's/v1/audio/transcriptionsinstead of loading a local model (defaultNone)
Usage:
5. VLMDescriptionStrategy¶
Purpose: Generate descriptions using a Vision-Language Model behind an OpenAI-compatible /v1 endpoint
Parameters:
vlm_endpoint: OpenAI-compatible/v1VLM endpoint URL (required)batch_size: Batch size for frame processing (default 500)timeout: Request timeout in seconds (default 10800 / 3 hours)vlm_concurrency: Keyframe describe-requests kept in flight (default 8). Concurrent requests are what feed vLLM's continuous batching — the chat API returns one completion per request, so there is no single-request multi-image describe. Raise it on a GPU that can serve a bigger batch; keep it low on a small one.
Usage:
strategy = VLMDescriptionStrategy(
vlm_endpoint="http://cogniverse-vllm-llm-student:8000/v1",
batch_size=500
)
Requires keyframe segmentation (FrameSegmentationStrategy): it describes keyframe images, so pairing it with a chunk-based segmentation is rejected by ProcessingStrategySet at construction.
6. MultiVectorEmbeddingStrategy¶
Purpose: Generate multi-vector embeddings (frame-by-frame)
Parameters:
model_name: Embedding model (default"TomoroAI/tomoro-colqwen3-embed-4b")inference_service: Optional named remote service for ColPali/ColQwen inference (defaultNone)
Usage:
Custom Method (shared by every embedding strategy: MultiVectorEmbeddingStrategy, SingleVectorEmbeddingStrategy, AudioEmbeddingStrategy, DocumentTextEmbeddingStrategy, DocumentVisualEmbeddingStrategy, CodeTextEmbeddingStrategy):
async generate_embeddings_with_processor(results: dict, pipeline_context) -> dict¶
Wraps results with video_id/video_path and delegates to pipeline_context.generate_embeddings(), which routes through EmbeddingGeneratorImpl. The processor manager is read off pipeline_context.processor_manager rather than passed in:
embeddings = await strategy.generate_embeddings_with_processor(
results={"keyframes": [...]},
pipeline_context=pipeline,
)
7. SingleVectorEmbeddingStrategy¶
Purpose: Generate single-vector embeddings (one per segment)
Parameters:
model_name: Embedding model (default"microsoft/xclip-large-patch14")inference_service: Optional named remote service for X-CLIP inference (defaultNone)
Usage:
8. NoDescriptionStrategy / NoTranscriptionStrategy¶
Purpose: No-op strategies for profiles that skip a stage entirely.
NoDescriptionStrategy—get_required_processors()returns{}; used when a profile has nodescriptionstage (e.g. non-VLM video profiles).NoTranscriptionStrategy—get_required_processors()returns{}; used for non-video content (images, documents, code) that has no audio track.
Usage:
strategy_set = ProcessingStrategySet(
segmentation=ImageSegmentationStrategy(),
transcription=NoTranscriptionStrategy(),
description=NoDescriptionStrategy(),
embedding=MultiVectorEmbeddingStrategy(),
)
9. ImageSegmentationStrategy¶
Purpose: Load images from a directory and present each one as a "keyframe" in the same shape FrameSegmentationStrategy produces, so MultiVectorEmbeddingStrategy works unchanged.
Parameters:
max_images: Maximum images to discover (default 10000)inference_service: Optional named remote inference service (defaultNone)
Usage:
strategy = ImageSegmentationStrategy(max_images=5000)
requirements = strategy.get_required_processors()
# Returns: {"image": {"max_images": 5000}}
Discovery for the image requirement key is handled inline by ProcessingStrategySet._process_segmentation — there is no ImageProcessor in ProcessorManager's auto-discovered processor set.
10. AudioFileSegmentationStrategy¶
Purpose: Discover audio files in a directory for standalone audio ingestion (analogous to ImageSegmentationStrategy).
Parameters:
max_files: Maximum audio files to discover (default 10000)
Usage:
strategy = AudioFileSegmentationStrategy(max_files=1000)
# get_required_processors() -> {"audio_file": {"max_files": 1000}}
11. AudioEmbeddingStrategy¶
Purpose: Generate CLAP acoustic (512-dim) + ColBERT semantic (128-dim multi-vector) dual embeddings for audio content.
Parameters:
clap_model: CLAP model for acoustic embeddings (default"laion/clap-htsat-unfused")colbert_model: ColBERT model for semantic embeddings (default"lightonai/LateOn")
Usage:
strategy = AudioEmbeddingStrategy(
clap_model="laion/clap-htsat-unfused",
colbert_model="lightonai/LateOn",
)
# get_required_processors() -> {"embedding": {"type": "audio", "clap_model": ..., "colbert_model": ...}}
12. DocumentSegmentationStrategy¶
Purpose: Discover document files (PDF via PyPDF2, plain text, markdown) in a directory for text-based ingestion.
Parameters:
max_files: Maximum document files to discover (default 10000)
Usage:
strategy = DocumentSegmentationStrategy(max_files=2000)
# get_required_processors() -> {"document_file": {"max_files": 2000}}
13. DocumentVisualSegmentationStrategy¶
Purpose: Render PDF pages to images (via pdf2image/poppler) for ColPali page-as-image ingestion — analogous to ImageSegmentationStrategy for images.
Parameters:
max_files: Maximum PDF files to discover (default 10000)dpi: Rendering resolution in DPI (default 150)
Usage:
strategy = DocumentVisualSegmentationStrategy(max_files=500, dpi=150)
# get_required_processors() -> {"document_page": {"max_files": 500, "dpi": 150}}
14. CodeSegmentationStrategy¶
Purpose: Parse source code files into AST-aware chunks using tree-sitter — function, method, class, and top-level block segments, each with name + signature + docstring + body (the same representation colgrep uses). Walks directories, respects .gitignore, filters by language extension.
Parameters:
languages: List of languages to parse (default["python", "typescript", "go", "javascript"])max_files: Maximum source files to discover (default 50000)
Usage:
strategy = CodeSegmentationStrategy(languages=["python", "go"], max_files=10000)
# get_required_processors() -> {"code_file": {"languages": [...], "max_files": 10000}}
Additional Methods: get_supported_extensions(), parse_file(file_path), and internal AST-walking helpers (_extract_segments, _extract_name, _extract_signature) that back the tree-sitter parsing.
15. DocumentTextEmbeddingStrategy¶
Purpose: Generate ColBERT multi-vector embeddings (128-dim per token) for document text extracted by DocumentSegmentationStrategy.
Parameters:
colbert_model: ColBERT model (default"lightonai/LateOn")inference_service: Optional named remote inference service (defaultNone)
Usage:
strategy = DocumentTextEmbeddingStrategy(colbert_model="lightonai/LateOn")
# get_required_processors() -> {"embedding": {"type": "document_text", "colbert_model": ...}}
16. DocumentVisualEmbeddingStrategy¶
Purpose: Generate Tomoro ColQwen3 multi-vector embeddings (320 dimensions per patch) for the PDF page images DocumentVisualSegmentationStrategy renders.
Parameters:
colpali_model: ColPali model (default"TomoroAI/tomoro-colqwen3-embed-4b")inference_service: Optional named remote inference service (defaultNone)
Usage:
strategy = DocumentVisualEmbeddingStrategy(colpali_model="TomoroAI/tomoro-colqwen3-embed-4b")
# get_required_processors() -> {"embedding": {"type": "document_visual", "colpali_model": ...}}
17. CodeTextEmbeddingStrategy¶
Purpose: Generate ColBERT multi-vector embeddings (128-dim per token) for source-code chunks produced by CodeSegmentationStrategy.
Parameters:
colbert_model: ColBERT model (default"lightonai/LateOn-Code-edge")
Usage:
strategy = CodeTextEmbeddingStrategy(colbert_model="lightonai/LateOn-Code-edge")
# get_required_processors() -> {"embedding": {"type": "code_text", "colbert_model": ...}}
Processors¶
1. KeyframeProcessor¶
Purpose: Extract representative keyframes using histogram comparison
Methods:
def extract_keyframes(video_path: Path, output_dir: Path = None) -> dict[str, Any]¶
Extract keyframes using histogram or FPS method.
processor = KeyframeProcessor(logger, threshold=0.999, max_frames=3000, fps=0.5)
result = processor.extract_keyframes(
video_path=Path("video.mp4"),
output_dir=Path("outputs/")
)
# Returns:
# {
# "keyframes": [
# {"frame_number": 0, "timestamp": 0.0, "filename": "...", "path": "..."},
# {"frame_number": 30, "timestamp": 1.0, "filename": "...", "path": "..."}
# ],
# "metadata": {...},
# "keyframes_dir": "/path/to/keyframes/",
# "video_id": "video"
# }
Extraction Modes:
- FPS Mode: Extract at regular intervals (e.g., 1 frame per second)
- Histogram Mode: Extract when scene changes (correlation < threshold)
Output:
- Saved keyframes:
{output_dir}/keyframes/{video_id}/{video_id}_keyframe_0000.jpg - Metadata JSON:
{output_dir}/metadata/{video_id}_keyframes.json
2. ChunkProcessor¶
Purpose: Extract video chunks using FFmpeg
Methods:
def extract_chunks(video_path: Path, output_dir: Path = None) -> dict[str, Any]¶
Extract video chunks with optional overlap.
processor = ChunkProcessor(logger, chunk_duration=30.0, chunk_overlap=0.0)
result = processor.extract_chunks(
video_path=Path("video.mp4"),
output_dir=Path("outputs/")
)
# Returns:
# {
# "chunks": [
# {"chunk_number": 0, "start_time": 0.0, "end_time": 30.0, "filename": "...", "path": "..."},
# {"chunk_number": 1, "start_time": 30.0, "end_time": 60.0, "filename": "...", "path": "..."}
# ],
# "metadata": {...},
# "chunks_dir": "/path/to/chunks/",
# "video_id": "video"
# }
FFmpeg Command:
ffmpeg -y -threads 4 -ss 0.0 -i video.mp4 -t 30.0 -map 0:v:0 -map 0:a? \
-c:v libx264 -threads 4 -preset ultrafast -pix_fmt yuv420p -c:a aac \
-avoid_negative_ts make_zero chunk_0000.mp4
Each chunk is re-encoded so it decodes on its own. Decoding and encoding run on 4 threads (FFMPEG_THREADS); left to ffmpeg, both size their thread pools from the node's cores rather than the container's CPU limit.
Output:
- Saved chunks:
{output_dir}/chunks/{video_id}/{video_id}_chunk_0000.mp4 - Metadata JSON:
{output_dir}/metadata/{video_id}_chunks.json
3. AudioProcessor¶
Purpose: Transcribe audio using Whisper
Methods:
def transcribe_audio(video_path: Path, output_dir: Path = None, cache=None) -> dict[str, Any]¶
Transcribe audio with caching support.
With an endpoint, the audio is decoded to 16 kHz mono and cut where vLLM's Whisper server cuts a long file (at most 30 s, at the quietest 0.1 s window of the chunk's last second). Each chunk is sent with timestamps (verbose_json); unless its timed segments run to the end of the chunk, it is also sent without (json), and the text comes from the json answer, timed by the verbose_json segments (whisper_transcription.align_text in core.md). The language the first chunk is answered in is sent with every later request. Segment times count from the start of the file. An answer that is a repetition loop, or empty for a chunk whose loudest 25 ms frame reaches -60 dBFS, is sent again at the next sampling temperature, up to six times; a chunk whose json answer never comes back usable fails the transcription with GarbledTranscriptError or EmptyTranscriptError naming the chunk and its time range; a silent chunk may come back empty. A failed transcription returns the error dict, which fails the pipeline's transcription stage.
processor = AudioProcessor(logger, model="whisper-large-v3", language="auto")
result = processor.transcribe_audio(
video_path=Path("video.mp4"),
output_dir=Path("outputs/"),
cache=cache
)
# Returns:
# {
# "video_id": "video",
# "model": "whisper-large-v3",
# "language": "en",
# "duration": 120.5,
# "transcription_time": 15.2,
# "full_text": "This is the full transcript...",
# "segments": [
# {"start": 0.0, "end": 5.2, "text": "Hello world"},
# {"start": 5.2, "end": 10.1, "text": "This is a test"}
# ]
# }
Model Mapping:
whisper-large-v3→large-v3whisper-medium→mediumwhisper-base→base
Output:
- Transcript JSON:
{output_dir}/transcripts/{video_id}_transcript.json
4. VLMProcessor¶
Purpose: Generate frame descriptions via a remote Vision-Language Model, delegating the OpenAI-compatible /v1 HTTP communication to VLMDescriptor.
Methods:
def generate_descriptions(frames_data: dict[str, Any]) -> dict[str, Any]¶
processor = VLMProcessor(
logger,
vlm_endpoint="http://cogniverse-vllm-llm-student:8000/v1",
batch_size=500,
)
result = processor.generate_descriptions(frames_data)
process(*args, **kwargs) (the BaseProcessor abstract method) forwards to generate_descriptions. cleanup() stops the underlying Modal service if VLMDescriptor auto-started it.
5. SingleVectorVideoProcessor (single_vector_processor.py)¶
Purpose: Process videos with sliding window segmentation for single-vector embeddings
Methods:
def process_video(video_path: Path, transcript_data: dict | None) -> dict[str, Any]¶
Process video into segments with transcript alignment.
processor = SingleVectorVideoProcessor(
logger,
strategy="sliding_window",
segment_duration=6.0,
segment_overlap=1.0,
sampling_fps=2.0
)
result = processor.process_video(
video_path=Path("video.mp4"),
transcript_data=transcript
)
# Returns:
# {
# "segments": [VideoSegment(...), VideoSegment(...), ...],
# "metadata": {...},
# "full_transcript": "...",
# "document_structure": {"type": "multi_document"}
# }
VideoSegment Structure:
@dataclass
class VideoSegment:
segment_id: int
start_time: float
end_time: float
frames: list[np.ndarray] # Actual frame data
frame_timestamps: list[float] # Timestamps for each frame
transcript_segments: list[dict[str, Any]] # Transcript segments
transcript_text: str = "" # Combined transcript text
metadata: dict[str, Any] = None # Additional metadata
Internal Helper Classes (not BaseProcessor subclasses)¶
These plain classes back the processors above and are not auto-discovered by ProcessorManager themselves:
AudioTranscriber(audio_transcriber.py) — Whisper model loading and the core transcription call; used byAudioProcessor.AudioEmbeddingGenerator(audio_embedding_generator.py) — lazy CLAP loading and acoustic embedding generation; used byEmbeddingGeneratorImpl._process_audio_segments(). Remote (clap_endpoint_url) calls reuse one pooledhttpx.Clientacross the instance instead of opening a connection per segment;close()releases that pooled client.VLMDescriptor(vlm_descriptor.py) — HTTP client for an OpenAI-compatible/v1vision chat endpoint; used byVLMProcessor.pool_document_tokens(embedding_generator/token_pooling.py) — hierarchical token pooling of a document's multi-vector before feed (Ward-linkage clusters of its tokens, cut atn_tokens // pool_factor, each mean-pooled and L2-renormalized), the method of colpali_engine'sHierarchicalTokenPoolerin numpy and scipy.EmbeddingGeneratorImplapplies it to frame, chunk and image multi-vectors when the profile setsmodel_config.token_pool_factor; queries are never pooled.
The four visual profiles (video_colpali_smol500_mv_frame, video_colqwen_omni_mv_chunk_30s, image_colpali_mv, document_visual_colpali) ship unpooled, model_config.token_pool_factor: 1, so every document token is stored (227 for a 640x360 frame). The generator honours the factor, so changing it takes effect at the next ingest. It ships at 1 because pooling was measured on the exported production corpora with the golden set (125 queries), paired 95% bootstrap CIs:
| Profile | Tokens unpooled → factor 2 → 3 | Result |
|---|---|---|
| frame (361 documents) | 153,947 → 76,851 → 51,146 (104.7 → 52.3 → 34.8 MB) | float strategies and default level; at factor 3 every binary-MaxSim strategy loses (binary_binary MRR −0.052, R@1 −0.088; hybrid_binary_bm25 −0.056, −0.104; hybrid_bm25_binary −0.051, −0.088); at factor 2 the two binary hybrids still lose MRR (−0.028, −0.030) |
| chunk_30s (34 documents) | 14,538 → 7,257 → 4,830 | level, with the binary strategies trending down at factor 3 (MRR −0.036) |
| image, document visual | — | no evaluation corpus |
Pooling cut the backend search p50 (frame default 45 → 26 → 22 ms), but not enough to outweigh the binary-strategy loss. - resolve_served_model_id(...) (served_model.py) — returns the model id an OpenAI-compatible /v1 endpoint serves, used by VLMDescriptor and AudioProcessor's remote path. The remote services scale to zero, so discovery retries a cold endpoint until the service spec's boot_deadline_seconds, caches the answer per process per endpoint, and raises ServedModelUnavailable naming the endpoint when the budget runs out or the endpoint answers a status waiting cannot repair. - EmbeddingGeneratorFactory (embedding_generator/embedding_generator_factory.py) — exposes create_embedding_generator(...), the factory function used to construct EmbeddingGeneratorImpl. - BackendFactory (embedding_generator/backend_factory.py) — BackendFactory.create(backend_type, tenant_id, config, ...) builds the IngestionBackend (Vespa) client fed to EmbeddingGeneratorImpl.
Data Flow¶
End-to-End Ingestion Flow¶
1. VIDEO INPUT (video.mp4)
↓
2. PIPELINE INITIALIZATION
• Load profile config: video_colpali_smol500_mv_frame
• Resolve strategy: FrameSegmentationStrategy
• Create strategy set: {segmentation, transcription, embedding}
• Initialize processors: KeyframeProcessor, AudioProcessor
• Initialize embedding generator with ColPali model
↓
3. PER-STRATEGY CACHE (shared multi-pod tier)
• Each strategy step below first checks the PipelineArtifactCache for
its artifact against the shared backend; on a hit it skips the
computation (keyframes additionally rehydrate their frame image
files to this pod's disk so downstream steps can open them), on a
miss it computes and stores the result.
↓
4. SEGMENTATION (FrameSegmentationStrategy)
• Cache hit → load cached keyframes + rehydrate frame files to disk
• Cache miss → KeyframeProcessor.extract_keyframes() (histogram scene
detection, ~150 keyframes saved to disk), then store in cache
• Return: {"keyframes": [{frame_number, timestamp, filename, path}, ...]}
↓
5. TRANSCRIPTION (AudioTranscriptionStrategy)
• Cache hit → load cached transcript
• Cache miss → AudioProcessor.transcribe_audio() (Whisper), then store
in cache
• Return: {"full_text": "...", "segments": [{start, end, text}, ...]}
↓
6. EMBEDDING GENERATION (MultiVectorEmbeddingStrategy)
• EmbeddingGeneratorImpl.generate_embeddings()
• Model: ColPali (TomoroAI/tomoro-colqwen3-embed-4b)
• Processing:
- Load keyframe images (150 images)
- Batch inference (batch_size=8)
- Generate embeddings per frame: [150 x 1024 x 320]
• Document Building:
- Per-frame documents: 150 documents
- Fields: video_id, frame_number, timestamp, embeddings (hex-encoded)
• Backend Feeding:
- Feed to Vespa via VespaPyClient
- feed_iterable() with batch_size=50
- 150 documents fed successfully
↓
7. PIPELINE RESULT
{
"video_id": "video",
"status": "completed",
"total_processing_time": 85.3,
"results": {
"keyframes": {...},
"transcript": {...},
"embeddings": {
"total_documents": 150,
"documents_fed": 150,
"processing_time": 45.2
}
}
}
Concurrent Multi-Video Processing¶
flowchart TB
subgraph Semaphore["<span style='color:#000'>Semaphore (limit=2)</span>"]
subgraph Active["<span style='color:#000'>Active Processing</span>"]
V1["<span style='color:#000'>Process v1.mp4<br/>├─ Segment<br/>├─ Transcribe<br/>├─ Embed<br/>└─ Feed (45s)</span>"]
V2["<span style='color:#000'>Process v2.mp4<br/>├─ Segment<br/>├─ Transcribe<br/>└─ Embed (30s)</span>"]
end
subgraph Queued["<span style='color:#000'>Queued</span>"]
V3["<span style='color:#000'>v3.mp4</span>"]
V4["<span style='color:#000'>v4.mp4</span>"]
V5["<span style='color:#000'>v5.mp4</span>"]
end
end
V1 -->|completes| V3
V2 -->|completes| V4
style Semaphore fill:#ce93d8,stroke:#7b1fa2,color:#000
style Active fill:#a5d6a7,stroke:#388e3c,color:#000
style Queued fill:#b0bec5,stroke:#546e7a,color:#000
style V1 fill:#a5d6a7,stroke:#388e3c,color:#000
style V2 fill:#a5d6a7,stroke:#388e3c,color:#000
style V3 fill:#b0bec5,stroke:#546e7a,color:#000
style V4 fill:#b0bec5,stroke:#546e7a,color:#000
style V5 fill:#b0bec5,stroke:#546e7a,color:#000 When v1 or v2 completes → v3 starts immediately
Usage Examples¶
Example 1: Basic Single Video Ingestion (ColPali)¶
import asyncio
from pathlib import Path
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline, PipelineConfig
async def ingest_video():
# Create pipeline for ColPali frame-based processing
pipeline = VideoIngestionPipeline(
tenant_id="your_org:production", # Required parameter
schema_name="video_colpali_smol500_mv_frame",
debug_mode=True
)
# Process single video
result = await pipeline.process_video_async(Path("data/videos/demo.mp4"))
print(f"Status: {result['status']}")
print(f"Documents fed: {result['results']['embeddings']['documents_fed']}")
print(f"Processing time: {result['total_processing_time']:.1f}s")
# Output:
# Status: completed
# Documents fed: 150
# Processing time: 85.3s
# Run the async function
asyncio.run(ingest_video())
Profile Configuration (configs/config.json, under backend.profiles):
video_colpali_smol500_mv_frame:
strategies:
segmentation:
class: "FrameSegmentationStrategy"
params:
fps: 0.5
threshold: 0.999
max_frames: 3000
transcription:
class: "AudioTranscriptionStrategy"
params:
model: "whisper-large-v3"
description:
class: "NoDescriptionStrategy"
params: {}
embedding:
class: "MultiVectorEmbeddingStrategy"
params:
model_name: "TomoroAI/tomoro-colqwen3-embed-4b"
Example 2: Batch Processing with Concurrency¶
from pathlib import Path
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline
# Create pipeline
pipeline = VideoIngestionPipeline(
tenant_id="your_org:production",
schema_name="video_colpali_smol500_mv_frame"
)
# Process directory with concurrent processing (3 videos at once)
results = pipeline.process_directory(
video_dir=Path("data/videos/"),
max_concurrent=3
)
print(f"Total videos: {results['total_videos']}")
print(f"Processed: {len(results['processed_videos'])}")
print(f"Failed: {len(results['failed_videos'])}")
print(f"Total time: {results['total_processing_time'] / 60:.1f} minutes")
print(f"Throughput: {results['total_videos'] / (results['total_processing_time'] / 60):.1f} videos/min")
# Output:
# Total videos: 50
# Processed: 48
# Failed: 2
# Total time: 25.3 minutes
# Throughput: 2.0 videos/min
Example 3: X-CLIP Single-Vector Processing¶
import asyncio
from pathlib import Path
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline
async def ingest_with_xclip():
# Create pipeline for X-CLIP single-vector embeddings
pipeline = VideoIngestionPipeline(
tenant_id="your_org:production",
schema_name="video_xclip_sv_chunk_6s"
)
# Process video
result = await pipeline.process_video_async(Path("data/videos/lecture.mp4"))
# Access single-vector processing results
sv_data = result['results']['single_vector_processing']
print(f"Segments: {len(sv_data['segments'])}")
print(f"Document structure: {sv_data['document_structure']['type']}")
print(f"Documents fed: {result['results']['embeddings']['documents_fed']}")
# Output:
# Segments: 20
# Document structure: multi_document
# Documents fed: 20
asyncio.run(ingest_with_xclip())
Profile Configuration:
video_xclip_sv_chunk_6s:
strategies:
segmentation:
class: "SingleVectorSegmentationStrategy"
params:
strategy: "sliding_window"
segment_duration: 6.0
segment_overlap: 1.0
sampling_fps: 2.0
max_frames_per_segment: 12
store_as_single_doc: false
transcription:
class: "AudioTranscriptionStrategy"
params:
model: "whisper-large-v3"
description:
class: "NoDescriptionStrategy"
params: {}
embedding:
class: "SingleVectorEmbeddingStrategy"
params:
model_name: "microsoft/xclip-large-patch14"
Example 4: ColQwen Chunk-Based Processing¶
import asyncio
from pathlib import Path
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline
async def ingest_with_colqwen():
# Create pipeline for ColQwen chunk processing
pipeline = VideoIngestionPipeline(
tenant_id="your_org:production",
schema_name="video_colqwen_omni_mv_chunk_30s"
)
# Process video
result = await pipeline.process_video_async(Path("data/videos/tutorial.mp4"))
# Access chunk processing results
chunks = result['results']['video_chunks']['chunks']
print(f"Chunks extracted: {len(chunks)}")
for chunk in chunks[:3]:
print(f" Chunk {chunk['chunk_number']}: {chunk['start_time']:.1f}s - {chunk['end_time']:.1f}s")
print(f"Documents fed: {result['results']['embeddings']['documents_fed']}")
# Output:
# Chunks extracted: 12
# Chunk 0: 0.0s - 30.0s
# Chunk 1: 30.0s - 60.0s
# Chunk 2: 60.0s - 90.0s
# Documents fed: 12
asyncio.run(ingest_with_colqwen())
Profile Configuration:
video_colqwen_omni_mv_chunk_30s:
strategies:
segmentation:
class: "ChunkSegmentationStrategy"
params:
chunk_duration: 30.0
chunk_overlap: 0.0
cache_chunks: true
transcription:
class: "AudioTranscriptionStrategy"
params:
model: "whisper-large-v3"
description:
class: "NoDescriptionStrategy"
params: {}
embedding:
class: "MultiVectorEmbeddingStrategy"
params:
model_name: "TomoroAI/tomoro-colqwen3-embed-4b"
Example 5: Custom Strategy Configuration¶
import asyncio
from pathlib import Path
from cogniverse_runtime.ingestion.strategy_factory import StrategyFactory
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline
async def ingest_with_custom_strategy():
# Define custom profile config
custom_profile = {
"strategies": {
"segmentation": {
"class": "FrameSegmentationStrategy",
"params": {"fps": 0.5, "max_frames": 500} # Fewer frames, lower FPS
},
"transcription": {
"class": "AudioTranscriptionStrategy",
"params": {"model": "whisper-medium"} # Faster model
},
"description": {
"class": "VLMDescriptionStrategy",
"params": {
"vlm_endpoint": "http://cogniverse-vllm-llm-student:8000/v1",
"batch_size": 500
}
},
"embedding": {
"class": "MultiVectorEmbeddingStrategy",
"params": {"model_name": "TomoroAI/tomoro-colqwen3-embed-4b"}
}
}
}
# Create strategy set
strategy_set = StrategyFactory.create_from_profile_config(custom_profile)
# Use in pipeline (manual initialization)
pipeline = VideoIngestionPipeline(tenant_id="your_org:production", app_config=app_config)
pipeline.strategy_set = strategy_set
pipeline.processor_manager.initialize_from_strategies(strategy_set)
# Process video
result = await pipeline.process_video_async(Path("video.mp4"))
return result
asyncio.run(ingest_with_custom_strategy())
Example 6: Production Batch Processing Script¶
#!/usr/bin/env python3
"""
Production ingestion script with monitoring and error handling.
"""
import asyncio
from pathlib import Path
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline
async def main():
profiles = [
"video_colpali_smol500_mv_frame",
"video_xclip_sv_chunk_6s",
"video_colqwen_omni_mv_chunk_30s"
]
video_dir = Path("data/production/videos/")
for profile in profiles:
print(f"\n{'='*60}")
print(f"Processing with profile: {profile}")
print(f"{'='*60}\n")
pipeline = VideoIngestionPipeline(tenant_id="your_org:production", schema_name=profile)
results = pipeline.process_directory(
video_dir=video_dir,
max_concurrent=2 # Conservative for production
)
# Save summary
summary_path = Path(f"outputs/ingestion_summary_{profile}.json")
import json
with open(summary_path, 'w') as f:
json.dump(results, f, indent=2)
print(f"\n✅ Profile {profile} completed:")
print(f" Processed: {len(results['processed_videos'])}")
print(f" Failed: {len(results['failed_videos'])}")
print(f" Summary: {summary_path}")
# Log failures
if results['failed_videos']:
print(f"\n⚠️ Failed videos:")
for failed in results['failed_videos']:
print(f" - {failed['video_path']}: {failed.get('error', 'Unknown error')}")
if __name__ == "__main__":
asyncio.run(main())
Production Considerations¶
1. Performance Optimization¶
Concurrent Processing:
- Use
max_concurrent=2-3for production to balance throughput and resource usage - Monitor memory usage (each video loads models + images into RAM)
- Consider GPU availability when setting concurrency limits
Caching Strategy:
pipeline_cache:
enabled: true
backends:
- type: "disk"
path: "outputs/cache/"
default_ttl: 0 # Infinite TTL for production
enable_compression: true
Model Loading:
- Load models lazily to reduce startup time
- Consider model quantization for faster inference (e.g., int8 embeddings)
- Use GPU when available (
export CUDA_VISIBLE_DEVICES=0)
2. Error Handling¶
Pipeline Exceptions:
from cogniverse_runtime.ingestion.exceptions import PipelineException
try:
result = await pipeline.process_video_async(video_path)
except PipelineException as e:
logger.error(f"Pipeline failed: {e}")
logger.error(f"Error context: {e.context}")
# Continue with next video
Per-Video Error Isolation:
- Concurrent processing isolates errors per video
- Failed videos don't affect successful ones
- All results include error details for debugging
3. Monitoring¶
Metrics to Track:
# Per-video metrics
- Processing time (segmentation, transcription, embedding)
- Document counts (total, processed, fed)
- Error rates and types
# Batch metrics
- Throughput (videos/minute)
- Success rate
- Cache hit rates
- Average processing time per profile
Logging:
# Profile-specific logs
outputs/logs/video_processing_{profile}_{timestamp}.log
# Example log entries
2025-10-07 10:30:15 - VideoIngestionPipeline_colpali - INFO - Starting async video processing: video_001
2025-10-07 10:30:45 - VideoIngestionPipeline_colpali - INFO - ✅ Extracted 150 keyframes using histogram method in 3.2s
2025-10-07 10:31:20 - VideoIngestionPipeline_colpali - INFO - ✅ Audio transcribed in 15.2s (45 segments)
2025-10-07 10:32:50 - VideoIngestionPipeline_colpali - INFO - ✅ Embeddings generated: 150 documents fed to backend
2025-10-07 10:32:50 - VideoIngestionPipeline_colpali - INFO - Async video processing completed in 155.3s
4. Storage Management¶
Output Directory Structure:
outputs/processing/
├── profile_video_colpali_smol500_mv_frame/
│ ├── keyframes/
│ │ ├── video_001/
│ │ │ ├── video_001_keyframe_0000.jpg
│ │ │ ├── video_001_keyframe_0001.jpg
│ │ │ └── ...
│ ├── metadata/
│ │ ├── video_001_keyframes.json
│ │ ├── video_001_chunks.json
│ │ └── ...
│ ├── transcripts/
│ │ └── video_001_transcript.json
│ └── pipeline_summary.json
└── profile_video_xclip_sv_chunk_6s/
└── ...
Disk Space Management:
- Keyframes: ~50KB per frame × 150 frames = ~7.5MB per video
- Chunks: ~1MB per 30s chunk × 12 chunks = ~12MB per video
- Cache: Enable compression to reduce storage (50% reduction typical)
- Clean up old profiles:
rm -rf outputs/processing/profile_old_*/
5. Scalability¶
Horizontal Scaling:
# Distribute videos across multiple machines
import socket
machine_id = socket.gethostname()
video_files = get_video_files(video_dir)
# Assign videos by hash
assigned_videos = [
v for i, v in enumerate(video_files)
if hash(str(v)) % num_machines == machine_id
]
pipeline.process_videos_concurrent(assigned_videos, max_concurrent=3)
Profile-Based Scaling:
# Process different profiles on different GPUs
profiles_gpu0 = ["video_colpali_smol500_mv_frame"]
profiles_gpu1 = ["video_xclip_sv_chunk_6s"]
# GPU 0
os.environ["CUDA_VISIBLE_DEVICES"] = "0"
for profile in profiles_gpu0:
process_profile(profile)
# GPU 1
os.environ["CUDA_VISIBLE_DEVICES"] = "1"
for profile in profiles_gpu1:
process_profile(profile)
6. Testing Strategy¶
Unit Tests:
- Test individual processors in isolation
- Mock video files with fixtures
- Verify output formats and metadata
Integration Tests:
- Test full pipeline with real videos (small samples)
- Verify Vespa document feeding
- Test cache functionality
End-to-End Tests:
- Test complete workflow: ingestion → search → retrieval
- Verify embedding quality through search results
- Test multiple profiles
Testing¶
Key Test Files¶
The tests/ingestion/ suite has 69 files (42 unit, 26 integration, 1 shared integration/conftest.py).
Unit Tests (tests/ingestion/unit/):¶
| File | Covers |
|---|---|
test_keyframe_processor.py | Keyframe extraction |
test_keyframe_processor_real.py | KeyframeProcessor against a real video file |
test_chunk_processor.py | Chunk extraction (mock logger) |
test_chunk_processor_basic.py | Basic ChunkProcessor behavior |
test_chunk_processor_real.py | ChunkProcessor against a real video file |
test_audio_processor.py | Real factory→manager wiring for the audio processor |
test_audio_processor_real.py | AudioProcessor against real audio |
test_audio_ingestion.py | Audio-file directory discovery (AudioFileSegmentationStrategy) |
test_audio_embedding_failure.py | Acoustic embedding failure raises instead of returning zeros |
test_document_ingestion.py | Document-file directory discovery (DocumentSegmentationStrategy) |
test_image_ingestion.py | Image directory discovery (ImageSegmentationStrategy) |
test_embedding_generator_impl.py | EmbeddingGeneratorImpl |
test_token_pooling.py | Token-pooling utilities |
test_single_vector_processor_basic.py | Basic SingleVectorVideoProcessor behavior |
test_single_vector_process_adapter.py | SingleVectorVideoProcessor.process forwards transcript_data/metadata |
test_vlm_descriptor.py | VLMDescriptor |
test_cold_inference_endpoint_discovery.py | Model-id discovery against a cold, failing or dead /v1 endpoint |
test_description_stage_named_on_failure.py | A failed description stage names itself in ContentProcessingError |
test_processor_base.py / test_processor_base_basic.py | BaseProcessor contract |
test_processor_manager.py | ProcessorManager auto-discovery and initialization |
test_strategy_factory_inference_services.py | StrategyFactory profile-level inference_services injection |
test_pipeline.py | PipelineConfig |
test_pipeline_orchestration.py | Strategy orchestration through ProcessingStrategySet |
test_pipeline_concurrency_wiring.py | with_concurrency() reaches the pipeline's concurrency control |
test_pipeline_locator_wiring.py | MediaLocator integration in VideoIngestionPipeline |
test_cache_key_collision_fix.py | URI-hash cache-key regression in PipelineArtifactCache |
test_source_url_in_documents.py | source_url threaded from the pipeline into every emitted Document |
test_run_ingestion.py | run_ingestion CLI script |
Integration Tests (tests/ingestion/integration/):¶
| File | Covers |
|---|---|
test_real_ingestion_pipeline.py | Real integration with actual video processing |
test_pipeline_real.py | Real-video pipeline smoke test |
test_end_to_end_processing.py | End-to-end processing with real processors |
test_backend_ingestion.py | Vespa document feeding |
test_multimodal_content_processing.py | VespaPyClient against the test Vespa instance |
test_audio_acoustic_text_embedding.py | CLAP-space 512-d text embedding from the cluster's CLAP service |
test_pipeline_cache_live_path.py | Live-path PipelineArtifactCache wiring |
test_pipeline_minio_round_trip.py | Pipeline reads from MinIO, writes to Vespa |
test_upload_via_queue.py | POST /ingestion/upload end-to-end |
test_backfill_source_url.py | scripts/backfill_source_url.py against real Vespa |
test_source_url_round_trip.py | source_url written at ingest comes back at search |
test_document_visual_ingestion_real.py | document_visual round-trip: PDF pages → ColPali → Vespa → search |
test_clap_remote_embedding.py | AudioEmbeddingGenerator remote path ↔ clap_embed sidecar |
test_colbert_remote_roundtrip.py | Remote ColBERT path round-trip |
test_lateon_real_model.py | LateOn real-model smoke test |
test_lateon_pylate_parity.py | PyLate-service LateOn per-token embeddings vs pylate oracle |
test_denseon_vllm_real.py | Real vLLM DenseOn vs sentence-transformers oracle |
test_tomoro_serving_real.py | Real vLLM Tomoro serving (320-dim normalization + pooling gate) |
test_whisper_remote_roundtrip.py | Remote Whisper transcription round-trip |
test_vllm_asr_real_sidecar.py | Real vLLM ASR sidecar |
test_vllm_colpali_real_sidecar.py | Real vLLM ColPali sidecar (RemoteColPaliLoader) |
Example Test Scenarios¶
# Test frame extraction
def test_keyframe_extraction_histogram():
processor = KeyframeProcessor(logger, threshold=0.999, max_frames=100)
result = processor.extract_keyframes(test_video_path)
assert result["keyframes"]
assert len(result["keyframes"]) > 0
assert all(Path(kf["path"]).exists() for kf in result["keyframes"])
# Test concurrent processing
@pytest.mark.asyncio
async def test_concurrent_video_processing():
pipeline = VideoIngestionPipeline(tenant_id="test", schema_name="test_profile")
video_files = [Path(f"test_video_{i}.mp4") for i in range(5)]
results = await pipeline.process_videos_concurrent(video_files, max_concurrent=2)
assert len(results) == 5
assert all(r["status"] == "completed" for r in results)
# Test embedding generation
def test_colpali_embedding_generation():
generator = create_embedding_generator(
config=test_config,
schema_name="video_colpali_mv_frame",
tenant_id="test",
config_manager=config_manager,
schema_loader=schema_loader,
)
result = generator.generate_embeddings(video_data, output_dir)
assert result.documents_fed == result.total_documents
assert result.processing_time > 0
assert len(result.errors) == 0
Real-Time Progress Notifications¶
Overview¶
The VideoIngestionPipeline reports progress on the EventQueue it is built with. A job started with POST /ingestion/start reports on its ingestion task in the shared task event store, so any runtime process streams, lists and cancels it:
- Live Progress:
GET /events/ingestion/{job_id}(task events) orGET /ingestion/{job_id}/events(status entries), from any process - Multiple Subscribers: the web client and CLI can watch the same job
- Graceful Cancellation:
POST /events/ingestion/{job_id}/cancelstops the pipeline before its next video - Job Tracking: the pipeline's
job_idis its queue's task id
Event Flow¶
With a queue, process_videos_concurrent emits:
- Job starts → StatusEvent("working", phase="starting")
- Each video → ProgressEvent(current=N, total=total_videos), then a StatusEvent at the video's start and a ProgressEvent at its end
- A video fails → ErrorEvent naming it
- Job ends → CompleteEvent with the results summary, or a StatusEvent("cancelled") when it was cancelled
On a job's ingestion task each event is stored as a status entry {"state": "running", "ingest_id", "event": ...}; the end-of-job event is held until /ingestion/start has recorded the job's outcome and is stored with the terminal status (complete, failed or cancelled).
Pipeline Outside the Runtime¶
from cogniverse_core.events import InMemoryEventQueue
from cogniverse_runtime.ingestion.pipeline import VideoIngestionPipeline, PipelineConfig
queue = InMemoryEventQueue(task_id="ingestion_job_123", tenant_id="acme:acme")
pipeline = VideoIngestionPipeline(
tenant_id="acme:acme",
config=PipelineConfig(video_dir=video_path, output_dir=output_path),
event_queue=queue,
)
result = await pipeline.process_videos_concurrent(video_files)
assert result["job_id"] == "ingestion_job_123"
# The pipeline checks the queue's cancellation token between videos
queue.cancel("User cancelled")
See Events Module for complete EventQueue documentation.
Related Modules¶
- Backends Module (
backends.md): Vespa search integration, document feeding - Common Module (
common.md): Model loading, configuration, output management - Events Module (
events.md): Task events, cancellation and the active-task listing - System Integration (
tests/ingestion/integration/test_end_to_end_processing.py,tests/e2e/test_ingestion_upload_e2e.py,tests/runtime/integration/test_ingestion_worker_e2e.py): End-to-end ingestion → search testing
Study Tip: Run the ingestion pipeline with debug_mode=True and follow the logs to understand each processing stage. Experiment with different profiles to see how strategies affect the output.
Production Checklist:
- ✅ Configure profiles for your video types
- ✅ Test with sample videos before production batch
- ✅ Enable caching to avoid reprocessing
- ✅ Monitor disk space and memory usage
- ✅ Set appropriate
max_concurrentbased on resources - ✅ Implement error handling and retry logic
- ✅ Save summaries for history
- ✅ Watch long jobs on
/events/ingestion/{job_id}
File References:
- Pipeline:
libs/runtime/cogniverse_runtime/ingestion/pipeline.py - StrategyFactory:
libs/runtime/cogniverse_runtime/ingestion/strategy_factory.py - Strategies:
libs/runtime/cogniverse_runtime/ingestion/strategies.py - EmbeddingGeneratorImpl:
libs/runtime/cogniverse_runtime/ingestion/processors/embedding_generator/embedding_generator_impl.py - KeyframeProcessor:
libs/runtime/cogniverse_runtime/ingestion/processors/keyframe_processor.py - ChunkProcessor:
libs/runtime/cogniverse_runtime/ingestion/processors/chunk_processor.py - AudioProcessor:
libs/runtime/cogniverse_runtime/ingestion/processors/audio_processor.py