Download src/knowledge_parsing/run.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 16.3 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_parsing/run.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/knowledge_parsing/run.py
-
curl -L -o run.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_parsing/run.py
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.") | |
| 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() | |