Skip to content

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

  1. Module Overview
  2. Architecture
  3. Core Components
  4. Processing Strategies
  5. Processors
  6. Data Flow
  7. Usage Examples
  8. Production Considerations
  9. 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:

  1. Video Segmentation - Extract frames, chunks, or sliding windows
  2. Audio Transcription - Whisper-based speech-to-text
  3. Visual Description - VLM-based frame descriptions
  4. Embedding Generation - ColPali, X-CLIP, ColQwen, ColBERT embeddings
  5. Image Processing - Directory of images presented as keyframes (ImageSegmentationStrategy) for ColPali embedding
  6. Document Processing - Text extraction and ColBERT semantic embeddings for PDFs/documents
  7. Document Visual Processing - PDF pages rendered to images (pdf2image) with ColPali multi-vector page embeddings (document_visual_colpali profile → document_visual schema)
  8. Audio Processing - CLAP acoustic + ColBERT semantic dual embeddings for audio content, plus standalone audio-file discovery (AudioFileSegmentationStrategy)
  9. Code Processing - tree-sitter AST-aware chunking of source files (CodeSegmentationStrategy) with ColBERT multi-vector embeddings (CodeTextEmbeddingStrategy)
  10. 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 by storage_mode (multi_doc vs single_doc), not a processing_type field
  • Document ColBERT (document_files present → _process_document_segments): ColBERT 128-dim per-token multi-vector for text documents
  • Document Visual ColPali (document_pages present → _process_document_visual_segments): ColPali (Tomoro ColQwen3) 320-dim per-patch multi-vector for PDF pages rendered to images (DocumentVisualSegmentationStrategy → DocumentVisualEmbeddingStrategy)
  • Code ColBERT (code_files present → _process_code_segments): LateOn-Code-edge 48-dim per-token multi-vector for source-code chunks (CodeSegmentationStrategy → CodeTextEmbeddingStrategy)
  • Audio Dual (audio_files present → _process_audio_segments): CLAP 512-dim acoustic single-vector + ColBERT 128-dim semantic multi-vector for audio content. The transcription result also supplies audio_language and audio_duration, which land in the schema fields of the same names. A profile that binds inference_services.acoustic_embedding fails 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. Each SegmentRecord carries text plus a segment_anchor: Mention with 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 calls DocExtractor.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 calls DocExtractor.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 accumulates ExtractionResult (anchored Node.mentions: List[Mention] + SPO Edge rows from ClaimExtractor) in segment order, runs CrossModalLinker.link() to add same_as edges across modalities, GraphManager.upsert()s the merged result, then PATCHes per-segment back-refs (entity_ids / relation_ids / claim_ids arrays) 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 config
  • config_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 logging
  • event_queue: Optional EventQueue for real-time progress notifications; a batch's job_id is its task id
  • max_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:

def __init__(self, **strategies)

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.

keyframe_proc = manager.get_processor("keyframe")
audio_proc = manager.get_processor("audio")

Supported Processors:

  • keyframe: KeyframeProcessor
  • chunk: ChunkProcessor
  • audio: AudioProcessor
  • vlm: VLMProcessor
  • single_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_pages present: ColPali multi-vector page-image embeddings for PDF pages
  • _process_document_segments() - document_files present: ColBERT 128-dim per-token multi-vector embeddings for text documents
  • _process_code_segments() - code_files present: LateOn-Code-edge 48-dim per-token multi-vector embeddings for source-code chunks
  • _process_audio_segments() - audio_files present: 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 by storage_mode

Required Profile Config Keys (enforced by EmbeddingGeneratorImpl.__init__):

  • embedding_model - Model identifier for the loader (required, raises ValueError if missing)
  • model_loader - Selects the loader class in ModelLoaderFactory (required, raises ValueError if missing)
  • semantic_model - Secondary model for the colbert model_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 dependencies
  • with_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 default max_concurrent (this is the "configured value" process_videos_concurrent falls back to when its own max_concurrent argument is omitted)
  • build() -> VideoIngestionPipeline — raises ValueError if tenant_id or config_manager was 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, AudioProcessor POSTs to the vLLM Whisper pod's /v1/audio/transcriptions instead of loading a local model (default None)

Usage:

strategy = AudioTranscriptionStrategy(model="whisper-large-v3", language="auto")

5. VLMDescriptionStrategy

Purpose: Generate descriptions using a Vision-Language Model behind an OpenAI-compatible /v1 endpoint

Parameters:

  • vlm_endpoint: OpenAI-compatible /v1 VLM 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 (default None)

Usage:

strategy = MultiVectorEmbeddingStrategy(model_name="TomoroAI/tomoro-colqwen3-embed-4b")

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 (default None)

Usage:

strategy = SingleVectorEmbeddingStrategy(model_name="microsoft/xclip-large-patch14")

8. NoDescriptionStrategy / NoTranscriptionStrategy

Purpose: No-op strategies for profiles that skip a stage entirely.

  • NoDescriptionStrategy — get_required_processors() returns {}; used when a profile has no description stage (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 (default None)

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 (default None)

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 (default None)

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-v3
  • whisper-medium → medium
  • whisper-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 by AudioProcessor.
  • AudioEmbeddingGenerator (audio_embedding_generator.py) — lazy CLAP loading and acoustic embedding generation; used by EmbeddingGeneratorImpl._process_audio_segments(). Remote (clap_endpoint_url) calls reuse one pooled httpx.Client across the instance instead of opening a connection per segment; close() releases that pooled client.
  • VLMDescriptor (vlm_descriptor.py) — HTTP client for an OpenAI-compatible /v1 vision chat endpoint; used by VLMProcessor.
  • 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 at n_tokens // pool_factor, each mean-pooled and L2-renormalized), the method of colpali_engine's HierarchicalTokenPooler in numpy and scipy. EmbeddingGeneratorImpl applies it to frame, chunk and image multi-vectors when the profile sets model_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

Videos: [v1.mp4, v2.mp4, v3.mp4, v4.mp4, v5.mp4]
max_concurrent: 2
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-3 for 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) or GET /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}/cancel stops the pipeline before its next video
  • Job Tracking: the pipeline's job_id is its queue's task id

Event Flow

With a queue, process_videos_concurrent emits:

  1. Job starts → StatusEvent("working", phase="starting")
  2. Each video → ProgressEvent(current=N, total=total_videos), then a StatusEvent at the video's start and a ProgressEvent at its end
  3. A video fails → ErrorEvent naming it
  4. 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.


  • 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_concurrent based 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