HyperSaaS
BackendDocuments & RAG

Ingestion Pipeline

Document parsing, chunking, and embedding via Celery tasks.

The ingestion pipeline runs as Celery tasks, processing documents asynchronously after upload or URL submission.

File Ingestion

Task: documents.process_document

@shared_task(
    bind=True,
    max_retries=3,
    soft_time_limit=600,   # 10 minutes (OCR on long PDFs is slow)
    time_limit=660,        # 11 minutes hard limit
    acks_late=True,
    reject_on_worker_lost=True,
)
def process_document(self, document_id, processing_task_id=None):

Pipeline Steps

1. Download from S3 → temp directory
    │
    ▼
2. Parse document
    ├─ Docling (primary) → structured output with sections, tables, headings,
    │                      in a separate process, for at most 240s
    └─ PyMuPDF (fallback) → page-by-page text, OCR'ing pages with no text layer
    │
    ▼
3. Chunk text
    ├─ Docling → HierarchicalChunker (max 512 tokens, respects section boundaries)
    └─ PyMuPDF → RecursiveCharacterTextSplitter (2048 chars, 100 overlap)
    │
    ▼
4. Embed chunks → OpenAI text-embedding-3-small (1536 dimensions), counted against credit
    │
    ▼
5. Bulk insert DocumentChunk rows (NUL characters removed; nothing extracted → failed)
    │
    ▼
6. Populate search_vector (GIN-indexed tsvector for keyword search)
    │
    ▼
7. Update Document status → "ready"

Parsing with Docling

Docling is the primary document parser, providing structured output with:

  • Section headings and hierarchy
  • Table extraction
  • Page numbers and bounding boxes
  • Element types (paragraph, heading, list, table)
from docling.document_converter import DocumentConverter
converter = DocumentConverter()
result = converter.convert(file_path)

Docling runs in a child process, not a thread, so it can be killed. If it takes longer than DOCUMENT_DOCLING_PARSE_TIMEOUT (240 seconds) or isn't installed, the pipeline falls back to PyMuPDF for page-level text extraction. On the 1,000-paper benchmark, Docling parsed 953 papers and the fallback the other 47, with none failing.

Docling only runs OCR on a PDF that has a page without a usable text layer, such as a scan. Running it on every PDF took one paper from 22 seconds to 172 for byte-identical text. OCR languages are DOCUMENT_RAPIDOCR_LANGUAGE for Docling and DOCUMENT_OCR_LANGUAGE for the fallback's Tesseract.

Parsing can leave a worker process several GB large. Set CELERY_WORKER_MAX_MEMORY_PER_CHILD (in KB) to have Celery replace a worker process after a task once it passes that size, before the kernel kills it mid-task.

Chunking Strategy

With Docling (HierarchicalChunker):

  • Respects document structure — never splits mid-table or mid-section
  • Max 512 tokens per chunk
  • Preserves section headings and page numbers in metadata

With PyMuPDF (RecursiveCharacterTextSplitter):

  • Character-based splitting with 100-character overlap
  • Chunk size: 2048 characters
  • Page numbers tracked per chunk

Each chunk produces:

{
    "text": "chunk content...",
    "page_number": 3,
    "section_heading": "Installation Guide",
    "metadata": {"element_type": "paragraph", ...},
    "chunk_index": 0
}

Embedding

Chunks are embedded in batches of 512 using OpenAI's API:

from langchain_openai import OpenAIEmbeddings

embeddings = OpenAIEmbeddings(
    model="text-embedding-3-small",
    dimensions=1536
)
vectors = embeddings.embed_documents(texts)  # Batched

Full-text search vector

After chunk insertion, a single UPDATE precomputes the tsvector used by the keyword leg of hybrid retrieval:

DocumentChunk.objects.filter(document=document).update(
    search_vector=SearchVector("content", config=FTS_LANGUAGE)
)

Chunks are immutable after ingestion, so this one-shot population keeps the GIN-indexed column in sync without database triggers. The config comes from DOCUMENT_FTS_LANGUAGE (default simple) and must match the config used at query time — both read the same setting. Both the file and URL pipelines run this step.

URL Ingestion

Task: documents.process_url_document

@shared_task(
    bind=True,
    max_retries=3,
    soft_time_limit=120,   # 2 minutes
    time_limit=180,        # 3 minutes hard limit
)
def process_url_document(self, document_id, processing_task_id=None):

Web Pages

Uses Trafilatura for article extraction:

from trafilatura import fetch_url, extract

downloaded = fetch_url(url)
text = extract(downloaded, include_tables=True)

Falls back to HTML stripping if Trafilatura yields nothing. Text is chunked with RecursiveCharacterTextSplitter.

YouTube Videos

Extracts transcripts using youtube-transcript-api:

from youtube_transcript_api import YouTubeTranscriptApi

transcript = YouTubeTranscriptApi.get_transcript(video_id)
# Returns: [{"start": 0.0, "text": "Hello...", "duration": 2.5}, ...]

Transcripts are formatted with timestamps ([MM:SS] text) and chunked by token limit while preserving timestamp boundaries. Video title is fetched via the oEmbed endpoint (no API key required).

A video with no transcript (captions disabled, or none in a usable language) fails with a message saying so, rather than being stored with placeholder text.

Error Handling

ScenarioBehavior
Network trouble (DNS, connection resets, S3 or OpenAI connection errors, rate limits)Retried up to 3 times, 30s, 60s, then 90s apart, then marked failed
Anything else (a parse error, a document over the time limit)Marked failed at once: it would only fail again
Nothing extractedMarked failed with "no text found", not ready with no chunks
Worker lost mid-taskTask is rejected and requeued (reject_on_worker_lost)
Soft time limitSoftTimeLimitExceeded raised, cleanup runs
Hard time limitWorker kills the task after 11 minutes

A failed document's processing_error, which the API returns, is a plain message for the user: "We couldn't process this file. Try again, or upload it in another format." (or "…this link…" for web pages and videos). The "no text found" messages are kept as they are. The full error and traceback go to the server log only.

A failed document can be retried with the /reprocess/ endpoint, which checks credit and the rate limit first.

Task Tracking

Each ingestion dispatches a DocumentProcessingTask record:

class DocumentProcessingTask(BaseModel):
    document = models.ForeignKey(Document)
    task_id = models.CharField(unique=True)  # Celery task ID
    status = models.CharField()  # PENDING, STARTED, SUCCESS, FAILURE, RETRY, REVOKED
    error_message = models.TextField(blank=True)

The frontend polls /processing-status/ to track progress without querying Celery directly.

On this page