Coverage for src/lilbee/app/ingest.py: 100%

201 statements  

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

1"""Register external source roots, and remove indexed documents durably.""" 

2 

3from __future__ import annotations 

4 

5import fnmatch 

6import logging 

7from collections.abc import Generator, Iterable 

8from contextlib import contextmanager 

9from dataclasses import dataclass, field 

10from pathlib import Path 

11 

12from lilbee.app.services import get_services 

13from lilbee.core import settings 

14from lilbee.core.config import active_config 

15from lilbee.data.ingest.discovery import ( 

16 excluded_extension_reasons, 

17 file_hash, 

18 resolve_source_path, 

19) 

20from lilbee.data.ingest.skip_marker import ( 

21 SkipRecords, 

22 held_out_names, 

23 mark_removed, 

24 update_skip_records, 

25) 

26from lilbee.data.store.types import RemoveResult 

27 

28 

29@dataclass 

30class RegisterResult: 

31 """Result of registering source roots into the knowledge base.""" 

32 

33 registered: list[str] = field(default_factory=list) # labels newly registered 

34 name_taken: list[str] = field(default_factory=list) 

35 """Labels held by a different live source or an owned entry; ``--force`` overwrites.""" 

36 overlapping: list[str] = field(default_factory=list) 

37 """Paths nesting under or over ``documents_dir`` or a live root; that source covers them.""" 

38 refused: list[str] = field(default_factory=list) 

39 """Files whose format lilbee does not index, as ``name: reason``.""" 

40 tracked: list[str] = field(default_factory=list) 

41 """Named sources the knowledge base already tracks, so nothing was registered. 

42 

43 Either the path already lives under ``documents_dir`` or this exact source is 

44 already registered under that label. Nothing is wrong and ``--force`` would 

45 change nothing: the sync that follows covers them. 

46 """ 

47 

48 @property 

49 def reached_corpus(self) -> bool: 

50 """Whether a named path is in the corpus, so a sync has something to index for it. 

51 

52 A refused or missing path is not, and neither is one whose label another 

53 source holds; a sync after one is a whole-vault pass whose summary would 

54 read as the outcome of the add. 

55 """ 

56 return bool(self.registered or self.tracked or self.overlapping) 

57 

58 

59def _resolve_label( 

60 base: str, roots: dict[str, str], docs_resolved: Path, *, force: bool 

61) -> str | None: 

62 """Choose the source-key label for a new root, or None when the name is taken. 

63 

64 An owned ``documents_dir`` top-level entry of the same name always wins and is 

65 never shadowed, even under ``force`` -- a label that shadows it would make 

66 resolve_source_path disagree with how discovery keyed the owned file. Reuses 

67 the label when a root of that name was registered before but its path has since 

68 vanished (the source moved: re-register in place, no ``--force`` needed) or 

69 when ``force`` overwrites a live registered root of the same name. 

70 """ 

71 if (docs_resolved / base).exists(): 

72 return None # an owned entry holds this name; never shadow it 

73 existing = roots.get(base) 

74 if existing is not None and not Path(existing).exists(): 

75 return base # dangling root; the source moved, re-point it to the new path 

76 if force: 

77 return base 

78 if base in roots: 

79 return None 

80 return base 

81 

82 

83def _overlaps_existing(src: Path, docs_resolved: Path, roots: dict[str, str]) -> bool: 

84 """Whether *src* overlaps ``documents_dir`` or a live registered root. 

85 

86 Two roots covering the same tree would walk the same file twice and index it 

87 under two keys (double-index). The caller already rejects *src* inside 

88 ``documents_dir``; this rejects *src* being an ANCESTOR of it, and *src* 

89 nesting under or over any live registered root. A vanished root cannot 

90 double-index, so it is ignored. 

91 """ 

92 if docs_resolved.is_relative_to(src): 

93 return True 

94 for target in roots.values(): 

95 root = Path(target) 

96 if not root.exists(): 

97 continue 

98 root = root.resolve() 

99 if src.is_relative_to(root) or root.is_relative_to(src): 

100 return True 

101 return False 

102 

103 

104def source_label_taken(name: str, target: Path | None = None) -> bool: 

