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
« 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."""
3from __future__ import annotations
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
28from rich.progress import (
29 BarColumn,
30 MofNCompleteColumn,
31 Progress,
32 SpinnerColumn,
33 TimeElapsedColumn,
34)
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
140log = logging.getLogger(__name__)
143def _max_concurrent() -> int:
144 """Files allowed in their compute phase at once.
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
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())
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.
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
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)
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.
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
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
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
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.
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
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
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.
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 )
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)
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())
323def _stat_unchanged(stored: SourceStat, current: SourceStat) -> bool:
324 """Whether the stored stat proves the file unchanged without hashing it.
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
337@dataclass(frozen=True)
338class _FileChangeVerdict:
339 """One file's sync verdict: process it, hold it out on its skip marker, or unchanged.
341 A held file is not in the index, so it is never counted as unchanged.
342 """
344 to_process: FileToProcess | None = None
345 backfill: SourceStatBackfill | None = None
346 is_update: bool = False
347 held: bool = False
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 )
394def _plan_workers() -> int:
395 """Worker count for the parallel planning pass: config override, else auto.
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()
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
412class _PlanProgress:
413 """Periodic progress for the plan/hash pass, with rate and ETA.
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 """
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
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 )
445class _StreamStop:
446 """Stop signal for a streamed plan: the caller's cancel, or the stream closing.
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 """
453 def __init__(self, cancel: CancelSignal | None) -> None:
454 self._cancel = cancel
455 self._closed = threading.Event()
457 def close(self) -> None:
458 self._closed.set()
460 def is_set(self) -> bool:
461 return self._closed.is_set() or (self._cancel is not None and self._cancel.is_set())
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
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.
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 """
509 def _classify(name: str, path: Path) -> _FileChangeVerdict:
510 return _classify_file_change(name, path, existing_sources.get(name), skip_markers)
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)
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 {})
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.
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.
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 )
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)
595@dataclass(frozen=True)
596class _Move:
597 """One relocated source: its old key, its new key, and the new file's stat."""
599 old: str
600 new: str
601 stat: SourceStat | None
604class _MovePool:
605 """Absent sources indexed by content hash, consumed as moves are paired.
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 """
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
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
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.
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
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.
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)
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.
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 ]
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
689def _plan_batch_bounds(total: int) -> Iterator[tuple[int, int]]:
690 """(start, stop) slices covering *total* files, doubling up to the cap.
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)
704@dataclass
705class _StreamedPlan:
706 """Bookkeeping a streamed plan accumulates across its batches.
708 ``added`` and ``updated`` are the dicts the ingest pass mutates as files land.
709 """
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
730 @property
731 def resolved(self) -> int:
732 """Files the plan disposed of without ingest, as it disposes of them.
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)
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.
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)
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 )
775 state.pending_hashes.update((entry.name, entry.file_hash) for entry in entries)
776 state.planned += len(entries)
777 return entries
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.
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")
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 )
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)
843def detect_pending() -> int:
844 """Count files in documents/ that are out of sync with the store.
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)
865# Refused files named in the log line before it falls back to a count.
866_EXCLUDED_LOG_SAMPLE = 5
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)
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.
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)
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.
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
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.
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))}
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))
937 update_skip_records(records_root, _merge)
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]
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.
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()
960def _force_rebuild_store(store: Any) -> None:
961 """Drop the store and re-embed the preserved memories table (blocking).
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))
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.
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)
995def _ignored_sources(sources: list[SourceRecord], rules: IgnoreRules) -> list[str]:
996 """Indexed sources a ``.lilbeeignore`` now excludes.
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
1014def _forget_ignored(sources: list[SourceRecord], rules: IgnoreRules) -> list[str]:
1015 """Drop sources the patterns now exclude from the index. Returns what went.
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
1028 removed = list(get_services().store.remove_documents(names).removed)
1029 forget_removed_from_wiki_index(removed)
1030 return removed
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
1042 removed = list(get_services().store.remove_documents(names).removed)
1043 forget_removed_from_wiki_index(removed)
1044 return removed
1047def _require_embedding_model() -> None:
1048 """Refuse ingest without an embedding model.
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}")
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.
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())
1092 from lilbee.retrieval.entities.lifecycle import ensure_entities
1094 await to_ingest_thread(ensure_entities, cancel)
1097async def _update_wiki(changed_sources: set[str], config: Config) -> None:
1098 """Refresh the wiki index, and regenerate pages when auto-update is on.
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.
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
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)
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)
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 )
1144def _merge_worker_shards(store: Store, specs: list[ShardSpec], touched: set[str]) -> None:
1145 """Fold every worker's shard into this index.
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
1153 scope = touched if store.has_chunks() else None
1154 merge_shards(store, [spec.config.lancedb_dir for spec in specs], sources=scope)
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
1206_SyncParams = ParamSpec("_SyncParams")
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."""
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)
1219 return _marked
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
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)
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 )
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})
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)
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)
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)
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]
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)
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
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
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
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 )
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 )
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 )
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
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.
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 """
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)
1440 return _callback
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 )
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
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.
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
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.
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)
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 )
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 )
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]
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
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.
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.
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)
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
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 )
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
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
1746class _ResultFeed:
1747 """Pull-based source of per-file ingest coroutines over a streamed plan.
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 """
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
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
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
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
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()
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()
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.
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))
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).
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
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.
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()
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)
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
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.
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 )
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.
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
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
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()
2090def _flush_batch(buffer: list[_IngestResult]) -> None:
2091 """Persist one flush unit in a single locked ``write_chunks_batch`` transaction.
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)
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 )
2140def _flush_concept_records(buffer: list[_IngestResult]) -> None:
2141 """Write the flush unit's buffered concept rows in one batched pass.
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)
2157def _flush_entity_rows(buffer: list[_IngestResult]) -> None:
2158 """Write the flush unit's buffered entity rows in one batched pass.
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)
2173def _purge_emptied_sources(names: list[str]) -> None:
2174 """Remove the prior index entry for files that now extract to nothing.
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)
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.
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()