Coverage for src/lilbee/data/ingest/discovery.py: 100%

174 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""File discovery, classification, hashing, and source-path resolution.""" 

2 

3from __future__ import annotations 

4 

5import hashlib 

6import logging 

7import os 

8import time 

9from collections.abc import Iterator, Mapping 

10from enum import StrEnum 

11from functools import cache 

12from pathlib import Path 

13from types import MappingProxyType 

14from typing import NamedTuple 

15 

16from lilbee.core.config import active_config 

17from lilbee.core.system import is_ignored_dir 

18from lilbee.data.extract.code_chunker import is_code_file 

19from lilbee.data.ingest.ignore import IgnoreRules 

20from lilbee.data.types import IMAGE_CONTENT_TYPE, PDF_CONTENT_TYPE, ShardId 

21 

22log = logging.getLogger(__name__) 

23 

24_PDF_MIME = "application/pdf" 

25 

26# Base MIME subtypes that name a container of other files. 

27_ARCHIVE_SUBTYPES = frozenset( 

28 { 

29 "7z-compressed", 

30 "bzip", 

31 "bzip2", 

32 "compress", 

33 "gtar", 

34 "gzip", 

35 "lzip", 

36 "lzma", 

37 "rar", 

38 "rar-compressed", 

39 "tar", 

40 "xz", 

41 "zip", 

42 "zip-compressed", 

43 "zstd", 

44 } 

45) 

46# Prefixes that mark a subtype as unregistered or vendor-specific, not part of its name. 

47_SUBTYPE_PREFIXES = ("x-", "vnd.", "prs.") 

48 

49 

50class ExclusionReason(StrEnum): 

51 """Why discovery refuses a file whose extension xberg could otherwise extract.""" 

52 

53 VECTOR_GRAPHIC = "vector graphic, not a document" 

54 NEEDS_TRANSCRIPTION = "audio or video, needs a transcription model lilbee does not run" 

55 

56 

57# Refusals the MIME type cannot express: image/svg+xml is a drawing, not a scan. 

58_DENIED_EXTENSIONS: dict[str, ExclusionReason] = {".svg": ExclusionReason.VECTOR_GRAPHIC} 

59# Refusals by MIME type: xberg errors on these without a transcription config. 

60_DENIED_MIME_PREFIXES: dict[str, ExclusionReason] = { 

61 "audio/": ExclusionReason.NEEDS_TRANSCRIPTION, 

62 "video/": ExclusionReason.NEEDS_TRANSCRIPTION, 

63} 

64 

65 

66def _content_type_for(ext: str, mime: str) -> str: 

67 """content_type for a xberg format: PDFs and images grouped, others keyed by extension.""" 

68 if mime == _PDF_MIME: 

69 return PDF_CONTENT_TYPE 

70 if mime.startswith("image/"): 

71 return IMAGE_CONTENT_TYPE 

72 return ext.lstrip(".") 

73 

74 

75def _normalized_ext(extension: str) -> str: 

76 """A xberg format's extension as a lowercase suffix with its leading dot.""" 

77 ext = extension.lower() 

78 return ext if ext.startswith(".") else f".{ext}" 

79 

80 

81def _is_archive_mime(mime: str) -> bool: 

82 """Whether the base subtype of *mime* names an archive. 

83 

84 ``application/epub+zip`` reads as ``epub``; ``application/x-tar`` reads as ``tar``. 

85 """ 

86 subtype = mime.partition("/")[2].partition("+")[0].strip().lower() 

87 for prefix in _SUBTYPE_PREFIXES: 

88 subtype = subtype.removeprefix(prefix) 

89 return subtype in _ARCHIVE_SUBTYPES 

90 

91 

92def _denied_mime_reason(mime: str) -> ExclusionReason | None: 

93 return next( 

94 (reason for prefix, reason in _DENIED_MIME_PREFIXES.items() if mime.startswith(prefix)), 

95 None, 

96 ) 

97 

98 

99@cache 

100def excluded_extension_reasons() -> Mapping[str, ExclusionReason]: 

101 """Extension -> why discovery refuses it, for xberg formats lilbee will not ingest.""" 

102 from xberg import list_supported_formats 

103 

104 refused = dict(_DENIED_EXTENSIONS) 

