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

843 statements  

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

1"""Top-level sync orchestration: discovery, dispatch, batching, post-sync hooks.""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6import contextlib 

7import functools 

8import logging 

9import os 

10import threading 

11import time 

12from collections import deque 

13from collections.abc import ( 

14 AsyncGenerator, 

15 Callable, 

16 Collection, 

17 Coroutine, 

18 Iterable, 

19 Iterator, 

20 Mapping, 

21) 

22from concurrent.futures import Future, ThreadPoolExecutor, as_completed 

23from dataclasses import dataclass, field 

24from itertools import count 

25from pathlib import Path 

26from typing import Any, ParamSpec, cast 

27 

28from rich.progress import ( 

29 BarColumn, 

30 MofNCompleteColumn, 

31 Progress, 

32 SpinnerColumn, 

33 TimeElapsedColumn, 

34) 

35 

36from lilbee.app.services import get_services 

37from lilbee.core.config import Config, active_config 

38from lilbee.data.extract.chunk import ChunkLimitError 

39from lilbee.data.extract.document import ( 

40 extract_batching, 

41 ingest_archive, 

42 ingest_document, 

43 ingest_markdown, 

44 ocr_backend, 

45 warn_if_table_model_ignored, 

46) 

47from lilbee.data.extract.trace import configure_from_env as configure_trace_from_env 

48from lilbee.data.ingest.adaptive import ( 

49 AdaptiveController, 

50 ResizableGate, 

51 enumerate_fleet_devices, 

52 make_signal_sampler, 

53 profile_for, 

54 resolve_mode, 

55) 

56from lilbee.data.ingest.code import ingest_code_sync 

57from lilbee.data.ingest.discovery import ( 

58 ExclusionReason, 

59 archive_content_types, 

60 classify_file, 

61 discover_corpus, 

62 discover_files, 

63 file_hash, 

64 resolve_source_root, 

65) 

66from lilbee.data.ingest.errors import error_reason 

67from lilbee.data.ingest.fanout import ( 

68 WORKER_LOG_NAME, 

69 ShardDone, 

70 ShardOptions, 

71 ShardSpec, 

72 aggregate_results, 

73 plan_fanout, 

74 run_workers, 

75) 

76from lilbee.data.ingest.ignore import IgnoreRules 

77from lilbee.data.ingest.skip_marker import ( 

78 SkipKind, 

79 SkipRecords, 

80 clear_failed_markers, 

81 clear_skip_markers, 

82 describe_skips, 

83 held_out_names, 

84 load_skip_markers, 

85 update_skip_records, 

86) 

87from lilbee.data.offload import ( 

88 embed_inflight_target, 

89 max_workers, 

90 to_executor, 

91 to_ingest_thread, 

92) 

93from lilbee.data.store import ( 

94 SOURCE_STAT_UNKNOWN, 

95 ChunkWrite, 

96 ConceptRecords, 

97 IndexMismatch, 

98 PageTextRecord, 

99 SourceMeta, 

100 SourceRecord, 

101 SourceStat, 

102 SourceStatBackfill, 

103 SourceType, 

104 Store, 

105 source_stat, 

106) 

107from lilbee.data.title import derive_title 

108from lilbee.data.types import ( 

109 ChunkRecord, 

110 DocumentRecords, 

111 FileChangePlan, 

112 FileToProcess, 

113 MemberRecords, 

114 OcrReport, 

115 ShardId, 

116 SyncResult, 

117 _IngestResult, 

118) 

119from lilbee.runtime.asyncio_loop import is_executor_shutdown 

120from lilbee.runtime.cancellation import CancelSignal, TaskCancelledError 

121from lilbee.runtime.cpu import available_cpu_count, cpu_quota 

122from lilbee.runtime.lock import LockTimeoutError, sync_running 

123from lilbee.runtime.progress import ( 

124 BatchProgressEvent, 

125 BatchStatus, 

126 DetailedProgressCallback, 

127 EmbedEvent, 

128 EventType, 

129 ExtractEvent, 

130 FileDoneEvent, 

131 FileStartEvent, 

132 OcrBackendUsed, 

133 OcrStartEvent, 

134 ProgressEvent, 

135 SyncDoneEvent, 

136 noop_callback, 

137) 

138from lilbee.runtime.progress.columns import literal_text_column 

139 

140log = logging.getLogger(__name__) 

141 

142 

143def _max_concurrent() -> int: 

144 """Files allowed in their compute phase at once. 

145 

