Spaces:
Sleeping
Sleeping
| import gc | |
| import inspect | |
| import logging | |
| import math | |
| import os | |
| import pickle | |
| import re | |
| import subprocess | |
| import sys | |
| import tempfile | |
| import time | |
| from pathlib import Path | |
| from typing import Any | |
| from docling.datamodel.accelerator_options import AcceleratorDevice, AcceleratorOptions | |
| from docling.datamodel.base_models import InputFormat | |
| from docling.datamodel.pipeline_options import PdfPipelineOptions | |
| from docling.datamodel.document import TextItem, TableItem, PictureItem | |
| from docling.document_converter import ( | |
| DocumentConverter, | |
| ExcelFormatOption, | |
| PdfFormatOption, | |
| WordFormatOption, | |
| ) | |
| try: | |
| from docling.chunking import HierarchicalChunker | |
| _HAS_HIERARCHICAL_CHUNKER = True | |
| except Exception: | |
| HierarchicalChunker = None | |
| _HAS_HIERARCHICAL_CHUNKER = False | |
| from app.config import PROJECT_ROOT, get_debug_chunks_dir, settings | |
| from app.services.callback import notify_update_sync | |
| log = logging.getLogger(__name__) | |
| # Docling progress_callback is in PR #3042; when merged we get per-phase/page updates | |
| _DOCLING_PHASES = ("BUILD", "ASSEMBLE", "ENRICH") | |
| _PHASE_MESSAGES = { | |
| "BUILD": "Building document...", | |
| "ASSEMBLE": "Assembling layout...", | |
| "ENRICH": "Enriching content...", | |
| } | |
| def _make_docling_progress_callback(update_context: dict[str, Any]) -> Any: | |
| """Build a progress callback for DocumentConverter when the API is available.""" | |
| update_url = update_context.get("update_url") or "" | |
| task_auth_token = update_context.get("task_auth_token") | |
| doc_id = update_context.get("doc_id") or "" | |
| project_id = update_context.get("project_id") | |
| if not update_url or not doc_id: | |
| return None | |
| total_phases = len(_DOCLING_PHASES) | |
| def on_progress(event: Any) -> None: | |
| try: | |
| # PhaseProgressEvent: phase name (BUILD/ASSEMBLE/ENRICH) | |
| phase = getattr(event, "phase", None) | |
| if phase is not None and phase in _DOCLING_PHASES: | |
| finished = _DOCLING_PHASES.index(phase) + 1 | |
| message = _PHASE_MESSAGES.get(phase, str(phase)) | |
| notify_update_sync( | |
| update_url, | |
| task_auth_token, | |
| doc_id, | |
| finished=finished, | |
| total=total_phases, | |
| message=message, | |
| project_id=project_id, | |
| ) | |
| return | |
| # PageProgressEvent: page_no, total (or total_pages) | |
| page_no = getattr(event, "page_no", None) | |
| total_pages = getattr(event, "total", None) or getattr(event, "total_pages", None) | |
| if page_no is not None and total_pages is not None and total_pages > 0: | |
| notify_update_sync( | |
| update_url, | |
| task_auth_token, | |
| doc_id, | |
| finished=page_no, | |
| total=total_pages, | |
| message=f"Page {page_no} of {total_pages}", | |
| project_id=project_id, | |
| ) | |
| except Exception as e: | |
| log.debug("Docling progress callback error: %s", e) | |
| return on_progress | |
| def _converter_supports_progress_callback() -> bool: | |
| sig = getattr(inspect, "signature", None) | |
| if sig is None: | |
| return False | |
| try: | |
| return "progress_callback" in inspect.signature(DocumentConverter.__init__).parameters | |
| except (ValueError, TypeError): | |
| return False | |
| _DEVICE_MAP = { | |
| "auto": AcceleratorDevice.AUTO, | |
| "cpu": AcceleratorDevice.CPU, | |
| "cuda": AcceleratorDevice.CUDA, | |
| "mps": AcceleratorDevice.MPS, | |
| } | |
| def _get_accelerator_options() -> AcceleratorOptions: | |
| raw = (settings.docling_device or "auto").strip().lower() | |
| device = _DEVICE_MAP.get(raw, AcceleratorDevice.AUTO) | |
| return AcceleratorOptions(device=device) | |
| def _normalize_for_search(text: str) -> str: | |
| if not text: | |
| return "" | |
| return " ".join(text.lower().split()) | |
| def _normalize_line_for_repeat_detection(text: str) -> str: | |
| if not text: | |
| return "" | |
| return " ".join(str(text).strip().split()) | |
| def _is_noise_chunk_text(text: str) -> bool: | |
| if not text: | |
| return False | |
| s = text.strip() | |
| if len(s) < 200: | |
| return False | |
| alnum = sum(ch.isalnum() for ch in s) | |
| if alnum >= 25: | |
| return False | |
| allowed = set("-|. \n\r\t") | |
| bad = sum(1 for ch in s if ch not in allowed) | |
| if bad > 0: | |
| return False | |
| return True | |
| def _looks_like_toc(text: str) -> bool: | |
| if not text: | |
| return False | |
| t = " ".join(text.lower().split()) | |
| if "innehållsförteckning" in t or "innehållsförteckning" in t: | |
| return True | |
| if "sida" in t and "§" in t and "." in t: | |
| return True | |
| dot_runs = len(re.findall(r"\.{4,}", text)) | |
| if dot_runs >= 6: | |
| return True | |
| return False | |
| def _clean_dot_leaders(text: str) -> str: | |
| if not text: | |
| return "" | |
| s = text | |
| s = re.sub(r"[.\u00B7·]{4,}", " — ", s) | |
| s = re.sub(r"[ \t]{2,}", " ", s) | |
| s = re.sub(r"\n{3,}", "\n\n", s) | |
| return s.strip() | |
| def _markdown_table_to_text(md: str) -> str: | |
| if not md: | |
| return "" | |
| out_lines: list[str] = [] | |
| for ln in md.splitlines(): | |
| s = ln.strip() | |
| if not s: | |
| continue | |
| if set(s) <= set("-|:"): | |
| continue | |
| if "|" in s: | |
| parts = [p.strip() for p in s.strip("|").split("|")] | |
| parts = [p for p in parts if p] | |
| if not parts: | |
| continue | |
| s = " ".join(parts) | |
| s = _clean_dot_leaders(s) | |
| if s: | |
| out_lines.append(s) | |
| return "\n".join(out_lines).strip() | |
| def _is_noise_markdown_table(md: str) -> bool: | |
| if not md or not _has_markdown_table(md): | |
| return False | |
| lines = [ln.rstrip() for ln in md.splitlines() if ln.strip()] | |
| if not lines: | |
| return False | |
| dash_walls = 0 | |
| dataish = 0 | |
| dot_runs = 0 | |
| for ln in lines: | |
| if len(ln) > 250 and sum(ch.isalnum() for ch in ln) < 5 and set(ln.strip()) <= set("-|"): | |
| dash_walls += 1 | |
| if "|" in ln and sum(ch.isalnum() for ch in ln) >= 5: | |
| dataish += 1 | |
| if re.search(r"\.{4,}", ln): | |
| dot_runs += 1 | |
| if dash_walls >= 1 and dataish <= 1: | |
| return True | |
| if dot_runs >= 3 and dataish <= 2: | |
| return True | |
| return False | |
| def _get_total_pages_hint(doc: Any) -> int | None: | |
| for attr in ("num_pages", "page_count", "n_pages"): | |
| val = getattr(doc, attr, None) | |
| try: | |
| if val is not None: | |
| return int(val) | |
| except Exception: | |
| pass | |
| pages = getattr(doc, "pages", None) | |
| if isinstance(pages, dict) and pages: | |
| try: | |
| keys = [int(k) for k in pages.keys()] | |
| return max(keys) if keys else None | |
| except Exception: | |
| return None | |
| return None | |
| def _pages_from_item(item: Any) -> set[int]: | |
| prov = getattr(item, "prov", None) | |
| if not prov: | |
| return set() | |
| pages: set[int] = set() | |
| for p in prov: | |
| try: | |
| page_no = getattr(p, "page_no", None) | |
| if page_no is not None: | |
| pages.add(int(page_no)) | |
| except Exception: | |
| pass | |
| return {p for p in pages if isinstance(p, int) and p > 0} | |
| def _item_text_for_scan(item: Any, doc: Any, *, light: bool) -> str: | |
| """Text for page-index / header scans. Light mode never calls export_to_markdown.""" | |
| if light: | |
| return (getattr(item, "text", "") or "").strip() | |
| if isinstance(item, (TableItem, PictureItem)): | |
| try: | |
| return (item.export_to_markdown(doc=doc) or "").strip() | |
| except Exception: | |
| pass | |
| return (getattr(item, "text", "") or "").strip() | |
| def _scan_doc_ancillary( | |
| doc: Any, | |
| total_pages_hint: int | None, | |
| *, | |
| collect_page_index: bool, | |
| collect_fragments: bool, | |
| collect_tables: bool, | |
| light: bool, | |
| ) -> tuple[dict[int, str], set[str], dict[int, list[str]]]: | |
| """Single pass over doc items for page index, repeated headers, and tables.""" | |
| page_parts: dict[int, list[str]] = {} | |
| fragments_to_pages: dict[str, set[int]] = {} | |
| tables_by_page: dict[int, list[str]] = {} | |
| for item, _level in doc.iterate_items(): | |
| pages = _pages_from_item(item) | |
| if not pages: | |
| continue | |
| t = _item_text_for_scan(item, doc, light=light) | |
| if not t: | |
| continue | |
| if collect_page_index: | |
| for pg in sorted(pages): | |
| page_parts.setdefault(pg, []).append(t) | |
| if collect_fragments: | |
| t_norm = _normalize_line_for_repeat_detection(t) | |
| if t_norm and 4 <= len(t_norm) <= 140: | |
| fragments_to_pages.setdefault(t_norm, set()).update(pages) | |
| if collect_tables and _has_markdown_table(t): | |
| for pg in sorted(pages): | |
| tables_by_page.setdefault(pg, []).append(t) | |
| page_index = { | |
| pg: _normalize_for_search("\n".join(parts)) for pg, parts in page_parts.items() | |
| } | |
| total_pages = int(total_pages_hint or 0) | |
| min_pages = 3 | |
| if total_pages > 0: | |
| min_pages = max(min_pages, int(math.ceil(total_pages * 0.35))) | |
| repeated_fragments = { | |
| frag for frag, pgs in fragments_to_pages.items() if len(pgs) >= min_pages | |
| } | |
| return page_index, repeated_fragments, tables_by_page | |
| def _build_page_text_index(doc: Any) -> dict[int, str]: | |
| page_index, _, _ = _scan_doc_ancillary( | |
| doc, | |
| _get_total_pages_hint(doc), | |
| collect_page_index=True, | |
| collect_fragments=False, | |
| collect_tables=False, | |
| light=False, | |
| ) | |
| return page_index | |
| def _detect_repeated_fragments(doc: Any, total_pages_hint: int | None) -> set[str]: | |
| _, repeated, _ = _scan_doc_ancillary( | |
| doc, | |
| total_pages_hint, | |
| collect_page_index=False, | |
| collect_fragments=True, | |
| collect_tables=False, | |
| light=False, | |
| ) | |
| return repeated | |
| _PAGE_MARKER_PATTERNS = [ | |
| re.compile(r"^\s*\d+\s*\(\s*\d+\s*\)\s*$"), | |
| re.compile(r"^\s*\d+\s*/\s*\d+\s*$"), | |
| re.compile(r"^\s*page\s+\d+\s+of\s+\d+\s*$", re.IGNORECASE), | |
| ] | |
| def _strip_headers_footers_for_inference(text: str, repeated_fragments: set[str]) -> tuple[str, int]: | |
| if not text: | |
| return "", 0 | |
| removed = 0 | |
| out_lines: list[str] = [] | |
| for ln in text.splitlines(): | |
| raw = ln.rstrip("\n") | |
| s = raw.strip() | |
| if not s: | |
| out_lines.append(raw) | |
| continue | |
| s_norm = _normalize_line_for_repeat_detection(s) | |
| if s_norm in repeated_fragments: | |
| removed += 1 | |
| continue | |
| if any(pat.match(s) for pat in _PAGE_MARKER_PATTERNS): | |
| removed += 1 | |
| continue | |
| out_lines.append(raw) | |
| return "\n".join(out_lines).strip(), removed | |
| def _infer_pages_from_text(text: str, page_index: dict[int, str], *, max_span_pages: int = 3) -> list[int]: | |
| if not text or not page_index: | |
| return [] | |
| lines = [ln.strip() for ln in text.splitlines() if ln.strip()] | |
| body = "\n".join(lines).strip() | |
| if not body: | |
| return [] | |
| candidates: list[str] = [] | |
| if len(lines) >= 6: | |
| mid_start = max(0, len(lines) // 2 - 2) | |
| mid_end = min(len(lines), mid_start + 5) | |
| candidates.append(" ".join(lines[mid_start:mid_end])) | |
| if len(body) > 240: | |
| mid = len(body) // 2 | |
| candidates.append(body[max(0, mid - 120): mid + 120]) | |
| candidates.append(body[:220]) | |
| candidates.append(body[-220:]) | |
| anchors: list[str] = [] | |
| for c in candidates: | |
| n = _normalize_for_search(c) | |
| if len(n) >= 25: | |
| anchors.append(n) | |
| if not anchors: | |
| return [] | |
| max_pages_for_anchor = max(2, int(math.ceil(len(page_index) * 0.35))) | |
| filtered: list[str] = [] | |
| for a in anchors: | |
| hits = 0 | |
| for pg_text in page_index.values(): | |
| if a and a in pg_text: | |
| hits += 1 | |
| if hits <= max_pages_for_anchor: | |
| filtered.append(a) | |
| if filtered: | |
| anchors = filtered | |
| best_score = 0 | |
| page_scores: dict[int, int] = {} | |
| for pg, pg_text in page_index.items(): | |
| score = 0 | |
| for a in anchors: | |
| if a and a in pg_text: | |
| score += 1 | |
| page_scores[pg] = score | |
| best_score = max(best_score, score) | |
| if best_score <= 0: | |
| return [] | |
| best_pages = sorted([pg for pg, sc in page_scores.items() if sc == best_score and sc > 0]) | |
| if len(best_pages) > 1 and best_score >= 2: | |
| span_min, span_max = best_pages[0], best_pages[-1] | |
| if span_max - span_min <= max_span_pages: | |
| return list(range(int(span_min), int(span_max) + 1)) | |
| return [int(best_pages[0])] | |
| def _filter_pages(pages: Any, total_pages: int | None) -> list[int]: | |
| if not pages: | |
| return [] | |
| out: list[int] = [] | |
| if isinstance(pages, int): | |
| pages = [pages] | |
| if isinstance(pages, tuple): | |
| pages = list(pages) | |
| if not isinstance(pages, list): | |
| return [] | |
| for p in pages: | |
| try: | |
| pi = int(p) | |
| except Exception: | |
| continue | |
| if pi <= 0: | |
| continue | |
| if total_pages is not None and pi > total_pages: | |
| continue | |
| out.append(pi) | |
| return sorted(set(out)) | |
| def _has_markdown_table(text: str) -> bool: | |
| if not text: | |
| return False | |
| pipe_lines = 0 | |
| for ln in text.splitlines(): | |
| s = ln.strip() | |
| if s.count("|") >= 2: | |
| pipe_lines += 1 | |
| if pipe_lines >= 2: | |
| return True | |
| return False | |
| def _looks_like_label(s: str, *, max_chars: int) -> bool: | |
| if not s: | |
| return False | |
| t = s.strip() | |
| if len(t) > max_chars or len(t) < 2: | |
| return False | |
| if t.count(" ") > 4: | |
| return False | |
| if t.endswith(".") or t.endswith(":"): | |
| return True | |
| digits = sum(ch.isdigit() for ch in t) | |
| if digits > 0 and digits / max(1, len(t)) > 0.7: | |
| return False | |
| return True | |
| def _split_blocks_by_blank_lines(text: str) -> list[str]: | |
| if not text: | |
| return [] | |
| segments: list[list[str]] = [] | |
| buf: list[str] = [] | |
| for ln in text.splitlines(): | |
| if ln.strip(): | |
| buf.append(ln.strip()) | |
| continue | |
| if buf: | |
| segments.append(buf) | |
| buf = [] | |
| if buf: | |
| segments.append(buf) | |
| max_label_chars = int(settings.docling_preserve_tables_max_label_chars or 40) | |
| blocks: list[str] = [] | |
| for seg in segments: | |
| if len(seg) > 1 and len(seg) <= 3: | |
| if all(_looks_like_label(s, max_chars=max_label_chars) for s in seg): | |
| blocks.extend([s.strip() for s in seg if s.strip()]) | |
| continue | |
| merged = " ".join(seg).strip() | |
| if merged: | |
| blocks.append(merged) | |
| return blocks | |
| def _kv_blocks_to_markdown_table(text: str) -> tuple[str, int] | None: | |
| if not text: | |
| return None | |
| if _has_markdown_table(text): | |
| return None | |
| blocks = _split_blocks_by_blank_lines(text) | |
| if len(blocks) < 6: | |
| return None | |
| max_label_chars = int(settings.docling_preserve_tables_max_label_chars or 40) | |
| min_pairs = int(settings.docling_preserve_tables_min_pairs or 3) | |
| section_headers = { | |
| "allmänt", | |
| "plan", | |
| "höjd", | |
| "markhöjd", | |
| "lägesbeskrivning", | |
| "distansbricka", | |
| "övrigt", | |
| "anmärkningar", | |
| "historik", | |
| } | |
| def is_section_header(s: str) -> bool: | |
| return (s or "").strip().lower() in section_headers | |
| best_offset: int | None = None | |
| best_pairs: list[tuple[str, str]] = [] | |
| best_prefix_blocks: list[str] = [] | |
| best_score: float | None = None | |
| best_end_index: int = 0 | |
| offsets = [0, 1, 2] | |
| if blocks and is_section_header(blocks[0]): | |
| offsets = [1, 0, 2] | |
| for offset in offsets: | |
| pairs: list[tuple[str, str]] = [] | |
| i = offset | |
| while i < len(blocks): | |
| label = blocks[i].strip() | |
| if is_section_header(label): | |
| i += 1 | |
| continue | |
| if not _looks_like_label(label, max_chars=max_label_chars): | |
| break | |
| j = i + 1 | |
| while j < len(blocks) and is_section_header(blocks[j].strip()): | |
| j += 1 | |
| value = "" | |
| next_i = j + 1 | |
| if j < len(blocks): | |
| cand = blocks[j].strip() | |
| cand_is_labelish = _looks_like_label(cand, max_chars=max_label_chars) | |
| if cand and cand_is_labelish: | |
| k = j + 1 | |
| while k < len(blocks) and is_section_header(blocks[k].strip()): | |
| k += 1 | |
| if k < len(blocks): | |
| cand2 = blocks[k].strip() | |
| cand2_is_labelish = _looks_like_label(cand2, max_chars=max_label_chars) | |
| cand2_is_valueish = ( | |
| (len(cand2) > max_label_chars) | |
| or ("(" in cand2) | |
| or (")" in cand2) | |
| or any(ch.isdigit() for ch in cand2) | |
| ) | |
| if cand2 and cand2_is_valueish and not cand2_is_labelish: | |
| value = cand2 | |
| next_i = k + 1 | |
| else: | |
| value = cand | |
| next_i = j + 1 | |
| else: | |
| value = cand | |
| next_i = j + 1 | |
| else: | |
| value = cand | |
| next_i = j + 1 | |
| pairs.append((label.rstrip(":").strip(), value)) | |
| i = max(next_i, i + 1) | |
| empty_values = sum(1 for _k, v in pairs if not v) | |
| suspicious_values = 0 | |
| for _k, v in pairs: | |
| if not v: | |
| continue | |
| v_is_labelish = _looks_like_label(v, max_chars=max_label_chars) | |
| v_is_valueish = (len(v) > max_label_chars) or ("(" in v) or (")" in v) or any(ch.isdigit() for ch in v) | |
| if v_is_labelish and not v_is_valueish: | |
| suspicious_values += 1 | |
| score = float(len(pairs)) - (empty_values * 0.75) - (suspicious_values * 0.5) | |
| if best_score is None or score > best_score: | |
| best_score = score | |
| best_pairs = pairs | |
| best_offset = offset | |
| best_prefix_blocks = blocks[:offset] | |
| best_end_index = i | |
| if best_offset is None or len(best_pairs) < min_pairs: | |
| return None | |
| empty_values = sum(1 for _k, v in best_pairs if not v) | |
| if empty_values / max(1, len(best_pairs)) > 0.65: | |
| return None | |
| header = "\n".join([b for b in best_prefix_blocks if b.strip()]).strip() | |
| lines: list[str] = [] | |
| if header: | |
| lines.append(header) | |
| lines.append("") | |
| lines.append("| Field | Value |") | |
| lines.append("|---|---|") | |
| for k, v in best_pairs: | |
| k2 = k.replace("\n", " ").strip() | |
| v2 = v.replace("\n", " ").strip() | |
| lines.append(f"| {k2} | {v2} |") | |
| remaining = [b for b in blocks[best_end_index:] if b.strip()] | |
| if remaining: | |
| lines.append("") | |
| lines.append("\n".join(remaining)) | |
| return "\n".join(lines).strip(), len(best_pairs) | |
| def _collect_table_markdown_by_page(doc: Any) -> dict[int, list[str]]: | |
| _, _, tables_by_page = _scan_doc_ancillary( | |
| doc, | |
| _get_total_pages_hint(doc), | |
| collect_page_index=False, | |
| collect_fragments=False, | |
| collect_tables=True, | |
| light=False, | |
| ) | |
| return tables_by_page | |
| def _hierarchical_light_enabled(requested_mode: str) -> bool: | |
| if requested_mode == "hierarchical_light": | |
| return True | |
| return requested_mode == "hierarchical" and bool(settings.docling_hierarchical_light) | |
| def _apply_table_preservation( | |
| text: str, | |
| *, | |
| covered_pages: list[int], | |
| tables_by_page: dict[int, list[str]], | |
| ) -> tuple[str, dict[str, Any]]: | |
| meta: dict[str, Any] = {} | |
| if not text: | |
| return text, meta | |
| if not settings.docling_preserve_tables: | |
| return text, meta | |
| if settings.docling_toc_cleanup and _looks_like_toc(text): | |
| meta["toc_cleaned"] = True | |
| cleaned_text = _clean_dot_leaders(text) | |
| extra_parts: list[str] = [] | |
| for pg in covered_pages: | |
| for md in tables_by_page.get(int(pg), []): | |
| extra = _markdown_table_to_text(md) | |
| if extra: | |
| extra_parts.append(extra) | |
| if extra_parts: | |
| merged = cleaned_text.rstrip() + "\n\n" + "\n\n".join(extra_parts) | |
| return merged.strip(), meta | |
| return cleaned_text, meta | |
| kv = _kv_blocks_to_markdown_table(text) | |
| if kv is not None: | |
| new_text, pairs = kv | |
| meta["table_mode"] = "kv_markdown" | |
| meta["table_pairs"] = pairs | |
| return new_text, meta | |
| if not covered_pages or not tables_by_page: | |
| return text, meta | |
| if _has_markdown_table(text): | |
| return text, meta | |
| collected: list[str] = [] | |
| for pg in covered_pages: | |
| for md in tables_by_page.get(int(pg), []): | |
| if settings.docling_toc_cleanup and _is_noise_markdown_table(md): | |
| continue | |
| collected.append(md) | |
| if not collected: | |
| return text, meta | |
| meta["table_mode"] = "tableitem_markdown" | |
| meta["table_count"] = len(collected) | |
| merged = text.rstrip() + "\n\n" + "\n\n".join(collected).strip() | |
| return merged, meta | |
| def _context_key(text: str, meta: dict[str, Any]) -> tuple[str, ...]: | |
| for key in ("headings", "heading_path", "heading", "headers"): | |
| val = meta.get(key) if isinstance(meta, dict) else None | |
| if not val: | |
| continue | |
| if isinstance(val, str) and val.strip(): | |
| return ("_headings", val.strip()) | |
| if isinstance(val, (list, tuple)): | |
| parts = [str(x).strip() for x in val if x is not None and str(x).strip()] | |
| if parts: | |
| return tuple(["_headings", *parts]) | |
| first_line = "" | |
| for ln in (text or "").splitlines(): | |
| ln = ln.strip() | |
| if ln: | |
| first_line = ln | |
| break | |
| if first_line and len(first_line) <= 140: | |
| return ("_firstline", first_line) | |
| return ("_none",) | |
| def _primary_page(meta: dict[str, Any]) -> int: | |
| pages = meta.get("covered_pages") or [] | |
| if not pages: | |
| return 0 | |
| try: | |
| return int(pages[0]) | |
| except (TypeError, ValueError): | |
| return 0 | |
| def _union_pages(a: list, b: list) -> list[int]: | |
| out: list[int] = [] | |
| for p in list(a) + list(b): | |
| try: | |
| pi = int(p) | |
| if pi > 0: | |
| out.append(pi) | |
| except (TypeError, ValueError): | |
| continue | |
| return sorted(set(out)) | |
| def _should_skip_item_chunk(text: str, label: str) -> bool: | |
| if label == "picture" and "Image not available" in text: | |
| return True | |
| if settings.docling_filter_noise_chunks and len(text) < 200 and _is_noise_chunk_text(text): | |
| return True | |
| return False | |
| def _collect_item_chunks( | |
| doc: Any, | |
| source_name: str, | |
| ) -> tuple[list[dict], int, int, int, dict[str, int]]: | |
| """One Docling item per chunk (same as legacy item mode).""" | |
| chunks: list[dict] = [] | |
| total_items = 0 | |
| kept_items = 0 | |
| export_errors = 0 | |
| item_type_counts: dict[str, int] = {} | |
| for item, _level in doc.iterate_items(): | |
| total_items += 1 | |
| item_type = type(item).__name__ | |
| item_type_counts[item_type] = item_type_counts.get(item_type, 0) + 1 | |
| try: | |
| if isinstance(item, (TextItem, TableItem, PictureItem)): | |
| if isinstance(item, (TableItem, PictureItem)): | |
| text_content = item.export_to_markdown(doc=doc).strip() | |
| else: | |
| text_content = (getattr(item, "text", None) or "").strip() | |
| if not text_content: | |
| continue | |
| label = str(item.label) if item.label is not None else "" | |
| if _should_skip_item_chunk(text_content, label): | |
| continue | |
| covered_pages: list[int] = [] | |
| if item.prov: | |
| covered_pages = [ | |
| int(p.page_no) | |
| for p in item.prov | |
| if getattr(p, "page_no", None) is not None | |
| ] | |
| chunks.append( | |
| { | |
| "text": text_content, | |
| "metadata": { | |
| "source": source_name, | |
| "covered_pages": covered_pages, | |
| "label": label, | |
| }, | |
| } | |
| ) | |
| kept_items += 1 | |
| except Exception as exc: | |
| export_errors += 1 | |
| log.warning( | |
| "DOCLING export_to_markdown failed [%s] item=%s: %s", | |
| source_name, | |
| type(item).__name__, | |
| exc, | |
| ) | |
| return chunks, total_items, kept_items, export_errors, item_type_counts | |
| def _merge_sequential_item_chunks(items: list[dict], *, max_chars: int) -> list[dict]: | |
| """ | |
| Merge adjacent item chunks in reading order (same page, under max_chars). | |
| Starts a new chunk on section_header labels. Low RAM vs HierarchicalChunker. | |
| """ | |
| if not items: | |
| return [] | |
| merged: list[dict] = [] | |
| current: dict | None = None | |
| def flush() -> None: | |
| nonlocal current | |
| if current and (current.get("text") or "").strip(): | |
| merged.append(current) | |
| current = None | |
| for it in items: | |
| text = (it.get("text") or "").strip() | |
| if not text: | |
| continue | |
| meta = dict(it.get("metadata") or {}) | |
| label = str(meta.get("label") or "") | |
| if label == "section_header": | |
| flush() | |
| current = {"text": text, "metadata": {**meta, "label": "section"}} | |
| continue | |
| if current is None: | |
| current = {"text": text, "metadata": meta} | |
| continue | |
| cur_meta = current.get("metadata") or {} | |
| combined = current["text"].rstrip() + "\n\n" + text | |
| cur_page = _primary_page(cur_meta) | |
| new_page = _primary_page(meta) | |
| # Same page when both have page numbers, or both lack pages (appendices, DOCX, etc.) | |
| same_page = cur_page == new_page and (cur_page > 0 or new_page == 0) | |
| if same_page and len(combined) <= max_chars: | |
| cur_meta = dict(cur_meta) | |
| cur_meta["covered_pages"] = _union_pages( | |
| cur_meta.get("covered_pages") or [], meta.get("covered_pages") or [] | |
| ) | |
| current = {"text": combined, "metadata": cur_meta} | |
| else: | |
| flush() | |
| current = {"text": text, "metadata": meta} | |
| flush() | |
| return merged | |
| def _apply_max_chunk_split(chunks: list[dict], *, max_chars: int, overlap: int) -> list[dict]: | |
| if max_chars <= 0: | |
| return chunks | |
| resized: list[dict] = [] | |
| for it in chunks: | |
| text = (it.get("text") or "").strip() | |
| if not text: | |
| continue | |
| if settings.docling_filter_noise_chunks and _is_noise_chunk_text(text): | |
| continue | |
| if len(text) <= max_chars: | |
| resized.append(it) | |
| continue | |
| for part in _split_text_keep_context(text, max_chars=max_chars, overlap=overlap): | |
| meta = dict(it.get("metadata") or {}) | |
| meta["split_from"] = True | |
| if settings.docling_filter_noise_chunks and _is_noise_chunk_text(part): | |
| continue | |
| resized.append({"text": part, "metadata": meta}) | |
| return resized | |
| def _merge_peer_chunks(items: list[dict], *, max_chars: int) -> list[dict]: | |
| merged: list[dict] = [] | |
| current: dict | None = None | |
| current_key: tuple[str, ...] | None = None | |
| def normalize_heading(s: str) -> str: | |
| if not s: | |
| return "" | |
| s = s.replace("\u00A0", " ").replace("\u200B", "") | |
| s = " ".join(s.strip().split()) | |
| return s.casefold() | |
| def first_nonempty_line(text: str) -> str: | |
| for ln in (text or "").splitlines(): | |
| s = ln.strip() | |
| if s: | |
| return s | |
| return "" | |
| def strip_leading_heading(text: str, *, heading_norm: str) -> str: | |
| if not text or not heading_norm: | |
| return text | |
| lines = text.splitlines() | |
| i = 0 | |
| while i < len(lines) and not lines[i].strip(): | |
| i += 1 | |
| if i < len(lines) and normalize_heading(lines[i]) == heading_norm: | |
| i += 1 | |
| while i < len(lines) and not lines[i].strip(): | |
| i += 1 | |
| return "\n".join(lines[i:]).strip() | |
| return text | |
| def flush() -> None: | |
| nonlocal current, current_key | |
| if current is not None: | |
| merged.append(current) | |
| current = None | |
| current_key = None | |
| for it in items: | |
| text = (it.get("text") or "").strip() | |
| if not text: | |
| continue | |
| meta = it.get("metadata") or {} | |
| if not isinstance(meta, dict): | |
| meta = {} | |
| key = _context_key(text, meta) | |
| if current is None: | |
| current = {"text": text, "metadata": dict(meta)} | |
| current_key = key | |
| continue | |
| can_merge = current_key == key and (len(current["text"]) + 2 + len(text)) <= max_chars | |
| if not can_merge: | |
| flush() | |
| current = {"text": text, "metadata": dict(meta)} | |
| current_key = key | |
| continue | |
| heading_norm = normalize_heading(first_nonempty_line(current["text"])) | |
| text_to_add = strip_leading_heading(text, heading_norm=heading_norm) | |
| if not text_to_add: | |
| continue | |
| current["text"] = current["text"].rstrip() + "\n\n" + text_to_add | |
| cur_pages = current["metadata"].get("covered_pages") or [] | |
| new_pages = meta.get("covered_pages") or [] | |
| try: | |
| pages_union = sorted({int(p) for p in list(cur_pages) + list(new_pages) if p is not None}) | |
| except Exception: | |
| pages_union = cur_pages or new_pages | |
| current["metadata"]["covered_pages"] = pages_union | |
| flush() | |
| return merged | |
| def _merge_kv_microchunks(items: list[dict]) -> list[dict]: | |
| if not items: | |
| return [] | |
| merged: list[dict] = [] | |
| i = 0 | |
| while i < len(items): | |
| cur = items[i] | |
| text = (cur.get("text") or "").strip() | |
| meta = cur.get("metadata") or {} | |
| if not isinstance(meta, dict): | |
| meta = {} | |
| if ( | |
| settings.docling_kv_merge_microchunks | |
| and settings.docling_preserve_tables | |
| and text.endswith(":") | |
| and "\n" not in text | |
| and len(text) <= 80 | |
| and i + 1 < len(items) | |
| ): | |
| nxt = items[i + 1] | |
| nxt_text = (nxt.get("text") or "").strip() | |
| nxt_meta = nxt.get("metadata") or {} | |
| if not isinstance(nxt_meta, dict): | |
| nxt_meta = {} | |
| if nxt_text and not nxt_text.endswith(":") and "\n" not in nxt_text and len(nxt_text) <= 120: | |
| cur_pages = meta.get("covered_pages") or [] | |
| nxt_pages = nxt_meta.get("covered_pages") or [] | |
| try: | |
| pages_union = sorted({int(p) for p in list(cur_pages) + list(nxt_pages) if p is not None}) | |
| except Exception: | |
| pages_union = cur_pages or nxt_pages | |
| new_meta = dict(meta) | |
| new_meta["covered_pages"] = pages_union | |
| merged.append({"text": f"{text} {nxt_text}".strip(), "metadata": new_meta}) | |
| i += 2 | |
| continue | |
| merged.append({"text": text, "metadata": dict(meta)}) | |
| i += 1 | |
| return merged | |
| def _split_text_keep_context(text: str, *, max_chars: int, overlap: int) -> list[str]: | |
| if not text: | |
| return [] | |
| if len(text) <= max_chars: | |
| return [text] | |
| lines = text.splitlines() | |
| header_lines: list[str] = [] | |
| remaining_lines: list[str] = [] | |
| seen_blank = False | |
| for ln in lines: | |
| if not ln.strip(): | |
| seen_blank = True | |
| remaining_lines.append(ln) | |
| continue | |
| if not seen_blank and len(header_lines) < 12: | |
| header_lines.append(ln) | |
| else: | |
| remaining_lines.append(ln) | |
| header = "\n".join(header_lines).strip() | |
| body = "\n".join(remaining_lines).strip() if remaining_lines else "" | |
| prefix = (header + "\n\n") if header else "" | |
| max_body = max(200, max_chars - len(prefix)) | |
| out: list[str] = [] | |
| start = 0 | |
| while start < len(body): | |
| end = min(len(body), start + max_body) | |
| piece = body[start:end] | |
| out.append((prefix + piece).strip()) | |
| if end >= len(body): | |
| break | |
| start = max(0, end - overlap) | |
| if start == end: | |
| break | |
| return [o for o in out if o] | |
| def _dedupe_repeated_heading_lines(text: str) -> str: | |
| if not text: | |
| return "" | |
| raw = text.strip() | |
| parts = [p for p in re.split(r"\n\s*\n", raw) if p.strip()] | |
| if len(parts) < 2: | |
| return raw | |
| def first_nonempty_line(p: str) -> str: | |
| for ln in p.splitlines(): | |
| s = ln.strip() | |
| if s: | |
| return s | |
| return "" | |
| def normalize_heading(s: str) -> str: | |
| if not s: | |
| return "" | |
| s = s.replace("\u00A0", " ").replace("\u200B", "") | |
| s = " ".join(s.strip().split()) | |
| return s.casefold() | |
| heading = first_nonempty_line(parts[0]) | |
| heading_norm = normalize_heading(heading) | |
| if not heading_norm or len(heading) > 140: | |
| return raw | |
| repeats = 0 | |
| for p in parts[1:]: | |
| if normalize_heading(first_nonempty_line(p)) == heading_norm: | |
| repeats += 1 | |
| if repeats == 0: | |
| return raw | |
| out_parts: list[str] = [parts[0].strip()] | |
| for p in parts[1:]: | |
| lines = p.splitlines() | |
| i = 0 | |
| while i < len(lines) and not lines[i].strip(): | |
| i += 1 | |
| if i < len(lines) and normalize_heading(lines[i]) == heading_norm: | |
| i += 1 | |
| while i < len(lines) and not lines[i].strip(): | |
| i += 1 | |
| rest = "\n".join(lines[i:]).strip() | |
| if rest: | |
| out_parts.append(rest) | |
| continue | |
| out_parts.append(p.strip()) | |
| return "\n\n".join([p for p in out_parts if p.strip()]) | |
| class ChunkingSubprocessError(RuntimeError): | |
| """Child process failed during PDF chunking (often OOM).""" | |
| def _should_use_chunking_subprocess(requested_mode: str) -> bool: | |
| if not settings.docling_chunking_subprocess: | |
| return False | |
| return requested_mode in ("hierarchical", "hierarchical_light") | |
| def _run_chunking_subprocess( | |
| file_path: str, | |
| source_name: str, | |
| *, | |
| force_mode: str | None = None, | |
| doc_id: str | None = None, | |
| ) -> list[dict]: | |
| """Run chunking in a separate Python process (stdlib subprocess, no mp semaphores).""" | |
| fd, out_path = tempfile.mkstemp(suffix=".chunks.pkl") | |
| os.close(fd) | |
| try: | |
| cmd = [ | |
| sys.executable, | |
| "-m", | |
| "app.chunking_worker", | |
| file_path, | |
| source_name, | |
| out_path, | |
| force_mode or "", | |
| doc_id or "", | |
| ] | |
| result = subprocess.run( | |
| cmd, | |
| cwd=str(PROJECT_ROOT), | |
| check=False, | |
| ) | |
| if result.returncode != 0: | |
| raise ChunkingSubprocessError( | |
| f"chunking subprocess exit code {result.returncode} " | |
| f"(negative often means OOM kill)" | |
| ) | |
| with open(out_path, "rb") as f: | |
| return pickle.load(f) | |
| finally: | |
| try: | |
| os.unlink(out_path) | |
| except OSError: | |
| pass | |
| def _process_pdf_to_chunks_impl( | |
| file_path: str, | |
| source_name: str, | |
| *, | |
| update_context: dict[str, Any] | None = None, | |
| force_mode: str | None = None, | |
| doc_id: str | None = None, | |
| ) -> list[dict]: | |
| accelerator_options = _get_accelerator_options() | |
| device_name = str(accelerator_options.device).replace("AcceleratorDevice.", "") | |
| log.info("DOCLING running on device: %s", device_name) | |
| requested_mode_log = (force_mode or settings.docling_chunking_mode or "item").strip().lower() | |
| input_format = Path(file_path).suffix.lower().lstrip(".") or "unknown" | |
| log.info( | |
| "DOCLING config: format=%s mode=%s light=%s preserve_tables=%s enable_table_structure_model=%s max_seconds=%s max_chars=%s overlap=%s", | |
| input_format, | |
| requested_mode_log, | |
| _hierarchical_light_enabled(requested_mode_log) | |
| if requested_mode_log in ("hierarchical", "hierarchical_light") | |
| else False, | |
| bool(settings.docling_preserve_tables), | |
| bool(settings.docling_enable_table_structure_model), | |
| settings.docling_max_seconds, | |
| settings.docling_max_chars_per_chunk, | |
| settings.docling_chunk_overlap, | |
| ) | |
| options = PdfPipelineOptions() | |
| options.do_ocr = False | |
| options.do_table_structure = bool(settings.docling_enable_table_structure_model) | |
| options.do_picture_classification = False | |
| options.accelerator_options = accelerator_options | |
| converter_kw: dict[str, Any] = { | |
| "allowed_formats": [InputFormat.PDF, InputFormat.DOCX, InputFormat.XLSX], | |
| "format_options": { | |
| InputFormat.PDF: PdfFormatOption(pipeline_options=options), | |
| InputFormat.DOCX: WordFormatOption(), | |
| InputFormat.XLSX: ExcelFormatOption(), | |
| }, | |
| } | |
| if update_context and _converter_supports_progress_callback(): | |
| cb = _make_docling_progress_callback(update_context) | |
| if cb is not None: | |
| converter_kw["progress_callback"] = cb | |
| converter = DocumentConverter(**converter_kw) | |
| t0 = time.perf_counter() | |
| result = converter.convert(file_path) | |
| doc = result.document | |
| del result | |
| gc.collect() | |
| requested_mode = (force_mode or settings.docling_chunking_mode or "item").strip().lower() | |
| mode_used = "item" | |
| chunks: list[dict] = [] | |
| total_items = 0 | |
| kept_items = 0 | |
| export_errors = 0 | |
| item_type_counts: dict[str, int] = {} | |
| if requested_mode == "merged_item" and doc is not None: | |
| raw_items, total_items, kept_items, export_errors, item_type_counts = _collect_item_chunks( | |
| doc, source_name | |
| ) | |
| max_chars = int(settings.docling_max_chars_per_chunk or 2000) | |
| overlap = int(settings.docling_chunk_overlap or 200) | |
| n_raw = len(raw_items) | |
| chunks = _merge_sequential_item_chunks(raw_items, max_chars=max_chars) | |
| chunks = _apply_max_chunk_split(chunks, max_chars=max_chars, overlap=overlap) | |
| mode_used = "merged_item" | |
| log.info( | |
| "DOCLING merged_item [%s]: %d items -> %d merged chunks (max_chars=%s)", | |
| source_name, | |
| n_raw, | |
| len(chunks), | |
| max_chars, | |
| ) | |
| del doc | |
| doc = None # type: ignore[assignment] | |
| gc.collect() | |
| elif requested_mode in ("hierarchical", "hierarchical_light") and _HAS_HIERARCHICAL_CHUNKER: | |
| try: | |
| light = _hierarchical_light_enabled(requested_mode) | |
| chunker = HierarchicalChunker() | |
| total_pages_hint = _get_total_pages_hint(doc) | |
| if light: | |
| log.info("DOCLING hierarchical LIGHT: skipping page scan, header strip, merges") | |
| page_index: dict[int, str] = {} | |
| repeated_fragments: set[str] = set() | |
| tables_by_page: dict[int, list[str]] = {} | |
| else: | |
| page_index, repeated_fragments, tables_by_page = _scan_doc_ancillary( | |
| doc, | |
| total_pages_hint, | |
| collect_page_index=True, | |
| collect_fragments=True, | |
| collect_tables=bool(settings.docling_preserve_tables), | |
| light=False, | |
| ) | |
| raw_chunks = chunker.chunk(doc) | |
| for hc in raw_chunks: | |
| try: | |
| try: | |
| text_content = chunker.contextualize(hc).strip() | |
| except Exception: | |
| text_content = (getattr(hc, "text", "") or "").strip() | |
| if not text_content: | |
| continue | |
| if not light: | |
| text_content = _dedupe_repeated_heading_lines(text_content) | |
| meta = getattr(hc, "metadata", {}) or {} | |
| if not isinstance(meta, dict): | |
| meta = {} | |
| if light: | |
| cleaned = text_content | |
| removed_lines = 0 | |
| else: | |
| cleaned, removed_lines = _strip_headers_footers_for_inference( | |
| text_content, repeated_fragments | |
| ) | |
| covered_pages = ( | |
| meta.get("covered_pages") | |
| or meta.get("page_numbers") | |
| or meta.get("pages") | |
| or meta.get("page_span") | |
| or [] | |
| ) | |
| covered_pages = _filter_pages(covered_pages, total_pages_hint) | |
| page_source = "chunk_meta" if covered_pages else "none" | |
| if not light and not covered_pages: | |
| inferred = _infer_pages_from_text(cleaned or text_content, page_index) | |
| covered_pages = _filter_pages(inferred, total_pages_hint) | |
| if covered_pages: | |
| page_source = "text_match" | |
| label = str(meta.get("label") or "hierarchical") | |
| chunk_meta: dict[str, Any] = { | |
| **meta, | |
| "source": source_name, | |
| "covered_pages": covered_pages, | |
| "label": label, | |
| "page_source": page_source, | |
| } | |
| if not light: | |
| chunk_meta["hf_removed_lines"] = int(removed_lines) | |
| chunks.append( | |
| { | |
| "text": (cleaned or text_content).strip(), | |
| "metadata": chunk_meta, | |
| } | |
| ) | |
| except Exception: | |
| export_errors += 1 | |
| continue | |
| if chunks and not light: | |
| merge_limit = int(settings.docling_max_chars_per_chunk or 2000) | |
| chunks = _merge_peer_chunks(chunks, max_chars=merge_limit) | |
| for it in chunks: | |
| it["text"] = _dedupe_repeated_heading_lines(it.get("text") or "") | |
| chunks = _merge_kv_microchunks(chunks) | |
| if settings.docling_preserve_tables: | |
| for it in chunks: | |
| meta = it.get("metadata") or {} | |
| if not isinstance(meta, dict): | |
| meta = {} | |
| raw_pages = meta.get("covered_pages") or [] | |
| covered = _filter_pages(raw_pages, total_pages_hint) | |
| new_text, table_meta = _apply_table_preservation( | |
| it.get("text") or "", | |
| covered_pages=covered, | |
| tables_by_page=tables_by_page, | |
| ) | |
| if table_meta: | |
| it["text"] = new_text | |
| meta.update(table_meta) | |
| it["metadata"] = meta | |
| if chunks: | |
| max_chars = int(settings.docling_max_chars_per_chunk or 2000) | |
| overlap = int(settings.docling_chunk_overlap or 200) | |
| if not light: | |
| for it in chunks: | |
| it["text"] = _dedupe_repeated_heading_lines(it.get("text") or "") | |
| chunks = _apply_max_chunk_split(chunks, max_chars=max_chars, overlap=overlap) | |
| elapsed_total = time.perf_counter() - t0 | |
| if settings.docling_max_seconds and elapsed_total > float(settings.docling_max_seconds): | |
| chunks = [] | |
| log.info( | |
| "DOCLING hierarchical guardrail triggered: elapsed_s=%.2f max_s=%s file=%s", | |
| elapsed_total, | |
| settings.docling_max_seconds, | |
| source_name, | |
| ) | |
| if settings.docling_fallback_on_document_chunk and len(chunks) == 1: | |
| only = chunks[0] | |
| meta0 = only.get("metadata") or {} | |
| if (meta0.get("label") == "document") or ( | |
| len((only.get("text") or "")) > int(settings.docling_max_chars_per_chunk or 2000) * 2 | |
| ): | |
| chunks = [] | |
| log.info("DOCLING hierarchical guardrail triggered: single giant chunk file=%s", source_name) | |
| total_pages = total_pages_hint | |
| if total_pages is None and page_index: | |
| total_pages = max(page_index.keys()) | |
| if ( | |
| not light | |
| and total_pages is not None | |
| and int(total_pages) >= int(settings.docling_min_chunks_page_threshold or 3) | |
| and len(chunks) < int(settings.docling_min_chunks or 3) | |
| ): | |
| chunk_count = len(chunks) | |
| chunks = [] | |
| log.info( | |
| "DOCLING hierarchical guardrail triggered: too_few_chunks chunks=%s pages=%s file=%s", | |
| chunk_count, | |
| total_pages, | |
| source_name, | |
| ) | |
| if chunks: | |
| mode_used = "hierarchical_light" if light else "hierarchical" | |
| total_items = len(chunks) | |
| kept_items = len(chunks) | |
| item_type_counts["HierarchicalChunk"] = len(chunks) | |
| try: | |
| del raw_chunks | |
| except NameError: | |
| pass | |
| try: | |
| del chunker | |
| except NameError: | |
| pass | |
| del doc | |
| doc = None # type: ignore[assignment] | |
| gc.collect() | |
| except Exception as e: | |
| chunks = [] | |
| log.warning("DOCLING hierarchical chunking failed [%s]: %s", source_name, e) | |
| if not chunks and doc is not None: | |
| chunks, total_items, kept_items, export_errors, item_type_counts = _collect_item_chunks( | |
| doc, source_name | |
| ) | |
| mode_used = "item" | |
| if not chunks and doc is not None: | |
| try: | |
| full_text = doc.export_to_markdown().strip() | |
| if full_text: | |
| chunks.append( | |
| { | |
| "text": full_text, | |
| "metadata": {"source": source_name, "covered_pages": [], "label": "document"}, | |
| } | |
| ) | |
| mode_used = "document" | |
| log.warning("DOCLING fallback used: exported full-document markdown as single chunk [%s]", source_name) | |
| except Exception as e: | |
| log.warning("DOCLING fallback export failed [%s]: %s", source_name, e) | |
| if total_items == 0: | |
| total_items = len(chunks) | |
| log.info( | |
| "DOCLING extraction stats [%s]: mode=%s total_items=%d kept_items=%d export_errors=%d item_types=%s", | |
| source_name, | |
| mode_used, | |
| total_items, | |
| kept_items, | |
| export_errors, | |
| item_type_counts, | |
| ) | |
| elapsed = time.perf_counter() - t0 | |
| log.info("DOCLING done in %.2fs, %d chunks [%s]", elapsed, len(chunks), source_name) | |
| _maybe_dump_chunks_debug( | |
| chunks, source_name, doc_id=doc_id, input_path=file_path | |
| ) | |
| return chunks | |
| def process_pdf_to_chunks( | |
| file_path: str, | |
| source_name: str, | |
| *, | |
| update_context: dict[str, Any] | None = None, | |
| doc_id: str | None = None, | |
| ) -> list[dict]: | |
| """ | |
| Chunk a PDF, DOCX, or XLSX. Hierarchical modes may run in a subprocess so OOM kills the worker, | |
| not the API server; on subprocess failure we retry once with item mode. | |
| """ | |
| mode = (settings.docling_chunking_mode or "item").strip().lower() | |
| use_subprocess = _should_use_chunking_subprocess(mode) and update_context is None | |
| if use_subprocess: | |
| log.info("DOCLING chunking in subprocess (mode=%s)", mode) | |
| try: | |
| return _run_chunking_subprocess(file_path, source_name, doc_id=doc_id) | |
| except ChunkingSubprocessError as exc: | |
| log.warning( | |
| "DOCLING subprocess chunking failed [%s]: %s — retrying with item mode", | |
| source_name, | |
| exc, | |
| ) | |
| return _run_chunking_subprocess( | |
| file_path, source_name, force_mode="merged_item", doc_id=doc_id | |
| ) | |
| return _process_pdf_to_chunks_impl( | |
| file_path, source_name, update_context=update_context, doc_id=doc_id | |
| ) | |
| def _sniff_office_openxml_kind(path: Path) -> str | None: | |
| """Distinguish docx vs xlsx inside ZIP Office Open XML (debug only).""" | |
| import zipfile | |
| try: | |
| with zipfile.ZipFile(path) as zf: | |
| ct = zf.read("[Content_Types].xml") | |
| except (OSError, KeyError, zipfile.BadZipFile): | |
| return None | |
| if b"spreadsheetml" in ct: | |
| return "xlsx" | |
| if b"wordprocessingml" in ct: | |
| return "docx" | |
| return None | |
| def _sniff_file_format(path: Path) -> str | None: | |
| """Best-effort format from magic bytes (debug only).""" | |
| try: | |
| with path.open("rb") as f: | |
| head = f.read(8) | |
| except OSError: | |
| return None | |
| if head.startswith(b"%PDF"): | |
| return "pdf" | |
| if head[:2] == b"PK": | |
| return _sniff_office_openxml_kind(path) or "docx" | |
| return None | |
| def _debug_chunks_output_path( | |
| out_dir: Path, | |
| source_name: str, | |
| doc_id: str | None, | |
| *, | |
| detected_format: str | None = None, | |
| ) -> Path: | |
| """One markdown file per job (doc_id); uses detected format when available.""" | |
| if doc_id: | |
| safe_id = re.sub(r"[^\w\-.]", "_", doc_id) | |
| if detected_format: | |
| return out_dir / f"{safe_id}.{detected_format}.chunks.md" | |
| safe_source = re.sub(r"[^\w\-.]", "_", source_name) or "document" | |
| return out_dir / f"{safe_id}_{safe_source}.chunks.md" | |
| safe_source = re.sub(r"[^\w\-.]", "_", source_name) or "document" | |
| return out_dir / f"{safe_source}.chunks.md" | |
| def _maybe_dump_chunks_debug( | |
| chunks: list[dict], | |
| source_name: str, | |
| *, | |
| doc_id: str | None = None, | |
| input_path: str | None = None, | |
| ) -> None: | |
| """If DEBUG_CHUNKS_DIR is set, write a markdown file with all chunks for inspection.""" | |
| out_dir = get_debug_chunks_dir() | |
| if out_dir is None: | |
| return | |
| try: | |
| out_dir.mkdir(parents=True, exist_ok=True) | |
| input_path_p = Path(input_path) if input_path else None | |
| input_suffix = input_path_p.suffix.lower() if input_path_p else "" | |
| sniffed = _sniff_file_format(input_path_p) if input_path_p else None | |
| out_path = _debug_chunks_output_path( | |
| out_dir, source_name, doc_id, detected_format=sniffed | |
| ) | |
| lines = [ | |
| "# Chunk debug\n", | |
| f"**upload_filename:** `{source_name}`\n", | |
| f"**doc_id:** `{doc_id or 'n/a'}`\n", | |
| ] | |
| if input_path_p: | |
| lines.append(f"**stored_as:** `{input_path_p.name}`\n") | |
| if input_suffix: | |
| lines.append(f"**extension_on_disk:** `{input_suffix.lstrip('.')}`\n") | |
| if sniffed: | |
| lines.append(f"**detected_format:** `{sniffed}`\n") | |
| if sniffed and input_suffix and sniffed != input_suffix.lstrip("."): | |
| lines.append( | |
| "**note:** upload extension does not match file contents " | |
| "(e.g. multipart `filename=file.pdf` with a `.docx` file).\n" | |
| ) | |
| lines.append(f"**Total chunks:** {len(chunks)}\n\n---\n") | |
| for i, chunk in enumerate(chunks): | |
| meta = chunk.get("metadata", {}) | |
| text = chunk.get("text", "") | |
| pages = meta.get("covered_pages", []) | |
| label = meta.get("label", "") | |
| lines.append(f"## Chunk {i + 1} / {len(chunks)}\n") | |
| lines.append(f"- **label:** `{label}`\n") | |
| lines.append(f"- **pages:** {pages}\n") | |
| if isinstance(meta, dict): | |
| if meta.get("page_source"): | |
| lines.append(f"- **page_source:** `{meta.get('page_source')}`\n") | |
| if meta.get("table_mode"): | |
| lines.append(f"- **table_mode:** `{meta.get('table_mode')}`\n") | |
| if meta.get("table_pairs") is not None: | |
| lines.append(f"- **table_pairs:** {meta.get('table_pairs')}\n") | |
| if meta.get("table_count") is not None: | |
| lines.append(f"- **table_count:** {meta.get('table_count')}\n") | |
| lines.append(f"- **length:** {len(text)} chars\n\n") | |
| lines.append("```\n") | |
| lines.append(text) | |
| lines.append("\n```\n\n---\n") | |
| out_path.write_text("".join(lines), encoding="utf-8") | |
| log.info( | |
| "DEBUG chunk dump written: %s (%d chunks)", | |
| out_path.resolve(), | |
| len(chunks), | |
| ) | |
| except Exception as e: | |
| log.warning("DEBUG chunk dump failed (dir=%s): %s", out_dir, e) | |
| process_document_to_chunks = process_pdf_to_chunks | |