105 for fmt in list_supported_formats(): 

106 reason = _denied_mime_reason(fmt.mime_type) 

107 if reason is not None: 

108 refused[_normalized_ext(fmt.extension)] = reason 

109 return MappingProxyType(refused) 

110 

111 

112@cache 

113def archive_content_types() -> frozenset[str]: 

114 """content_types of the containers whose members ingest as their own sources. 

115 

116 Found by the MIME type xberg reports, so ``application/epub+zip`` stays a book. 

117 """ 

118 from xberg import list_supported_formats 

119 

120 return frozenset( 

121 _content_type_for(_normalized_ext(fmt.extension), fmt.mime_type) 

122 for fmt in list_supported_formats() 

123 if _is_archive_mime(fmt.mime_type) 

124 ) 

125 

126 

127@cache 

128def supported_extension_map() -> dict[str, str]: 

129 """Extension -> content_type for every format lilbee ingests. 

130 

131 Built from ``xberg.list_supported_formats()`` so lilbee covers the full set 

132 without a hand-maintained list, minus the formats ``excluded_extension_reasons`` 

133 refuses. Source-code files are routed separately (their extensions are absent 

134 here), so ``classify_file`` falls through to the code path. 

135 """ 

136 from xberg import list_supported_formats 

137 

138 excluded = excluded_extension_reasons() 

139 out: dict[str, str] = {} 

140 for fmt in list_supported_formats(): 

141 ext = _normalized_ext(fmt.extension) 

142 if ext not in excluded: 

143 out[ext] = _content_type_for(ext, fmt.mime_type) 

144 return out 

145 

146 

147# How often the discovery walk logs progress. The walk runs before the file count 

148# is known (it is what produces the count), so it cannot show an ETA; it just 

149# proves the run is alive. A large or NFS-backed tree can take minutes to walk, 

150# during which the plan pass has not started and nothing else logs. 

151_SCAN_LOG_INTERVAL_S = 10.0 

152 

153 

154class _ScanProgress: 

155 """Periodic progress for the pre-plan discovery walk. 

156 

157 Emitted at warning level, not info: the default LILBEE_LOG_LEVEL is WARNING, 

158 so an info line would be filtered before any handler and a headless 

159 ``lilbee sync`` would show nothing while the tree is walked. Interval-gated, 

160 so a fast walk (the common case) stays silent -- the first line appears only 

161 once the walk has run longer than the interval. 

162 """ 

163 

164 def __init__(self) -> None: 

165 self._examined = 0 

166 self._matched = 0 

167 self._started = time.monotonic() 

168 self._last = self._started 

169 

170 def tick(self, *, matched: bool) -> None: 

171 self._examined += 1 

172 if matched: 

173 self._matched += 1 

174 now = time.monotonic() 

175 if now - self._last < _SCAN_LOG_INTERVAL_S: 

176 return 

177 self._last = now 

178 elapsed = now - self._started 

179 rate = self._examined / elapsed if elapsed > 0 else 0.0 

180 log.warning( 

181 "Scanning for files: examined %d, matched %d (%.0f files/s, %.0fs elapsed)", 

182 self._examined, 

183 self._matched, 

184 rate, 

185 elapsed, 

186 ) 

187 

188 

189def file_hash(path: Path) -> str: 

190 """Compute SHA-256 hex digest of a file.""" 

191 with open(path, "rb") as f: 

192 return hashlib.file_digest(f, "sha256").hexdigest() 

193 

194 

195def classify_file(path: Path) -> str | None: 

196 """Classify a file by extension: a xberg content_type, "code", or None. 

197 

198 Ingestable xberg formats win; source code (not in xberg's set) routes to the 

199 code chunker; a refused container and anything else is unsupported. 

200 """ 

201 doc_type = supported_extension_map().get(path.suffix.lower()) 

202 if doc_type is not None: 

203 return doc_type 

204 if is_code_file(path): 

205 return "code" 

206 return None 

207 

208 

209def resolve_source_path(filename: str) -> Path: 

