"""Stages 1-3 wired together: a folder of documents in, chunks + manifest out. py -m src.knowledge_parsing.run --input data/knowledge_docs/ py -m src.knowledge_parsing.run --input data/knowledge_docs/ --backend pipeline Deliberate properties: - one document failing does NOT stop the rest (recorded in failures.jsonl) - resumable: documents already finished are skipped (--no-resume to force) - MinerU's own output is never touched, only pointed at from runs/ """ from __future__ import annotations import argparse import json import sys import time import traceback from dataclasses import dataclass from datetime import datetime from pathlib import Path from typing import Any from .assets import AssetSink, LocalAssetSink, NullAssetSink from .checks import check_items, vocabulary_per_page from .config import PipelineConfig from .contracts import ParsedDocument from .manifest import DocumentRecord, Manifest from .normalize import normalise, normaliser_version from .normalize_mistral import normalise_mistral from .normalize_mistral import normaliser_version as mistral_normaliser_version from .normalize_paddle import normalise_paddle, unknown_labels from .normalize_paddle import normaliser_version as paddle_normaliser_version from .parse import ParseResult, file_hash, mineru_version, parse_document EXTENSIONS = {".pdf", ".docx", ".pptx", ".doc", ".ppt"} def collect_documents(folder: Path) -> list[Path]: return sorted(p for p in folder.rglob("*") if p.suffix.lower() in EXTENSIONS) def _count(n: int, word: str) -> str: """`3 chunks` / `1 chunk` — console output only, never parsed by anything.""" return f"{n} {word}" if n == 1 else f"{n} {word}s" def _reject_duplicate_names(documents: list[Path]) -> None: """Two files with one name is ambiguous, and the ambiguity is silent. `doc_id` is the file name, `chunk_id` is built from `doc_id`, and every `term_id` / `formula_id` / `rule_id` downstream is derived from it in turn. Two documents sharing a name therefore claim the same identity — and because `resume` skips a document whose `chunks.json` already exists, the second one is quietly reported as "already present" and never parsed at all. The run ends `ok / failed : 2 / 0` with one document missing. This fails instead, before any parsing is paid for. It deliberately does NOT disambiguate automatically: renaming one to `Laporan_2` would silently change the identifier an expert's approve/reject decisions hang on, which is a worse failure than stopping and asking. """ by_stem: dict[str, list[Path]] = {} for src in documents: by_stem.setdefault(src.stem, []).append(src) clashes = {stem: paths for stem, paths in by_stem.items() if len(paths) > 1} if not clashes: return lines = ["Two or more documents share one name, and doc_id is built from it:"] for stem, paths in sorted(clashes.items()): lines.append(f" '{stem}'") lines.extend(f" {p}" for p in sorted(paths)) lines.append(" -> rename one, or parse them in separate runs.") raise SystemExit("\n".join(lines)) def _warn_identical_content(documents: list[Path]) -> None: """Same bytes under two names: legal, wasteful, and worth saying out loud. A document downloaded twice arrives as `X.pdf` and `X (2).pdf`. Both are real files with different names, so neither this module nor extraction may silently drop one — which of them is authoritative is a curation decision. But both are parsed (the second free, from the content-addressed cache) and both reach extraction as separate documents, so the expert reviews the same terms twice with nothing indicating why. Detection is cheap: identical content means an identical file hash. """ by_hash: dict[str, list[Path]] = {} for src in documents: by_hash.setdefault(file_hash(src), []).append(src) for paths in by_hash.values(): if len(paths) > 1: print(f"⚠ {len(paths)} documents have identical content:") for p in sorted(paths): print(f" {p.name}") print(" Both are parsed and both reach extraction as separate documents.") @dataclass class _ParseOutcome: """What `_parse_one` produces: the artifact plus the run-level numbers. `parse_one` (the library entry point) returns only `.artifact`; `run()` reads the rest — parse timing, cache hit, quality checks — for its per-document manifest entry. Those live on `ParseResult`/`checks`, not on the artifact, so surfacing them here keeps `run()`'s manifest exactly as rich as before without recomputing anything on a second, drift-prone path. """ artifact: ParsedDocument parse_result: ParseResult checks: dict[str, Any] normalise_seconds: float def _parse_one( source: Path, cfg: PipelineConfig, doc_id: str, asset_sink: AssetSink | None = None, ) -> _ParseOutcome: """The per-document core: parse, normalise, assemble the `ParsedDocument`. The single place a `ParsedDocument` is built. Both `parse_one` and `run()` go through here, so the two callers cannot drift — a change to the artifact is a change in one spot, and the seam tests catch it because the shape is versioned. """ result = parse_document(source, cfg) if cfg.backend == "paddleocr": # Added 2026-09-24, reached from v2 document processing only (the CLI's # --backend choices are unchanged). Unknown layout labels are kept as prose # and listed here so a new label gets noticed. response = json.loads(result.content_list.read_text(encoding="utf-8")) checks: dict[str, Any] = { "backend": "paddleocr-vl", "unknown_labels": unknown_labels(response), } t0 = time.perf_counter() chunks = normalise_paddle(response, doc_id, asset_sink=asset_sink) normalise_seconds = time.perf_counter() - t0 elif cfg.backend == "mistral": response = json.loads(result.content_list.read_text(encoding="utf-8")) # MinerU's content-drift check and the word-boundary vocabulary are both # MinerU-specific — Mistral emits neither a content_list nor # character-spaced formulas, so neither applies. checks = {"backend": "mistral-ocr"} t0 = time.perf_counter() chunks = normalise_mistral(response, doc_id, asset_sink=asset_sink) normalise_seconds = time.perf_counter() - t0 else: items = json.loads(result.content_list.read_text(encoding="utf-8")) # The source PDF goes in too: without something to compare against, a # number that changed silently (the EOQ bug) would never be detected. checks = check_items( items, cfg.min_latex_len, source_pdf=source, page_offset=cfg.start_page, ) # The source page vocabulary, used to restore the word boundaries lost # when MinerU writes formulas one character at a time. Re-keyed to # MinerU's page indices: page_idx 0 means page `start_page` in the PDF. raw_vocabulary = vocabulary_per_page(source) vocabulary = ( {p - cfg.start_page: words for p, words in raw_vocabulary.items() if p >= cfg.start_page} if raw_vocabulary else None ) t0 = time.perf_counter() chunks = normalise(items, doc_id, vocabulary) normalise_seconds = time.perf_counter() - t0 artifact = ParsedDocument( doc_id=doc_id, # The human-readable label stays the file stem even once `doc_id` becomes # a stable key — see the field's note in contracts.py. source_title=source.stem, chunks=chunks, source_path=str(source), content_hash=file_hash(source, length=64), n_pages=result.pages, parser_name={"mistral": "mistral-ocr", "paddleocr": "paddleocr-vl"}.get( cfg.backend, "mineru" ), parser_version=result.mineru_version or mineru_version(), # backend as the parser actually recorded it, not the config parser_backend=result.backend_recorded or cfg.backend, parser_config=cfg.fingerprint(), # Which normalisation code produced these chunks — the mistral and # MinerU normalisers are fingerprinted separately. normaliser_version=( mistral_normaliser_version() if cfg.backend == "mistral" else paddle_normaliser_version() if cfg.backend == "paddleocr" else normaliser_version() ), created_at=datetime.now().isoformat(timespec="seconds"), raw_output_dir=str(result.cache_dir), ) return _ParseOutcome( artifact=artifact, parse_result=result, checks=checks, normalise_seconds=normalise_seconds, ) def parse_one( source: Path, cfg: PipelineConfig, doc_id: str | None = None, asset_sink: AssetSink | None = None, ) -> ParsedDocument: """Parse one document and return the artifact in memory. Raises on failure. The single-document entry point for callers that persist the artifact themselves — the HTTP ingestion service, which parses one `documents` row at a time. No folder scan, no manifest, no file written here; `run()` owns all of that and calls the same core for its parse, so there is one artifact builder. `doc_id` is the identity every downstream `entity_id` derives from. Passing it explicitly is what lets ingestion supply a STABLE key (`documents.id`) so a file rename no longer mints all-new ids and orphans the expert approvals already recorded against the old ones. `None` keeps today's behaviour — the file stem — which is what `run()` passes. """ if doc_id is None: doc_id = source.stem return _parse_one(source, cfg, doc_id, asset_sink).artifact def run(cfg: PipelineConfig) -> Manifest: cfg.validate() documents = collect_documents(cfg.input_dir) if not documents: raise SystemExit(f"No documents found in {cfg.input_dir}") # Both run before anything is parsed: a 16-hour batch must not discover a # name clash on its last document. _reject_duplicate_names(documents) _warn_identical_content(documents) manifest = Manifest.new(cfg.to_dict()) manifest.mineru_version = mineru_version() run_dir = cfg.runs_dir / manifest.run_id run_dir.mkdir(parents=True, exist_ok=True) failures = run_dir / "failures.jsonl" print(f"Run {manifest.run_id} | backend={cfg.backend} | {_count(len(documents), 'document')}") print(f"Output to: {run_dir}\n") for n, source in enumerate(documents, 1): doc_id = source.stem target = run_dir / doc_id chunk_file = target / "chunks.json" if cfg.resume and chunk_file.exists(): print(f"[{n}/{len(documents)}] {doc_id} — skipped (already present)") continue print(f"[{n}/{len(documents)}] {doc_id} … ", end="", flush=True) record = DocumentRecord(doc_id=doc_id, source=str(source), status="ok") try: # Same core the ingestion service calls — the artifact is built once, # here and there. run() adds only the folder-batch concerns around it: # the manifest record and writing the result to disk. # Assets land beside the artifact they belong to, so a run folder is # self-describing and movable. `--no-assets` keeps the metadata and # drops only the bytes (see assets.NullAssetSink). sink: AssetSink = ( NullAssetSink() if cfg.no_assets else LocalAssetSink(target) ) outcome = _parse_one(source, cfg, doc_id, asset_sink=sink) artifact = outcome.artifact result = outcome.parse_result record.pages = result.pages record.parse_seconds = result.seconds record.from_cache = result.from_cache record.backend_recorded = result.backend_recorded record.checks = outcome.checks record.normalise_seconds = outcome.normalise_seconds record.chunk_count = len(artifact.chunks) target.mkdir(parents=True, exist_ok=True) chunk_file.write_text( artifact.model_dump_json(indent=2), encoding="utf-8", ) # A pointer to MinerU's own output — not copied, so it is not doubled. (target / "mineru-source.txt").write_text( str(result.cache_dir), encoding="utf-8" ) mark = " (cache)" if result.from_cache else f" {result.seconds:.1f}s" warnings = record.checks.get("warning_count", 0) mark += f" | {_count(len(artifact.chunks), 'chunk')}" if warnings: mark += f" | ⚠ {_count(warnings, 'warning')}" print("ok" + mark) except Exception as e: # one failure must not stop the rest record.status = "failed" record.error = f"{type(e).__name__}: {e}" print(f"FAILED — {record.error}") with failures.open("a", encoding="utf-8") as f: f.write(json.dumps({ "doc_id": doc_id, "source": str(source), "error": record.error, "trace": traceback.format_exc(), }, ensure_ascii=False) + "\n") manifest.documents.append(record) manifest.save(run_dir / "manifest.json") s = manifest.summary() print("\n--- summary ---") print(f" ok / failed : {s['documents_ok']} / {s['documents_failed']}") print(f" total pages : {s['total_pages']}") if s["seconds_per_page"]: print(f" seconds/page : {s['seconds_per_page']} " f"(from {s['pages_measured']} page(s) actually parsed)") if s["total_warnings"]: print(f" ⚠ quality warnings: {s['total_warnings']} — see manifest.json") print(f" manifest : {run_dir / 'manifest.json'}") return manifest def main() -> None: # A Windows console defaults to cp1252, which has no "⚠". Without this, # printing a quality warning raises UnicodeEncodeError INSIDE the per-document # try block — so a document carrying warnings is marked FAILED even though its # parse succeeded and was written. The inverse of the intent: the documents # whose quality needs attention are the ones thrown away. for stream in (sys.stdout, sys.stderr): try: stream.reconfigure(encoding="utf-8", errors="replace") # type: ignore[union-attr] except (AttributeError, OSError): pass p = argparse.ArgumentParser(description="Document parsing pipeline (MinerU)") p.add_argument("--input", type=Path, help="folder of documents") p.add_argument("--backend", choices=["pipeline", "vlm", "hybrid", "mistral"], help="mistral=API, no GPU (default); pipeline=CPU MinerU; " "vlm/hybrid=GPU MinerU") p.add_argument("--lang", default=None) # `hybrid` only. Not merely fast-vs-thorough: MinerU FORCE-DISABLES image # analysis at `medium`, so descriptions exist only at `high` — and always # on `vlm`. p.add_argument("--effort", choices=["medium", "high"], default=None) p.add_argument("--no-resume", action="store_true", help="reprocess everything") p.add_argument("--no-assets", action="store_true", help="keep asset metadata but do not write image bytes to disk") p.add_argument("--start-page", type=int, default=None) p.add_argument("--end-page", type=int, default=None) a = p.parse_args() cfg = PipelineConfig() if a.input: cfg.input_dir = a.input if a.backend: cfg.backend = a.backend if a.lang: cfg.lang = a.lang if a.effort: cfg.effort = a.effort if a.no_resume: cfg.resume = False if a.no_assets: cfg.no_assets = True if a.start_page is not None: cfg.start_page = a.start_page if a.end_page is not None: cfg.end_page = a.end_page run(cfg) if __name__ == "__main__": main()