From 2b2b475279c5c08fad9a5f059bfa2edba86ab82b Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Wed, 24 Jun 2026 10:40:09 +0300 Subject: [PATCH] Collapse docling compression to a single function --- CHANGELOG.md | 1 + haiku_rag_slim/haiku/rag/client/rebuild.py | 2 +- haiku_rag_slim/haiku/rag/store/compression.py | 22 ++---------- .../haiku/rag/store/models/document.py | 4 +-- .../haiku/rag/store/upgrades/v0_38_0.py | 2 +- tests/store/test_compression.py | 9 ++--- tests/store/test_document_items.py | 5 ++- tests/store/test_v0_40_0_migration.py | 2 +- tests/store/test_v0_48_0_migration.py | 4 +-- tests/test_document.py | 35 +++++++------------ 10 files changed, 29 insertions(+), 57 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a9e5f3ed..9612b7fb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ ### Fixed - zstd compression uses a fresh `zstandard` compressor/decompressor per call instead of shared module-level singletons, fixing process segfaults when ingester workers compress documents concurrently (Python < 3.14). +- Ingestion compresses the DoclingDocument structure directly from the in-memory dict, removing one full-size serialized copy from per-document peak memory. ## [0.61.0] - 2026-06-23 diff --git a/haiku_rag_slim/haiku/rag/client/rebuild.py b/haiku_rag_slim/haiku/rag/client/rebuild.py index f7702f45..0af3cd83 100644 --- a/haiku_rag_slim/haiku/rag/client/rebuild.py +++ b/haiku_rag_slim/haiku/rag/client/rebuild.py @@ -643,7 +643,7 @@ def _apply_descriptions_sync( pic.meta = PictureMeta() pic.meta.description = DescriptionMetaField(text=text) - structure_bytes, _ = compress_docling_split(docling_doc.model_dump_json()) + 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) diff --git a/haiku_rag_slim/haiku/rag/store/compression.py b/haiku_rag_slim/haiku/rag/store/compression.py index 97edf1e9..989f1fbd 100644 --- a/haiku_rag_slim/haiku/rag/store/compression.py +++ b/haiku_rag_slim/haiku/rag/store/compression.py @@ -34,20 +34,7 @@ def decompress_json(data: bytes) -> str: return _zstd_decompress(data).decode("utf-8") -def compress_docling_split(json_str: str) -> tuple[bytes, bytes | None]: - """Parse a DoclingDocument JSON string and compress it. - - Thin wrapper over :func:`compress_docling_data` for callers that only hold - the serialized string — store migrations and rebuild-from-blob, neither of - which is speed-sensitive. The ingestion hot path should call - ``compress_docling_data`` with ``DoclingDocument.model_dump(mode="json")`` - instead, to avoid serializing the document to a full JSON string only to - parse it straight back into a dict. - """ - return compress_docling_data(json.loads(json_str)) - - -def compress_docling_data(data: dict) -> tuple[bytes, bytes | None]: +def compress_docling_split(data: dict) -> tuple[bytes, bytes | None]: """Split a DoclingDocument dict into structure and pages, compress both with zstd. Picture image URIs are stripped from the structure blob — they are stored on @@ -70,10 +57,7 @@ def compress_docling_data(data: dict) -> tuple[bytes, bytes | None]: if isinstance(picture, dict): picture["image"] = None - structure_bytes = _zstd_compress(json.dumps(data).encode("utf-8")) - - pages_bytes = None - if pages: - pages_bytes = _zstd_compress(json.dumps(pages).encode("utf-8")) + structure_bytes = compress_json(json.dumps(data)) + pages_bytes = compress_json(json.dumps(pages)) if pages else None return structure_bytes, pages_bytes diff --git a/haiku_rag_slim/haiku/rag/store/models/document.py b/haiku_rag_slim/haiku/rag/store/models/document.py index 48dbf432..f045b939 100644 --- a/haiku_rag_slim/haiku/rag/store/models/document.py +++ b/haiku_rag_slim/haiku/rag/store/models/document.py @@ -4,7 +4,7 @@ from typing import TYPE_CHECKING from pydantic import BaseModel, Field -from haiku.rag.store.compression import compress_docling_data, decompress_json +from haiku.rag.store.compression import compress_docling_split, decompress_json if TYPE_CHECKING: from docling_core.types.doc.document import DoclingDocument, PageItem @@ -32,7 +32,7 @@ class Document(BaseModel): Sets docling_document (zstd-compressed structure without pages), docling_pages (zstd-compressed page images), and docling_version. """ - structure, pages = compress_docling_data(docling_doc.model_dump(mode="json")) + structure, pages = compress_docling_split(docling_doc.model_dump(mode="json")) self.docling_document = structure self.docling_pages = pages self.docling_version = docling_doc.version diff --git a/haiku_rag_slim/haiku/rag/store/upgrades/v0_38_0.py b/haiku_rag_slim/haiku/rag/store/upgrades/v0_38_0.py index 7cc2e2ff..4fe74222 100644 --- a/haiku_rag_slim/haiku/rag/store/upgrades/v0_38_0.py +++ b/haiku_rag_slim/haiku/rag/store/upgrades/v0_38_0.py @@ -58,7 +58,7 @@ async def _apply_split_pages_zstd(store: Store) -> None: # pragma: no cover json_str = docling_blob.decode("utf-8") # Split structure and pages, re-compress with zstd - structure_bytes, pages_bytes = compress_docling_split(json_str) + structure_bytes, pages_bytes = compress_docling_split(json.loads(json_str)) metadata_raw = row.get("metadata") metadata_str = ( diff --git a/tests/store/test_compression.py b/tests/store/test_compression.py index 45d3f4b2..a50fd0c2 100644 --- a/tests/store/test_compression.py +++ b/tests/store/test_compression.py @@ -54,8 +54,7 @@ class TestDoclingCompressionSplit: "texts": [{"text": "hello"}], "pages": {"1": {"image": "base64data"}, "2": {"image": "more"}}, } - json_str = json.dumps(data) - structure_bytes, pages_bytes = compress_docling_split(json_str) + structure_bytes, pages_bytes = compress_docling_split(data) assert structure_bytes is not None assert pages_bytes is not None @@ -73,8 +72,7 @@ class TestDoclingCompressionSplit: def test_split_without_pages(self): data = {"name": "test_doc", "texts": []} - json_str = json.dumps(data) - structure_bytes, pages_bytes = compress_docling_split(json_str) + structure_bytes, pages_bytes = compress_docling_split(data) assert structure_bytes is not None assert pages_bytes is None @@ -84,8 +82,7 @@ class TestDoclingCompressionSplit: def test_split_with_empty_pages(self): data = {"name": "test_doc", "texts": [], "pages": {}} - json_str = json.dumps(data) - structure_bytes, pages_bytes = compress_docling_split(json_str) + structure_bytes, pages_bytes = compress_docling_split(data) assert structure_bytes is not None assert pages_bytes is None diff --git a/tests/store/test_document_items.py b/tests/store/test_document_items.py index 45a7c6a4..37700f01 100644 --- a/tests/store/test_document_items.py +++ b/tests/store/test_document_items.py @@ -397,8 +397,7 @@ class TestDocumentItemMigration: from haiku.rag.store.upgrades.v0_40_0 import _apply_populate_document_items docling_doc = _make_docling_doc() - json_str = docling_doc.model_dump_json() - structure, pages = compress_docling_split(json_str) + structure, pages = compress_docling_split(docling_doc.model_dump(mode="json")) # Create a database at a pre-migration version with a document async with Store(temp_db_path, create=True, skip_migration_check=True) as store: @@ -687,7 +686,7 @@ class TestCompressDoclingSplitStripsPictureUris: "pages": {}, } - structure_bytes, pages_bytes = compress_docling_split(json.dumps(doc_json)) + structure_bytes, pages_bytes = compress_docling_split(doc_json) decoded = json.loads(decompress_json(structure_bytes)) for pic in decoded["pictures"]: diff --git a/tests/store/test_v0_40_0_migration.py b/tests/store/test_v0_40_0_migration.py index e11f630e..c4d3274b 100644 --- a/tests/store/test_v0_40_0_migration.py +++ b/tests/store/test_v0_40_0_migration.py @@ -35,7 +35,7 @@ async def test_populate_handles_extra_columns_on_items_table(temp_db_path): accepted even though it only writes the original 6 columns. """ docling_doc = _simple_docling_doc() - structure, pages = compress_docling_split(docling_doc.model_dump_json()) + structure, pages = compress_docling_split(docling_doc.model_dump(mode="json")) async with Store(temp_db_path, create=True, skip_migration_check=True) as store: # _init_tables already created document_items with the latest schema. diff --git a/tests/store/test_v0_48_0_migration.py b/tests/store/test_v0_48_0_migration.py index 94cdd142..4e73dbf1 100644 --- a/tests/store/test_v0_48_0_migration.py +++ b/tests/store/test_v0_48_0_migration.py @@ -26,7 +26,7 @@ class TestV0_48_0Migration: async def test_backfill_populates_levels(self, temp_db_path): docling_doc = _docling_with_levels() - structure, pages = compress_docling_split(docling_doc.model_dump_json()) + structure, pages = compress_docling_split(docling_doc.model_dump(mode="json")) async with Store(temp_db_path, create=True, skip_migration_check=True) as store: await store.set_haiku_version("0.45.0") @@ -79,7 +79,7 @@ class TestV0_48_0Migration: async def test_backfill_idempotent(self, temp_db_path): docling_doc = _docling_with_levels() - structure, pages = compress_docling_split(docling_doc.model_dump_json()) + structure, pages = compress_docling_split(docling_doc.model_dump(mode="json")) async with Store(temp_db_path, create=True, skip_migration_check=True) as store: await store.set_haiku_version("0.45.0") diff --git a/tests/test_document.py b/tests/test_document.py index 317aeb46..f850915c 100644 --- a/tests/test_document.py +++ b/tests/test_document.py @@ -281,14 +281,11 @@ def test_set_docling_with_page_images(): assert "1" in pages -def test_set_docling_dict_path_matches_string_path(): - """The ingestion dict path must produce the same stored bytes as the - legacy string path. - - ``set_docling`` feeds ``compress_docling_data`` the dict from - ``model_dump(mode="json")`` instead of round-tripping through - ``model_dump_json()`` + ``json.loads``. The on-disk format must not change - — in particular int-keyed ``pages`` must still serialize to string keys. +def test_compress_docling_split_dict_sources_match(): + """Both ways of building the compression input dict must yield identical + stored bytes: ``model_dump(mode="json")`` (used at ingest) and + ``json.loads(model_dump_json())`` (used when migrating a stored string). + In particular int-keyed ``pages`` must serialize to string keys either way. """ import json @@ -296,28 +293,22 @@ def test_set_docling_dict_path_matches_string_path(): from docling_core.types.doc.document import DoclingDocument, PageItem from docling_core.types.doc.labels import DocItemLabel - from haiku.rag.store.compression import ( - compress_docling_data, - compress_docling_split, - decompress_json, - ) + from haiku.rag.store.compression import compress_docling_split docling_doc = DoclingDocument(name="equivalence_test") docling_doc.add_text(label=DocItemLabel.PARAGRAPH, text="Hello world") docling_doc.pages[1] = PageItem(size=Size(width=612, height=792), page_no=1) - struct_dict, pages_dict = compress_docling_data( + struct_dump, pages_dump = compress_docling_split( docling_doc.model_dump(mode="json") ) - struct_str, pages_str = compress_docling_split(docling_doc.model_dump_json()) + struct_str, pages_str = compress_docling_split( + json.loads(docling_doc.model_dump_json()) + ) - assert json.loads(decompress_json(struct_dict)) == json.loads( - decompress_json(struct_str) - ) - assert pages_dict is not None and pages_str is not None - assert json.loads(decompress_json(pages_dict)) == json.loads( - decompress_json(pages_str) - ) + assert struct_dump == struct_str + assert pages_dump is not None and pages_str is not None + assert pages_dump == pages_str def test_get_page_images():