Coverage for src/lilbee/data/extract/document.py: 100%
320 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"""Document extraction: one xberg pass that natively extracts text and OCRs
2scanned pages/images through the registered backend; chunk + embed the result."""
4from __future__ import annotations
6import contextvars
7import logging
8import time
9from collections.abc import AsyncGenerator, Generator, Sequence
10from contextlib import asynccontextmanager, contextmanager
11from dataclasses import dataclass, replace
12from pathlib import Path
13from typing import TYPE_CHECKING, Any, Protocol
15from lilbee.app.services import get_services
16from lilbee.core.config import active_config
17from lilbee.core.config.enums import OcrPageStrategy
18from lilbee.data.offload import to_ingest_thread
19from lilbee.data.store import ChunkType, PageTextRecord, SourceMeta
20from lilbee.data.title import derive_title, source_meta_from_extraction
21from lilbee.data.types import (
22 IMAGE_CONTENT_TYPE,
23 MARKDOWN_OUTPUT,
24 PDF_CONTENT_TYPE,
25 ChunkRecord,
26 DocumentRecords,
27 ExtractMode,
28 MemberRecords,
29 OcrBackendName,
30 OcrReport,
31)
32from lilbee.providers.base import aux_options
33from lilbee.runtime.progress import (
34 DetailedProgressCallback,
35 EventType,
36 ExtractEvent,
37 OcrBackendUsed,
38 OcrStartEvent,
39 noop_callback,
40)
42from .backends.vision_ocr import backend_options_for, ocr_request
43from .batch import active_extract_batcher
44from .chunk import ChunkLimitError, build_chunking_config, chunk_text, enforce_chunk_limit
45from .trace import ExtractionTrace, trace_extraction, trace_log
47if TYPE_CHECKING:
48 from xberg import (
49 ExtractedDocument,
50 ExtractionConfig,
51 LayoutDetectionConfig,
52 OcrConfig,
53 OcrStrategy,
54 PdfConfig,
55 )
57 from .batch import ExtractBatcher
59log = logging.getLogger(__name__)
62class _ExtractedTable(Protocol):
63 """The table fields lilbee indexes as dedicated chunks.
65 Structural, so it is satisfied by both xberg's public ``Table`` and the native
66 type that ``ExtractedDocument.tables`` actually yields.
67 """
69 @property
70 def markdown(self) -> str: ...
72 @property
73 def page_number(self) -> int: ...
76def content_type_to_mode(content_type: str) -> ExtractMode:
77 """Map a content_type to the extraction mode (paginated for PDFs and images)."""
78 if content_type in (PDF_CONTENT_TYPE, IMAGE_CONTENT_TYPE):
79 return ExtractMode.PAGINATED
80 return ExtractMode.MARKDOWN
83def _page_text_record(source: str, page: int, text: str, content_type: str) -> PageTextRecord:
84 """Build one per-page text row for the export dataset."""
85 return PageTextRecord(source=source, page=page, text=text, content_type=content_type)
88_ocr_enable_override: contextvars.ContextVar[bool | None] = contextvars.ContextVar(
89 "lilbee_ocr_enable_override", default=None
90)
91_ocr_timeout_override: contextvars.ContextVar[float | None] = contextvars.ContextVar(
92 "lilbee_ocr_timeout_override", default=None
93)
96# Per-file title for embedding input; a ContextVar (like the OCR overrides) so the
97# call chains need no signature threading and concurrent ingests don't leak titles.
98_embed_title: contextvars.ContextVar[str] = contextvars.ContextVar("lilbee_embed_title", default="")
101@contextmanager
102def _title_scope(title: str) -> Generator[None, None, None]:
103 token = _embed_title.set(title or "")
104 try:
105 yield
106 finally:
107 _embed_title.reset(token)
110def _embed_inputs(texts: list[str], title: str | None = None) -> list[str]:
111 """Embedding inputs, title-prefixed when ``cfg.embed_titles`` is on.
113 Only the vector sees the title; the stored chunk text is unchanged.
114 ``None`` falls back to the scoped per-file title (the OCR chains).
115 """
116 effective = title if title is not None else _embed_title.get()
117 if not effective or not active_config().embed_titles:
118 return texts
119 return [f"{effective}\n{text}" for text in texts]
122# Contextual enrichment: characters of document head shown to the model, and
123# the reply budget for the one situating sentence.
124_ENRICH_HEAD_CHARS = 2000
125_ENRICH_CHUNK_CHARS = 2000
126_ENRICH_MAX_TOKENS = 60
127_ENRICH_PROMPT = (
128 "Document beginning:\n{head}\n\nChunk from the same document:\n{chunk}\n\n"
129 "Write one short sentence situating this chunk within the document, to "
130 "improve search retrieval of the chunk. Answer with only the sentence."
131)
134def _enrich_texts(texts: list[str], doc_head: str, source_name: str) -> list[str]:
135 """Embedding inputs with one LLM-written situating sentence per chunk.
137 Anthropic-style contextual retrieval, opt-in (``cfg.contextual_enrichment``):
138 one generation per chunk, so ingest slows accordingly. Only the vector sees
139 the sentence; stored chunk text and citations stay verbatim. Any failure
140 keeps that chunk's bare text.
141 """
142 if not active_config().contextual_enrichment or not texts:
143 return texts
144 from lilbee.retrieval.reasoning import strip_reasoning
146 provider = get_services().provider
147 head = doc_head[:_ENRICH_HEAD_CHARS]
148 enriched: list[str] = []
149 failed = 0
150 for text in texts:
151 prompt = _ENRICH_PROMPT.format(head=head, chunk=text[:_ENRICH_CHUNK_CHARS])
152 try:
153 response = provider.chat(
154 [{"role": "user", "content": prompt}],
155 stream=False,
156 options=aux_options(_ENRICH_MAX_TOKENS),
157 )
158 lines = strip_reasoning(response.text).strip().splitlines()
159 sentence = lines[0].strip() if lines else ""
160 except Exception:
161 failed += 1
162 sentence = ""
163 enriched.append(f"{sentence}\n{text}" if sentence else text)
164 if failed:
165 log.warning(
166 "Contextual enrichment failed for %d of %d chunks in %s; those embed bare",
167 failed,
168 len(texts),
169 source_name,
170 )
171 return enriched
174def _effective_enable_ocr() -> bool | None:
175 """``cfg.enable_ocr`` unless a per-request OCR override is active.
177 The override is a ContextVar, not a global cfg mutation, so concurrent
178 ingests on the shared HTTP daemon each see their own setting.
179 """
180 override = _ocr_enable_override.get()
181 return active_config().enable_ocr if override is None else override
184def _extraction_timeout_secs() -> int | None:
185 """``cfg.extraction_timeout`` as xberg's per-file cap; None when uncapped.
187 xberg defaults this to 600s. Passing lilbee's own value on every call keeps
188 the cap something a user can see and raise instead of an inherited default.
189 """
190 timeout = active_config().extraction_timeout
191 return timeout if timeout > 0 else None
194def _effective_ocr_timeout() -> float:
195 """``cfg.ocr_timeout`` unless a per-request OCR timeout override is active."""
196 override = _ocr_timeout_override.get()
197 return active_config().ocr_timeout if override is None else override
200@contextmanager
201def ocr_override(
202 enable_ocr: bool | None = None, ocr_timeout: float | None = None
203) -> Generator[None, None, None]:
204 """Scope per-request OCR settings without mutating the global cfg.
206 A ``None`` argument leaves that setting at its cfg default. Each override is
207 isolated to the entering context, so overlapping ingests never clobber one
208 another's OCR config.
209 """
210 tokens: list[tuple[contextvars.ContextVar[Any], contextvars.Token[Any]]] = []
211 try:
212 if enable_ocr is not None:
213 tokens.append((_ocr_enable_override, _ocr_enable_override.set(enable_ocr)))
214 if ocr_timeout is not None:
215 tokens.append((_ocr_timeout_override, _ocr_timeout_override.set(ocr_timeout)))
216 yield
217 finally:
218 for var, token in reversed(tokens):
219 var.reset(token)
222def ocr_backend() -> OcrBackendUsed:
223 """The OCR backend for this extraction; ``enable_ocr`` False wins over a vision model."""
224 return OcrBackendUsed.chosen(_effective_enable_ocr(), active_config().vision_model)
227def _ocr_config(ocr_token: str | None) -> OcrConfig | None:
228 """xberg's OcrConfig for the backend ``ocr_backend`` picks, or None when OCR is off.
230 xberg auto-OCRs only the pages that lack a text layer.
231 """
232 from xberg import OcrConfig
234 config = active_config()
235 backend = ocr_backend()
236 if backend is OcrBackendUsed.NONE:
237 return None
238 if backend is OcrBackendUsed.VISION:
239 options = backend_options_for(ocr_token) if ocr_token else None
240 return OcrConfig(
241 backend=OcrBackendName.LILBEE_VISION,
242 backend_options=options,
243 )
244 # xberg requires a non-empty language list (4.x defaulted to English;
245 # xberg 1.0 errors on an empty one). cfg.ocr_language is validated non-empty.
246 return OcrConfig(
247 backend=OcrBackendName.TESSERACT,
248 language=list(config.ocr_language),
249 )
252def _ocr_strategy() -> OcrStrategy:
253 """xberg's page-selection strategy for cfg.ocr_strategy; auto when OCR is off."""
254 from xberg import OcrStrategy
256 config = active_config()
257 # xberg rejects scanned_pages when OCR is disabled.
258 if _effective_enable_ocr() is False or config.ocr_strategy is OcrPageStrategy.AUTO:
259 return OcrStrategy(OcrPageStrategy.AUTO.value)
260 return OcrStrategy.scanned_pages(config.ocr_scan_confidence)
263def _force_ocr_pages() -> list[int] | None:
264 """cfg.force_ocr_pages for xberg; None when empty or OCR is off."""
265 if _effective_enable_ocr() is False:
266 return None
267 return list(active_config().force_ocr_pages) or None
270def _ocr_force_requested() -> bool:
271 """Whether LILBEE_OCR_FORCE forces vision OCR on every page (targeted re-ingest lever)."""
272 import os
274 return os.environ.get("LILBEE_OCR_FORCE", "").strip().lower() in {"1", "true", "yes"}
277# Header/footer band stripped when layout detection is on: outermost 5%.
278_TOP_MARGIN_FRACTION = 0.05
279_BOTTOM_MARGIN_FRACTION = 0.05
282def _pdf_options() -> PdfConfig | None:
283 """PdfConfig for the enabled opt-in features (tables, layout), or None when all off."""
284 config = active_config()
285 if not (config.table_extraction or config.layout_detection):
286 return None
287 from xberg import PdfConfig
289 kwargs: dict[str, Any] = {}
290 if config.table_extraction:
291 kwargs["extract_tables"] = True
292 if config.layout_detection:
293 kwargs.update(
294 reading_order=True,
295 top_margin_fraction=_TOP_MARGIN_FRACTION,
296 bottom_margin_fraction=_BOTTOM_MARGIN_FRACTION,
297 )
298 return PdfConfig(**kwargs)
301def warn_if_table_model_ignored() -> None:
302 """Warn when table extraction runs with layout_detection off: xberg only
303 applies table_model inside layout detection, so the model is silently ignored.
304 """
305 config = active_config()
306 if config.table_extraction and not config.layout_detection:
307 log.warning(
308 "table_model=%s is ignored while layout_detection is off: tables use "
309 "the native extractor, not the structure model. Enable layout_detection "
310 "to apply the table model.",
311 config.table_model.value,
312 )
315def _layout_config() -> LayoutDetectionConfig | None:
316 """AUTO-strategy layout config when enabled, else None."""
317 config = active_config()
318 if not config.layout_detection:
319 return None
320 from xberg import LayoutDetectionConfig, LayoutStrategy
322 return LayoutDetectionConfig(strategy=LayoutStrategy.AUTO, table_model=config.table_model)
325def extraction_config(mode: ExtractMode, *, ocr_token: str | None = None) -> ExtractionConfig:
326 """Build ExtractionConfig for the given extraction mode."""
327 from xberg import ExtractionConfig, PageConfig
329 # Files are extracted one per call; xberg parallelizes OCR across a document's
330 # pages internally, and cross-file concurrency is the pipeline's semaphore.
331 chunking = build_chunking_config()
332 ocr = _ocr_config(ocr_token)
333 # OCR off sends no OCR block (xberg 1.2.7 OCRs page images under any block, even
334 # disabled; 1.2.9 fixes that) and sets disable_ocr (without it, xberg auto-OCRs a
335 # PDF that has no text layer).
336 disable_ocr = ocr is None
337 # Defeats xberg's text-layer short-circuit; vision path only (GPU re-OCR lever).
338 force_ocr = (
339 _ocr_force_requested() and ocr is not None and ocr.backend == OcrBackendName.LILBEE_VISION
340 )
341 if mode is ExtractMode.PAGINATED:
342 paginated = ExtractionConfig(
343 chunking=chunking,
344 pages=PageConfig(extract_pages=True, insert_page_markers=False),
345 ocr=ocr,
346 disable_ocr=disable_ocr,
347 force_ocr=force_ocr,
348 pdf_options=_pdf_options(),
349 extraction_timeout_secs=_extraction_timeout_secs(),
350 ocr_strategy=_ocr_strategy(),
351 force_ocr_pages=_force_ocr_pages(),
352 )
353 # The layout fields keep xberg's defaults when layout detection is off.
354 layout = _layout_config()
355 if layout is None:
356 return paginated
357 return replace(paginated, layout=layout, use_layout_for_markdown=True)
358 return ExtractionConfig(
359 chunking=chunking,
360 output_format=MARKDOWN_OUTPUT,
361 ocr=ocr,
362 disable_ocr=disable_ocr,
363 force_ocr=force_ocr,
364 extraction_timeout_secs=_extraction_timeout_secs(),
365 ocr_strategy=_ocr_strategy(),
366 force_ocr_pages=_force_ocr_pages(),
367 )
370def make_extract_batcher() -> ExtractBatcher | None:
371 """The extraction batcher for this ingest run, or None when batching is off."""
372 config = active_config()
373 if not config.batch_extraction:
374 return None
375 from .batch import ExtractBatcher
376 from .xberg import aextract_batch
378 return ExtractBatcher(
379 size=config.batch_extraction_size,
380 config_fn=extraction_config,
381 ocr_fn=_ocr_config,
382 batch_fn=aextract_batch,
383 )
386@asynccontextmanager
387async def extract_batching() -> AsyncGenerator[None]:
388 """Activate extraction batching for the enclosed ingest, when the toggle is on.
390 The batcher is set before the block runs so the ingest tasks created inside it
391 inherit it in their copied context; off (the default) is a no-op.
392 """
393 batcher = make_extract_batcher()
394 if batcher is None:
395 yield
396 return
397 from .batch import reset_active_batcher, set_active_batcher
399 token = set_active_batcher(batcher)
400 try:
401 yield
402 finally:
403 await batcher.close()
404 reset_active_batcher(token)
407def _chunk_pages(page_texts: Sequence[tuple[int, str]]) -> list[tuple[int, str]]:
408 """Chunk each page's text. Semantic chunking is off: a single page rarely spans
409 multiple topics, so the semantic round-trip is not worth it."""
410 return [
411 (page_num, chunk)
412 for page_num, text in page_texts
413 for chunk in chunk_text(text, use_semantic=False)
414 ]
417async def chunk_and_embed_pages(
418 page_texts: Sequence[tuple[int, str]],
419 source_name: str,
420 content_type: str,
421 on_progress: DetailedProgressCallback,
422) -> list[ChunkRecord]:
423 """Chunk per-page text and embed every chunk. Used by the dataset import path."""
424 if not page_texts:
425 return []
427 # chunk_text runs xberg's synchronous extractor; offload it so a long
428 # document does not stall sibling files sharing this event loop.
429 all_chunks = await to_ingest_thread(_chunk_pages, page_texts)
430 if not all_chunks:
431 return []
432 texts = [c for _, c in all_chunks]
433 embed_texts = await to_ingest_thread(_enrich_texts, texts, page_texts[0][1], source_name)
434 vectors = await to_ingest_thread(
435 get_services().embedder.embed_batch,
436 _embed_inputs(embed_texts),
437 source=source_name,
438 on_progress=on_progress,
439 )
440 return [
441 ChunkRecord(
442 source=source_name,
443 content_type=content_type,
444 chunk_type=ChunkType.RAW,
445 page_start=page_num,
446 page_end=page_num,
447 line_start=0,
448 line_end=0,
449 chunk=text,
450 chunk_index=i,
451 vector=vec,
452 )
453 for i, ((page_num, text), vec) in enumerate(zip(all_chunks, vectors, strict=True))
454 ]
457def _capture_result_page_texts(
458 doc: ExtractedDocument,
459 source_name: str,
460 content_type: str,
461 page_texts_out: list[PageTextRecord] | None,
462) -> None:
463 """Append an extraction's page texts to the export accumulator.
465 Paginated documents yield one row per ``doc.pages`` entry; others have no
466 page split, so the full ``doc.content`` is recorded as page 0.
467 """
468 if page_texts_out is None:
469 return
470 if doc.pages:
471 page_texts_out.extend(
472 _page_text_record(source_name, page.page_number, page.content, content_type)
473 for page in doc.pages
474 )
475 elif doc.content.strip():
476 page_texts_out.append(_page_text_record(source_name, 0, doc.content, content_type))
479def _document_tables(doc: ExtractedDocument) -> list[_ExtractedTable]:
480 """The result tables to index as dedicated chunks, when table extraction is on.
482 Each table becomes its own chunk carrying xberg's markdown serialization and
483 page metadata. The table's flattened text also stays inside the content
484 chunks: stripping it there would tear holes in reading-order prose and the
485 page-text export, so both are indexed deliberately. The dedicated table
486 chunk adds the structured serialization for targeted retrieval.
487 """
488 if not active_config().table_extraction:
489 return []
490 return [t for t in (doc.tables or []) if t.markdown and t.markdown.strip()]
493def _warn_empty_ocr(source_name: str, media: str, backend: OcrBackendUsed) -> None:
494 """Warn that extraction yielded no text, with advice that matches the backend that ran."""
495 if backend is OcrBackendUsed.NONE:
496 log.warning(
497 "Skipped %s: text extraction produced no usable text. "
498 "OCR is off (enable_ocr = false); set it to true to OCR %s.",
499 source_name,
500 media,
501 )
502 return
503 if backend is OcrBackendUsed.VISION:
504 log.warning(
505 "Skipped %s: the vision model returned no usable text for %s.",
506 source_name,
507 media,
508 )
509 return
510 log.warning(
511 "Skipped %s: text extraction produced no usable text. "
512 "For better results on %s, configure a vision model via PUT /api/models/vision.",
513 source_name,
514 media,
515 )
518async def ingest_document(
519 path: Path,
520 source_name: str,
521 content_type: str,
522 *,
523 quiet: bool = False,
524 on_progress: DetailedProgressCallback = noop_callback,
525 page_texts_out: list[PageTextRecord] | None = None,
526) -> DocumentRecords:
527 """Extract, chunk, and embed a document in a single xberg pass, with its metadata.
529 xberg extracts native text and, where a page has none, OCRs it through the
530 registered backend (lilbee's vision model, or tesseract). Vision OCR progress
531 is streamed per page via ``ocr_request``; Tesseract reports no pages, so a
532 scanned file gets one OCR_START event before its extraction. ``quiet`` is
533 accepted for pipeline call compatibility. The returned metadata carries the document's
534 extraction title/authors/date and is derived even when extraction yields nothing;
535 the OCR report says which backend the extraction ran and how many pages it OCR'd.
536 """
537 del quiet
538 doc, ocr = await _extract_document(
539 path, source_name, content_type, content_type_to_mode(content_type), on_progress
540 )
541 records, meta = await _records_from_document(
542 doc,
543 source_name,
544 content_type,
545 on_progress=on_progress,
546 page_texts_out=page_texts_out,
547 ocr_backend=ocr.backend,
548 )
549 return DocumentRecords(records, meta, ocr)
552async def ingest_archive(
553 path: Path,
554 source_name: str,
555 content_type: str,
556 *,
557 on_progress: DetailedProgressCallback = noop_callback,
558) -> list[MemberRecords]:
559 """Extract an archive once and build records for every member, nested archives included.
561 The archive itself contributes no chunks. Each member is its own source named
562 ``<archive>/<member path>``. xberg unpacks to ``max_archive_depth`` under its
563 zip-bomb limits, so depth and size are enforced before this runs.
564 """
565 doc, ocr = await _extract_document(
566 path, source_name, content_type, ExtractMode.PAGINATED, on_progress
567 )
568 members: list[MemberRecords] = []
569 await _collect_members(doc, source_name, members, on_progress, ocr.backend)
570 return members
573async def _collect_members(
574 doc: ExtractedDocument,
575 prefix: str,
576 members: list[MemberRecords],
577 on_progress: DetailedProgressCallback,
578 ocr_backend: OcrBackendUsed,
579) -> None:
580 # circular: document -> ingest.discovery via the ingest package's pipeline import
581 from lilbee.data.ingest.discovery import archive_content_types, classify_file
583 unsupported: list[str] = []
584 for entry in doc.children or []:
585 name = f"{prefix}/{entry.path}"
586 # The same gate discovery applies on disk: a member of an unsupported format
587 # is never chunked, so an archive of machine output cannot flood the index.
588 content_type = classify_file(Path(entry.path))
589 if content_type is None:
590 unsupported.append(entry.path)
591 continue
592 if content_type in archive_content_types():
593 await _collect_members(entry.result, name, members, on_progress, ocr_backend)
594 continue
595 page_texts: list[PageTextRecord] = []
596 try:
597 records, meta = await _records_from_document(
598 entry.result,
599 name,
600 content_type,
601 on_progress=on_progress,
602 page_texts_out=page_texts,
603 ocr_backend=ocr_backend,
604 )
605 except ChunkLimitError as exc:
606 raise ChunkLimitError(exc.count, exc.limit, member=name) from None
607 members.append(MemberRecords(name, content_type, records, page_texts, meta))
608 if unsupported:
609 log.info(
610 "Skipped %d member(s) of %s, unsupported format: %s",
611 len(unsupported),
612 prefix,
613 ", ".join(sorted(unsupported)),
614 )
617def _page_count_config() -> ExtractionConfig:
618 """A metadata-only ExtractionConfig: OCR and page bodies off, structure kept."""
619 from xberg import ExtractionConfig, PageConfig
621 # No OCR block at all: xberg 1.2.7 OCRs a PDF's page images whenever one is present,
622 # even disabled (fixed in 1.2.9).
623 return ExtractionConfig(
624 pages=PageConfig(extract_pages=False, insert_page_markers=False),
625 disable_ocr=True,
626 enable_quality_processing=False,
627 )
630@dataclass(frozen=True)
631class _PageProbe:
632 """What the metadata-only pass read: the page count and whether any page is a scan."""
634 pages: int = 0
635 has_scanned_pages: bool = False
638def _has_scanned_pages(doc: ExtractedDocument) -> bool:
639 """Whether xberg found a PDF page with no usable text layer."""
640 fmt = doc.metadata.format
641 pdf = fmt.pdf if fmt is not None else None
642 return bool(pdf is not None and pdf.scanned_pages)
645async def _probe_pages(data: bytes, filename: str) -> _PageProbe:
646 """Read *data*'s page count and scanned pages in a metadata-only xberg pass.
648 Costs a second structural parse of *data* with OCR and page bodies both off.
649 Any failure here returns an empty probe rather than raising, so the caller's
650 real extraction still runs and reports its own error.
651 """
652 from .xberg import aextract_document
654 try:
655 doc = await aextract_document(data, filename=filename, config=_page_count_config())
656 except Exception:
657 log.debug("Page-count probe failed for %s; OCR progress total stays unknown", filename)
658 return _PageProbe()
659 return _PageProbe(pages=doc.counts.pages, has_scanned_pages=_has_scanned_pages(doc))
662def _announce_tesseract_ocr(
663 probe: _PageProbe, source_name: str, on_progress: DetailedProgressCallback
664) -> None:
665 """Emit OCR_START when Tesseract will OCR this file; it reports no per-page progress."""
666 if probe.has_scanned_pages and ocr_backend() is OcrBackendUsed.TESSERACT:
667 on_progress(EventType.OCR_START, OcrStartEvent(file=source_name, total_pages=probe.pages))
670async def _extract_document(
671 path: Path,
672 source_name: str,
673 content_type: str,
674 mode: ExtractMode,
675 on_progress: DetailedProgressCallback,
676) -> tuple[ExtractedDocument, OcrReport]:
677 """Run one xberg pass over *path*, with per-page OCR progress and the extraction trace."""
678 from .xberg import aextract_document
680 data = path.read_bytes()
681 probe = _PageProbe()
682 # content_type_to_mode(content_type), not the *mode* argument: ingest_archive
683 # always requests PAGINATED regardless of the archive's own content_type, and
684 # an archive has no single page count to probe for.
685 if content_type_to_mode(content_type) is ExtractMode.PAGINATED:
686 probe = await _probe_pages(data, path.name)
687 _announce_tesseract_ocr(probe, source_name, on_progress)
689 page_seen = 0
691 def _tick() -> None:
692 nonlocal page_seen
693 page_seen += 1
694 on_progress(
695 EventType.EXTRACT,
696 ExtractEvent(
697 file=source_name,
698 page=page_seen,
699 total_pages=probe.pages,
700 ocr_backend=OcrBackendUsed.VISION,
701 ),
702 )
704 trace_log.debug("extract-start source=%r type=%s", source_name, content_type)
705 backend = ocr_backend()
706 started = time.perf_counter()
707 with ocr_request(on_page=_tick, timeout=_effective_ocr_timeout()) as token:
708 batcher = active_extract_batcher()
709 if batcher is not None:
710 doc = await batcher.submit(mode, data, path.name, token)
711 else:
712 config = extraction_config(mode, ocr_token=token)
713 # xberg's extract is async; awaiting it keeps the OCR page loop off this thread.
714 doc = await aextract_document(data, filename=path.name, config=config)
715 elapsed = time.perf_counter() - started
716 ocr = OcrReport(backend=backend, pages=_ocr_page_count(doc))
718 # One trace line per extraction (filename, timing, counts, OCR pages), plus a
719 # vision line for scanned files. Emitted for empty results too (a slow file
720 # that yields nothing is worth surfacing).
721 trace_extraction(
722 ExtractionTrace(
723 source=source_name,
724 content_type=content_type,
725 elapsed_s=elapsed,
726 page_count=len(doc.pages or []) or len(doc.chunks or []),
727 chunk_count=len(doc.chunks or []),
728 ocr=ocr,
729 )
730 )
732 return doc, ocr
735def _ocr_page_count(doc: ExtractedDocument) -> int:
736 """Pages xberg OCR'd, which are the pages it attaches an OCR confidence to."""
737 return sum(1 for page in doc.pages or [] if page.ocr_confidence is not None)
740async def _records_from_document(
741 doc: ExtractedDocument,
742 source_name: str,
743 content_type: str,
744 *,
745 on_progress: DetailedProgressCallback,
746 page_texts_out: list[PageTextRecord] | None,
747 ocr_backend: OcrBackendUsed,
748) -> tuple[list[ChunkRecord], SourceMeta]:
749 """Chunk-cap, page-capture, and embed one extracted document into its records."""
750 # Derived before the empty-result return so a scan's title/authors survive zero chunks.
751 meta = source_meta_from_extraction(doc.metadata, source_name)
753 tables = _document_tables(doc)
754 if not doc.chunks and not tables:
755 if content_type in (PDF_CONTENT_TYPE, IMAGE_CONTENT_TYPE):
756 _warn_empty_ocr(source_name, "scanned documents", ocr_backend)
757 return [], meta
759 enforce_chunk_limit(len(doc.chunks or []) + len(tables))
760 _capture_result_page_texts(doc, source_name, content_type, page_texts_out)
762 # One EXTRACT event per file so progress subscribers show "extracted N pages"
763 # before embedding; result.pages, or the chunk count for non-paginated docs.
764 page_count = len(doc.pages or []) or len(doc.chunks or [])
765 ran = ocr_backend if _ocr_page_count(doc) else OcrBackendUsed.NONE
766 on_progress(
767 EventType.EXTRACT,
768 ExtractEvent(file=source_name, page=page_count, total_pages=page_count, ocr_backend=ran),
769 )
771 # Content chunks and table serializations share one embed batch; the vector
772 # list is split back apart below by position.
773 texts = [chunk.content for chunk in doc.chunks or []]
774 table_texts = [table.markdown for table in tables]
775 embed_texts = await to_ingest_thread(
776 _enrich_texts, texts + table_texts, texts[0] if texts else "", source_name
777 )
778 vectors = await to_ingest_thread(
779 get_services().embedder.embed_batch,
780 _embed_inputs(embed_texts, meta.title),
781 source=source_name,
782 on_progress=on_progress,
783 )
784 records = [
785 ChunkRecord(
786 source=source_name,
787 content_type=content_type,
788 chunk_type=ChunkType.RAW,
789 page_start=chunk.metadata.first_page or 0,
790 page_end=chunk.metadata.last_page or 0,
791 line_start=0,
792 line_end=0,
793 chunk=text,
794 chunk_index=chunk.metadata.chunk_index,
795 vector=vec,
796 )
797 for chunk, text, vec in zip(doc.chunks or [], texts, vectors[: len(texts)], strict=True)
798 ]
799 # Table chunk indices continue after the content chunks so a source's
800 # (source, chunk_index) pairs stay unique.
801 records.extend(
802 ChunkRecord(
803 source=source_name,
804 content_type=content_type,
805 chunk_type=ChunkType.TABLE,
806 page_start=table.page_number,
807 page_end=table.page_number,
808 line_start=0,
809 line_end=0,
810 chunk=text,
811 chunk_index=len(texts) + i,
812 vector=vec,
813 )
814 for i, (table, text, vec) in enumerate(
815 zip(tables, table_texts, vectors[len(texts) :], strict=True)
816 )
817 )
818 return records, meta
821def _markdown_h1(text: str) -> str | None:
822 """The document's leading ``# Heading``, the best title a note carries.
824 Only a top-level ATX heading counts; ``##`` and deeper are sections, not the
825 document title. None when the note opens without one.
826 """
827 for line in text.splitlines():
828 stripped = line.strip()
829 if not stripped:
830 continue
831 if stripped.startswith("# "):
832 return stripped[2:].strip() or None
833 return None
834 return None
837async def ingest_markdown(
838 path: Path,
839 source_name: str,
840 on_progress: DetailedProgressCallback = noop_callback,
841 page_texts_out: list[PageTextRecord] | None = None,
842) -> tuple[list[ChunkRecord], SourceMeta]:
843 """Chunk a markdown file with heading context prepended to each chunk.
845 Each chunk gets the heading hierarchy path (e.g. "# Setup > ## Install")
846 prepended for better retrieval context. When ``page_texts_out`` is given,
847 the full text is appended as page 0 for export. The returned metadata's
848 title is the note's leading ``# Heading`` when it has one, else the stem.
849 """
850 raw_text = await to_ingest_thread(path.read_text, encoding="utf-8", errors="replace")
851 meta = SourceMeta(title=derive_title(source_name, _markdown_h1(raw_text)))
852 if not raw_text.strip():
853 return [], meta
855 # chunk_text runs xberg's synchronous extractor; offload it so a large
856 # markdown doc does not stall sibling files sharing this event loop.
857 texts = await to_ingest_thread(
858 chunk_text, raw_text, mime_type="text/markdown", heading_context=True
859 )
860 if not texts:
861 return [], meta
863 enforce_chunk_limit(len(texts))
864 if page_texts_out is not None:
865 page_texts_out.append(_page_text_record(source_name, 0, raw_text, "text"))
867 embed_texts = await to_ingest_thread(_enrich_texts, texts, raw_text, source_name)
868 vectors = await to_ingest_thread(
869 get_services().embedder.embed_batch,
870 _embed_inputs(embed_texts, meta.title),
871 source=source_name,
872 on_progress=on_progress,
873 )
874 records = [
875 ChunkRecord(
876 source=source_name,
877 content_type="text",
878 chunk_type=ChunkType.RAW,
879 page_start=0,
880 page_end=0,
881 line_start=0,
882 line_end=0,
883 chunk=t,
884 chunk_index=idx,
885 vector=vec,
886 )
887 for idx, (t, vec) in enumerate(zip(texts, vectors, strict=True))
888 ]
889 return records, meta