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
« 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."""
3from __future__ import annotations
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
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
22log = logging.getLogger(__name__)
24_PDF_MIME = "application/pdf"
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.")
50class ExclusionReason(StrEnum):
51 """Why discovery refuses a file whose extension xberg could otherwise extract."""
53 VECTOR_GRAPHIC = "vector graphic, not a document"
54 NEEDS_TRANSCRIPTION = "audio or video, needs a transcription model lilbee does not run"
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}
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(".")
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}"
81def _is_archive_mime(mime: str) -> bool:
82 """Whether the base subtype of *mime* names an archive.
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
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 )
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
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)
112@cache
113def archive_content_types() -> frozenset[str]:
114 """content_types of the containers whose members ingest as their own sources.
116 Found by the MIME type xberg reports, so ``application/epub+zip`` stays a book.
117 """
118 from xberg import list_supported_formats
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 )
127@cache
128def supported_extension_map() -> dict[str, str]:
129 """Extension -> content_type for every format lilbee ingests.
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
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
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
154class _ScanProgress:
155 """Periodic progress for the pre-plan discovery walk.
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 """
164 def __init__(self) -> None:
165 self._examined = 0
166 self._matched = 0
167 self._started = time.monotonic()
168 self._last = self._started
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 )
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()
195def classify_file(path: Path) -> str | None:
196 """Classify a file by extension: a xberg content_type, "code", or None.
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
209def resolve_source_path(filename: str) -> Path:
210 """Map a stored source key back to the file it tracks on disk.
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.
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
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.
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
255def resolve_source_path_checked(filename: str) -> Path | None:
256 """Resolve *filename*, returning None if it escapes its owning root.
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
275class ScannedFile(NamedTuple):
276 """One file the corpus walk kept: its source key, its path, and why it is refused.
278 ``excluded`` is None for a file that will be ingested.
279 """
281 key: str
282 path: Path
283 excluded: ExclusionReason | None = None
286class CorpusScan(NamedTuple):
287 """One walk's outcome: the files to ingest, and the refused ones keyed to their reason."""
289 files: dict[str, Path]
290 excluded: dict[str, ExclusionReason]
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)
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*).
312 A refused format is yielded with its reason; an unknown format is left out.
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.
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
346def _walk_corpus(rules: IgnoreRules | None = None) -> Iterator[ScannedFile]:
347 """Yield every file lilbee knows in the owned tree and in each registered root.
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
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`.
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)
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}.
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).
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.
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
404def corpus_has_at_least(count: int) -> bool:
405 """Whether the corpus holds at least *count* ingestable files.
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))