ishaq101's picture
/fix planner and unstructured pipeline (#23)
7124acf
Raw History Blame Contribute Delete
14.3 kB
"""Stage 2 — parse a document with MinerU, cached by content.
Why the cache is keyed to document CONTENT rather than to a run id: parsing is
by far the most expensive stage (19 s/page on CPU). A per-run cache would
re-parse everything each time the normaliser changes and the pipeline is run
again, for nothing. Keyed by content + settings + version, the same document is
parsed once — however many times the normaliser is revised.
Cache key = sha256(file bytes) + settings fingerprint + MinerU version. The
version is in there so figures from an older MinerU never quietly mix with
figures from a newer one — which matters, because those figures are used to
justify spend.
"""
from __future__ import annotations
import hashlib
import json
import shutil
import time
from dataclasses import dataclass
from pathlib import Path
from .config import PipelineConfig
@dataclass
class ParseResult:
doc_id: str
source: Path
cache_dir: Path # folder holding MinerU's output untouched
content_list: Path # *_content_list.json
middle_json: Path | None
pages: int
seconds: float
from_cache: bool
backend_recorded: str | None # read from _middle.json, NOT from the config
mineru_version: str | None
def mineru_version() -> str:
try:
from mineru.version import __version__ # type: ignore
return str(__version__)
except Exception:
try:
from importlib.metadata import version
return version("mineru")
except Exception:
return "unknown"
def file_hash(path: Path, length: int = 16) -> str:
h = hashlib.sha256()
with path.open("rb") as f:
for block in iter(lambda: f.read(1 << 20), b""):
h.update(block)
return h.hexdigest()[:length]
def cache_key(source: Path, cfg: PipelineConfig) -> str:
# The parser version goes in the key so figures from two parser versions never
# mix. For Mistral that version is the model id, not the MinerU build.
parser_ver = (
cfg.mistral_model if cfg.backend == "mistral"
else cfg.paddle_model if cfg.backend == "paddleocr"
else mineru_version()
)
return f"{file_hash(source)}-{cfg.fingerprint()}-{parser_ver}"
def _mineru_name(doc_id: str, source: Path) -> str:
"""A length-bounded folder name for MinerU's own output.
MinerU names its output folder after whatever is passed as
`pdf_file_names`, and then writes images into it under a full SHA-256
filename. Passing the real document name blows the Windows MAX_PATH limit
of 260: measured on this repo, the fixed part of the path costs 217
characters, so the name may not exceed 43. `STD_2026_006_MNO_ Rev.0.0 -
Production Parameter and ECA` is 56 and fails.
**Cosmetic only.** Uniqueness is already guaranteed by the parent folder,
which is `cache_key` — the full SHA-256 of the file contents. Two different
documents can never share a cache folder however alike their names. The
readable part exists so the cache can be browsed while debugging, nothing
more.
Head AND tail are kept, because the names that actually collide share a
long prefix and differ at the end — "… and ECA" versus "… and ECA (2)".
Keeping only the head would render exactly those indistinguishable.
**Temporary.** Once `doc_id` becomes the content hash rather than the file
name (see the README), every path here is short by construction and this
function goes away. It is deliberately kept simple for that reason.
"""
safe = "".join(c if c.isalnum() or c in " ._-()" else "_" for c in doc_id).strip()
short = safe if len(safe) <= 20 else f"{safe[:12]}~{safe[-8:]}"
return f"{short}-{file_hash(source, 8)}"
def require_gpu(cfg: PipelineConfig) -> None:
"""`hybrid` and `vlm` need a GPU — checked only when MinerU is about to run.
Deliberately not a config-time check. A cache hit needs no GPU, and
rebuilding artifacts from MinerU's saved output is the normal way to pick up
a normalisation fix without paying for a re-parse: the expensive stage
(PDF -> content_list) is cached, and only the cheap CPU stage re-runs.
Refusing at startup would block exactly that.
"""
# RuntimeError, not SystemExit: the batch may hold a mix of cached and
# uncached documents, and one that needs a GPU must not abort the ones that
# do not. It is recorded as a per-document failure like any other.
if cfg.backend in {"hybrid", "vlm"} and shutil.which("nvidia-smi") is None:
raise RuntimeError(
f"Backend '{cfg.backend}' needs a GPU to parse, and nvidia-smi was not found.\n"
" -> This document is not in the parse cache, so MinerU has to run.\n"
" -> To parse on a laptop without one, ask for pipeline:\n"
" python -m src.knowledge_parsing.run "
"--input data/knowledge_docs/ --backend pipeline\n"
" -> That output is valid, but has NO image descriptions and the\n"
" multiplication sign is copied as the letter 'x'. Do not use\n"
" it as a reference measurement."
)
def _find_output(root: Path) -> tuple[Path | None, Path | None]:
"""Locate content_list & middle json anywhere under root.
Deliberately a glob rather than a fixed path: MinerU's subfolder layout
differs per backend ('auto' for pipeline, something else for vlm). Globbing
means this module needs no change when the backend does.
"""
content = next(iter(sorted(root.rglob("*_content_list.json"))), None)
middle = next(iter(sorted(root.rglob("*_middle.json"))), None)
return content, middle
def _read_middle(middle: Path | None) -> tuple[str | None, str | None]:
"""The backend & version that ACTUALLY ran, taken from MinerU's own output."""
if not middle or not middle.exists():
return None, None
try:
d = json.loads(middle.read_text(encoding="utf-8"))
return d.get("_backend"), d.get("_version_name")
except Exception:
return None, None
def _parse_mistral(
source: Path, cfg: PipelineConfig, doc_id: str, target: Path, marker: Path
) -> ParseResult:
"""Parse one document with Mistral OCR, cached exactly like the MinerU path.
The cached output is a single `mistral_ocr.json` (the API response) rather
than MinerU's `*_content_list.json`; `ParseResult.content_list` points at it
and run.py dispatches on `cfg.backend` to the matching normaliser.
"""
out = target / "mistral_ocr.json"
if marker.exists() and out.exists():
meta = json.loads(marker.read_text(encoding="utf-8"))
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=out, middle_json=None,
pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0),
from_cache=True, backend_recorded="mistral-ocr",
mineru_version=meta.get("model") or cfg.mistral_model,
)
from .mistral_ocr import run_ocr
started = time.perf_counter()
response = run_ocr(source, cfg) # no GPU; needs MISTRAL_API_KEY
seconds = time.perf_counter() - started
pages = len(response.get("pages", []))
model = response.get("model") or cfg.mistral_model
staging = target.with_name(target.name + ".partial")
shutil.rmtree(staging, ignore_errors=True)
staging.mkdir(parents=True, exist_ok=True)
(staging / "mistral_ocr.json").write_text(
json.dumps(response, ensure_ascii=False, indent=2), encoding="utf-8"
)
staging.rename(target)
marker.write_text(
json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id, "model": model}),
encoding="utf-8",
)
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=target / "mistral_ocr.json", middle_json=None,
pages=pages, seconds=seconds, from_cache=False,
backend_recorded="mistral-ocr", mineru_version=model,
)
def _parse_paddle(
source: Path, cfg: PipelineConfig, doc_id: str, target: Path, marker: Path
) -> ParseResult:
"""Parse one document with PaddleOCR-VL, cached exactly like the Mistral path.
The cached `paddle_ocr.json` is `paddle_ocr.run_ocr`'s output — the page results
with every figure's bytes already embedded, so a cache hit needs no network.
"""
out = target / "paddle_ocr.json"
if marker.exists() and out.exists():
meta = json.loads(marker.read_text(encoding="utf-8"))
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=out, middle_json=None,
pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0),
from_cache=True, backend_recorded="paddleocr-vl",
mineru_version=meta.get("model") or cfg.paddle_model,
)
from .paddle_ocr import run_ocr
started = time.perf_counter()
response = run_ocr(source, cfg) # hosted API; needs PADDLEOCR_ACCESS_TOKEN
seconds = time.perf_counter() - started
pages = len(response.get("pages", []))
model = response.get("model") or cfg.paddle_model
staging = target.with_name(target.name + ".partial")
shutil.rmtree(staging, ignore_errors=True)
staging.mkdir(parents=True, exist_ok=True)
(staging / "paddle_ocr.json").write_text(
json.dumps(response, ensure_ascii=False), encoding="utf-8"
)
shutil.rmtree(target, ignore_errors=True)
staging.rename(target)
marker.write_text(
json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id, "model": model}),
encoding="utf-8",
)
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=out, middle_json=None,
pages=pages, seconds=seconds, from_cache=False,
backend_recorded="paddleocr-vl", mineru_version=model,
)
def parse_document(source: Path, cfg: PipelineConfig) -> ParseResult:
"""Parse one document. A cache hit is not re-run."""
doc_id = source.stem
target = cfg.cache_dir / "parse" / cache_key(source, cfg)
marker = target / ".done"
if cfg.backend == "mistral":
return _parse_mistral(source, cfg, doc_id, target, marker)
if cfg.backend == "paddleocr":
return _parse_paddle(source, cfg, doc_id, target, marker)
if marker.exists():
content, middle = _find_output(target)
if content:
meta = json.loads(marker.read_text(encoding="utf-8"))
backend_recorded, version = _read_middle(middle)
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=content, middle_json=middle,
pages=meta.get("pages", 0), seconds=meta.get("seconds", 0.0),
from_cache=True,
backend_recorded=backend_recorded, mineru_version=version,
)
shutil.rmtree(target, ignore_errors=True) # cache corrupt, redo it
from mineru.cli.common import do_parse, read_fn
# MinerU dispatches on prefixes: "pipeline", then `startswith("vlm-")` and
# `startswith("hybrid-")`. The bare names "vlm"/"hybrid" match NONE of those
# branches — MinerU simply does nothing, without error, and the failure only
# surfaces here as "content_list not found", which points at the wrong place
# entirely. The short names stay in the config because they are part of the
# cache key and read better; the translation happens here, right before
# MinerU is touched.
mineru_backend = {
"pipeline": "pipeline",
"vlm": "vlm-engine",
"hybrid": "hybrid-engine",
}[cfg.backend]
# The playground does this before parsing; aliases like "latin"/"en" are
# resolved to a canonical code. Matched here so results are identical.
# Cache missed, so MinerU really is about to run — this is the point where a
# GPU stops being optional.
require_gpu(cfg)
lang = cfg.canonical_lang()
# with_suffix() REPLACES the last suffix, and the folder name ends in the
# MinerU version ("...-3.4.4") — so ".4" would be thrown away and two
# different versions would collide on one temp directory.
staging = target.with_name(target.name + ".partial")
shutil.rmtree(staging, ignore_errors=True)
staging.mkdir(parents=True, exist_ok=True)
pdf_bytes = read_fn(source)
started = time.perf_counter()
do_parse(
output_dir=str(staging),
# NOT doc_id: this becomes a folder name, and a long document name
# exceeds MAX_PATH once MinerU appends a SHA-256 image filename to it.
pdf_file_names=[_mineru_name(doc_id, source)],
pdf_bytes_list=[pdf_bytes],
p_lang_list=[lang],
backend=mineru_backend,
formula_enable=cfg.formula_enable,
table_enable=cfg.table_enable,
f_draw_layout_bbox=cfg.write_debug_pdf,
f_draw_span_bbox=cfg.write_debug_pdf,
start_page_id=cfg.start_page,
end_page_id=cfg.end_page,
**({"effort": cfg.effort} if cfg.backend == "hybrid" else {}),
)
seconds = time.perf_counter() - started
content, middle = _find_output(staging)
if content is None:
raise RuntimeError(
f"MinerU finished but no *_content_list.json was found in {staging}"
)
pages = _count_pages(source, cfg)
staging.rename(target)
content = target / content.relative_to(staging)
middle = target / middle.relative_to(staging) if middle else None
marker.write_text(
json.dumps({"pages": pages, "seconds": seconds, "doc_id": doc_id}),
encoding="utf-8",
)
backend_recorded, version = _read_middle(middle)
return ParseResult(
doc_id=doc_id, source=source, cache_dir=target,
content_list=content, middle_json=middle,
pages=pages, seconds=seconds, from_cache=False,
backend_recorded=backend_recorded, mineru_version=version,
)
def _count_pages(source: Path, cfg: PipelineConfig) -> int:
if cfg.end_page is not None:
return cfg.end_page - cfg.start_page + 1
try:
from pypdf import PdfReader
return len(PdfReader(str(source)).pages)
except Exception:
return 0