210 """Map a stored source key back to the file it tracks on disk. 

211 

212 A key's first segment is a registered root label when ``add`` recorded that 

213 root; the file then lives at ``linked_roots[label]/<rest>`` (or at the root 

214 itself for a single-file root, where there is no rest). Every other key 

215 belongs to a file lilbee owns under ``documents_dir`` and resolves there. 

216 The path is returned whether or not it still exists: a source whose file was 

217 moved or deleted keeps its index entry, and the dead path surfaces only when 

218 something tries to open it. 

219 

220 A registered label owns its whole key namespace: if an owned ``documents_dir`` 

221 subtree of the same top-level name is created after the root is registered, 

222 its files resolve to the root, not the owned copy. ``discover_files`` walks 

223 the root after the owned tree and so keys the same file identically, keeping 

224 resolution and discovery in agreement; ``add`` blocks the reverse collision 

225 (registering a label that shadows an existing owned entry). 

226 """ 

227 config = active_config() 

228 first, _, rest = filename.partition("/") 

229 root = config.linked_roots.get(first) 

230 if root is not None: 

231 base = Path(root) 

232 return base / rest if rest else base 

233 return config.documents_dir / filename 

234 

235 

236def resolve_source_root(filename: str) -> tuple[Path, Path] | None: 

237 """The walked root and the resolved path for *filename*, or None if none walks it. 

238 

239 Pairs a source key with the base its patterns are written relative to, so the 

240 index can be reconciled against ``.lilbeeignore`` without a second walk. A 

241 single-file root is the file the user named, never a tree, so nothing walks 

242 it and no pattern applies. 

243 """ 

244 config = active_config() 

245 first, _, rest = filename.partition("/") 

246 root = config.linked_roots.get(first) 

247 if root is None: 

248 return config.documents_dir, config.documents_dir / filename 

249 if not rest: 

250 return None 

251 base = Path(root) 

252 return base, base / rest 

253 

254 

255def resolve_source_path_checked(filename: str) -> Path | None: 

256 """Resolve *filename*, returning None if it escapes its owning root. 

257 

258 Guards a surface that resolves a caller-supplied source key (the HTTP 

259 document-serving endpoint): a key with ``..`` that would climb out of 

260 ``documents_dir`` or a registered root is rejected. Keys produced by 

261 discovery never contain ``..``; this defends against a crafted request, not 

262 stored data. 

263 """ 

264 config = active_config() 

265 resolved = resolve_source_path(filename).resolve(strict=False) 

266 roots = [ 

267 config.documents_dir.resolve(), 

268 *(Path(root).resolve() for root in config.linked_roots.values()), 

269 ] 

270 if any(resolved == root or root in resolved.parents for root in roots): 

271 return resolved 

272 return None 

273 

274 

275class ScannedFile(NamedTuple): 

276 """One file the corpus walk kept: its source key, its path, and why it is refused. 

277 

278 ``excluded`` is None for a file that will be ingested. 

279 """ 

280 

281 key: str 

282 path: Path 

283 excluded: ExclusionReason | None = None 

284 

285 

286class CorpusScan(NamedTuple): 

287 """One walk's outcome: the files to ingest, and the refused ones keyed to their reason.""" 

288 

289 files: dict[str, Path] 

290 excluded: dict[str, ExclusionReason] 

291 

292 

293def _scan_entry(path: Path, key: str) -> ScannedFile | None: 

294 """The walk's verdict for one file, or None when nothing here can be ingested.""" 

295 reason = excluded_extension_reasons().get(path.suffix.lower()) 

296 if reason is not None: 

297 return ScannedFile(key, path, reason) 

298 if classify_file(path) is None: 

299 return None 

300 return ScannedFile(key, path) 

301 

302 

303def _walk_root( 

304 base: Path, 

305 label: str | None, 

306 ignore_dirs: frozenset[str], 

307 progress: _ScanProgress, 

308 rules: IgnoreRules, 

309) -> Iterator[ScannedFile]: 

310 """Yield the files under *base* lilbee knows, keyed relative to it (prefixed by *label*). 

311 

312 A refused format is yielded with its reason; an unknown format is left out. 

313 

314 Symlinks are not followed (``followlinks=False``): each root is walked as the 

315 real tree it names, so there is no traversal loop and no path can escape the 

316 root it was registered under. 

317 

318 A directory ``.lilbeeignore`` excludes is pruned rather than filtered per 

319 file, so an excluded tree costs nothing to skip and no pattern beneath it can 

320 re-include a file -- git's rule, holding here because the walk never descends. 

321 """ 