146 ``cpu_quota()`` (cpu_count // 2) keeps worker storms from starving the TUI's asyncio 

147 main thread, and is the cap for text/code ingest. Vision OCR is different: every file 

148 in compute holds a continuous-batching slot on a vision server, so an OCR run is bounded 

149 by the vision slot capacity. Once the fleet is up that capacity is the servers' real 

150 fitted ``--parallel`` slots (a memory-constrained card fits fewer than requested); until 

151 then it is estimated from ``replicas x per-server pages``. Sizing to real capacity keeps 

152 the OCR queue shallow instead of piling pages behind the dispatcher with their deadlines 

153 ticking, and still keeps every slot fed. 

154 """ 

155 from lilbee.providers.fleet.replicas import gpu_device_count, resolve_replica_count 

156 from lilbee.providers.roles import WorkerRole 

157 

158 config = active_config() 

159 if ocr_backend() is OcrBackendUsed.VISION: 

160 fitted = get_services().provider.vision_slot_capacity() 

161 if fitted is not None: 

162 return fitted 

163 replicas = resolve_replica_count(WorkerRole.VISION, gpu_device_count()) 

164 return max(1, replicas * config.vision_ocr_concurrency) 

165 if config.ingest_max_inflight > 0: 

166 return config.ingest_max_inflight # explicit override 

167 # Auto: keep every embed replica fed. The CPU-bound quota alone leaves a 

168 # many-core multi-GPU box starved (~4 files/card), so scale admission with 

169 # the detected fleet size -- no manual cap needed. 

170 return max(cpu_quota(), embed_inflight_target()) 

171 

172 

173async def _rebuild_concept_clusters( 

174 added: Collection[str] = (), 

175 updated: Collection[str] = (), 

176 removed: Collection[str] = (), 

177) -> None: 

178 """Re-run Leiden clustering after sync. No-op if disabled. 

179 

180 Removals force a full pass: their rows are already gone, so the touched 

181 communities are unrecoverable. Adds and updates recluster incrementally. 

182 """ 

183 if not active_config().concept_graph: 

184 return 

185 from lilbee.retrieval.concepts import concepts_available 

186 

187 if not concepts_available(): 

188 return 

189 try: 

190 cg = get_services().concepts 

191 if not cg.get_graph(): 

192 return 

193 if removed: 

194 await to_ingest_thread(cg.rebuild_clusters) 

195 else: 

196 await to_ingest_thread(cg.rebuild_clusters, set(added), set(updated)) 

197 except Exception: 

198 log.warning("Concept cluster rebuild failed", exc_info=True) 

199 

200 

201async def build_entity_records(records: list[ChunkRecord], source_name: str) -> list[dict] | None: 

202 """Extract typed entities for ingested chunks. None when the mode is off. 

203 

204 Gated twice: the ``entity_extraction`` config flag, and a schema already 

205 induced into the index; absent either, syncs cost nothing. Extraction 

206 failures degrade to no rows for the file, mirroring concept extraction. 

207 """ 

208 config = active_config() 

209 if not config.entity_extraction or not records: 

210 return None 

211 from lilbee.retrieval.entities import ExtractorKind, extract_entities, load_schema 

212 

213 schema = load_schema(get_services().store) 

214 if schema is None: 

215 return None 

216 nlp = None 

217 if any(t.kind is ExtractorKind.SPACY for t in schema.types): 

218 from lilbee.retrieval.concepts import concepts_available 

219 from lilbee.retrieval.concepts.nlp import load_spacy_pipeline 

220 

221 if concepts_available(): 

222 try: 

223 nlp = load_spacy_pipeline() 

224 except ImportError: 

225 log.warning("spaCy model unavailable; spacy-kind entity types skipped") 

226 provider = None 

227 if any(t.kind is ExtractorKind.LLM for t in schema.types): 

228 provider = get_services().provider 

229 try: 

230 return await to_ingest_thread( 

231 extract_entities, 

232 cast("list[Mapping[str, Any]]", records), 

233 schema, 

234 provider=provider, 

235 nlp=nlp, 

236 ) 

237 except Exception: 

238 log.warning("Entity extraction failed for %s", source_name, exc_info=True) 

239 return None 

240 

241 

242async def build_concept_records( 

243 records: list[ChunkRecord], source_name: str 

244) -> ConceptRecords | None: 

245 """Extract concepts for ingested chunks and build their table rows. None if disabled. 

246 

247 Pure record building, no store access: the rows are buffered on the file's 

248 ingest result and written once per flush (see :func:`_flush_concept_records`), 

249 so a large sync pays one concept-table write per flush, not per file. 

250 """ 

251 if not active_config().concept_graph or not records: 

252 return None 

253 from lilbee.retrieval.concepts import concepts_available 

254 

255 if not concepts_available(): 

256 return None 

257 try: 

258 cg = get_services().concepts 

259 texts = [r["chunk"] for r in records] 

260 concept_lists = await to_ingest_thread(cg.extract_concepts_batch, texts) 

261 chunk_ids = [(source_name, r["chunk_index"]) for r in records] 

262 return await to_ingest_thread(cg.build_concept_records, chunk_ids, concept_lists) 

263 except Exception: 

264 log.warning("Concept extraction failed for %s", source_name, exc_info=True) 

265 return None 

266 

267 

268async def produce_records( 

269 path: Path, 

270 source_name: str, 

271 content_type: str, 

272 *, 

273 quiet: bool = False, 

274 on_progress: DetailedProgressCallback = noop_callback, 

275 page_texts_out: list[PageTextRecord] | None = None, 

276) -> DocumentRecords: 

277 """Extract, chunk, and embed a single file into its records, metadata and OCR report. 

278 

279 The LanceDB write is deferred: records are returned to the caller and written 

280 in a batched flush (see :func:`_flush_writes`), so bulk ingest pays one 

281 write-lock acquisition per batch instead of one per file. The per-page text 

282 dataset rows land in ``page_texts_out`` and are written by the same flush. 

283 The returned metadata (extraction-provided when available, stem-derived title 

284 otherwise) stamps every record's ``title`` and updates the source row. The OCR 

285 report is ``None`` for code and markdown, which never reach OCR. 

286 """ 

287 records: list[ChunkRecord] 

288 ocr: OcrReport | None = None 

289 page_texts: list[PageTextRecord] = page_texts_out if page_texts_out is not None else [] 

290 if content_type == "code": 

291 records = await to_ingest_thread(ingest_code_sync, path, source_name, on_progress) 

292 meta = SourceMeta(title=derive_title(source_name)) 

293 elif path.suffix.lower() == ".md": 

294 records, meta = await ingest_markdown( 

295 path, source_name, on_progress, page_texts_out=page_texts 

296 ) 

297 else: 

298 records, meta, ocr = await ingest_document( 

299 path, 

300 source_name, 

301 content_type, 

302 quiet=quiet, 

303 on_progress=on_progress, 

304 page_texts_out=page_texts, 

305 ) 

306 

307 for record in records: 

308 # NULL (not "") for an absent title, so chunk rows match the migration 

309 # and the _sources table, which both persist absence as NULL. 

310 record["title"] = meta.title or None 

311 return DocumentRecords(records, meta, ocr) 

312 

313 

314def _disk_stat(path: Path) -> SourceStat | None: 

315 """Current size/mtime of *path* stamped with now, or None when it cannot be stat'd.""" 

316 try: 

317 st = path.stat() 

318 except OSError: 

319 return None 

320 return SourceStat(st.st_size, st.st_mtime_ns, time.time_ns()) 

321 

322 

323def _stat_unchanged(stored: SourceStat, current: SourceStat) -> bool: 

324 """Whether the stored stat proves the file unchanged without hashing it. 

325 

326 Git-style racily-clean guard: a matching (size, mtime) only counts when the 

327 mtime is strictly older than the time the stat was recorded; a same-size 

328 edit landing in the same mtime tick is otherwise missed forever. 

329 """ 

330 if (stored.size_bytes, stored.mtime_ns) != (current.size_bytes, current.mtime_ns): 

331 return False 

332 # Unknown capture or mtime >= capture hashes anyway; clock skew can only widen 

333 # hashing, never widen skipping past the pre-existing same-tick window. 

334 return stored.captured_ns != SOURCE_STAT_UNKNOWN and current.mtime_ns < stored.captured_ns 

335 

336 

337@dataclass(frozen=True) 

338class _FileChangeVerdict: 

339 """One file's sync verdict: process it, hold it out on its skip marker, or unchanged. 

340 

341 A held file is not in the index, so it is never counted as unchanged. 

342 """ 

343 

344 to_process: FileToProcess | None = None 

345 backfill: SourceStatBackfill | None = None 

346 is_update: bool = False 

347 held: bool = False 

348 

349 

350def _classify_file_change( 

351 name: str, 

352 path: Path, 

353 record: SourceRecord | None, 

354 skip_markers: dict[str, str], 

355) -> _FileChangeVerdict: 

356 """Decide one file's verdict: stat-unchanged, hash-unchanged, skip-marked, or process.""" 

357 content_type = classify_file(path) 

358 if content_type is None: 

359 raise ValueError(f"Unsupported file slipped through discovery: {name}") 

360 stored_stat = source_stat(record) if record is not None else None 

361 current_stat = _disk_stat(path) 

362 if ( 

363 record is not None 

364 and stored_stat is not None 

365 and current_stat is not None 

366 and _stat_unchanged(stored_stat, current_stat) 

367 ): 

368 return _FileChangeVerdict() 

369 old_hash = record["file_hash"] if record is not None else None 

370 current_hash = file_hash(path) 

371 if old_hash == current_hash: 

372 # Content verified unchanged; persist the stat pair so the next 

373 # sync skips the hash entirely. 

374 backfill = ( 

375 SourceStatBackfill(record, current_stat) 

376 if record is not None and current_stat is not None 

377 else None 

378 ) 

379 return _FileChangeVerdict(backfill=backfill) 

380 if skip_markers.get(name) == current_hash: 

381 # Failed last sync at this exact hash; skip the retry. 

382 return _FileChangeVerdict(held=True) 

383 # needs_cleanup=True unconditionally: delete_by_source is idempotent, 

384 # and this closes the race where a prior ingest wrote chunks but died 

385 # before upsert_source, leaving orphaned chunks that would duplicate. 

386 return _FileChangeVerdict( 

387 to_process=FileToProcess( 

388 name, path, content_type, current_hash, needs_cleanup=True, stat=current_stat 

389 ), 

390 is_update=old_hash is not None, 

391 ) 

392 

393 

394def _plan_workers() -> int: 

395 """Worker count for the parallel planning pass: config override, else auto. 

396 

397 ``config.ingest_workers`` (also set per run by ``add --max-cpus``) wins when 

398 positive; otherwise size to the container-aware CPU budget so a big corpus 

399 hashes on every core the pod actually has, not the host's vCPU count. 

400 """ 

401 configured = active_config().ingest_workers 

402 return configured if configured > 0 else available_cpu_count() 

403 

404 

405# How often the plan pass logs progress. The pass can run for tens of minutes on 

406# a multi-million-file corpus while the Rich bar renders nothing without a TTY 

407# and stdout is block-buffered when piped; a periodic line (which logging flushes 

408# per record) keeps a headless run observable instead of looking hung. 

409_PLAN_LOG_INTERVAL_S = 10.0 

410 

411 

412class _PlanProgress: 

413 """Periodic progress for the plan/hash pass, with rate and ETA. 

414 

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

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

417 ``lilbee sync`` would show nothing during the plan pass and still look hung. 

418 """ 

419 

420 def __init__(self, total: int) -> None: 

421 self._total = total 

422 self._done = 0 

423 self._started = time.monotonic() 

424 self._last = self._started 

425 

426 def tick(self) -> None: 

427 self._done += 1 

428 now = time.monotonic() 

429 if now - self._last < _PLAN_LOG_INTERVAL_S: 

430 return 

431 self._last = now 

432 elapsed = now - self._started 

433 rate = self._done / elapsed if elapsed > 0 else 0.0 

434 remaining = (self._total - self._done) / rate if rate > 0 else 0.0 

435 log.warning( 

436 "Planning: examined %d/%d files (%.0f%%, %.0f files/s, ~%.0fs left)", 

437 self._done, 

438 self._total, 

439 100.0 * self._done / self._total, 

440 rate, 

441 remaining, 

442 ) 

443 

444 

445class _StreamStop: 

446 """Stop signal for a streamed plan: the caller's cancel, or the stream closing. 

447 

448 Closing the stream has to reach the plan batch in flight, not just the next one. 

449 Build-vs-buy: the hashers run in a thread pool, so the flag must be a 

450 thread-visible ``threading.Event``; ``anyio.CancelScope`` is async-only. 

451 """ 

452 

453 def __init__(self, cancel: CancelSignal | None) -> None: 

454 self._cancel = cancel 

455 self._closed = threading.Event() 

456 

457 def close(self) -> None: 

458 self._closed.set() 

459 

460 def is_set(self) -> bool: 

461 return self._closed.is_set() or (self._cancel is not None and self._cancel.is_set()) 

462 

463 

464def _classify_pooled( 

465 pool: ThreadPoolExecutor, 

466 items: list[tuple[str, Path]], 

467 classify: Callable[[str, Path], _FileChangeVerdict], 

468 cancel: CancelSignal | None, 

469 progress: _PlanProgress, 

470) -> dict[str, _FileChangeVerdict]: 

471 """Fan *items* across *pool*, returning name -> verdict for what completed.""" 

472 verdicts: dict[str, _FileChangeVerdict] = {} 

473 futures: dict[Future[_FileChangeVerdict], str] = {} 

474 for name, path in items: 

475 if cancel and cancel.is_set(): 

476 break 

477 futures[pool.submit(classify, name, path)] = name 

478 for future in as_completed(futures): 

479 if cancel and cancel.is_set(): 

480 # Drop queued-but-unstarted work; running tasks drain on their own. 

481 for pending in futures: 

482 pending.cancel() 

483 break 

484 verdicts[futures[future]] = future.result() 

485 progress.tick() 

486 return verdicts 

487 

488 

489def _classify_changes( 

490 items: list[tuple[str, Path]], 

491 existing_sources: dict[str, SourceRecord], 

492 skip_markers: dict[str, str], 

493 cancel: CancelSignal | None, 

494 *, 

495 progress: _PlanProgress | None = None, 

496 pool: ThreadPoolExecutor | None = None, 

497) -> dict[str, _FileChangeVerdict]: 

498 """Classify each file (stat + hash) by name, fanning across a thread pool. 

499 

500 Independent per file and side-effect-free, so pooled classification matches a 

501 serial pass; ``hashlib`` releases the GIL during digest, giving real speedup 

502 on a large corpus. Returns name -> verdict. A set ``cancel`` stops promptly: 

503 submission halts and queued-but-unstarted work is cancelled, so a mid-pass 

504 cancel over a huge corpus does not hash every remaining file. A *pool* and 

505 *progress* passed in are shared across the batches of a streamed plan, so the 

506 ETA covers the whole corpus and the workers are spun up once. 

507 """ 

508 

509 def _classify(name: str, path: Path) -> _FileChangeVerdict: 

510 return _classify_file_change(name, path, existing_sources.get(name), skip_markers) 

511 

512 total = len(items) 

513 progress = progress or _PlanProgress(total) 

514 if pool is not None: 

515 return _classify_pooled(pool, items, _classify, cancel, progress) 

516 workers = _plan_workers() 

517 if workers <= 1 or total <= 1: 

518 verdicts: dict[str, _FileChangeVerdict] = {} 

519 for name, path in items: 

520 if cancel and cancel.is_set(): 

521 break 

522 verdicts[name] = _classify(name, path) 

523 progress.tick() 

524 return verdicts 

525 with ThreadPoolExecutor(max_workers=workers, thread_name_prefix="lilbee-plan") as owned: 

526 return _classify_pooled(owned, items, _classify, cancel, progress) 

527 

528 

529def _plan_file_changes( 

530 disk_files: dict[str, Path], 

531 existing_sources: dict[str, SourceRecord], 

532 cancel: CancelSignal | None, 

533 skip_markers: dict[str, str] | None = None, 

534) -> FileChangePlan: 

535 """Diff every disk file against the store in one pass (see :func:`_plan_items`).""" 

536 return _plan_items(sorted(disk_files.items()), existing_sources, cancel, skip_markers or {}) 

537 

538 

539def _plan_items( 

540 items: list[tuple[str, Path]], 

541 existing_sources: dict[str, SourceRecord], 

542 cancel: CancelSignal | None, 

543 skip_markers: dict[str, str], 

544 *, 

545 progress: _PlanProgress | None = None, 

546 pool: ThreadPoolExecutor | None = None, 

547) -> FileChangePlan: 

548 """Diff *items* (sorted name -> path) against the store, hashing only drifted files. 

549 

550 A tracked file whose stored (size, mtime) matches the disk stat, and whose 

551 mtime predates the stat capture (see :func:`_stat_unchanged`), is unchanged 

552 without reading its bytes; everything else is SHA-256 hashed. A file whose 

553 current hash matches a marker in ``skip_markers`` (set by a prior failed 

554 attempt) is held out rather than retried every sync, and is reported as held 

555 out rather than counted unchanged. Edit the file or run 

556 ``/sync --force-rebuild`` to clear the marker and try again. 

557 

558 Classification fans across a thread pool (see :func:`_classify_changes`); the 

559 plan is assembled from the results in the original sorted order, so a partial 

560 or reordered completion never yields a wrong or reordered plan -- only a 

561 shorter one when cancelled mid-pass. 

562 """ 

563 verdicts = _classify_changes( 

564 items, existing_sources, skip_markers, cancel, progress=progress, pool=pool 

565 ) 

566 

567 files_to_process: list[FileToProcess] = [] 

568 added: dict[str, None] = {} 

569 updated: dict[str, None] = {} 

570 stat_backfills: list[SourceStatBackfill] = [] 

571 held_out: list[str] = [] 

572 unchanged = 0 

573 # Assemble in the original sorted order so a partial (cancelled) or reordered 

574 # completion never changes the plan the serial pass would have produced. 

575 for name, _path in items: 

576 verdict = verdicts.get(name) 

577 if verdict is None: 

578 continue # cancelled before this file was classified 

579 if verdict.to_process is None: 

580 if verdict.held: 

581 held_out.append(name) 

582 else: 

583 unchanged += 1 

584 if verdict.backfill is not None: 

585 stat_backfills.append(verdict.backfill) 

586 continue 

587 files_to_process.append(verdict.to_process) 

588 if verdict.is_update: 

589 updated[name] = None 

590 else: 

591 added[name] = None 

592 return FileChangePlan(files_to_process, added, updated, unchanged, stat_backfills, held_out) 

593 

594 

595@dataclass(frozen=True) 

596class _Move: 

597 """One relocated source: its old key, its new key, and the new file's stat.""" 

598 

599 old: str 

600 new: str 

601 stat: SourceStat | None 

602 

603 

604class _MovePool: 

605 """Absent sources indexed by content hash, consumed as moves are paired. 

606 

607 Built once per sync and drained across the streamed plan's batches, so a file 

608 that moved matches exactly the one old key it would have matched in a 

609 single-pass plan however the corpus is sharded. 

610 """ 

611 

612 def __init__(self, absent: list[str], existing_sources: dict[str, SourceRecord]) -> None: 

613 by_hash: dict[str, list[str]] = {} 

614 for name in absent: 

615 record = existing_sources.get(name) 

616 if record is not None: 

617 by_hash.setdefault(record["file_hash"], []).append(name) 

618 for candidates in by_hash.values(): 

619 candidates.sort() 

620 self._by_hash = by_hash 

621 

622 def take(self, file_hash: str) -> str | None: 

623 """The next absent source with this content hash, or None.""" 

624 matches = self._by_hash.get(file_hash) 

625 return matches.pop(0) if matches else None 

626 

627 

628def _detect_moves( 

629 files_to_process: list[FileToProcess], 

630 added: dict[str, None], 

631 pool: _MovePool, 

632) -> list[_Move]: 

633 """Pair brand-new files with absent sources of the same content hash. 

634 

635 Only additions (files with a new name) can be moves; an update keeps its name. 

636 When several absent sources share a hash, pairing is deterministic (sorted) 

637 and one-to-one, so a duplicated file that moved matches exactly one old key 

638 and any leftovers stay indexed under their old key. 

639 """ 

640 moves: list[_Move] = [] 

641 for entry in files_to_process: 

642 if entry.name not in added: 

643 continue 

644 old = pool.take(entry.file_hash) 

645 if old is not None: 

646 moves.append(_Move(old, entry.name, entry.stat)) 

647 return moves 

648 

649 

650def _apply_moves( 

651 moves: list[_Move], 

652 files_to_process: list[FileToProcess], 

653 added: dict[str, None], 

654) -> tuple[list[FileToProcess], list[str]]: 

655 """Fold detected moves out of the add set after they were relocated. 

656 

657 Drops each moved file from the ingest list and the added set: its chunks were 

658 re-keyed onto the new source name, not rebuilt. Returns the trimmed 

659 ``(files_to_process, relocated)``. 

660 """ 

661 moved_new = {m.new for m in moves} 

662 for name in moved_new: 

663 added.pop(name, None) 

664 remaining = [e for e in files_to_process if e.name not in moved_new] 

665 return remaining, sorted(moved_new) 

666 

667 

668def _absent_sources(sources: list[SourceRecord], disk_files: dict[str, Path]) -> list[str]: 

669 """Document sources whose backing file is not on disk this sync. 

670 

671 A vanished file is never removed on its own (its chunks stay searchable, a 

672 dead path-link discovered at open time); this set exists only to pair a 

673 reappeared identical file to its old key in move detection. Imported sources 

674 have no backing file, so they are excluded. 

675 """ 

676 return [ 

677 s["filename"] 

678 for s in sources 

679 if s["filename"] not in disk_files and s["source_type"] != SourceType.IMPORTED 

680 ] 

681 

682 

683# A plan batch is only done when its slowest hash is, so the first one is small (work 

684# reaches the fleet within a second) and later ones amortize that barrier. 

685_PLAN_SHARD_MIN_FILES = 256 

686_PLAN_SHARD_MAX_FILES = 8192 

687 

688 

689def _plan_batch_bounds(total: int) -> Iterator[tuple[int, int]]: 

690 """(start, stop) slices covering *total* files, doubling up to the cap. 

691 

692 Build-vs-buy: ``itertools.batched`` is the stock slicer but is 3.12+ (floor is 

693 3.11) and fixed-size, so it cannot ramp the batch size. 

694 """ 

695 start = 0 

696 size = _PLAN_SHARD_MIN_FILES 

697 while start < total: 

698 stop = min(start + size, total) 

699 yield start, stop 

700 start = stop 

701 size = min(size * 2, _PLAN_SHARD_MAX_FILES) 

702 

703 

704@dataclass 

705class _StreamedPlan: 

706 """Bookkeeping a streamed plan accumulates across its batches. 

707 

708 ``added`` and ``updated`` are the dicts the ingest pass mutates as files land. 

709 """ 

710 

711 added: dict[str, None] = field(default_factory=dict) 

712 updated: dict[str, None] = field(default_factory=dict) 

713 # Processed files' content hashes, for the skip markers written after the run. 

714 pending_hashes: dict[str, str] = field(default_factory=dict) 

715 relocated: list[str] = field(default_factory=list) 

716 # Old keys of relocated sources. The wiki index is keyed by source name, so 

717 # without these a move leaves the old name in it forever: its mentions 

718 # double-count and its dead chunk refs occupy the per-subject cap. 

719 relocated_from: list[str] = field(default_factory=list) 

720 unchanged: int = 0 

721 # Files a skip marker held out of this run, in plan order (ordered set). 

722 held_out: dict[str, None] = field(default_factory=dict) 

723 planned: int = 0 

724 # Files this pass's slice holds, from the discovery walk. Fixed before the 

725 # first batch is planned, so it is what progress is measured against: the 

726 # plan's own running totals grow as batches land and cannot say how much of 

727 # the corpus is left. Summed across a fan-out's workers it is the corpus. 

728 corpus_total: int = 0 

729 

730 @property 

731 def resolved(self) -> int: 

732 """Files the plan disposed of without ingest, as it disposes of them. 

733 

734 Unchanged files, skip-marker held-out files and repointed moves are done 

735 as far as the corpus is concerned, and nothing downstream ingests them, 

736 so an incremental sync would otherwise show a handful of changed files 

737 against the whole corpus. Files a cancelled plan never classified are 

738 absent from every count and so are not claimed as done. 

739 """ 

740 return self.unchanged + len(self.held_out) + len(self.relocated) 

741 

742 

743async def _absorb_plan_batch( 

744 plan: FileChangePlan, state: _StreamedPlan, moves: _MovePool 

745) -> list[FileToProcess]: 

746 """Fold one batch's plan into *state* and return the files it queues for ingest. 

747 

748 Relocations and stat backfills are written per batch, so they contend with the 

749 batch flushes and take the same one-shot lock retry. 

750 """ 

751 store = get_services().store 

752 entries = plan.files_to_process 

753 state.unchanged += plan.unchanged 

754 state.held_out.update(dict.fromkeys(plan.held_out)) 

755 state.added.update(plan.added) 

756 state.updated.update(plan.updated) 

757 

758 detected = _detect_moves(entries, plan.added, moves) 

759 if detected: 

760 relocations = [(m.old, m.new, m.stat) for m in detected] 

761 await to_ingest_thread( 

762 _retry_after_lock_timeout, lambda: store.relocate_sources(relocations) 

763 ) 

764 entries, relocated = _apply_moves(detected, entries, plan.added) 

765 state.relocated_from.extend(m.old for m in detected) 

766 for name in relocated: 

767 state.added.pop(name, None) 

768 state.relocated.extend(relocated) 

769 if plan.stat_backfills: 

770 backfills = plan.stat_backfills 

771 await to_ingest_thread( 

772 _retry_after_lock_timeout, lambda: store.update_source_stats(backfills) 

773 ) 

774 

775 state.pending_hashes.update((entry.name, entry.file_hash) for entry in entries) 

776 state.planned += len(entries) 

777 return entries 

778 

779 

780async def _plan_batches( 

781 disk_files: dict[str, Path], 

782 existing_sources: dict[str, SourceRecord], 

783 skip_markers: dict[str, str], 

784 absent: list[str], 

785 state: _StreamedPlan, 

786 cancel: CancelSignal | None, 

787) -> AsyncGenerator[list[FileToProcess]]: 

788 """Plan the corpus batch by batch, yielding each batch's files to ingest. 

789 

790 Shards are contiguous slices of one sorted item list consumed in order, so 

791 this delivers the single-pass plan of :func:`_plan_items` in pieces. The next 

792 batch is planned while the current one ingests, on a dedicated thread: the 

793 shared ingest pool is saturated by extraction and would stall the stream it 

794 feeds. Empty batches are not yielded, so the first yield means there is work. 

795 """ 

796 items = sorted(disk_files.items()) 

797 moves = _MovePool(absent, existing_sources) 

798 progress = _PlanProgress(len(items)) 

799 stop = _StreamStop(cancel) 

800 workers = _plan_workers() 

801 hashers = ThreadPoolExecutor(max_workers=workers, thread_name_prefix="lilbee-plan") 

802 driver = ThreadPoolExecutor(max_workers=1, thread_name_prefix="lilbee-plan-driver") 

803 

804 def _plan(lo: int, hi: int) -> FileChangePlan: 

805 return _plan_items( 

806 items[lo:hi], 

807 existing_sources, 

808 stop, 

809 skip_markers, 

810 progress=progress, 

811 pool=hashers if workers > 1 else None, 

812 ) 

813 

814 bounds = list(_plan_batch_bounds(len(items))) 

815 try: 

816 ahead = asyncio.ensure_future(to_executor(driver, _plan, *bounds[0])) if bounds else None 

817 for index, _bound in enumerate(bounds): 

818 if ahead is None: 

819 break 

820 plan = await ahead 

821 ahead = ( 

822 asyncio.ensure_future(to_executor(driver, _plan, *bounds[index + 1])) 

823 if index + 1 < len(bounds) and not stop.is_set() 

824 else None 

825 ) 

826 entries = await _absorb_plan_batch(plan, state, moves) 

827 if entries: 

828 yield entries 

829 finally: 

830 stop.close() 

831 if ahead is not None: 

832 ahead.cancel() 

833 # Retrieve the outcome: a prefetch that had already failed would 

834 # otherwise surface as an unretrieved task exception. 

835 with contextlib.suppress(Exception, asyncio.CancelledError): 

836 await ahead 

837 # No wait: the running batch notices the stop and exits on its own, and a 

838 # generator being closed must not block the event loop on a hash in flight. 

839 hashers.shutdown(wait=False) 

840 driver.shutdown(wait=False) 

841 

842 

843def detect_pending() -> int: 

844 """Count files in documents/ that are out of sync with the store. 

845 

846 Cheap operation: filesystem walk + stat-gated SHA-256 hashing + a single 

847 sources-table read. No embedding, no writes. Returns the count of files that 

848 would be ingested (added + updated), which is what the TaskBar hint surfaces. 

849 A vanished file is not pending work: sync leaves it indexed, so it is not 

850 counted. Reuses ``_plan_file_changes`` so the diff logic stays single-sourced. 

851 Honors skip markers: a file that failed last time at this hash does 

852 not show up as pending. Blocking: callers on the event loop run it via 

853 ``asyncio.to_thread``. 

854 """ 

855 config = active_config() 

856 disk_files = discover_files() 

857 if not disk_files: 

858 return 0 

859 existing_sources = {s["filename"]: s for s in get_services().store.get_sources()} 

860 skip_markers = load_skip_markers(config.data_root) 

861 plan = _plan_file_changes(disk_files, existing_sources, cancel=None, skip_markers=skip_markers) 

862 return len(plan.files_to_process) 

863 

864 

865# Refused files named in the log line before it falls back to a count. 

866_EXCLUDED_LOG_SAMPLE = 5 

867 

868 

869def _log_excluded(excluded: dict[str, ExclusionReason]) -> None: 

870 """Warn once per exclusion reason, naming at most ``_EXCLUDED_LOG_SAMPLE`` files.""" 

871 by_reason: dict[ExclusionReason, list[str]] = {} 

872 for name, why in excluded.items(): 

873 by_reason.setdefault(why, []).append(name) 

874 for why, names in by_reason.items(): 

875 shown = sorted(names)[:_EXCLUDED_LOG_SAMPLE] 

876 more = f" and {len(names) - len(shown)} more" if len(names) > len(shown) else "" 

877 log.warning("Skipped %d file(s), %s: %s%s", len(names), why.value, ", ".join(shown), more) 

878 

879 

880def _clear_skip_records(records_root: Path, *, force_rebuild: bool, retry_skipped: bool) -> None: 

881 """Clear every skip record for a rebuild, or only the failed ones for a retry. 

882 

883 Entries are otherwise kept whether or not a pass discovers the file. A 

884 marker is what holds a removed or unextractable file out of the next sync, 

885 so a pass that cannot see the file -- its root is unmounted or moved, or a 

886 worker owns only a shard of the corpus -- must not erase the record and 

887 re-offer the file the moment it comes back. A marker is dropped when the 

888 file ingests cleanly, by ``rebuild``, or for a failure by ``retry-skipped``. 

889 """ 

890 if force_rebuild: 

891 clear_skip_markers(records_root) 

892 elif retry_skipped: 

893 clear_failed_markers(records_root) 

894 

895 

896def _prepare_skip_records( 

897 shard: ShardId | None, *, force_rebuild: bool, retry_skipped: bool 

898) -> Path: 

899 """The data root whose skip records this sync reads and writes. 

900 

901 Only a sync that is not a worker clears them, so a fan-out clears once: a 

902 worker clearing the shared records would erase a sibling's verdicts. 

903 """ 

904 if shard is not None: 

905 return shard.records_root 

906 records_root = active_config().data_root 

907 _clear_skip_records(records_root, force_rebuild=force_rebuild, retry_skipped=retry_skipped) 

908 return records_root 

909 

910 

911def _persist_skip_records( 

912 records_root: Path, 

913 pending_hashes: dict[str, str], 

914 reasons: dict[str, str], 

915 *, 

916 succeeded: Iterable[str], 

917 failed: Iterable[str], 

918) -> None: 

919 """Merge this sync's verdicts into the skip records as they are on disk now. 

920 

921 Only the files this sync decided on change: a clean ingest drops its 

922 record, a file that produced no chunks gains one with its reason and the 

923 FAILED kind. Records written or cleared since the sync started (a reset, a 

924 removal, a rolled back add) keep their reason and kind. 

925 """ 

926 dropped = list(succeeded) 

927 held = list(failed) 

928 marked = {name: fhash for name in held if (fhash := pending_hashes.get(name))} 

929 

930 def _merge(records: SkipRecords) -> None: 

931 for name in dropped: 

932 records.markers.pop(name, None) 

933 records.markers.update(marked) 

934 records.reasons.update({name: reasons[name] for name in held if name in reasons}) 

935 records.kinds.update(dict.fromkeys(marked, SkipKind.FAILED)) 

936 

937 update_skip_records(records_root, _merge) 

938 

939 

940def _failures_among(records_root: Path, held: Iterable[str]) -> list[str]: 

941 """The files in *held* that an ingestion failure holds out, in order; removals are left out.""" 

942 failed = set(held_out_names(records_root)) 

943 return [name for name in held if name in failed] 

944 

945 

946def _report_index_mismatch(store: Store) -> IndexMismatch | None: 

947 """Name an index built with another embedder than cfg; the sync leaves it as it is. 

948 

949 An unchanged corpus never reaches the write gate that refuses such an index, 

950 so without this a sync finishes green over an index search refuses. Wiping 

951 it is the caller's decision, never the sync's. 

952 """ 

953 mismatch = store.index_mismatch() 

954 if mismatch is None: 

955 return None 

956 log.warning("Sync left the index as it is: %s", mismatch) 

957 return mismatch.describe() 

958 

959 

960def _force_rebuild_store(store: Any) -> None: 

961 """Drop the store and re-embed the preserved memories table (blocking). 

962 

963 Run off the event loop by ``sync``. ``drop_all`` keeps the memories table, so 

964 its vectors are refreshed under the (possibly changed) embedding model; 

965 a no-op when empty or no embedder. 

966 """ 

967 store.drop_all() 

968 embedder = get_services().embedder 

969 if embedder.embedding_available(): 

970 store.rebuild_memory_embeddings(lambda texts: embedder.embed_batch(texts)) 

971 

972 

973def _reconcile_missing( 

974 disk_files: dict[str, Path], 

975 sources: list[SourceRecord], 

976 failed: Iterable[str], 

977 skipped: Iterable[str], 

978 held: Iterable[str], 

979) -> list[str]: 

980 """On-disk document files absent from the store that no mechanism accounts for. 

981 

982 A file discovery found and classified that ended up in neither the sources table 

983 nor any of the accounting sets was dropped with no signal -- the silent 

984 data-loss case (a scanned PDF that never made it into the index yet reported no 

985 error). Everything legitimately not indexed is excluded: a failed extraction is in 

986 ``failed``, a zero-text file this run attempted is in ``skipped``, a file this run 

987 held out on its skip marker is in ``held``, and an unsupported type was never 

988 returned by discovery in the first place. ``held`` is the run's held-out set, not 

989 the marker file: a stale marker (file edited since it failed) must not hide a drop. 

990 """ 

991 accounted = {s["filename"] for s in sources} | set(failed) | set(skipped) | set(held) 

992 return sorted(name for name in disk_files if name not in accounted) 

993 

994 

995def _ignored_sources(sources: list[SourceRecord], rules: IgnoreRules) -> list[str]: 

996 """Indexed sources a ``.lilbeeignore`` now excludes. 

997 

998 Asked of the patterns rather than of discovery, which cannot tell an excluded 

999 file from a lost one -- both simply stop being yielded, and the two need 

1000 opposite handling. Reading every source rather than only the undiscovered 

1001 ones keeps the answer independent of which slice of the corpus this pass 

1002 walked. Imported sources have no file to match a pattern against. 

1003 """ 

1004 ignored = [] 

1005 for source in sources: 

1006 if source["source_type"] == SourceType.IMPORTED: 

1007 continue 

1008 resolved = resolve_source_root(source["filename"]) 

1009 if resolved is not None and rules.excludes_path(resolved[1], base=resolved[0]): 

1010 ignored.append(source["filename"]) 

1011 return ignored 

1012 

1013 

1014def _forget_ignored(sources: list[SourceRecord], rules: IgnoreRules) -> list[str]: 

1015 """Drop sources the patterns now exclude from the index. Returns what went. 

1016 

1017 A corpus-wide pass: it reads every source, so one worker of a fan-out must 

1018 not run it. No skip marker is written either, because the ignore file is 

1019 itself the durable statement -- a marker could only outlive it and hold a 

1020 source out after its pattern was deleted. Deleting the pattern brings the 

1021 source back on the next sync. 

1022 """ 

1023 names = _ignored_sources(sources, rules) 

1024 if not names: 

1025 return [] 

1026 from lilbee.app.ingest import forget_removed_from_wiki_index 

1027 

1028 removed = list(get_services().store.remove_documents(names).removed) 

1029 forget_removed_from_wiki_index(removed) 

1030 return removed 

1031 

1032 

1033def _forget_refused( 

1034 excluded: Mapping[str, ExclusionReason], existing: Mapping[str, SourceRecord] 

1035) -> list[str]: 

1036 """Drop indexed sources that discovery now refuses. Returns what went.""" 

1037 names = [name for name in excluded if name in existing] 

1038 if not names: 

1039 return [] 

1040 from lilbee.app.ingest import forget_removed_from_wiki_index 

1041 

1042 removed = list(get_services().store.remove_documents(names).removed) 

1043 forget_removed_from_wiki_index(removed) 

1044 return removed 

1045 

1046 

1047def _require_embedding_model() -> None: 

1048 """Refuse ingest without an embedding model. 

1049 

1050 Ingest has no degraded mode: without one, every chunk fails to embed after 

1051 the run has already paid the parse and OCR cost. Search and chat fall back 

1052 to keyword via embedding_available() and carry on. 

1053 """ 

1054 if not get_services().embedder.validate_model(): 

1055 ref = active_config().embedding_model 

1056 detail = ( 

1057 f"{ref!r} is not available: pull it, or set a different embedding_model." 

1058 if ref 

1059 else "None is configured: pull one and set embedding_model." 

1060 ) 

1061 raise RuntimeError(f"Ingest needs an embedding model. {detail}") 

1062 

1063 

1064async def _run_post_ingest_passes( 

1065 store: Any, 

1066 *, 

1067 indexed_anything: bool, 

1068 cluster_added: Collection[str], 

1069 cluster_updated: Collection[str], 

1070 cluster_removed: Collection[str], 

1071 touched: set[str], 

1072 cancel: CancelSignal | None, 

1073) -> None: 

1074 """Index maintenance, concept clusters, the wiki hook, and the entity lifecycle. 

1075 

1076 The index and wiki passes only run when this sync indexed something. The 

1077 concept clusters are recomputed only when a source was added, updated or 

1078 removed: a relocation changes no co-occurrence. The entity lifecycle runs 

1079 every sync (a cheap no-op when off or already current) so turning the 

1080 setting on takes effect without a separate operation. 

1081 """ 

1082 if indexed_anything: 

1083 store.ensure_fts_index() 

1084 store.ensure_scalar_indexes() 

1085 store.ensure_vector_index() 

1086 store.optimize_sources() 

1087 if cluster_added or cluster_updated or cluster_removed: 

1088 await _rebuild_concept_clusters(cluster_added, cluster_updated, cluster_removed) 

1089 if indexed_anything: 

1090 await _update_wiki(touched, active_config()) 

1091 

1092 from lilbee.retrieval.entities.lifecycle import ensure_entities 

1093 

1094 await to_ingest_thread(ensure_entities, cancel) 

1095 

1096 

1097async def _update_wiki(changed_sources: set[str], config: Config) -> None: 

1098 """Refresh the wiki index, and regenerate pages when auto-update is on. 

1099 

1100 The index refresh runs whenever the wiki is enabled, because it spends no 

1101 LLM call and it is what lets the browse tree list a page the moment its 

1102 document lands. Regeneration is the expensive half and stays behind 

1103 wiki_auto_update. 

1104 

1105 Best effort: the ingest itself already succeeded and `lilbee wiki update` 

1106 re-runs the regeneration, so a wiki failure must not skip the post-ingest 

1107 entity pass or the reconciliation guard. 

1108 """ 

1109 if not config.wiki: 

1110 return 

1111 # circular: lilbee.wiki imports lilbee.data.ingest.file_hash, so the 

1112 # post-ingest hook stays function-local at this boundary. 

1113 from lilbee.wiki.ingest import incremental_update 

1114 from lilbee.wiki.stubs import refresh_stub_index 

1115 

1116 try: 

1117 await to_ingest_thread( 

1118 refresh_stub_index, get_services().store, config, sources=changed_sources 

1119 ) 

1120 except Exception: 

1121 log.warning("Wiki index refresh failed after sync", exc_info=True) 

1122 

1123 if not config.wiki_auto_update: 

1124 return 

1125 try: 

1126 await incremental_update(changed_sources, config) 

1127 except Exception: 

1128 log.warning("Wiki auto-update failed after sync", exc_info=True) 

1129 

1130 

1131def _worker_failure_message(failures: list[ShardDone], specs: list[ShardSpec]) -> str: 

1132 """What a failed fan-out reports; its workers' shards are kept for the re-run.""" 

1133 detail = "; ".join( 

1134 f"worker {failure.index} ({specs[failure.index].config.data_root / WORKER_LOG_NAME}): " 

1135 f"{failure.error}" 

1136 for failure in failures 

1137 ) 

1138 return ( 

1139 f"{len(failures)} ingest worker(s) failed, so the index was not updated: {detail}. " 

1140 "Their work is kept: re-running sync continues from where they stopped." 

1141 ) 

1142 

1143 

1144def _merge_worker_shards(store: Store, specs: list[ShardSpec], touched: set[str]) -> None: 

1145 """Fold every worker's shard into this index. 

1146 

1147 A store with no chunks of its own takes the shards whole; one that already 

1148 holds a corpus takes only the sources this run touched, so a re-sync replaces 

1149 those rows instead of appending a second copy of everything. 

1150 """ 

1151 from lilbee.data.store.shard_merge import merge_shards 

1152 

1153 scope = touched if store.has_chunks() else None 

1154 merge_shards(store, [spec.config.lancedb_dir for spec in specs], sources=scope) 

1155 

1156 

1157async def _sync_across_workers( 

1158 specs: list[ShardSpec], 

1159 store: Store, 

1160 *, 

1161 prune_ignored: bool = False, 

1162 options: ShardOptions, 

1163 quiet: bool, 

1164 on_progress: DetailedProgressCallback, 

1165 cancel: CancelSignal | None, 

1166) -> SyncResult: 

1167 """Ingest on one worker per GPU, then fold their shards into the one index.""" 

1168 verdicts = await run_workers( 

1169 specs, options=options, quiet=quiet, on_progress=on_progress, cancel=cancel 

1170 ) 

1171 if cancel is not None and cancel.is_set(): 

1172 raise asyncio.CancelledError 

1173 if failures := [verdict for verdict in verdicts if verdict.error is not None]: 

1174 raise RuntimeError(_worker_failure_message(failures, specs)) 

1175 result = aggregate_results(verdicts) 

1176 touched = set(result.added) | set(result.updated) | set(result.relocated) 

1177 await to_ingest_thread(_merge_worker_shards, store, specs, touched) 

1178 # No worker sees the whole corpus, so each one leaves this pass to the parent. 

1179 if prune_ignored: 

1180 result.removed = await to_ingest_thread( 

1181 _forget_ignored, store.get_sources(), IgnoreRules.for_corpus() 

1182 ) 

1183 await _run_post_ingest_passes( 

1184 store, 

1185 indexed_anything=bool(touched), 

1186 cluster_added=result.added, 

1187 cluster_updated=result.updated, 

1188 cluster_removed=result.removed, 

1189 touched=touched, 

1190 cancel=cancel, 

1191 ) 

1192 on_progress( 

1193 EventType.SYNC_DONE, 

1194 SyncDoneEvent( 

1195 added=len(result.added), 

1196 updated=len(result.updated), 

1197 removed=len(result.removed), 

1198 failed=len(result.failed), 

1199 skipped=len(result.skipped), 

1200 relocated=len(result.relocated), 

1201 ), 

1202 ) 

1203 return result 

1204 

1205 

1206_SyncParams = ParamSpec("_SyncParams") 

1207 

1208 

1209def _marks_sync_running( 

1210 run: Callable[_SyncParams, Coroutine[Any, Any, SyncResult]], 

1211) -> Callable[_SyncParams, Coroutine[Any, Any, SyncResult]]: 

1212 """Hold the data root's sync mark for the whole run, so a reset refuses meanwhile.""" 

1213 

1214 @functools.wraps(run) 

1215 async def _marked(*args: _SyncParams.args, **kwargs: _SyncParams.kwargs) -> SyncResult: 

1216 async with sync_running(active_config().data_root): 

1217 return await run(*args, **kwargs) 

1218 

1219 return _marked 

1220 

1221 

1222@_marks_sync_running 

1223async def sync( 

1224 force_rebuild: bool = False, 

1225 quiet: bool = False, 

1226 *, 

1227 on_progress: DetailedProgressCallback = noop_callback, 

1228 cancel: CancelSignal | None = None, 

1229 retry_skipped: bool = False, 

1230 prune_ignored: bool = False, 

1231 shard: ShardId | None = None, 

1232) -> SyncResult: 

1233 """Sync documents/ with the vector store. 

1234 Returns a SyncResult with the added/updated/removed/unchanged/failed/skipped lists. 

1235 When *quiet* is True, the Rich progress bar is suppressed (for JSON output). 

1236 When *cancel* is set mid-run, planning and processing stop between files 

1237 without data loss (completed work is flushed) and CancelledError is raised; 

1238 a cancel already set on entry returns an empty result instead. 

1239 When *retry_skipped* is set, the failed-file skip markers are cleared so this 

1240 sync attempts those files again; *force_rebuild* clears removals too. 

1241 When *prune_ignored* is set, sources a ``.lilbeeignore`` now excludes are 

1242 dropped from the index. Off by default: the patterns govern what sync takes 

1243 in, and removing what a past sync already indexed is the caller's decision. 

1244 A *shard* runs this sync as one worker of a multi-GPU fan-out: it sees only 

1245 its slice of the corpus and leaves the corpus-wide passes to the parent. 

1246 """ 

1247 config = active_config() 

1248 _store = get_services().store 

1249 

1250 if force_rebuild: 

1251 # drop_all + memory re-embedding are heavy blocking store work; run them 

1252 # off the event loop so a rebuild doesn't stall other admitted requests. 

1253 await to_ingest_thread(_force_rebuild_store, _store) 

1254 

1255 config.documents_dir.mkdir(parents=True, exist_ok=True) 

1256 index_mismatch = _report_index_mismatch(_store) 

1257 records_root = _prepare_skip_records( 

1258 shard, force_rebuild=force_rebuild, retry_skipped=retry_skipped 

1259 ) 

1260 

1261 if shard is None and (specs := plan_fanout()): 

1262 merged = await _sync_across_workers( 

1263 specs, 

1264 _store, 

1265 prune_ignored=prune_ignored, 

1266 options=ShardOptions( 

1267 parent_pid=os.getpid(), 

1268 force_rebuild=force_rebuild, 

1269 ), 

1270 quiet=quiet, 

1271 on_progress=on_progress, 

1272 cancel=cancel, 

1273 ) 

1274 return merged.model_copy(update={"index_mismatch": index_mismatch}) 

1275 

1276 rules = IgnoreRules.for_corpus() 

1277 scan = discover_corpus(shard, rules) 

1278 disk_files = scan.files 

1279 sources = _store.get_sources() 

1280 existing_sources = {s["filename"]: s for s in sources} 

1281 skip_markers = load_skip_markers(records_root) 

1282 

1283 failed: dict[str, None] = {} 

1284 # Refused formats start the run skipped; they get no skip marker (no planned hash). 

1285 skipped: dict[str, OcrReport | None] = dict.fromkeys(scan.excluded) 

1286 # filename → why it was skipped/failed (for reporting) 

1287 reasons: dict[str, str] = {name: why.value for name, why in scan.excluded.items()} 

1288 flush_failed: set[str] = set() 

1289 _log_excluded(scan.excluded) 

1290 

1291 # Opt-in, and corpus-wide: a worker sees one slice but the whole sources 

1292 # table, so it leaves this pass to the parent rather than racing its siblings. 

1293 ignored = _forget_ignored(sources, rules) if prune_ignored and shard is None else [] 

1294 # Shard-safe: a worker removes only keys from its own slice. 

1295 refused = _forget_refused(scan.excluded, existing_sources) 

1296 

1297 # Sources whose backing file is not on disk this pass. A vanished file is NOT 

1298 # removed: it stays indexed and searchable, a dead path-link the user 

1299 # discovers only when they try to open it, and the set pairs a reappeared 

1300 # identical file to its old key below. What was just removed leaves the set, 

1301 # where it could otherwise capture a real move. 

1302 gone = set(ignored) | set(refused) 

1303 absent = [name for name in _absent_sources(sources, disk_files) if name not in gone] 

1304 

1305 # The planning pass stats (and where needed hashes) every file on disk, batch by 

1306 # batch off the event loop and overlapped with ingest. A brand-new file whose 

1307 # content hash matches an absent source is folded in per batch as a move, not an 

1308 # add: repointed in place so its chunks and embeddings are reused, not rebuilt. 

1309 state = _StreamedPlan(corpus_total=len(disk_files)) 

1310 added, updated, pending_hashes = state.added, state.updated, state.pending_hashes 

1311 plan_batches = _plan_batches(disk_files, existing_sources, skip_markers, absent, state, cancel) 

1312 

1313 # Snapshot the cumulative truncation counter so the delta over this sync can 

1314 # surface "N chunks truncated" instead of being lost in per-chunk debug logs. 

1315 truncated_before = get_services().embedder.truncated_total 

1316 

1317 # Ingest files (with optional progress bar). Only non-empty batches are yielded, 

1318 # so the first one arriving is what proves there is work to do. 

1319 try: 

1320 first = await anext(plan_batches, None) 

1321 if first is not None: 

1322 # Hold the embed fleet resident for the whole batch: an unevenly loaded 

1323 # replica must not idle-unload and reload cold mid-run (which snowballs 

1324 # into a fleet collapse). The ContextVar propagates into the ingest 

1325 # thread pool, where the fleet actually spawns on the first embed. 

1326 from lilbee.providers.fleet.ingest_warmth import keep_fleet_warm 

1327 

1328 with keep_fleet_warm(): 

1329 _require_embedding_model() 

1330 await ingest_stream( 

1331 _chain_plan_batches(first, plan_batches), 

1332 added, 

1333 updated, 

1334 failed, 

1335 skipped, 

1336 plan=state, 

1337 quiet=quiet, 

1338 on_progress=on_progress, 

1339 cancel=cancel, 

1340 flush_failed=flush_failed, 

1341 reasons=reasons, 

1342 ) 

1343 if cancel is not None and cancel.is_set(): 

1344 # The stream stops feeding on cancel, so ingest can drain its 

1345 # admitted files and return without raising. A cancelled run must 

1346 # not go on to write skip markers or reconcile an unplanned corpus. 

1347 raise asyncio.CancelledError 

1348 finally: 

1349 # Idempotent, and the only close when the stream is never consumed. 

1350 await plan_batches.aclose() 

1351 relocated = state.relocated 

1352 

1353 # A flush failure is a transient store-side problem, not a verdict on the 

1354 # file: leaving it unmarked re-plans it next sync instead of skipping it. 

1355 marker_failed = [name for name in (*failed, *skipped) if name not in flush_failed] 

1356 _persist_skip_records( 

1357 records_root, pending_hashes, reasons, succeeded=[*added, *updated], failed=marker_failed 

1358 ) 

1359 

1360 if shard is None: 

1361 # A worker's shard is merged before the indexes are built, so the passes 

1362 # run once corpus-wide in the parent instead of once per shard. 

1363 await _run_post_ingest_passes( 

1364 _store, 

1365 indexed_anything=bool(state.planned or relocated), 

1366 cluster_added=added, 

1367 cluster_updated=updated, 

1368 cluster_removed=[*ignored, *refused], 

1369 # The old names of relocated sources ride along so the wiki index 

1370 # subtracts them in the same pass that merges their new ones. 

1371 touched=set(added) | set(updated) | set(relocated) | set(state.relocated_from), 

1372 cancel=cancel, 

1373 ) 

1374 

1375 # Reconciliation guard against silent data loss: any on-disk document file that 

1376 # ended up in neither the index nor an accounting set was dropped without a 

1377 # signal. Surface it loudly instead of letting a whole dataset vanish quietly. 

1378 if missing := _reconcile_missing( 

1379 disk_files, _store.get_sources(), failed, skipped, state.held_out 

1380 ): 

1381 log.warning( 

1382 "Sync reconciliation: %d document file(s) on disk are absent from the index " 

1383 "with no failure reported (possible silent drop): %s", 

1384 len(missing), 

1385 ", ".join(missing[:20]), 

1386 ) 

1387 

1388 result = SyncResult( 

1389 added=list(added), 

1390 updated=list(updated), 

1391 removed=ignored + refused, 

1392 unchanged=state.unchanged, 

1393 relocated=relocated, 

1394 failed=list(failed), 

1395 skipped=list(skipped), 

1396 skipped_ocr={name: ocr for name, ocr in skipped.items() if ocr is not None}, 

1397 held_out=describe_skips(records_root, _failures_among(records_root, state.held_out)), 

1398 truncated=get_services().embedder.truncated_total - truncated_before, 

1399 index_mismatch=index_mismatch, 

1400 ) 

1401 on_progress( 

1402 EventType.SYNC_DONE, 

1403 SyncDoneEvent( 

1404 added=len(result.added), 

1405 updated=len(result.updated), 

1406 removed=len(result.removed), 

1407 failed=len(result.failed), 

1408 skipped=len(result.skipped), 

1409 relocated=len(result.relocated), 

1410 ), 

1411 ) 

1412 return result 

1413 

1414 

1415def _phase_progress_callback( 

1416 progress: Progress, ptask: Any, chain: DetailedProgressCallback 

1417) -> DetailedProgressCallback: 

1418 """Wrap *chain*, updating the bar's description on per-page / per-chunk events. 

1419 

1420 EXTRACT (page i/N, named for the OCR backend that ran), OCR_START (Tesseract 

1421 running on a file) and EMBED (chunk i/N) events would otherwise leave the bar 

1422 frozen between file completions; surfacing them on the spinner description 

1423 keeps a single large file's row visibly moving. All events still forward to 

1424 *chain* so the caller's own callback (TUI / JSON) is unaffected. 

1425 """ 

1426 

1427 def _callback(event_type: EventType, data: ProgressEvent) -> None: 

1428 if event_type is EventType.EXTRACT and isinstance(data, ExtractEvent): 

1429 progress.update( 

1430 ptask, description=f"{data.step} {data.file} (page {data.page}/{data.total_pages})" 

1431 ) 

1432 elif event_type is EventType.OCR_START and isinstance(data, OcrStartEvent): 

1433 progress.update(ptask, description=data.status_text) 

1434 elif event_type is EventType.EMBED and isinstance(data, EmbedEvent): 

1435 progress.update( 

1436 ptask, description=f"Embedding {data.file} ({data.chunk}/{data.total_chunks})" 

1437 ) 

1438 chain(event_type, data) 

1439 

1440 return _callback 

1441 

1442 

1443def _ingest_progress_bar(quiet: bool) -> Progress: 

1444 """The transient bar an in-process ingest reports on, disabled when *quiet*.""" 

1445 return Progress( 

1446 SpinnerColumn(), 

1447 literal_text_column("{task.description}"), 

1448 BarColumn(), 

1449 MofNCompleteColumn(), 

1450 TimeElapsedColumn(), 

1451 transient=True, 

1452 disable=quiet, 

1453 ) 

1454 

1455 

1456# In-flight task cap, as a multiple of _max_concurrent(): enough queued tasks to 

1457# keep every compute slot fed, without materializing one task object per file. 

1458_TASK_WINDOW_MULTIPLIER = 2 

1459 

1460 

1461def _build_admission( 

1462 baseline: int, pages_done: list[int] 

1463) -> tuple[asyncio.Semaphore | ResizableGate, int, asyncio.Task[None] | None]: 

1464 """The batch's admission control, plus its task-window size and controller task. 

1465 

1466 Static mode (the default) returns a fixed semaphore and no controller. Adaptive 

1467 mode, when a GPU fleet is present to feed, returns a resizable gate and a running 

1468 :class:`AdaptiveController` that tunes it toward this box's throughput knee; with 

1469 no fleet it falls back to the static path so a GPU-less host is never affected. 

1470 """ 

1471 profile = profile_for(resolve_mode()) 

1472 devices = enumerate_fleet_devices() if profile is not None else [] 

1473 if profile is None or not devices: 

1474 return asyncio.Semaphore(baseline), baseline * _TASK_WINDOW_MULTIPLIER, None 

1475 permit_max = max_workers() 

1476 gate = ResizableGate(min(baseline, permit_max)) 

1477 controller = AdaptiveController( 

1478 gate, 

1479 profile, 

1480 make_signal_sampler(devices), 

1481 lambda: pages_done[0], 

1482 permit_min=1, 

1483 permit_max=permit_max, 

1484 ) 

1485 task = asyncio.ensure_future(controller.run()) 

1486 # warning, not info: the default LILBEE_LOG_LEVEL is WARNING, so the 

1487 # auto-chosen concurrency would otherwise never surface on a headless sync. 

1488 log.warning( 

1489 "Adaptive ingest concurrency (%s): start %d, max %d", profile.name, gate.limit, permit_max 

1490 ) 

1491 return gate, permit_max * _TASK_WINDOW_MULTIPLIER, task 

1492 

1493 

1494def _failed_result( 

1495 exc: Exception, 

1496 entry: FileToProcess, 

1497 *, 

1498 pages_done: list[int], 

1499 on_progress: DetailedProgressCallback, 

1500 cancel: CancelSignal | None, 

1501) -> _IngestResult: 

1502 """A file's failure as a result, or cancellation when the run is stopping. 

1503 

1504 During shutdown, worker pools raise RuntimeError from submit(). Those are 

1505 cancellation, not ingest failures: the cancel flag is the source of truth, 

1506 and the executor's shutdown message covers the race where cancel was set 

1507 after the submit. 

1508 """ 

1509 if (cancel and cancel.is_set()) or is_executor_shutdown(exc): 

1510 raise asyncio.CancelledError from exc 

1511 # Suppress TaskCancelledError on the FILE_DONE notice: the user already 

1512 # cancelled, and re-raising here would strand sibling tasks awaiting in 

1513 # _collect_results. 

1514 with contextlib.suppress(TaskCancelledError): 

1515 on_progress(EventType.FILE_DONE, FileDoneEvent(file=entry.name, status="error", chunks=0)) 

1516 pages_done[0] += 1 # cleared the gate (as a failure); still a throughput tick 

1517 return _IngestResult(entry.name, entry.path, 0, error=exc) 

1518 

1519 

1520async def _archive_result( 

1521 entry: FileToProcess, on_progress: DetailedProgressCallback, pages_done: list[int] 

1522) -> _IngestResult: 

1523 """Ingest an archive: its members become sources, the archive row keeps the disk stat.""" 

1524 members = await ingest_archive( 

1525 entry.path, entry.name, entry.content_type, on_progress=on_progress 

1526 ) 

1527 concept_batches = [await build_concept_records(m.records, m.name) for m in members] 

1528 entity_rows = [ 

1529 row for m in members for row in (await build_entity_records(m.records, m.name) or []) 

1530 ] 

1531 chunk_total = sum(len(m.records) for m in members) 

1532 on_progress( 

1533 EventType.FILE_DONE, FileDoneEvent(file=entry.name, status="ok", chunks=chunk_total) 

1534 ) 

1535 pages_done[0] += max(1, sum(len(m.page_texts) for m in members)) 

1536 found = [batch for batch in concept_batches if batch is not None] 

1537 return _IngestResult( 

1538 entry.name, 

1539 entry.path, 

1540 chunk_total, 

1541 error=None, 

1542 file_hash=entry.file_hash, 

1543 records=[], 

1544 needs_cleanup=entry.needs_cleanup, 

1545 page_texts=[], 

1546 stat=entry.stat, 

1547 concept_records=ConceptRecords.merged(found) if found else None, 

1548 entity_rows=entity_rows or None, 

1549 meta=SourceMeta(title=derive_title(entry.name)), 

1550 members=members, 

1551 ) 

1552 

1553 

1554def _over_limit_result( 

1555 exc: ChunkLimitError, 

1556 entry: FileToProcess, 

1557 *, 

1558 pages_done: list[int], 

1559 on_progress: DetailedProgressCallback, 

1560) -> _IngestResult: 

1561 """A file over the per-file chunk limit, recorded as skipped and skip-marked at its hash.""" 

1562 log.warning("Skipped %s: %s", entry.name, exc) 

1563 with contextlib.suppress(TaskCancelledError): 

1564 on_progress(EventType.FILE_DONE, FileDoneEvent(file=entry.name, status="skipped", chunks=0)) 

1565 pages_done[0] += 1 

1566 return _IngestResult( 

1567 entry.name, 

1568 entry.path, 

1569 0, 

1570 error=None, 

1571 file_hash=entry.file_hash, 

1572 skip_reason=str(exc), 

1573 needs_cleanup=entry.needs_cleanup, 

1574 ) 

1575 

1576 

1577async def _stream_tasks( 

1578 plan_batches: AsyncGenerator[list[FileToProcess]], 

1579 make_task: Callable[[FileToProcess, int], Coroutine[Any, Any, _IngestResult]], 

1580) -> AsyncGenerator[list[Coroutine[Any, Any, _IngestResult]]]: 

1581 """Each plan batch's per-file coroutines, in plan order.""" 

1582 index = count(1) 

1583 async for plan_batch in plan_batches: 

1584 yield [make_task(entry, next(index)) for entry in plan_batch] 

1585 

1586 

1587async def _chain_plan_batches( 

1588 first: list[FileToProcess], rest: AsyncGenerator[list[FileToProcess]] 

1589) -> AsyncGenerator[list[FileToProcess]]: 

1590 """Yield an already-pulled plan batch, then the remainder of its stream.""" 

1591 yield first 

1592 async for plan_batch in rest: 

1593 yield plan_batch 

1594 

1595 

1596async def ingest_stream( 

1597 plan_batches: AsyncGenerator[list[FileToProcess]], 

1598 added: dict[str, None], 

1599 updated: dict[str, None], 

1600 failed: dict[str, None], 

1601 skipped: dict[str, OcrReport | None], 

1602 *, 

1603 plan: _StreamedPlan | None = None, 

1604 quiet: bool = False, 

1605 on_progress: DetailedProgressCallback = noop_callback, 

1606 cancel: CancelSignal | None = None, 

1607 flush_failed: set[str] | None = None, 

1608 reasons: dict[str, str] | None = None, 

1609) -> None: 

1610 """Ingest a stream of planned file batches, optionally showing a Rich progress bar. 

1611 

1612 Files are admitted as their batch is planned, so ingest starts on the first 

1613 batch instead of waiting for the whole corpus to be diffed. Old chunks are 

1614 deleted in the same transaction as the new write, so the two are atomic per 

1615 file. When *cancel* is set, pending files raise CancelledError before starting. 

1616 

1617 *plan* is the bookkeeping the batches were planned into; it carries the corpus 

1618 the run is measured against. Without it progress is reported with no total, 

1619 since a bare stream of batches does not say what corpus it came from. 

1620 """ 

1621 # Honor LILBEE_INGEST_TRACE once per batch: it raises the trace loggers above 

1622 # the default WARNING so per-file extraction lines actually surface. 

1623 configure_trace_from_env() 

1624 warn_if_table_model_ignored() 

1625 # Throughput is measured in OCR pages, not documents: a document's cost scales 

1626 # with its page count (a 500-page scan is 500x a memo), so pages are the unbiased 

1627 # unit of GPU-feeding work for the adaptive controller to hill-climb on. 

1628 pages_done = [0] 

1629 # Sized off the files with no source row yet, not the planned count: the plan 

1630 # streams in batches and its total is unknown until the stream drains, but the 

1631 # pool has to be decided before the first batch is dispatched. Unindexed files 

1632 # are the one part of the plan that is known without diffing, exact for a first 

1633 # ingest or a rebuild and near zero for an incremental sync, so a small sync 

1634 # over a large corpus stays in-process. Undercounting only keeps a run 

1635 # in-process, which is the safe direction. 

1636 admission, window, controller_task = _build_admission(_max_concurrent(), pages_done) 

1637 progress = _ingest_progress_bar(quiet) 

1638 # The corpus the discovery walk found, known before the first batch is planned. 

1639 corpus_total = plan.corpus_total if plan is not None else 0 

1640 ptask = progress.add_task("Ingesting documents...", total=corpus_total or None) 

1641 # Rebound once, so every event the run emits, per file or per batch, passes the 

1642 # bar: a single large file's OCR and embed phases move it between completions. 

1643 on_progress = _phase_progress_callback(progress, ptask, on_progress) 

1644 

1645 async def _process_one(entry: FileToProcess, file_index: int) -> _IngestResult: 

1646 name = entry.name 

1647 async with admission: 

1648 if cancel and cancel.is_set(): 

1649 raise asyncio.CancelledError 

1650 

1651 try: 

1652 on_progress( 

1653 EventType.FILE_START, 

1654 FileStartEvent(file=name, total_files=feed.planned, current_file=file_index), 

1655 ) 

1656 except TaskCancelledError as exc: 

1657 # FILE_START itself can raise the cooperative cancel signal; 

1658 # normalize so _collect_results can drain siblings cleanly. 

1659 raise asyncio.CancelledError from exc 

1660 try: 

1661 # The source's old chunks are deleted in the same locked 

1662 # transaction as the new write (see _flush_writes), so cleanup is 

1663 # carried on the result rather than run eagerly here. 

1664 if entry.content_type in archive_content_types(): 

1665 return await _archive_result(entry, on_progress, pages_done) 

1666 page_texts: list[PageTextRecord] = [] 

1667 records, meta, ocr = await produce_records( 

1668 entry.path, 

1669 name, 

1670 entry.content_type, 

1671 quiet=quiet, 

1672 on_progress=on_progress, 

1673 page_texts_out=page_texts, 

1674 ) 

1675 concept_records = await build_concept_records(records, name) 

1676 entity_rows = await build_entity_records(records, name) 

1677 on_progress( 

1678 EventType.FILE_DONE, 

1679 FileDoneEvent(file=name, status="ok", chunks=len(records)), 

1680 ) 

1681 pages_done[0] += max(1, len(page_texts)) # OCR pages cleared: the throughput signal 

1682 return _IngestResult( 

1683 name, 

1684 entry.path, 

1685 len(records), 

1686 error=None, 

1687 file_hash=entry.file_hash, 

1688 records=records, 

1689 needs_cleanup=entry.needs_cleanup, 

1690 page_texts=page_texts, 

1691 stat=entry.stat, 

1692 concept_records=concept_records, 

1693 entity_rows=entity_rows, 

1694 meta=meta, 

1695 ocr=ocr, 

1696 ) 

1697 except ChunkLimitError as exc: 

1698 return _over_limit_result( 

1699 exc, entry, pages_done=pages_done, on_progress=on_progress 

1700 ) 

1701 except (asyncio.CancelledError, TaskCancelledError) as exc: 

1702 # TaskCancelledError is the TUI's cooperative cancel signal raised 

1703 # by reporter.check_cancelled() inside on_progress; treat it as 

1704 # asyncio cancellation so _collect_results can drain siblings 

1705 # cleanly instead of orphaning their pending exceptions. 

1706 raise asyncio.CancelledError from exc 

1707 except Exception as exc: 

1708 return _failed_result( 

1709 exc, entry, pages_done=pages_done, on_progress=on_progress, cancel=cancel 

1710 ) 

1711 

1712 feed = _ResultFeed(_stream_tasks(plan_batches, _process_one), plan) 

1713 try: 

1714 # extract_batching coalesces extractions into xberg batch calls when the 

1715 # toggle is on (off by default); the per-file collect contract is unchanged. 

1716 with progress: 

1717 async with extract_batching(): 

1718 await _collect_results( 

1719 feed, 

1720 added, 

1721 updated, 

1722 failed, 

1723 skipped, 

1724 window=window, 

1725 on_progress=on_progress, 

1726 progress=progress, 

1727 ptask=ptask, 

1728 flush_failed=flush_failed, 

1729 reasons=reasons, 

1730 ) 

1731 finally: 

1732 # Stop the adaptive controller (if any) before returning: its background 

1733 # loop must not outlive the batch it was tuning. 

1734 if controller_task is not None: 

1735 controller_task.cancel() 

1736 with contextlib.suppress(asyncio.CancelledError): 

1737 await controller_task 

1738 

1739 

1740# Accumulate roughly this many chunks across documents before one batched 

1741# LanceDB write. Bounds buffered-vector memory while amortizing the write lock 

1742# and per-transaction overhead over many documents instead of one write per file. 

1743_WRITE_FLUSH_CHUNKS = 2000 

1744 

1745 

1746class _ResultFeed: 

1747 """Pull-based source of per-file ingest coroutines over a streamed plan. 

1748 

1749 ``take(wait=False)`` hands back an already-planned file without blocking, so 

1750 the collector waits on the planner only when it has nothing left to run. 

1751 Build-vs-buy: an ``asyncio.Queue`` is the stock bounded channel, but it would 

1752 need a separate producer task to pump the plan-batch generator into it and a 

1753 sentinel to close it; pulling ``anext`` on demand keeps the plan stream the 

1754 single driver and needs neither. ``planned`` is the file count seen so far: 

1755 the run's total once the stream is drained. 

1756 """ 

1757 

1758 def __init__( 

1759 self, 

1760 plan_batches: AsyncGenerator[list[Coroutine[Any, Any, _IngestResult]]], 

1761 plan: _StreamedPlan | None = None, 

1762 ) -> None: 

1763 self._plan_batches = plan_batches 

1764 self._buffer: deque[Coroutine[Any, Any, _IngestResult]] = deque() 

1765 self._pull: asyncio.Task[list[Coroutine[Any, Any, _IngestResult]] | None] | None = None 

1766 self._drained = False 

1767 self._plan = plan if plan is not None else _StreamedPlan() 

1768 self.planned = 0 

1769 

1770 @property 

1771 def corpus_total(self) -> int: 

1772 """Files the run's slice holds, or 0 when the caller declared no corpus.""" 

1773 return self._plan.corpus_total 

1774 

1775 @property 

1776 def resolved(self) -> int: 

1777 """Files already disposed of by the plan, which never reach this feed.""" 

1778 return self._plan.resolved 

1779 

1780 def pull(self) -> asyncio.Task[list[Coroutine[Any, Any, _IngestResult]] | None] | None: 

1781 """The in-flight plan-batch prefetch, so a waiting collector wakes when it lands.""" 

1782 return self._pull 

1783 

1784 async def take(self, *, wait: bool) -> Coroutine[Any, Any, _IngestResult] | None: 

1785 """The next planned file, or None when the stream is drained (or, with 

1786 *wait* False, when the next batch is not planned yet).""" 

1787 while not self._buffer: 

1788 if self._drained: 

1789 return None 

1790 if self._pull is None: 

1791 self._pull = asyncio.ensure_future(anext(self._plan_batches, None)) 

1792 if not wait and not self._pull.done(): 

1793 return None 

1794 plan_batch = await self._pull 

1795 self._pull = None 

1796 if plan_batch is None: 

1797 self._drained = True 

1798 return None 

1799 self._buffer.extend(plan_batch) 

1800 self.planned += len(plan_batch) 

1801 return self._buffer.popleft() 

1802 

1803 async def aclose(self) -> None: 

1804 """Close the plan stream and discard files that were never started.""" 

1805 if self._pull is not None: 

1806 self._pull.cancel() 

1807 with contextlib.suppress(asyncio.CancelledError): 

1808 # A batch that landed before the cancel took effect still owns 

1809 # coroutines; recover it so they are closed rather than leaked. 

1810 plan_batch = await self._pull 

1811 if plan_batch: 

1812 self._buffer.extend(plan_batch) 

1813 self._pull = None 

1814 for coro in self._buffer: 

1815 coro.close() 

1816 self._buffer.clear() 

1817 await self._plan_batches.aclose() 

1818 

1819 

1820async def _refill_window( 

1821 in_flight: set[asyncio.Task[_IngestResult]], 

1822 feed: _ResultFeed, 

1823 window: int, 

1824) -> None: 

1825 """Top up the in-flight task set from *feed*, capped at *window* tasks. 

1826 

1827 Waits on the planner only when nothing is running, so a slow batch never 

1828 stalls files that are already planned. 

1829 """ 

1830 while len(in_flight) < window: 

1831 coro = await feed.take(wait=not in_flight) 

1832 if coro is None: 

1833 return 

1834 in_flight.add(asyncio.ensure_future(coro)) 

1835 

1836 

1837async def _next_completions( 

1838 in_flight: set[asyncio.Task[_IngestResult]], prefetch: asyncio.Future[Any] | None 

1839) -> tuple[Iterable[asyncio.Future[Any]], set[asyncio.Task[_IngestResult]]]: 

1840 """Wait for the next file to finish, returning (completed, still running). 

1841 

1842 The feed's plan-batch prefetch waits alongside the running files, so work is 

1843 admitted as soon as it is planned rather than on the next file completion, 

1844 and it is filtered out of the still-running set here since it is not a file. 

1845 """ 

1846 waiting: set[asyncio.Future[Any]] = set(in_flight) 

1847 if prefetch is not None: 

1848 waiting.add(prefetch) 

1849 done, still_running = await asyncio.wait(waiting, return_when=asyncio.FIRST_COMPLETED) 

1850 # Explicit loop, like _cancel_in_flight: Nuitka miscompiled the comprehension 

1851 # form of this task-set filtering. 

1852 remaining: set[asyncio.Task[_IngestResult]] = set() 

1853 for task in still_running: 

1854 if task is not prefetch: 

1855 remaining.add(cast("asyncio.Task[_IngestResult]", task)) 

1856 return done, remaining 

1857 

1858 

1859async def _collect_results( 

1860 feed: _ResultFeed, 

1861 added: dict[str, None], 

1862 updated: dict[str, None], 

1863 failed: dict[str, None], 

1864 skipped: dict[str, OcrReport | None], 

1865 *, 

1866 window: int, 

1867 on_progress: DetailedProgressCallback = noop_callback, 

1868 progress: Progress | None = None, 

1869 ptask: Any = None, 

1870 flush_failed: set[str] | None = None, 

1871 reasons: dict[str, str] | None = None, 

1872) -> None: 

1873 """Run *feed* through a bounded task window, batching writes and progress. 

1874 

1875 At most *window* tasks exist at once: results are consumed as they complete 

1876 and the window is refilled from the feed, so memory stays flat however many 

1877 files a sync covers. Successful files are buffered and flushed to LanceDB in 

1878 batches (one locked transaction per batch) rather than one write per file. 

1879 The buffer is flushed on the way out too -- even on cancel -- so 

1880 completed-but-unwritten work is persisted. On exception (typically 

1881 asyncio.CancelledError from a user cancel), cancel every in-flight sibling 

1882 and await them with ``return_exceptions=True`` so their pending 

1883 CancelledErrors don't surface as "Task exception was never retrieved". 

1884 """ 

1885 buffer: list[_IngestResult] = [] 

1886 buffered_chunks = 0 

1887 completed_count = 0 

1888 to_purge: list[str] = [] 

1889 in_flight: set[asyncio.Task[_IngestResult]] = set() 

1890 try: 

1891 await _refill_window(in_flight, feed, window) 

1892 while in_flight: 

1893 prefetch = feed.pull() 

1894 done, in_flight = await _next_completions(in_flight, prefetch) 

1895 saw_cancel = False 

1896 for fut in done: 

1897 if fut is prefetch: 

1898 continue # a planned batch landing, not a file result 

1899 try: 

1900 result = fut.result() 

1901 except asyncio.CancelledError: 

1902 # A user cancel completes several futures together. Flag it but 

1903 # keep draining `done` so a sibling that genuinely finished in 

1904 # the same batch is still buffered and flushed (the 

1905 # cancel-persists contract), then propagate after the loop. A 

1906 # non-cancel exception still propagates immediately, as before, 

1907 # so a genuine ingest bug surfaces and cancels the siblings. 

1908 saw_cancel = True 

1909 continue 

1910 completed_count += 1 

1911 status = _classify_result(result, added, updated, failed, skipped, reasons) 

1912 if status is BatchStatus.INGESTED: 

1913 buffered_chunks = await _buffer_and_maybe_flush( 

1914 result, 

1915 buffer, 

1916 buffered_chunks, 

1917 added, 

1918 updated, 

1919 failed, 

1920 skipped, 

1921 flush_failed, 

1922 ) 

1923 elif status is BatchStatus.SKIPPED and result.needs_cleanup: 

1924 # Zero-text result is never buffered; collect it for the 

1925 # purge pass (see _purge_emptied_sources). 

1926 to_purge.append(result.name) 

1927 _report_file_progress( 

1928 result, 

1929 status, 

1930 feed.resolved + completed_count, 

1931 feed.corpus_total, 

1932 on_progress, 

1933 progress, 

1934 ptask, 

1935 ) 

1936 if saw_cancel: 

1937 # Completed siblings in this batch are now buffered; propagate the 

1938 # cancel so the finally flushes them and cancels still-running work. 

1939 raise asyncio.CancelledError 

1940 await _refill_window(in_flight, feed, window) 

1941 finally: 

1942 # The inner finally guarantees the sibling cancel even if the flush 

1943 # itself raises (e.g. a cancellation landing on the to_thread await). 

1944 try: 

1945 await to_ingest_thread( 

1946 _flush_writes, buffer, added, updated, failed, skipped, flush_failed 

1947 ) 

1948 await to_ingest_thread(_purge_emptied_sources, to_purge) 

1949 finally: 

1950 try: 

1951 await _cancel_in_flight(in_flight) 

1952 finally: 

1953 # Closing the feed stops the planner behind it, so a cancelled 

1954 # sync does not keep hashing the rest of the corpus. 

1955 await feed.aclose() 

1956 

1957 

1958async def _cancel_in_flight(in_flight: set[asyncio.Task[_IngestResult]]) -> None: 

1959 """Cancel still-running tasks and await them so their CancelledErrors are retrieved.""" 

1960 # Explicit loop: Nuitka miscompiled the comprehension form of this cleanup. 

1961 still_pending = [] 

1962 for t in in_flight: 

1963 if not t.done(): 

1964 still_pending.append(t) 

1965 for task in still_pending: 

1966 task.cancel() 

1967 if still_pending: 

1968 await asyncio.gather(*still_pending, return_exceptions=True) 

1969 

1970 

1971async def _buffer_and_maybe_flush( 

1972 result: _IngestResult, 

1973 buffer: list[_IngestResult], 

1974 buffered_chunks: int, 

1975 added: dict[str, None], 

1976 updated: dict[str, None], 

1977 failed: dict[str, None], 

1978 skipped: dict[str, OcrReport | None], 

1979 flush_failed: set[str] | None, 

1980) -> int: 

1981 """Buffer one ingested file, flushing at the chunk threshold; returns the new count.""" 

1982 buffer.append(result) 

1983 # Zero-chunk files count one unit so the buffer stays bounded. 

1984 buffered_chunks += max(result.chunk_count, 1) 

1985 if buffered_chunks >= _WRITE_FLUSH_CHUNKS: 

1986 await to_ingest_thread(_flush_writes, buffer, added, updated, failed, skipped, flush_failed) 

1987 buffered_chunks = 0 

1988 return buffered_chunks 

1989 

1990 

1991def _report_file_progress( 

1992 result: _IngestResult, 

1993 status: BatchStatus, 

1994 completed_count: int, 

1995 total: int, 

1996 on_progress: DetailedProgressCallback, 

1997 progress: Progress | None, 

1998 ptask: Any, 

1999) -> None: 

2000 """Advance the Rich bar (when present) and emit one BATCH_PROGRESS event. 

2001 

2002 *completed_count* counts every file the pass has disposed of and *total* is 

2003 the corpus the discovery walk found, so the pair answers how much of the 

2004 corpus is done rather than how much of the plan so far is. 

2005 """ 

2006 if progress is not None and ptask is not None: 

2007 desc = f"Ingested {result.name}" if result.error is None else f"Failed {result.name}" 

2008 # Set, not advanced: files the plan resolved without ingest produce no 

2009 # result of their own and would otherwise never reach the bar. 

2010 progress.update(ptask, description=desc, completed=completed_count) 

2011 with contextlib.suppress(TaskCancelledError): 

2012 on_progress( 

2013 EventType.BATCH_PROGRESS, 

2014 BatchProgressEvent( 

2015 file=result.name, 

2016 status=status, 

2017 current=completed_count, 

2018 total=total, 

2019 ), 

2020 ) 

2021 

2022 

2023def _classify_result( 

2024 result: _IngestResult, 

2025 added: dict[str, None], 

2026 updated: dict[str, None], 

2027 failed: dict[str, None], 

2028 skipped: dict[str, OcrReport | None], 

2029 reasons: dict[str, str] | None = None, 

2030) -> BatchStatus: 

2031 """Record a completed file's outcome and return its batch status. 

2032 

2033 Failures, refusals and zero-chunk files are tracked here; a successful file is 

2034 reported as ``INGESTED`` and its chunks are persisted by the batched flush, so 

2035 it stays in ``added`` / ``updated`` until then. When *reasons* is given, the 

2036 human-readable cause is recorded there (filename → reason) for reporting. 

2037 """ 

2038 if result.skip_reason is not None: 

2039 added.pop(result.name, None) 

2040 updated.pop(result.name, None) 

2041 skipped[result.name] = None 

2042 if reasons is not None: 

2043 reasons[result.name] = result.skip_reason 

2044 return BatchStatus.SKIPPED 

2045 if result.error is not None: 

2046 # A traceback here would bleed into the TUI chat pane; the full trace stays at DEBUG. 

2047 log.warning("Failed to ingest %s: %s", result.name, result.error) 

2048 log.debug("Traceback for failed ingest of %s", result.name, exc_info=result.error) 

2049 added.pop(result.name, None) 

2050 updated.pop(result.name, None) 

2051 failed[result.name] = None 

2052 if reasons is not None: 

2053 reasons[result.name] = error_reason(result.error) 

2054 return BatchStatus.FAILED 

2055 if result.chunk_count == 0: 

2056 # No searchable chunks: never report it as added/updated. With no page 

2057 # texts either, nothing is persisted and the file retries next sync. With 

2058 # page texts, it stays INGESTED so its pages persist (export/recon) and it 

2059 # stops replanning, but it is reported as skipped since search can't see it. 

2060 added.pop(result.name, None) 

2061 updated.pop(result.name, None) 

2062 skipped[result.name] = result.ocr 

2063 if reasons is not None: 

2064 reasons[result.name] = ( 

2065 "no text extracted (0 chunks)" 

2066 if not result.page_texts 

2067 else "stored page text only (0 searchable chunks)" 

2068 ) 

2069 return BatchStatus.SKIPPED if not result.page_texts else BatchStatus.INGESTED 

2070 return BatchStatus.INGESTED 

2071 

2072 

2073# Back off briefly before the single flush retry: the usual contender is a 

2074# search-triggered FTS optimize holding the store lock past its 30s timeout. 

2075_FLUSH_RETRY_DELAY_SECONDS = 2.0 

2076 

2077 

2078def _retry_after_lock_timeout(write: Callable[[], object]) -> None: 

2079 """Run one store write, retrying once after a lock timeout.""" 

2080 try: 

2081 write() 

2082 except LockTimeoutError: 

2083 log.warning( 

2084 "Store write lock busy; retrying batch flush in %.0fs", _FLUSH_RETRY_DELAY_SECONDS 

2085 ) 

2086 time.sleep(_FLUSH_RETRY_DELAY_SECONDS) 

2087 write() 

2088 

2089 

2090def _flush_batch(buffer: list[_IngestResult]) -> None: 

2091 """Persist one flush unit in a single locked ``write_chunks_batch`` transaction. 

2092 

2093 Page texts travel inside each :class:`ChunkWrite` so the store writes them 

2094 after the cleanup delete (which clears the source's old page-text rows) and 

2095 before the source row: a page-text failure leaves the row stale and the file 

2096 replans next sync instead of losing its pages forever behind the stat 

2097 short-circuit. The write retries once on a lock timeout. 

2098 """ 

2099 store = get_services().store 

2100 items: list[ChunkWrite] = [] 

2101 stale: list[str] = [] 

2102 for r in buffer: 

2103 digest = r.file_hash or file_hash(r.path) 

2104 items.append( 

2105 ChunkWrite( 

2106 source=r.name, 

2107 file_hash=digest, 

2108 records=cast(list[dict], r.records or []), 

2109 needs_cleanup=r.needs_cleanup, 

2110 stat=r.stat, 

2111 page_texts=cast(list[dict], r.page_texts or []), 

2112 meta=r.meta, 

2113 ) 

2114 ) 

2115 if r.members is None: 

2116 continue 

2117 items.extend(_member_write(m, digest, r.stat) for m in r.members) 

2118 current = {m.name for m in r.members} 

2119 stale.extend(n for n in store.member_sources(r.name) if n not in current) 

2120 if stale: 

2121 store.remove_documents(stale) 

2122 _retry_after_lock_timeout(lambda: store.write_chunks_batch(items)) 

2123 _flush_concept_records(buffer) 

2124 _flush_entity_rows(buffer) 

2125 

2126 

2127def _member_write(member: MemberRecords, digest: str, stat: SourceStat | None) -> ChunkWrite: 

2128 """A member's write item: its own source row, keyed to the archive's hash and stat.""" 

2129 return ChunkWrite( 

2130 source=member.name, 

2131 file_hash=digest, 

2132 records=cast(list[dict], member.records), 

2133 needs_cleanup=True, 

2134 stat=stat, 

2135 page_texts=cast(list[dict], member.page_texts), 

2136 meta=member.meta, 

2137 ) 

2138 

2139 

2140def _flush_concept_records(buffer: list[_IngestResult]) -> None: 

2141 """Write the flush unit's buffered concept rows in one batched pass. 

2142 

2143 Runs after the chunk write so a failed flush (files moved to ``failed`` 

2144 and replanned) never lands concept rows for unwritten chunks. A concept 

2145 write failure is logged and never fails the files, matching the 

2146 per-file extraction failure semantics. 

2147 """ 

2148 batches = [r.concept_records for r in buffer if r.concept_records is not None] 

2149 if not batches: 

2150 return 

2151 try: 

2152 get_services().concepts.write_concept_records(ConceptRecords.merged(batches)) 

2153 except Exception: 

2154 log.warning("Concept indexing failed for %d-file batch", len(batches), exc_info=True) 

2155 

2156 

2157def _flush_entity_rows(buffer: list[_IngestResult]) -> None: 

2158 """Write the flush unit's buffered entity rows in one batched pass. 

2159 

2160 Runs after the chunk write, which also performed the per-source deletes, 

2161 so replacement never leaves a source's stale entity rows behind. A write 

2162 failure is logged and never fails the files, matching concept semantics. 

2163 """ 

2164 rows = [row for r in buffer if r.entity_rows for row in r.entity_rows] 

2165 if not rows: 

2166 return 

2167 try: 

2168 get_services().store.add_entities(rows) 

2169 except Exception: 

2170 log.warning("Entity indexing failed for a %d-row batch", len(rows), exc_info=True) 

2171 

2172 

2173def _purge_emptied_sources(names: list[str]) -> None: 

2174 """Remove the prior index entry for files that now extract to nothing. 

2175 

2176 An already-indexed file edited to yield zero chunks and zero page texts is 

2177 classified SKIPPED and never buffered, so the batched cleanup delete never 

2178 runs and its old chunks and source row would linger in search results. Full 

2179 removal here keeps the index consistent; ``remove_documents`` is a no-op for 

2180 never-indexed (brand-new empty) files, so unindexed inputs cost nothing. 

2181 """ 

2182 if not names: 

2183 return 

2184 get_services().store.remove_documents(names) 

2185 

2186 

2187def _flush_writes( 

2188 buffer: list[_IngestResult], 

2189 added: dict[str, None], 

2190 updated: dict[str, None], 

2191 failed: dict[str, None], 

2192 skipped: dict[str, OcrReport | None], 

2193 flush_failed: set[str] | None = None, 

2194) -> None: 

2195 """Flush the buffered documents to the store; track a write failure. 

2196 

2197 Each buffered file's page texts, chunks, cleanup delete, and source upsert 

2198 are written by :func:`_flush_batch`. If that fails, every file in the batch 

2199 is moved to ``failed`` since its source row did not land, and recorded in 

2200 *flush_failed* so the caller replans them next sync instead of skip-marking 

2201 them; the exception never escapes, so the caller's sibling-cancel and 

2202 skip-marker path always runs. The buffer is cleared either way. 

2203 """ 

2204 if not buffer: 

2205 return 

2206 try: 

2207 _flush_batch(buffer) 

2208 except Exception as exc: 

2209 for r in buffer: 

2210 log.warning("Failed to write %s: %s", r.name, exc) 

2211 added.pop(r.name, None) 

2212 updated.pop(r.name, None) 

2213 # A page-text-only file was pre-marked skipped at classification; on a 

2214 # flush failure it belongs in failed only, never both. 

2215 skipped.pop(r.name, None) 

2216 failed[r.name] = None 

2217 if flush_failed is not None: 

2218 flush_failed.add(r.name) 

2219 finally: 

2220 buffer.clear()