Coverage for src/lilbee/data/types.py: 100%
129 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""Shared ingest types and constants."""
3from __future__ import annotations
5import hashlib
6from dataclasses import dataclass
7from enum import StrEnum
8from pathlib import Path
9from typing import NamedTuple, NotRequired, TypedDict
11from pydantic import BaseModel
12from rich.highlighter import ReprHighlighter
13from rich.text import Text
15from lilbee.core.vectors import Vector
16from lilbee.data.store import (
17 ChunkType,
18 ConceptRecords,
19 IndexMismatch,
20 PageTextRecord,
21 SourceMeta,
22 SourceStat,
23 SourceStatBackfill,
24)
25from lilbee.runtime.progress import OcrBackendUsed
27# PDF and image content types route to paginated extraction; every other format
28# routes to markdown extraction. content_type is derived per-file in
29# discovery.classify_file (PDFs and images grouped; others keyed by extension).
30PDF_CONTENT_TYPE = "pdf"
31IMAGE_CONTENT_TYPE = "image"
32MARKDOWN_OUTPUT = "markdown"
33MARKDOWN_MIME = "text/markdown"
34# Sync summary note for a skipped document whose extraction ran with OCR off.
35SKIPPED_OCR_OFF_NOTE = ": OCR is off (enable_ocr = false)"
38@dataclass(frozen=True)
39class ShardId:
40 """Which slice of the corpus one ingest worker owns.
42 *records_root* is the data root holding the corpus's skip records, which
43 every worker reads and writes in place of its own private data root.
44 """
46 index: int
47 count: int
48 records_root: Path
50 def owns(self, key: str) -> bool:
51 """Whether source *key* belongs to this slice.
53 Hashed with blake2b, not ``hash()``, which is salted per process: two
54 runs would deal the same corpus differently and every resume would
55 re-embed what a sibling already holds.
56 """
57 digest = hashlib.blake2b(key.encode("utf-8"), digest_size=8).digest()
58 return int.from_bytes(digest, "big") % self.count == self.index
61class FileToProcess(NamedTuple):
62 """A file queued for ingestion with its metadata."""
64 name: str
65 path: Path
66 content_type: str
67 file_hash: str
68 needs_cleanup: bool
69 stat: SourceStat | None = None
72class MemberRecords(NamedTuple):
73 """One archive member's records, written as its own source."""
75 name: str
76 content_type: str
77 records: list[ChunkRecord]
78 page_texts: list[PageTextRecord]
79 meta: SourceMeta
82class DocumentRecords(NamedTuple):
83 """One file's records, its source metadata, and the OCR its extraction ran with."""
85 records: list[ChunkRecord]
86 meta: SourceMeta
87 ocr: OcrReport | None = None
90class SkippedSource(BaseModel):
91 """One file a skip marker holds out of the index, and why."""
93 filename: str
94 reason: str
97class FileChangePlan(NamedTuple):
98 """Outcome of diffing disk files against the tracked sources."""
100 files_to_process: list[FileToProcess]
101 added: dict[str, None]
102 updated: dict[str, None]
103 unchanged: int
104 stat_backfills: list[SourceStatBackfill]
105 # Files a skip marker holds out; not in the index, so never counted as unchanged.
106 held_out: list[str]
109class OcrBackendName(StrEnum):
110 """OCR backends lilbee selects in OcrConfig: xberg's tesseract or lilbee's vision plugin."""
112 TESSERACT = "tesseract"
113 LILBEE_VISION = "lilbee-vision"
116class OcrReport(BaseModel, frozen=True):
117 """Which OCR backend one extraction ran and how many pages it OCR'd."""
119 backend: OcrBackendUsed
120 pages: int = 0
123class EmbeddingBackendName(StrEnum):
124 """Embedding backends registered with xberg. lilbee registers its own embedder
125 as a plugin so the semantic chunker detects boundaries with the same model that
126 vectorizes chunks."""
128 LILBEE = "lilbee"
131class TokenizerBackendName(StrEnum):
132 """Tokenizer backends registered with xberg. lilbee registers its embedder's
133 tokenizer so ChunkSizing counts chunk budgets in the same tokens the embedder
134 consumes, instead of a chars-per-token heuristic. Separate registry from the
135 embedding backend, so sharing the ``lilbee`` name is fine."""
137 LILBEE = "lilbee"
140class ExtractMode(StrEnum):
141 """Extraction topology: paginated (PDFs/images) vs markdown output (text formats)."""
143 MARKDOWN = "markdown"
144 PAGINATED = "paginated"
147class ChunkRecord(TypedDict):
148 """A single store-ready chunk record matching store.CHUNKS_SCHEMA."""
150 source: str
151 content_type: str
152 chunk_type: ChunkType
153 page_start: int
154 page_end: int
155 line_start: int
156 line_end: int
157 chunk: str
158 chunk_index: int
159 vector: Vector
160 # Stamped once per document by the pipeline (see produce_records); None
161 # when the title is empty, so chunk rows persist NULL like the _sources table.
162 title: NotRequired[str | None]
165class SyncResult(BaseModel):
166 """Summary of a sync operation."""
168 added: list[str] = []
169 updated: list[str] = []
170 removed: list[str] = []
171 unchanged: int = 0
172 # Sources recognized as moved (same content hash, new location): re-keyed to
173 # the new name in place, so their chunks and embeddings were reused, not rebuilt.
174 relocated: list[str] = []
175 failed: list[str] = []
176 skipped: list[str] = []
177 # The OCR each skipped document ran with; files that never reach OCR are absent.
178 skipped_ocr: dict[str, OcrReport] = {}
179 # Files an earlier sync skip-marked, so this run did not attempt them.
180 held_out: list[SkippedSource] = []
181 # Chunks whose text exceeded the embedder's char budget and were truncated
182 # before embedding. Non-zero means some tail content did not reach the index.
183 truncated: int = 0
184 # Set when the index was built with another embedder than the one configured:
185 # the sync left it as it is, and search refuses it until a rebuild or a switch back.
186 index_mismatch: IndexMismatch | None = None
188 def _lines(self) -> list[list[tuple[str, str]]]:
189 """The summary as lines of ``(text, style)`` segments; ``""`` means unstyled."""
190 lines: list[list[tuple[str, str]]] = [
191 [(f"Added: {len(self.added)}", "")],
192 [(f"Updated: {len(self.updated)}", "")],
193 [(f"Removed: {len(self.removed)}", "")],
194 [(f"Unchanged: {self.unchanged}", "")],
195 ]
196 if self.index_mismatch is not None:
197 lines.append([("Index mismatch:", "red"), (f" {self.index_mismatch.message}", "")])
198 if self.relocated:
199 lines.append([(f"Relocated: {len(self.relocated)}", "")])
200 lines += [
201 [(f"Held out: {len(self.held_out)}", "")],
202 [(f"Skipped: {len(self.skipped)}", "")],
203 [(f"Failed: {len(self.failed)}", "")],
204 [(f"Truncated: {self.truncated}", "")],
205 ]
206 lines += [
207 [(" ", ""), (h.filename, "yellow"), (f": {h.reason}", "")] for h in self.held_out
208 ]
209 lines += [[(" ", ""), (name, "yellow"), self._skip_note(name)] for name in self.skipped]
210 lines += [[(" ", ""), (name, "red")] for name in self.failed]
211 return lines
213 def _skip_note(self, name: str) -> tuple[str, str]:
214 """The OCR-off note for a skipped file, or an empty segment."""
215 report = self.skipped_ocr.get(name)
216 if report is not None and report.backend is OcrBackendUsed.NONE:
217 return (SKIPPED_OCR_OFF_NOTE, "")
218 return ("", "")
220 def __str__(self) -> str:
221 return "\n".join(
222 "".join(f"[{style}]{text}[/{style}]" if style else text for text, style in line)
223 for line in self._lines()
224 )
226 def __repr__(self) -> str:
227 return (
228 f"SyncResult(added={len(self.added)}, updated={len(self.updated)}, "
229 f"removed={len(self.removed)}, unchanged={self.unchanged}, "
230 f"held_out={len(self.held_out)}, skipped={len(self.skipped)}, "
231 f"failed={len(self.failed)}, truncated={self.truncated})"
232 )
234 def __rich__(self) -> Text:
235 """Render the summary with every filename and reason as literal text."""
236 rendered = Text("\n").join(Text.assemble(*line) for line in self._lines())
237 return ReprHighlighter()(rendered)
240@dataclass
241class _IngestResult:
242 """Outcome of a single file ingestion attempt.
244 ``records`` carries the produced (extracted + embedded) chunks until the
245 batched flush writes them; ``None`` on a failed file. ``needs_cleanup``
246 travels with the records so the flush can delete the source's old chunks in
247 the same transaction. ``page_texts`` carries the per-page text dataset rows
248 and ``concept_records`` the file's concept-table rows, and ``entity_rows``
249 the file's typed-entity rows, all written by the same flush. ``meta``
250 carries the document's extraction-time metadata for the source row.
251 ``skip_reason`` is set when the file was refused rather than attempted, and
252 it decides the outcome ahead of the chunk count. ``ocr`` is the OCR the
253 document extraction ran with; ``None`` for files that never reach OCR.
254 """
256 name: str
257 path: Path
258 chunk_count: int
259 error: Exception | None
260 file_hash: str = ""
261 skip_reason: str | None = None
262 records: list[ChunkRecord] | None = None
263 needs_cleanup: bool = True
264 page_texts: list[PageTextRecord] | None = None
265 stat: SourceStat | None = None
266 concept_records: ConceptRecords | None = None
267 entity_rows: list[dict] | None = None
268 meta: SourceMeta | None = None
269 members: list[MemberRecords] | None = None
270 ocr: OcrReport | None = None