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

1"""Shared ingest types and constants.""" 

2 

3from __future__ import annotations 

4 

5import hashlib 

6from dataclasses import dataclass 

7from enum import StrEnum 

8from pathlib import Path 

9from typing import NamedTuple, NotRequired, TypedDict 

10 

11from pydantic import BaseModel 

12from rich.highlighter import ReprHighlighter 

13from rich.text import Text 

14 

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 

26 

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)" 

36 

37 

38@dataclass(frozen=True) 

39class ShardId: 

40 """Which slice of the corpus one ingest worker owns. 

41 

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 """ 

45 

46 index: int 

47 count: int 

48 records_root: Path 

49 

50 def owns(self, key: str) -> bool: 

51 """Whether source *key* belongs to this slice. 

52 

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 

59 

60 

61class FileToProcess(NamedTuple): 

62 """A file queued for ingestion with its metadata.""" 

63 

64 name: str 

65 path: Path 

66 content_type: str 

67 file_hash: str 

68 needs_cleanup: bool 

69 stat: SourceStat | None = None 

70 

71 

72class MemberRecords(NamedTuple): 

73 """One archive member's records, written as its own source.""" 

74 

75 name: str 

76 content_type: str 

77 records: list[ChunkRecord] 

78 page_texts: list[PageTextRecord] 

79 meta: SourceMeta 

80 

81 

82class DocumentRecords(NamedTuple): 

83 """One file's records, its source metadata, and the OCR its extraction ran with.""" 

84 

85 records: list[ChunkRecord] 

86 meta: SourceMeta 

87 ocr: OcrReport | None = None 

88 

89 

90class SkippedSource(BaseModel): 

91 """One file a skip marker holds out of the index, and why.""" 

92 

93 filename: str 

94 reason: str 

95 

96 

97class FileChangePlan(NamedTuple): 

98 """Outcome of diffing disk files against the tracked sources.""" 

99 

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] 

107 

108 

109class OcrBackendName(StrEnum): 

110 """OCR backends lilbee selects in OcrConfig: xberg's tesseract or lilbee's vision plugin.""" 

111 

112 TESSERACT = "tesseract" 

113 LILBEE_VISION = "lilbee-vision" 

114 

115 

116class OcrReport(BaseModel, frozen=True): 

117 """Which OCR backend one extraction ran and how many pages it OCR'd.""" 

118 

119 backend: OcrBackendUsed 

120 pages: int = 0 

121 

122 

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.""" 

127 

128 LILBEE = "lilbee" 

129 

130 

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.""" 

136 

137 LILBEE = "lilbee" 

138 

139 

140class ExtractMode(StrEnum): 

141 """Extraction topology: paginated (PDFs/images) vs markdown output (text formats).""" 

142 

143 MARKDOWN = "markdown" 

144 PAGINATED = "paginated" 

145 

146 

147class ChunkRecord(TypedDict): 

148 """A single store-ready chunk record matching store.CHUNKS_SCHEMA.""" 

149 

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] 

163 

164 

165class SyncResult(BaseModel): 

166 """Summary of a sync operation.""" 

167 

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 

187 

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 

212 

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 ("", "") 

219 

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 ) 

225 

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 ) 

233 

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) 

238 

239 

240@dataclass 

241class _IngestResult: 

242 """Outcome of a single file ingestion attempt. 

243 

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 """ 

255 

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