105 """Whether registering *target* under *name* would collide with a different source. 

106 

107 The confirm-before-overwrite affordance in the TUI reads this. The label 

108 rule is delegated to :func:`_resolve_label` (the authority register_sources 

109 itself applies) rather than mirrored, and a *target* whose exact path is 

110 already registered under *name* is not a collision: re-adding the same 

111 source is idempotent, matching register_sources' by-target no-op. 

112 """ 

113 config = active_config() 

114 roots = dict(config.linked_roots) 

115 if target is not None and roots.get(name) == str(target.resolve()): 

116 return False 

117 return _resolve_label(name, roots, config.documents_dir.resolve(), force=False) is None 

118 

119 

120def register_sources(paths: list[Path], *, force: bool = False) -> RegisterResult: 

121 """Register each path as a root lilbee indexes where it already lives. 

122 

123 A prepared corpus is already on local disk, so ``add`` records where it is 

124 rather than copying or linking it: discovery walks the registered root and 

125 keys its files under the root's label (its basename). A path already inside 

126 ``documents_dir`` is left to the owned-files walk; a path already registered 

127 under the same target is a no-op; a label already taken by a different live 

128 root or an owned entry is skipped unless ``force``. The registry is persisted 

129 so later processes index the same roots. 

130 """ 

131 config = active_config() 

132 documents_dir = config.documents_dir 

133 documents_dir.mkdir(parents=True, exist_ok=True) 

134 docs_resolved = documents_dir.resolve() 

135 result = RegisterResult() 

136 if not paths: 

137 return result 

138 

139 def _mutate(persisted: dict[str, str] | None) -> tuple[dict[str, str], RegisterResult]: 

140 # Read the registry from config.toml INSIDE the lock (not the possibly 

141 # stale in-memory copy) so two processes registering roots concurrently 

142 # cannot lose each other's entry. 

143 roots = dict(persisted or {}) 

144 by_target = {target: label for label, target in roots.items()} 

145 refused = excluded_extension_reasons() 

146 for p in paths: 

147 src = p.resolve() 

148 reason = refused.get(src.suffix.lower()) if src.is_file() else None 

149 if reason is not None: 

150 result.refused.append(f"{p.name}: {reason}") 

151 continue 

152 if src == docs_resolved or docs_resolved in src.parents: 

153 result.tracked.append(p.name) # already owned by the knowledge base 

154 continue 

155 already = by_target.get(str(src)) 

156 if already is not None: 

157 result.tracked.append(already) # this exact source is already registered 

158 continue 

159 if _overlaps_existing(src, docs_resolved, roots): 

160 result.overlapping.append(p.name) # would walk the same files twice 

161 continue 

162 label = _resolve_label(src.name, roots, docs_resolved, force=force) 

163 if label is None: 

164 result.name_taken.append(src.name) 

165 continue 

166 roots[label] = str(src) 

167 by_target[str(src)] = label 

168 result.registered.append(label) 

169 config.linked_roots = roots # refresh the in-process view (picks up merges) 

170 return roots, result 

171 

172 result = settings.mutate_value(config.data_root, "linked_roots", _mutate) 

173 unmark_sources_under(paths) 

174 return result 

175 

176 

177def unmark_sources_under(paths: list[Path]) -> None: 

178 """Drop the skip records (marker, reason and kind) of every source *paths* covers. 

179 

180 A marker exists to stop *discovery* from resurrecting a source the user 

181 removed, or from re-paying the extract cost on a file that yielded nothing. 

182 Naming the path outranks it: ``add`` is the user asking for that source 

183 back, so the marker goes and the sync that follows ingests the file again. 

184 Without this a removal would be permanent, undoable only by ``rebuild``, 

185 which the user has no reason to reach for after typing the path they want. 

186 

187 Each root in *paths* must be registered when this runs: marker keys resolve 

188 to files through the live registry. 

189 """ 

190 

191 def _drop_covered(records: SkipRecords) -> None: 

192 for name in _markers_covering(records.markers, paths): 

193 records.markers.pop(name) 

194 

195 update_skip_records(active_config().data_root, _drop_covered) 

196 

197 