322 for root, dirs, filenames in os.walk(base, topdown=True, followlinks=False): 

323 here = Path(root) 

324 dirs[:] = [ 

325 d 

326 for d in dirs 

327 if not is_ignored_dir(d, ignore_dirs) 

328 and not rules.excludes_entry(here / d, base=base, is_dir=True) 

329 ] 

330 for fname in filenames: 

331 if fname.startswith("."): 

332 continue 

333 path = here / fname 

334 if rules.excludes_entry(path, base=base, is_dir=False): 

335 progress.tick(matched=False) 

336 continue 

337 rel = path.relative_to(base).as_posix() 

338 entry = _scan_entry(path, f"{label}/{rel}" if label else rel) 

339 # tick per file visited, not per match: a skip-heavy tree still walks 

340 # slowly and must still show a heartbeat. 

341 progress.tick(matched=entry is not None and entry.excluded is None) 

342 if entry is not None: 

343 yield entry 

344 

345 

346def _walk_corpus(rules: IgnoreRules | None = None) -> Iterator[ScannedFile]: 

347 """Yield every file lilbee knows in the owned tree and in each registered root. 

348 

349 A single-file root is the file the user named at ``add`` time, so no ignore 

350 pattern is consulted for it: naming a file is a stronger statement than a 

351 pattern that would have swept it up. 

352 """ 

353 config = active_config() 

354 progress = _ScanProgress() 

355 rules = rules if rules is not None else IgnoreRules.for_corpus() 

356 if config.documents_dir.exists(): 

357 yield from _walk_root(config.documents_dir, None, config.ignore_dirs, progress, rules) 

358 for label, root in config.linked_roots.items(): 

359 root_path = Path(root) 

360 if root_path.is_dir(): 

361 yield from _walk_root(root_path, label, config.ignore_dirs, progress, rules) 

362 elif root_path.is_file() and (entry := _scan_entry(root_path, label)) is not None: 

363 yield entry 

364 

365 

366def discover_corpus(shard: ShardId | None = None, rules: IgnoreRules | None = None) -> CorpusScan: 

367 """Scan the owned documents dir and every registered root into a :class:`CorpusScan`. 

368 

369 A refused file (see ``excluded_extension_reasons``) lands in ``excluded``, not ``files``. 

370 """ 

371 files: dict[str, Path] = {} 

372 excluded: dict[str, ExclusionReason] = {} 

373 for entry in _walk_corpus(rules): 

374 if shard is not None and not shard.owns(entry.key): 

375 continue 

376 if entry.excluded is not None: 

377 excluded[entry.key] = entry.excluded 

378 else: 

379 files[entry.key] = entry.path 

380 return CorpusScan(files, excluded) 

381 

382 

383def discover_files( 

384 shard: ShardId | None = None, rules: IgnoreRules | None = None 

385) -> dict[str, Path]: 

386 """Scan the owned documents dir and every registered root, return {key: path}. 

387 

388 Files lilbee owns under ``documents_dir`` (crawl and upload output) are keyed 

389 by their path relative to it. Each root ``add`` registered is indexed where it 

390 lives: a directory root contributes its files keyed under the root's label; a 

391 single-file root contributes one entry keyed by the label alone. A root whose 

392 path has since vanished contributes nothing this pass, and its already-indexed 

393 sources are left in place (a dead path-link, not a removal). 

394 

395 A *shard* keeps only the keys that slice owns, so one worker of a multi-GPU 

396 ingest holds the paths of its own slice and not the whole corpus. 

397 

398 A caller that also reconciles the index passes the *rules* it will reconcile 

399 with, so the walk and that pass read one set of compiled patterns. 

400 """ 

401 return discover_corpus(shard, rules).files 

402 

403 

404def corpus_has_at_least(count: int) -> bool: 

405 """Whether the corpus holds at least *count* ingestable files. 

406 

407 Stops at the threshold. The answer gates the multi-GPU ingest fan-out, and 

408 walking a million-file tree to learn "yes, more than a few thousand" would 

409 cost minutes before any work starts. 

410 """ 

411 ingestable = (entry for entry in _walk_corpus() if entry.excluded is None) 

412 return any(seen >= count for seen, _ in enumerate(ingestable, start=1))