Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 37 additions & 5 deletions apps/worker/app/services/document_parser/formats/pdf/parser.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from loguru import logger

from shared.core.config import settings
from shared.core.exceptions.domain_exceptions import StorageServiceException
from shared.services.storage.job_file_storage import JobFileStorage


Expand Down Expand Up @@ -98,10 +99,11 @@ def _parse_pdf_via_shards(
1. PROFILE → shard plan + TOC
2. map_agent_shards → 1:1 MinerU shards
3. split_pdf (exclude TOC pages)
4. MinerU per shard (parallel)
5. **Per-shard heading prediction** (parallel)
6. Merge lines_with_heading + images
7. parse_md Phase B (skip TOC detection + heading prediction)
4. Upload split shards to S3 and confirm HeadObject succeeds
5. MinerU per shard (parallel)
6. **Per-shard heading prediction** (parallel)
7. Merge lines_with_heading + images
8. parse_md Phase B (skip TOC detection + heading prediction)
"""
from concurrent.futures import ThreadPoolExecutor, as_completed
from dataclasses import dataclass
Expand Down Expand Up @@ -200,7 +202,11 @@ def _parse_pdf_via_shards(
for shard_index, _shard_pdf_path in enumerate(shard_pdf_paths)
]

# 5. Parse each shard via MinerU (parallel)
with stage_timer("pdf.upload_shards", filename=filename):
_upload_temp_shard_pdfs(shard_pdf_paths, temp_shard_s3_keys)

# 5. Parse each shard via MinerU (parallel) only after S3 has the
# split objects. URL mode HeadObject then reuses those keys.
shard_output_dirs = [None] * len(shard_pdf_paths)

def _parse_single_shard(shard_idx, shard_pdf):
Expand Down Expand Up @@ -355,6 +361,32 @@ def _predict_shard_headings_with_tracking(
_cleanup_temp_shard_s3_assets(temp_shard_s3_keys)
_cleanup_local_shard_workspace(work_dir)


def _upload_temp_shard_pdfs(
shard_pdf_paths: list[str],
shard_s3_keys: list[str],
) -> None:
"""Upload split shards to S3 and confirm each object exists before MinerU."""
storage = JobFileStorage()
for shard_index, (shard_path, shard_s3_key) in enumerate(
zip(shard_pdf_paths, shard_s3_keys, strict=True)
):
logger.info(
f" ☁️ Uploading MinerU shard_{shard_index} to S3 ({shard_s3_key})"
)
storage.upload_source_file(shard_path, shard_s3_key)
existing_file = storage.verify_upload_exists(shard_s3_key)
if not existing_file.get("exists"):
raise StorageServiceException(
internal_message=(
"Temporary MinerU shard was uploaded but is not visible "
f"in S3: {shard_s3_key}"
),
operation="verify_source_object",
)
logger.info(f" ✅ Uploaded MinerU shard_{shard_index}: {shard_s3_key}")


def _build_temp_shard_s3_key(
*,
source_s3_key: str | None,
Expand Down
11 changes: 11 additions & 0 deletions apps/worker/tests/contract/test_parse_task_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -982,8 +982,18 @@ def test_oversized_pdf_happy_path_uses_shard_pipeline_without_external_services(
calls: dict[str, object] = {}
parse_s3_keys: list[str | None] = []
deleted_s3_keys: list[str] = []
uploaded_s3_keys: list[str] = []

class _FakeJobFileStorage:
def upload_source_file(self, local_file_path: str, storage_key: str) -> dict[str, str]:
assert Path(local_file_path).exists()
uploaded_s3_keys.append(storage_key)
return {"etag": "test"}

def verify_upload_exists(self, storage_key: str) -> dict[str, bool]:
assert storage_key in uploaded_s3_keys
return {"exists": True}

def delete_upload_file(self, storage_key: str) -> bool:
deleted_s3_keys.append(storage_key)
return True
Expand Down Expand Up @@ -1100,6 +1110,7 @@ def _identity_eval_md_headings(
"tmp/mineru-shards/job-oversized/shard_1.pdf",
]
assert parse_s3_keys == expected_s3_keys
assert uploaded_s3_keys == expected_s3_keys
assert deleted_s3_keys == expected_s3_keys
assert not (output_dir / "_shards").exists()
assert list(df["type"]) == ["PTXT", "PTXT"]
Expand Down
22 changes: 22 additions & 0 deletions apps/worker/tests/unit/test_mineru_pdf_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -279,3 +279,25 @@ def test_url_poll_other_pdf_error_does_not_direct_upload(
poll_mineru_task.assert_called_once()
request_upload_target.assert_not_called()
upload_file.assert_not_called()


def test_resolve_reuses_existing_object_without_reupload(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
) -> None:
monkeypatch.setattr(pdf_service.settings, "MINERU_UPLOAD_MODE_ENABLED", False)
storage = Mock()
storage.verify_upload_exists.return_value = {"exists": True, "size": 12}
monkeypatch.setattr(pdf_service, "JobFileStorage", Mock(return_value=storage))
local_pdf = tmp_path / "source.pdf"
local_pdf.write_bytes(b"%PDF-1.4")
s3_key = "tmp/mineru-shards/job_abc/shard_0.pdf"

resolved = pdf_service.resolve_mineru_source_s3_key(
s3_key,
local_file_path=str(local_pdf),
)

assert resolved == s3_key
storage.upload_source_file.assert_not_called()

237 changes: 237 additions & 0 deletions apps/worker/tests/unit/test_pdf_shard_s3_upload.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,237 @@
"""Split shards must be uploaded to S3 before MinerU URL mode starts."""

from __future__ import annotations

import os
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import Mock, call

import pytest

os.environ.setdefault("DATABASE_URL", "postgresql+asyncpg://test:test@localhost/test")
os.environ.setdefault("TMP_PATH", "/tmp/knowhere-test")
os.environ.setdefault("S3_BUCKET_NAME", "test-uploads")
os.environ.setdefault("S3_ACCESS_KEY_ID", "test")
os.environ.setdefault("S3_SECRET_ACCESS_KEY", "test")
os.environ.setdefault("S3_TEMP_PATH", "/tmp")

from app.services.document_agent.manifest import (
PageAnatomyMap,
Shard,
ShardPlan,
TocResult,
)
from app.services.document_parser.formats.pdf import parser as pdf_parser
from shared.core.exceptions.domain_exceptions import StorageServiceException


def test_upload_temp_shard_pdfs_writes_each_local_file(
monkeypatch: pytest.MonkeyPatch,
) -> None:
storage = Mock()
storage.verify_upload_exists.return_value = {"exists": True}
monkeypatch.setattr(pdf_parser, "JobFileStorage", Mock(return_value=storage))

pdf_parser._upload_temp_shard_pdfs(
["/tmp/shard_0.pdf", "/tmp/shard_1.pdf"],
[
"tmp/mineru-shards/job_abc/shard_0.pdf",
"tmp/mineru-shards/job_abc/shard_1.pdf",
],
)

assert storage.upload_source_file.call_args_list == [
call("/tmp/shard_0.pdf", "tmp/mineru-shards/job_abc/shard_0.pdf"),
call("/tmp/shard_1.pdf", "tmp/mineru-shards/job_abc/shard_1.pdf"),
]
assert storage.verify_upload_exists.call_args_list == [
call("tmp/mineru-shards/job_abc/shard_0.pdf"),
call("tmp/mineru-shards/job_abc/shard_1.pdf"),
]


def _toc_excluded_profile(tmp_path: Path, job_id: str) -> SimpleNamespace:
return SimpleNamespace(
anatomy=PageAnatomyMap(
job_id=job_id,
file_path=str(tmp_path / "doc.pdf"),
page_count=2,
page_features=[],
page_labels=[],
toc_result=TocResult(toc_pages=[1], method="vlm_batch"),
shard_plan=ShardPlan(
enabled=True,
reason="too_large",
shards=[
Shard(
shard_index=0,
page_start=1,
page_end=2,
page_offset=0,
anchor_type="toc_leaf_boundary",
anchor_evidence="Intro",
)
],
),
)
)


def test_shard_pipeline_uploads_before_mineru(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
) -> None:
events: list[str] = []

def fake_split_pdf(pdf_path, shards, work_dir, exclude_pages=None):
events.append("split")
shard_path = Path(work_dir) / "shard_0.pdf"
shard_path.write_bytes(b"%PDF-1.4")
return [str(shard_path)], None

class _FakeStorage:
def upload_source_file(
self, local_file_path: str, storage_key: str
) -> dict[str, str]:
events.append(f"upload:{storage_key}")
assert Path(local_file_path).exists()
return {"etag": "ok"}

def verify_upload_exists(self, storage_key: str) -> dict[str, bool]:
events.append(f"verify:{storage_key}")
return {"exists": True}

def delete_upload_file(self, storage_key: str) -> bool:
events.append(f"delete:{storage_key}")
return True

def fake_parse_via_full(shard_pdf, shard_filename, shard_out, s3_key=None):
events.append(f"parse:{s3_key}")
Path(shard_out).mkdir(parents=True, exist_ok=True)
(Path(shard_out) / "full.md").write_text("# Intro\nBody\n", encoding="utf-8")

monkeypatch.setattr(
"app.services.document_parser.formats.pdf.shard_splitter.split_pdf",
fake_split_pdf,
)
monkeypatch.setattr(pdf_parser, "JobFileStorage", _FakeStorage)
monkeypatch.setattr(pdf_parser, "parse_via_full", fake_parse_via_full)
monkeypatch.setattr(
"app.services.document_parser.formats.markdown.parser.eval_md_headings",
lambda md_lines, *args, **kwargs: list(md_lines),
)
monkeypatch.setattr(
pdf_parser,
"parse_md",
lambda *_args, **kwargs: {"ok": True, "lines": kwargs["lines_with_heading"]},
)
monkeypatch.setattr(
"app.services.document_parser.formats.pdf.shard_merger.merge_images",
lambda *_args, **_kwargs: None,
)

pdf_parser._parse_pdf_via_shards(
str(tmp_path / "doc.pdf"),
"doc.pdf",
str(tmp_path / "out"),
{"smart_title_parse": False, "model_name": "test-model"},
profile=_toc_excluded_profile(tmp_path, "job-upload-first"),
s3_key="uploads/job-upload-first.pdf",
job_id="job-upload-first",
)

assert events[:4] == [
"split",
"upload:tmp/mineru-shards/job-upload-first/shard_0.pdf",
"verify:tmp/mineru-shards/job-upload-first/shard_0.pdf",
"parse:tmp/mineru-shards/job-upload-first/shard_0.pdf",
]


def test_failed_shard_upload_does_not_start_mineru(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
) -> None:
parse_via_full = Mock()

def fake_split_pdf(pdf_path, shards, work_dir, exclude_pages=None):
shard_path = Path(work_dir) / "shard_0.pdf"
shard_path.write_bytes(b"%PDF-1.4")
return [str(shard_path)], None

class _FailingStorage:
def upload_source_file(
self, local_file_path: str, storage_key: str
) -> dict[str, str]:
raise StorageServiceException(
internal_message="S3 upload failed",
operation="upload_local_file",
)

def delete_upload_file(self, storage_key: str) -> bool:
return False

monkeypatch.setattr(
"app.services.document_parser.formats.pdf.shard_splitter.split_pdf",
fake_split_pdf,
)
monkeypatch.setattr(pdf_parser, "JobFileStorage", _FailingStorage)
monkeypatch.setattr(pdf_parser, "parse_via_full", parse_via_full)

with pytest.raises(StorageServiceException):
pdf_parser._parse_pdf_via_shards(
str(tmp_path / "doc.pdf"),
"doc.pdf",
str(tmp_path / "out"),
{"smart_title_parse": False, "model_name": "test-model"},
profile=_toc_excluded_profile(tmp_path, "job-upload-fail"),
s3_key="uploads/job-upload-fail.pdf",
job_id="job-upload-fail",
)

parse_via_full.assert_not_called()


def test_unverified_shard_upload_does_not_start_mineru(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
) -> None:
parse_via_full = Mock()

def fake_split_pdf(pdf_path, shards, work_dir, exclude_pages=None):
shard_path = Path(work_dir) / "shard_0.pdf"
shard_path.write_bytes(b"%PDF-1.4")
return [str(shard_path)], None

class _InvisibleStorage:
def upload_source_file(
self, local_file_path: str, storage_key: str
) -> dict[str, str]:
return {"etag": "ok"}

def verify_upload_exists(self, storage_key: str) -> dict[str, bool]:
return {"exists": False}

def delete_upload_file(self, storage_key: str) -> bool:
return True

monkeypatch.setattr(
"app.services.document_parser.formats.pdf.shard_splitter.split_pdf",
fake_split_pdf,
)
monkeypatch.setattr(pdf_parser, "JobFileStorage", _InvisibleStorage)
monkeypatch.setattr(pdf_parser, "parse_via_full", parse_via_full)

with pytest.raises(StorageServiceException, match="not visible"):
pdf_parser._parse_pdf_via_shards(
str(tmp_path / "doc.pdf"),
"doc.pdf",
str(tmp_path / "out"),
{"smart_title_parse": False, "model_name": "test-model"},
profile=_toc_excluded_profile(tmp_path, "job-upload-invisible"),
s3_key="uploads/job-upload-invisible.pdf",
job_id="job-upload-invisible",
)

parse_via_full.assert_not_called()
Loading