198def _markers_covering(markers: dict[str, str], paths: list[Path]) -> set[str]: 

199 """Marker keys whose file is one of *paths* or lives beneath one. 

200 

201 Each key is resolved back to the file it tracks -- the same mapping 

202 discovery keyed it by -- so an owned ``documents_dir`` entry, a file under a 

203 registered root, and a single-file root are all matched by the one rule 

204 instead of three shape-specific ones. 

205 """ 

206 named = [p.resolve() for p in paths] 

207 covered = set() 

208 for name in markers: 

209 tracked = resolve_source_path(name).resolve(strict=False) 

210 if any(tracked == path or path in tracked.parents for path in named): 

211 covered.add(name) 

212 return covered 

213 

214 

215_GLOB_CHARS = frozenset("*?[") 

216 

217 

218def _is_glob(name: str) -> bool: 

219 """Whether *name* should be matched as a glob rather than a literal source.""" 

220 return any(char in _GLOB_CHARS for char in name) 

221 

222 

223def folder_members(name: str, known: Iterable[str]) -> list[str]: 

224 """Known sources under folder *name*, matched on whole path segments. 

225 

226 ``myrepo`` covers ``myrepo/a.py`` but never ``myrepo-2/x``. Empty when *name* 

227 is not a parent directory of any known source. 

228 """ 

229 prefix = name.rstrip("/") + "/" 

230 return [source for source in known if source.startswith(prefix)] 

231 

232 

233def removable_names(indexed: list[str] | None = None) -> list[str]: 

234 """What a remove can name: *indexed* (the store's sources when not given), then failures.""" 

235 if indexed is None: 

236 indexed = [s["filename"] for s in get_services().store.get_sources()] 

237 seen = set(indexed) 

238 failed = held_out_names(active_config().data_root) 

239 return indexed + [name for name in failed if name not in seen] 

240 

241 

242def expand_remove_targets(names: list[str], known: list[str] | None = None) -> list[str]: 

243 """Expand folder names and glob patterns to the known sources they cover. 

244 

245 An exact source name is kept. A folder name (a parent directory of known 

246 sources) expands to every source beneath it. A glob (a name containing 

247 ``* ? [``) expands to every source it fnmatches. A name matching none of 

248 these is kept unchanged so the caller reports it not-found. Order and 

249 de-duplication are preserved. *known* is ``removable_names()`` when not 

250 supplied; a caller that already has it passes it to avoid a second read. 

251 """ 

252 if known is None: 

253 known = removable_names() 

254 known_set = set(known) 

255 expanded: list[str] = [] 

256 seen: set[str] = set() 

257 

258 def _add(candidate: str) -> None: 

259 if candidate not in seen: 

260 seen.add(candidate) 

261 expanded.append(candidate) 

262 

263 for name in names: 

264 if name in known_set: 

265 _add(name) 

266 continue 

267 if _is_glob(name): 

268 matches = [source for source in known if fnmatch.fnmatchcase(source, name)] 

269 else: 

270 matches = folder_members(name, known) 

271 if matches: 

272 for match in matches: 

273 _add(match) 

274 else: 

275 _add(name) # not-found; reported by the store 

276 return expanded 

277 

278 

279def unregister_roots(names: Iterable[str]) -> list[str]: 

280 """Un-register any top-level source root named in *names*. Returns removed labels. 

281 

282 ``add`` registers a source root; removing it by its label drops the registry 

283 entry so discovery stops finding its files, which then need no skip marker. 

284 The source bytes on disk are never touched. Nested names (``corpus/a.txt``) 

285 are not roots and are left alone. 

286 """ 

287 config = active_config() 

288 labels = list(names) 

289 removed: list[str] = [] 

290 if not labels: 

291 return removed 

292 

293 def _mutate(persisted: dict[str, str] | None) -> tuple[dict[str, str], list[str]]: 

294 roots = dict(persisted or {}) 

295 for name in labels: 

296 label = name.strip("/") 

297 if "/" in label or label not in roots: 

298 continue 

299 del roots[label] 

300 removed.append(label) 

301 config.linked_roots = roots # refresh the in-process view 

302 return roots, removed 

