check_source_accessible narrowed its handler to ValueError, but Path.exists re-raises errno values outside its ignored set (EACCES, ENAMETOOLONG). Those were swallowed before and now escaped into the rebuild sweep the guard exists to protect. Catch OSError too. Restore the arity guard in _common_path_prefix: without it an empty list raises from min() and a single label yields a prefix covering the whole path. Two tests would have hung rather than failed on regression (the vacuum skip and the protected-wait cancellation); both are now bounded. The import vacuum test raced against the done-callback that discards the task, and now spies on the call instead, with a negative control. Replace assertions that could not fail: blank-query search against an empty corpus, a batch flush counted against an empty table, a picture description asserting its own input state, and an FS scheme check with nothing on disk to resolve. The get_model matrix asserted only the returned type across 26 cases and now pins the per-provider settings. The three batching tests now count flushes, which revealed embed-only writes through chunks_table.add rather than _flush_rebuild_batch.
855 lines
32 KiB
Python
855 lines
32 KiB
Python
import asyncio
|
|
import json
|
|
import logging
|
|
from collections.abc import AsyncGenerator
|
|
from datetime import datetime
|
|
from typing import TYPE_CHECKING
|
|
|
|
from docling_core.types.doc.document import DescriptionMetaField, PictureMeta
|
|
from lancedb.pydantic import LanceModel
|
|
|
|
from haiku.rag.client.documents import check_source_accessible
|
|
from haiku.rag.converters import get_converter
|
|
from haiku.rag.store.compression import compress_docling_split
|
|
from haiku.rag.store.engine import ChunkRecordBase
|
|
from haiku.rag.store.models.chunk import Chunk
|
|
from haiku.rag.store.models.document import Document
|
|
from haiku.rag.store.models.document_item import extract_items
|
|
from haiku.rag.store.repositories.settings import SettingsRepository
|
|
|
|
if TYPE_CHECKING:
|
|
from docling_core.types.doc.document import DoclingDocument
|
|
|
|
from haiku.rag.client import HaikuRAG, RebuildMode
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
_REBUILD_BATCH_SIZE = 50
|
|
_STAGING_TABLE_NAME = "chunks_rebuild_staging"
|
|
_STAGING_MARKER_TABLE_NAME = "chunks_rebuild_marker"
|
|
_STAGING_COPY_BATCH_SIZE = 1000
|
|
|
|
|
|
class _StagingChunkRecord(LanceModel):
|
|
"""Non-vector copy of a chunk row, used by ``_rebuild_embed_only``.
|
|
|
|
The staging table holds the original chunks' identity and content while
|
|
the live ``chunks`` table is dropped and recreated with a potentially
|
|
different vector dimension. The vector itself is omitted — re-embedding
|
|
is the whole point — and ``content_fts`` is regenerated by
|
|
``contextualize`` during phase 2.
|
|
|
|
Mirrors ``ChunkRecordBase`` minus ``content_fts`` and ``vector``. Keep
|
|
in sync: ``test_staging_chunk_record_mirrors_chunk_record_schema``
|
|
enforces parity so a new column on ``ChunkRecordBase`` can't silently
|
|
get dropped on every embed-only rebuild.
|
|
"""
|
|
|
|
id: str
|
|
document_id: str
|
|
content: str
|
|
metadata: str
|
|
order: int
|
|
|
|
|
|
class _StagingMarkerRecord(LanceModel):
|
|
"""Sentinel marking the staging table as complete.
|
|
|
|
The marker table is created only after ``_populate_staging_table`` writes
|
|
every chunk into staging. Its presence at the top of a rebuild means
|
|
phase 2 (the embed loop) was interrupted by an earlier crash, so staging
|
|
is the authoritative source for the original chunk identities and we
|
|
must resume from it instead of rerunning phase 1.
|
|
"""
|
|
|
|
id: str
|
|
|
|
|
|
async def rebuild_database(
|
|
client: "HaikuRAG", mode: "RebuildMode"
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Rebuild the database with the specified mode.
|
|
|
|
Yields the ID of each document as it is processed.
|
|
|
|
Holds the store's rebuild lock for the whole run so tag operations fail
|
|
fast instead of snapshotting a half-rebuilt database. The lock is held
|
|
across yields; an abandoned generator releases it when closed or
|
|
garbage-collected.
|
|
"""
|
|
from haiku.rag.client import RebuildMode
|
|
|
|
async with client.store._rebuild_lock:
|
|
if mode == RebuildMode.SET_EMBEDDER:
|
|
await _set_embedder(client)
|
|
return
|
|
|
|
async for doc_id in _rebuild_locked(client, mode):
|
|
yield doc_id
|
|
|
|
|
|
async def _rebuild_locked(
|
|
client: "HaikuRAG", mode: "RebuildMode"
|
|
) -> AsyncGenerator[str, None]:
|
|
from haiku.rag.client import RebuildMode
|
|
|
|
# Resolve any leftover staging/marker tables from a previously
|
|
# interrupted rebuild. Returns True only when phase 1 was already
|
|
# complete and the current mode is EMBED_ONLY, in which case we resume
|
|
# phase 2 from the existing staging table instead of recopying.
|
|
resume_from_staging = await _resolve_rebuild_recovery(client, mode)
|
|
|
|
# Wait for any already-scheduled background vacuum before the destructive
|
|
# table operations at the top of RECHUNK / FULL. Rebuild drops and
|
|
# recreates tables (and creates indices); a concurrent optimize on the
|
|
# same table fails with "CreateIndex transaction was preempted" from
|
|
# lance. Note: FULL calls create_document_from_source inside its loop,
|
|
# which may schedule *new* background vacuums — those run after the
|
|
# destructive phase and are fine.
|
|
await client._await_vacuum_tasks()
|
|
|
|
# Update settings to current config
|
|
settings_repo = SettingsRepository(client.store)
|
|
await settings_repo.save_current_settings()
|
|
|
|
# Light listing — id/uri/title/metadata only. Each rebuild function
|
|
# fetches full content (including the multi-MB docling_pages blob) one
|
|
# document at a time so a 1000-doc database doesn't pull ~15 GB of
|
|
# blobs into memory before the loop starts.
|
|
documents = await client.list_documents(include_content=False)
|
|
|
|
if mode == RebuildMode.TITLE_ONLY:
|
|
async for doc_id in _rebuild_title_only(client, documents):
|
|
yield doc_id
|
|
elif mode == RebuildMode.EMBED_ONLY:
|
|
async for doc_id in _rebuild_embed_only(
|
|
client, documents, resume_from_staging=resume_from_staging
|
|
):
|
|
yield doc_id
|
|
elif mode == RebuildMode.RECHUNK:
|
|
await client.chunk_repository.delete_all()
|
|
await client.store.recreate_embeddings_table()
|
|
async for doc_id in _rebuild_rechunk(client, documents):
|
|
yield doc_id
|
|
elif mode == RebuildMode.DESCRIPTIONS:
|
|
await client.chunk_repository.delete_all()
|
|
await client.store.recreate_embeddings_table()
|
|
async for doc_id in _rebuild_descriptions(client, documents):
|
|
yield doc_id
|
|
else: # FULL
|
|
await client.chunk_repository.delete_all()
|
|
await client.store.recreate_embeddings_table()
|
|
async for doc_id in _rebuild_full(client, documents):
|
|
yield doc_id
|
|
|
|
# Final maintenance if auto_vacuum enabled. Swallowing only so that a
|
|
# failed post-rebuild optimize doesn't mask a successful rebuild — but
|
|
# log it so the failure is visible in the output.
|
|
if client._config.storage.auto_vacuum:
|
|
try:
|
|
await client.store.vacuum()
|
|
except Exception:
|
|
logger.warning("Post-rebuild vacuum failed", exc_info=True)
|
|
|
|
|
|
async def _set_embedder(client: "HaikuRAG") -> None:
|
|
"""Adopt the current embedder identity without re-embedding.
|
|
|
|
Only valid when the vector dimension is unchanged — the stored vectors stay
|
|
usable, so just the recorded provider/name are updated. A changed dimension
|
|
requires regenerating every embedding via a full rebuild.
|
|
"""
|
|
from haiku.rag.store.repositories.settings import ConfigMismatchError
|
|
|
|
settings_repo = SettingsRepository(client.store)
|
|
stored = await settings_repo.get_current_settings()
|
|
stored_dim = stored.get("embeddings", {}).get("model", {}).get("vector_dim")
|
|
current_dim = client._config.embeddings.model.vector_dim
|
|
|
|
if stored_dim is not None and current_dim != stored_dim:
|
|
raise ConfigMismatchError(
|
|
f"Stored vector dimension {stored_dim} differs from current "
|
|
f"{current_dim}; embeddings must be regenerated. Run 'haiku-rag rebuild'."
|
|
)
|
|
|
|
await settings_repo.save_current_settings()
|
|
|
|
|
|
async def _hydrate(
|
|
client: "HaikuRAG", light_docs: list[Document]
|
|
) -> AsyncGenerator[Document, None]:
|
|
"""Yield fully-loaded documents one at a time from a light listing.
|
|
|
|
The light listing in ``rebuild_database`` skips the multi-MB
|
|
``docling_document``/``docling_pages`` blobs; this helper fetches each
|
|
full record on demand so peak memory stays at ~one document. Documents
|
|
that disappeared between listing and processing are silently skipped.
|
|
"""
|
|
for light_doc in light_docs:
|
|
assert light_doc.id is not None
|
|
doc = await client.get_document_by_id(light_doc.id)
|
|
if doc is None:
|
|
continue
|
|
assert doc.id is not None
|
|
yield doc
|
|
|
|
|
|
async def _rebuild_title_only(
|
|
client: "HaikuRAG", documents: list[Document]
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Generate titles for documents that don't have one."""
|
|
untitled = [d for d in documents if d.title is None]
|
|
async for doc in _hydrate(client, untitled):
|
|
try:
|
|
title = await client.generate_title(doc)
|
|
except Exception:
|
|
logger.warning(
|
|
"Failed to generate title for document %s", doc.id, exc_info=True
|
|
)
|
|
continue
|
|
if title is not None:
|
|
doc.title = title
|
|
await client.document_repository.update_meta(doc)
|
|
assert doc.id is not None
|
|
yield doc.id
|
|
|
|
|
|
async def _resolve_rebuild_recovery(client: "HaikuRAG", mode: "RebuildMode") -> bool:
|
|
"""Resolve any partially-completed rebuild state from a previous crash.
|
|
|
|
Returns ``True`` if ``_rebuild_embed_only`` should resume from the
|
|
existing staging table (phase 1 was already complete). In all other
|
|
cases stale recovery tables are dropped and the rebuild starts fresh.
|
|
|
|
State at entry → action
|
|
--------------------------------
|
|
no staging, no marker → return False (normal start)
|
|
staging only → drop staging (phase 1 was interrupted; ``chunks`` is intact)
|
|
marker only → drop marker (corrupted state)
|
|
staging + marker, embed → return True (resume phase 2 from staging)
|
|
staging + marker, other → drop both (staging is for embed-only; user picked a different mode)
|
|
"""
|
|
from haiku.rag.client import RebuildMode
|
|
|
|
db = client.store.db
|
|
tables = (await db.list_tables()).tables
|
|
has_staging = _STAGING_TABLE_NAME in tables
|
|
has_marker = _STAGING_MARKER_TABLE_NAME in tables
|
|
|
|
if not has_staging and not has_marker:
|
|
return False
|
|
|
|
if has_marker and not has_staging:
|
|
logger.warning(
|
|
"Found '%s' without staging table; dropping orphaned marker.",
|
|
_STAGING_MARKER_TABLE_NAME,
|
|
)
|
|
await db.drop_table(_STAGING_MARKER_TABLE_NAME)
|
|
return False
|
|
|
|
if not has_marker:
|
|
logger.warning(
|
|
"Dropping incomplete '%s' from an interrupted phase 1.",
|
|
_STAGING_TABLE_NAME,
|
|
)
|
|
await db.drop_table(_STAGING_TABLE_NAME)
|
|
return False
|
|
|
|
# has_staging and has_marker
|
|
if mode == RebuildMode.EMBED_ONLY:
|
|
logger.warning(
|
|
"Resuming interrupted embed-only rebuild: phase 2 will run from "
|
|
"existing '%s'.",
|
|
_STAGING_TABLE_NAME,
|
|
)
|
|
return True
|
|
|
|
logger.warning(
|
|
"Dropping staging tables from a prior embed-only rebuild — current "
|
|
"mode (%s) does not consume them.",
|
|
mode.name,
|
|
)
|
|
await db.drop_table(_STAGING_MARKER_TABLE_NAME)
|
|
await db.drop_table(_STAGING_TABLE_NAME)
|
|
return False
|
|
|
|
|
|
async def _populate_staging_table(client: "HaikuRAG") -> None:
|
|
"""Stream the non-vector columns of the chunks table into staging.
|
|
|
|
Uses ``to_batches`` for a single streaming read (no offset/limit
|
|
pagination drift), so peak memory stays bounded regardless of corpus
|
|
size. The vector column is omitted — the point of embed-only rebuild is
|
|
to regenerate it.
|
|
|
|
Requires ``_resolve_rebuild_recovery`` to have cleared any leftover
|
|
staging table first: ``create_table`` raises if the name is already taken.
|
|
"""
|
|
db = client.store.db
|
|
tables = (await db.list_tables()).tables
|
|
|
|
staging = await db.create_table(_STAGING_TABLE_NAME, schema=_StagingChunkRecord)
|
|
if "chunks" not in tables:
|
|
return
|
|
|
|
stream = (
|
|
await client.store.chunks_table.query()
|
|
.select(["id", "document_id", "content", "metadata", "order"])
|
|
.to_batches(max_batch_length=_STAGING_COPY_BATCH_SIZE)
|
|
)
|
|
async for batch in stream:
|
|
rows = batch.to_pylist()
|
|
records = [
|
|
_StagingChunkRecord(
|
|
id=r["id"],
|
|
document_id=r["document_id"],
|
|
content=r["content"],
|
|
metadata=r["metadata"],
|
|
order=r["order"],
|
|
)
|
|
for r in rows
|
|
]
|
|
await staging.add(records)
|
|
|
|
|
|
async def _mark_phase1_complete(client: "HaikuRAG") -> None:
|
|
"""Create the marker table that designates staging as authoritative.
|
|
|
|
Called after ``_populate_staging_table`` finishes. On crash recovery the
|
|
marker's presence flips ``_rebuild_embed_only`` into resume mode.
|
|
"""
|
|
db = client.store.db
|
|
if _STAGING_MARKER_TABLE_NAME in (await db.list_tables()).tables:
|
|
return
|
|
marker = await db.create_table(
|
|
_STAGING_MARKER_TABLE_NAME, schema=_StagingMarkerRecord
|
|
)
|
|
await marker.add([_StagingMarkerRecord(id="phase1_complete")])
|
|
|
|
|
|
async def _drop_staging_tables(client: "HaikuRAG") -> None:
|
|
"""Drop the marker first, then the staging table.
|
|
|
|
Ordering matters: if a crash interrupts cleanup between the two drops,
|
|
the next rebuild sees ``staging`` without ``marker`` and treats it as a
|
|
partial phase 1 → drops staging harmlessly. The reverse order would
|
|
leak a marker pointing at nothing.
|
|
"""
|
|
db = client.store.db
|
|
tables = (await db.list_tables()).tables
|
|
if _STAGING_MARKER_TABLE_NAME in tables:
|
|
await db.drop_table(_STAGING_MARKER_TABLE_NAME)
|
|
if _STAGING_TABLE_NAME in tables:
|
|
await db.drop_table(_STAGING_TABLE_NAME)
|
|
|
|
|
|
async def _read_chunks_from_staging(staging_table, document_id: str) -> list[Chunk]:
|
|
"""Read chunks for one document from the staging table.
|
|
|
|
Only non-vector columns are selected: the staging table may have a
|
|
different vector dimension than the new chunks table (during a dim
|
|
migration), and we re-embed anyway.
|
|
"""
|
|
rows = (
|
|
await staging_table.query()
|
|
.where(f"document_id = '{document_id}'")
|
|
.select(["id", "document_id", "content", "metadata", "order"])
|
|
.to_arrow()
|
|
).to_pylist()
|
|
chunks: list[Chunk] = []
|
|
for row in rows:
|
|
chunks.append(
|
|
Chunk(
|
|
id=row["id"],
|
|
document_id=row["document_id"],
|
|
content=row["content"],
|
|
metadata=json.loads(row["metadata"]),
|
|
order=row["order"],
|
|
)
|
|
)
|
|
chunks.sort(key=lambda c: c.order)
|
|
return chunks
|
|
|
|
|
|
async def _rebuild_embed_only(
|
|
client: "HaikuRAG",
|
|
documents: list[Document],
|
|
*,
|
|
resume_from_staging: bool = False,
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Re-embed all chunks without changing chunk boundaries.
|
|
|
|
Two-phase pattern that keeps peak memory bounded regardless of corpus
|
|
size and is idempotent across crashes:
|
|
|
|
1. Stream-copy the chunks table's non-vector columns into a staging
|
|
table, then write a marker row that designates staging as complete.
|
|
LanceDB OSS does not support ``rename_table``, so the staging copy
|
|
is the only safe way to preserve chunk identity while the live
|
|
``chunks`` table is dropped and recreated.
|
|
2. Drop-and-recreate ``chunks`` with the current schema, then stream
|
|
from staging one document at a time, re-embed in batches of
|
|
``embeddings.batch_size``, and flush to the new chunks table every
|
|
``_REBUILD_BATCH_SIZE`` documents.
|
|
|
|
Cleanup runs only on success: a crash anywhere in phase 2 leaves both
|
|
staging and marker in place so the next rebuild can re-enter phase 2
|
|
via ``resume_from_staging=True``. The order of the success cleanup —
|
|
drop marker before staging — keeps an interruption between the two
|
|
drops recoverable: the next rebuild sees staging without marker and
|
|
treats it as a partial phase 1, which is harmless because phase 2 has
|
|
already finished writing the new chunks table.
|
|
"""
|
|
from haiku.rag.embeddings import contextualize, embed_chunks
|
|
|
|
db = client.store.db
|
|
embedder = client.chunk_repository.embedder
|
|
|
|
if not resume_from_staging:
|
|
# Phase 1: copy chunks into staging, then mark it complete. After the
|
|
# marker exists, a crash will resume phase 2 from staging.
|
|
await _populate_staging_table(client)
|
|
await _mark_phase1_complete(client)
|
|
|
|
# Recreate the chunks table fresh (idempotent; handles vector-dim
|
|
# changes and discards any partial new chunks from a prior crashed
|
|
# phase 2).
|
|
await client.store.recreate_embeddings_table()
|
|
|
|
staging_table = await db.open_table(_STAGING_TABLE_NAME)
|
|
|
|
pending_records: list[ChunkRecordBase] = []
|
|
yielded_docs: set[str] = set()
|
|
|
|
for doc in documents:
|
|
assert doc.id is not None
|
|
chunks = await _read_chunks_from_staging(staging_table, doc.id)
|
|
if not chunks:
|
|
continue
|
|
|
|
# Re-attach picture bytes stripped by the staging copy so picture
|
|
# chunks route through embed_image rather than being text-embedded.
|
|
# Bytes live in document_items (embed-only never touches that table).
|
|
if embedder.supports_images:
|
|
picture_data = await client.document_item_repository.get_all_picture_data(
|
|
doc.id
|
|
)
|
|
for chunk in chunks:
|
|
if "picture" not in (chunk.metadata.get("labels") or []):
|
|
continue
|
|
refs = chunk.metadata.get("doc_item_refs") or []
|
|
data = next((picture_data[r] for r in refs if r in picture_data), None)
|
|
if data is not None:
|
|
chunk._picture_data = data
|
|
else:
|
|
logger.warning(
|
|
"Document %s picture chunk %s has no recoverable bytes; "
|
|
"embedding its caption as text.",
|
|
doc.id,
|
|
chunk.id,
|
|
)
|
|
|
|
content_fts_list = contextualize(chunks)
|
|
embedded_chunks = await embed_chunks(chunks, embedder, client._config)
|
|
|
|
for chunk, content_fts, embedded in zip(
|
|
chunks, content_fts_list, embedded_chunks
|
|
):
|
|
assert chunk.id is not None
|
|
assert chunk.document_id is not None
|
|
assert embedded.embedding is not None
|
|
pending_records.append(
|
|
client.store.ChunkRecord(
|
|
id=chunk.id,
|
|
document_id=chunk.document_id,
|
|
content=chunk.content,
|
|
content_fts=content_fts,
|
|
metadata=json.dumps(chunk.metadata),
|
|
order=chunk.order,
|
|
vector=embedded.embedding,
|
|
)
|
|
)
|
|
|
|
yielded_docs.add(doc.id)
|
|
# Yield per-doc for progress reporting; the actual write batches up
|
|
# to _REBUILD_BATCH_SIZE docs. If the process is interrupted between
|
|
# yield and the next flush, the next rebuild resumes phase 2 from
|
|
# the staging table and redoes the batch (see _rebuild_rechunk for
|
|
# the original comment on the yield/flush gap).
|
|
yield doc.id
|
|
|
|
if len(yielded_docs) % _REBUILD_BATCH_SIZE == 0 and pending_records:
|
|
await client.store.chunks_table.add(pending_records)
|
|
pending_records = []
|
|
|
|
if pending_records:
|
|
await client.store.chunks_table.add(pending_records)
|
|
|
|
# Phase 2 finished. Drop the recovery state — marker first so a crash
|
|
# between the two drops leaves only staging behind, which the next
|
|
# rebuild discards harmlessly.
|
|
await _drop_staging_tables(client)
|
|
|
|
# Yield docs with no chunks
|
|
for doc in documents:
|
|
if doc.id and doc.id not in yielded_docs:
|
|
yield doc.id
|
|
|
|
|
|
async def _flush_rebuild_batch(
|
|
client: "HaikuRAG", documents: list[Document], chunks: list[Chunk]
|
|
) -> None:
|
|
"""Batch write documents and chunks during rebuild.
|
|
|
|
Performs two writes: one for all document updates (via merge_insert), one
|
|
for all chunks. Also repopulates document items from the stored docling
|
|
document. Used by RECHUNK and FULL modes after the chunks table has been
|
|
cleared.
|
|
"""
|
|
from haiku.rag.store.engine import DocumentMetaRecord, DocumentRecord
|
|
|
|
if not documents:
|
|
return
|
|
|
|
now = datetime.now().isoformat()
|
|
|
|
# Batch update documents and document_meta using merge_insert (one LanceDB
|
|
# version per table). Content+blobs go to documents; mutable attributes go
|
|
# to document_meta.
|
|
doc_records = []
|
|
meta_records = []
|
|
for doc in documents:
|
|
assert doc.id is not None
|
|
doc_records.append(
|
|
DocumentRecord(
|
|
id=doc.id,
|
|
content=doc.content,
|
|
docling_document=doc.docling_document,
|
|
docling_pages=doc.docling_pages,
|
|
docling_version=doc.docling_version,
|
|
)
|
|
)
|
|
meta_records.append(
|
|
DocumentMetaRecord(
|
|
id=doc.id,
|
|
uri=doc.uri,
|
|
title=doc.title,
|
|
metadata=json.dumps(doc.metadata),
|
|
created_at=doc.created_at.isoformat() if doc.created_at else now,
|
|
updated_at=now,
|
|
)
|
|
)
|
|
|
|
await (
|
|
client.store.documents_table.merge_insert("id")
|
|
.when_matched_update_all()
|
|
.execute(doc_records)
|
|
)
|
|
await (
|
|
client.store.document_meta_table.merge_insert("id")
|
|
.when_matched_update_all()
|
|
.when_not_matched_insert_all()
|
|
.execute(meta_records)
|
|
)
|
|
|
|
# Batch create all chunks (single LanceDB version)
|
|
if chunks:
|
|
await client.chunk_repository.create(chunks)
|
|
|
|
# Repopulate document items from stored docling data. The stored docling
|
|
# blob has had its picture URIs stripped (compress_docling_split), so
|
|
# re-extracting from it would lose picture_data — snapshot the existing
|
|
# bytes per document and merge them back.
|
|
for doc in documents:
|
|
assert doc.id is not None
|
|
docling_doc = doc.get_docling_document()
|
|
if docling_doc is not None:
|
|
existing_picture_data = (
|
|
await client.document_item_repository.get_all_picture_data(doc.id)
|
|
)
|
|
await client.document_item_repository.delete_by_document_id(doc.id)
|
|
items = extract_items(
|
|
doc.id,
|
|
docling_doc,
|
|
existing_picture_data=existing_picture_data,
|
|
)
|
|
await client.document_item_repository.create_items(doc.id, items)
|
|
|
|
|
|
async def _rebuild_rechunk(
|
|
client: "HaikuRAG", documents: list[Document]
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Re-chunk and re-embed each document from its stored docling blob."""
|
|
from haiku.rag.embeddings import embed_chunks
|
|
|
|
pending_chunks: list[Chunk] = []
|
|
pending_docs: list[Document] = []
|
|
embedder = client.embedder
|
|
|
|
async for doc in _hydrate(client, documents):
|
|
assert doc.id is not None
|
|
docling_document = doc.get_docling_document()
|
|
if docling_document is None:
|
|
raise ValueError(
|
|
f"Document {doc.id} has no stored docling document; rechunk "
|
|
"requires it. Run a full rebuild (without --rechunk) instead."
|
|
)
|
|
|
|
# Stored blob has stripped picture URIs; pass the snapshot so
|
|
# build_picture_chunks (inside chunk()) can recover the bytes.
|
|
existing_picture_data = (
|
|
await client.document_item_repository.get_all_picture_data(doc.id)
|
|
if embedder.supports_images
|
|
else None
|
|
)
|
|
chunks = await client.chunk(
|
|
docling_document,
|
|
existing_picture_data=existing_picture_data,
|
|
document_id=doc.id,
|
|
)
|
|
embedded_chunks = await embed_chunks(chunks, embedder, client._config)
|
|
|
|
# Prepare chunks with document_id and order
|
|
for order, chunk in enumerate(embedded_chunks):
|
|
chunk.document_id = doc.id
|
|
chunk.order = order
|
|
|
|
pending_chunks.extend(embedded_chunks)
|
|
pending_docs.append(doc)
|
|
# Yield per-doc so progress reporting moves immediately. The actual
|
|
# write batches up to _REBUILD_BATCH_SIZE for throughput; if the
|
|
# process is interrupted between yield and flush, up to one
|
|
# batch's worth of trailing yields aren't persisted, which is
|
|
# consistent with the rebuild already being non-atomic.
|
|
yield doc.id
|
|
|
|
# Flush batch when size reached
|
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
pending_chunks = []
|
|
pending_docs = []
|
|
|
|
# Flush remaining
|
|
if pending_docs:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
|
|
|
|
def _apply_descriptions_sync(
|
|
docling_doc: "DoclingDocument", doc: Document, descriptions: dict[str, str]
|
|
) -> int:
|
|
"""Patch picture descriptions into the docling document and re-compress.
|
|
|
|
Updates only docling_document — set_docling would also overwrite
|
|
docling_pages by routing through compress_docling_split, which
|
|
extracts pages from the in-memory JSON and finds none (the pages
|
|
blob is stored separately and is not loaded by get_docling_document).
|
|
That would silently destroy page rasters for every doc with at
|
|
least one undescribed picture.
|
|
"""
|
|
for pic in docling_doc.pictures:
|
|
text = descriptions.get(pic.self_ref)
|
|
if not text:
|
|
continue
|
|
if pic.meta is None:
|
|
pic.meta = PictureMeta()
|
|
pic.meta.description = DescriptionMetaField(text=text)
|
|
|
|
structure_bytes, _ = compress_docling_split(docling_doc.model_dump(mode="json"))
|
|
doc.docling_document = structure_bytes
|
|
doc.docling_version = docling_doc.version
|
|
return len(descriptions)
|
|
|
|
|
|
async def _patch_picture_descriptions(client: "HaikuRAG", doc: Document) -> int:
|
|
"""Run the VLM against pictures lacking a description, patch the docling
|
|
blob in-place. Returns the number of newly described pictures.
|
|
Pictures that already carry ``meta.description.text`` are skipped, so the
|
|
operation is safe to re-run after a partial failure.
|
|
"""
|
|
from haiku.rag.providers.picture_description import describe_pictures
|
|
|
|
assert doc.id is not None
|
|
docling_doc = doc.get_docling_document()
|
|
if docling_doc is None or not docling_doc.pictures:
|
|
return 0
|
|
|
|
needs_description: list[str] = []
|
|
for pic in docling_doc.pictures:
|
|
existing = (
|
|
pic.meta.description.text if pic.meta and pic.meta.description else None
|
|
)
|
|
if not (existing and existing.strip()):
|
|
needs_description.append(pic.self_ref)
|
|
|
|
if not needs_description:
|
|
return 0
|
|
|
|
bytes_by_ref = await client.document_item_repository.get_pictures_for_chunk(
|
|
doc.id, needs_description
|
|
)
|
|
if not bytes_by_ref:
|
|
logger.warning(
|
|
"Document %s has %d pictures missing descriptions but no stored "
|
|
"picture bytes — skipping. Run a full rebuild from source to "
|
|
"recover the bytes.",
|
|
doc.id,
|
|
len(needs_description),
|
|
)
|
|
return 0
|
|
|
|
descriptions = await describe_pictures(bytes_by_ref, config=client._config)
|
|
|
|
if not descriptions:
|
|
return 0
|
|
|
|
return await asyncio.to_thread(
|
|
_apply_descriptions_sync, docling_doc, doc, descriptions
|
|
)
|
|
|
|
|
|
async def _rebuild_descriptions(
|
|
client: "HaikuRAG", documents: list[Document]
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Run the VLM over already-stored picture bytes, patch descriptions into
|
|
the docling blob, then re-chunk + re-embed.
|
|
|
|
Skips the docling parse entirely (the blob is already there); only the VLM
|
|
cost remains. Idempotent: pictures whose ``meta.description.text`` is
|
|
already populated are not re-described.
|
|
"""
|
|
from haiku.rag.embeddings import embed_chunks
|
|
|
|
if client._config.processing.pictures != "description":
|
|
raise ValueError(
|
|
"rebuild --descriptions requires processing.pictures = 'description' "
|
|
"in your config."
|
|
)
|
|
|
|
pending_chunks: list[Chunk] = []
|
|
pending_docs: list[Document] = []
|
|
embedder = client.embedder
|
|
|
|
described_total = 0
|
|
async for doc in _hydrate(client, documents):
|
|
assert doc.id is not None
|
|
docling_document = doc.get_docling_document()
|
|
if docling_document is None:
|
|
raise ValueError(
|
|
f"Document {doc.id} has no stored docling document; "
|
|
"rebuild --descriptions requires it. Run a full rebuild instead."
|
|
)
|
|
|
|
n = await _patch_picture_descriptions(client, doc)
|
|
described_total += n
|
|
# Use the (possibly patched) docling document for chunking.
|
|
docling_document = doc.get_docling_document()
|
|
assert docling_document is not None
|
|
|
|
existing_picture_data = (
|
|
await client.document_item_repository.get_all_picture_data(doc.id)
|
|
if embedder.supports_images
|
|
else None
|
|
)
|
|
chunks = await client.chunk(
|
|
docling_document,
|
|
existing_picture_data=existing_picture_data,
|
|
document_id=doc.id,
|
|
)
|
|
embedded_chunks = await embed_chunks(chunks, embedder, client._config)
|
|
|
|
for order, chunk in enumerate(embedded_chunks):
|
|
chunk.document_id = doc.id
|
|
chunk.order = order
|
|
|
|
pending_chunks.extend(embedded_chunks)
|
|
pending_docs.append(doc)
|
|
yield doc.id
|
|
|
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
pending_chunks = []
|
|
pending_docs = []
|
|
|
|
if pending_docs:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
|
|
logger.info(
|
|
"rebuild --descriptions: %d new picture descriptions added across %d documents",
|
|
described_total,
|
|
len(documents),
|
|
)
|
|
|
|
|
|
async def _rebuild_full(
|
|
client: "HaikuRAG", documents: list[Document]
|
|
) -> AsyncGenerator[str, None]:
|
|
"""Full rebuild: re-convert from source, re-chunk, re-embed."""
|
|
from haiku.rag.embeddings import embed_chunks
|
|
|
|
pending_chunks: list[Chunk] = []
|
|
pending_docs: list[Document] = []
|
|
converter = get_converter(client._config)
|
|
embedder = client.embedder
|
|
|
|
for light_doc in documents:
|
|
assert light_doc.id is not None
|
|
|
|
# Try to rebuild from source if available — uses the light listing
|
|
# directly, no need to load the stored content/blobs first.
|
|
if light_doc.uri and check_source_accessible(light_doc.uri):
|
|
try:
|
|
# Flush pending batch before source rebuild (creates new doc)
|
|
if pending_docs:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
pending_chunks = []
|
|
pending_docs = []
|
|
|
|
await client.delete_document(light_doc.id)
|
|
new_doc = await client.create_document_from_source(
|
|
source=light_doc.uri, metadata=light_doc.metadata or {}
|
|
)
|
|
assert isinstance(new_doc, Document)
|
|
assert new_doc.id is not None
|
|
yield new_doc.id
|
|
continue
|
|
except Exception as e:
|
|
logger.error(
|
|
"Error recreating document from source %s: %s",
|
|
light_doc.uri,
|
|
e,
|
|
)
|
|
continue
|
|
|
|
# Fallback: rebuild from stored content. Now we need the full
|
|
# record (content + docling_pages for the round-trip write).
|
|
doc = await client.get_document_by_id(light_doc.id)
|
|
if doc is None:
|
|
continue
|
|
assert doc.id is not None
|
|
if doc.uri:
|
|
logger.warning("Source missing for %s, re-embedding from content", doc.uri)
|
|
|
|
docling_document = await converter.convert_text(doc.content, format="md")
|
|
chunks = await client.chunk(docling_document)
|
|
embedded_chunks = await embed_chunks(chunks, embedder, client._config)
|
|
|
|
doc.set_docling(docling_document)
|
|
|
|
# Prepare chunks with document_id and order
|
|
for order, chunk in enumerate(embedded_chunks):
|
|
chunk.document_id = doc.id
|
|
chunk.order = order
|
|
|
|
pending_chunks.extend(embedded_chunks)
|
|
pending_docs.append(doc)
|
|
yield doc.id
|
|
|
|
# Flush batch when size reached
|
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
|
pending_chunks = []
|
|
pending_docs = []
|
|
|
|
# Flush remaining
|
|
if pending_docs:
|
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|