ishaq101's picture
/fix planner and unstructured pipeline (#23)
7124acf
Raw History Blame Contribute Delete
16.3 kB
"""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()