303 

304 return settings.mutate_value(config.data_root, "linked_roots", _mutate) 

305 

306 

307log = logging.getLogger(__name__) 

308 

309 

310def remove_documents_durably(names: list[str], targets: list[str] | None = None) -> RemoveResult: 

311 """Remove documents from the index (folders and globs expand) and make it stick. 

312 

313 Never deletes source bytes. A name is an indexed source, a file an ingestion 

314 failure holds out, a folder or glob covering either, or a registered root. 

315 Each removed file is held out of every later sync as a removal, at its 

316 current hash, so a held-out failure stays out instead of being retried. 

317 Removing a registered root un-registers it and drops the skip records under 

318 it: discovery can no longer find its files. Editing the source (new hash), 

319 ``rebuild`` or adding the path again restores it; ``retry-skipped`` does not. 

320 *targets* (the expanded names) is computed when not supplied; a caller that 

321 already expanded for a confirmation prompt passes it to avoid re-expanding. 

322 """ 

323 if targets is None: 

324 targets = expand_remove_targets(names) 

325 result = get_services().store.remove_documents(targets) 

326 failed = set(held_out_names(active_config().data_root)) 

327 held = [name for name in result.not_found if name in failed] 

328 roots = forget_roots(names) 

329 _hold_out_removed([*result.removed, *held], roots) 

330 forget_removed_from_wiki_index(list(result.removed)) 

331 missing = [name for name in result.not_found if name not in failed] 

332 emptied = [name for name in missing if name.strip("/") in roots] 

333 return RemoveResult( 

334 removed=[*result.removed, *held, *emptied], 

335 not_found=[name for name in missing if name not in emptied], 

336 ) 

337 

338 

339def forget_roots(names: list[str]) -> list[str]: 

340 """Un-register every root named in *names* and drop the skip records under it.""" 

341 roots = active_config().linked_roots 

342 named = [Path(roots[label]) for label in (name.strip("/") for name in names) if label in roots] 

343 unmark_sources_under(named) # records resolve through the registry, so before un-registering 

344 return unregister_roots(names) 

345 

346 

347def _hold_out_removed(names: list[str], roots: list[str]) -> None: 

348 """Hold each of *names* out of later syncs as a removal, except under an un-registered root. 

349 

350 The marker takes the file's current hash. An imported source has no file and 

351 needs no marker; a held-out file that is not reachable keeps its marker's hash. 

352 """ 

353 hashes: dict[str, str] = {} 

354 unreachable: list[str] = [] 

355 for name in names: 

356 if any(name == root or name.startswith(root + "/") for root in roots): 

357 continue # the root is gone; discovery won't resurrect these 

358 path = resolve_source_path(name) 

359 if path.exists(): 

360 hashes[name] = file_hash(path) 

361 else: 

362 unreachable.append(name) 

363 mark_removed(active_config().data_root, hashes, unreachable) 

364 

365 

366def forget_removed_from_wiki_index(removed: list[str]) -> None: 

367 """Drop removed documents from the wiki's browse index. 

368 

369 Their skip markers keep them out of later syncs, so no refresh would ever 

370 revisit their entries and the tree would keep offering pages the library 

371 can no longer support. Best effort: the removal itself already succeeded. 

372 """ 

373 if not active_config().wiki or not removed: 

374 return 

375 from lilbee.wiki.stubs import drop_sources_from_index 

376 

377 try: 

378 drop_sources_from_index(set(removed)) 

379 except Exception: 

380 log.warning("Failed to drop removed documents from the wiki index", exc_info=True) 

381 

382 

383@contextmanager 

384def temporary_ocr_config( 

385 enable_ocr: bool | None = None, 

386 ocr_timeout: float | None = None, 

387) -> Generator[None, None, None]: 

388 """Override OCR config for the duration of the block, per request. 

389 

390 Backed by a ContextVar rather than a global ``cfg`` mutation, so concurrent 

391 ingests on the shared HTTP daemon do not clobber one another's OCR settings. 

392 """ 

393 from lilbee.data.extract.document import ocr_override 

394 

395 with ocr_override(enable_ocr, ocr_timeout): 

396 yield