from __future__ import annotations
import html
import json
import logging
import os
import re
import shutil
from datetime import datetime, timezone
from json import JSONDecodeError
from pathlib import Path
import fitz
import requests
from PIL import Image
from paperforge.config import paperforge_paths
from paperforge.worker import sync as _sync
from paperforge.worker._progress import progress_bar
def _read_dotenv(vault: Path, key: str) -> str:
"""Read *key* from vault/.env or System/PaperForge/.env, return empty if not found."""
for dotenv_path in (vault / ".env", paperforge_paths(vault).get("pipeline", Path()) / ".env"):
try:
for line in dotenv_path.read_text(encoding="utf-8").splitlines():
line = line.strip()
if line.startswith("#") or "=" not in line:
continue
k, _, v = line.partition("=")
if k.strip() == key:
return v.strip().strip('"').strip("'")
except (OSError, UnicodeDecodeError):
continue
return ""
try:
from paperforge.worker.ocr_roles import assign_block_role # noqa: F401
except ImportError:
pass
from paperforge.worker._retry import retry_with_meta
from paperforge.worker._utils import (
pipeline_paths,
read_json,
write_json,
)
from paperforge.worker.asset_index import refresh_index_entry
from paperforge.worker.ocr_artifacts import (
artifact_paths_for_key,
build_version_payload,
compute_json_hash,
compute_pdf_fingerprint,
)
from paperforge.worker.ocr_blocks import (
build_raw_blocks_for_result_lines,
build_structured_blocks,
write_raw_blocks_jsonl,
write_structured_blocks_jsonl,
)
from paperforge.worker.ocr_figures import (
build_figure_inventory,
write_figure_inventory,
)
from paperforge.worker.ocr_metadata import (
extract_frontmatter_candidates,
resolve_metadata,
write_resolved_metadata,
)
from paperforge.worker.ocr_tables import (
build_table_inventory,
write_table_inventory,
)
from paperforge.worker.sync import (
load_control_actions,
load_export_rows,
)
logger = logging.getLogger(__name__)
def ensure_ocr_meta(vault: Path, row: dict) -> dict:
paths = pipeline_paths(vault)
key = row["zotero_key"]
meta_path = paths["ocr"] / key / "meta.json"
meta_path.parent.mkdir(parents=True, exist_ok=True)
try:
meta = read_json(meta_path) if meta_path.exists() else {}
except Exception:
meta = {}
meta.setdefault("zotero_key", key)
meta.setdefault("source_pdf", row.get("pdf_path", ""))
meta.setdefault("ocr_provider", "PaddleOCR-VL-1.6")
meta.setdefault("mode", "async")
meta.setdefault("ocr_status", "pending")
meta.setdefault("ocr_job_id", "")
meta.setdefault("ocr_started_at", "")
meta.setdefault("ocr_finished_at", "")
meta.setdefault("page_count", 0)
meta.setdefault("markdown_path", "")
meta.setdefault("json_path", "")
try:
assets_path = str((paths["ocr"] / key / "assets").relative_to(vault)).replace("\\", "/")
except ValueError:
assets_path = str(paths["ocr"] / key / "assets")
meta.setdefault("assets_path", assets_path)
meta.setdefault("fulltext_md_path", "")
meta.setdefault("error", "")
meta.setdefault("retry_count", 0)
meta.setdefault("last_error", None)
meta.setdefault("last_attempt_at", None)
return meta
def _read_meta_or_empty(meta_path: Path) -> dict:
try:
return read_json(meta_path) if meta_path.exists() else {}
except Exception:
return {}
def validate_ocr_meta(paths: dict[str, Path], meta: dict) -> tuple[str, str]:
status = str(meta.get("ocr_status", "pending") or "pending").strip().lower()
if status != "done":
return (status, str(meta.get("error", "") or ""))
key = str(meta.get("zotero_key", "") or "").strip()
if not key:
return ("done_incomplete", "Missing zotero_key in OCR meta")
ocr_root = paths["ocr"] / key
fulltext_path = ocr_root / "fulltext.md"
json_path = ocr_root / "json" / "result.json"
page_count = int(meta.get("page_count", 0) or 0)
if not fulltext_path.exists():
return ("done_incomplete", "OCR fulltext.md missing")
if not json_path.exists():
return ("done_incomplete", "OCR result.json missing")
fulltext_size = fulltext_path.stat().st_size
json_size = json_path.stat().st_size
if page_count < 1:
return ("done_incomplete", "OCR page_count invalid")
if fulltext_size < 500:
return ("done_incomplete", "OCR fulltext.md too small")
if json_size < 1000:
return ("done_incomplete", "OCR result.json too small")
try:
rendered_pages = fulltext_path.read_text(encoding="utf-8").count(""]
reference_blocks: list[dict] = []
reference_continuations: list[dict] = []
footnotes: list[str] = []
footnote_counter = 0
references_heading_seen = False
references_section_active = False
references_emitted = False
references_insert_index: int | None = None
rendered_cluster_ids: set[int] = set()
rendered_caption_media_ids: set[int] = set()
first_page_meta_done = page_index != 1
affiliation_buffer: list[str] = []
deferred_meta: list[str] = []
for block in blocks:
label = block.get("block_label", "")
content = block.get("block_content", "")
bbox = block.get("block_bbox", [0, 0, 0, 0])
composite_region = composite_by_block_id.get(block.get("block_id"))
if composite_region:
if not composite_region.get("rendered"):
composite_region["rendered"] = True
region_bbox = composite_region["bbox"]
asset_path = (
images_dir
/ "blocks"
/ f"page_{page_index:03d}_figure_{region_bbox[0]}_{region_bbox[1]}_{region_bbox[2]}_{region_bbox[3]}.jpg"
)
if page_image and crop_block_asset(page_image, region_bbox, asset_path):
rendered.append(_image_embed_for_obsidian(asset_path.relative_to(vault)))
continue
if label in {"header", "header_image", "footer", "footer_image", "number"}:
continue
if label == "doc_title":
rendered.append(f"# {clean_block_text(content)}")
continue
if label == "paragraph_title":
title = clean_block_text(content)
if (
title.lower() == "check for updates"
or is_subfigure_label(title)
or is_reference_tail_noise_line(title)
or is_embedded_figure_text_block(block, blocks, page_width=ocr_width, page_height=ocr_height)
):
continue
if page_index == 1 and affiliation_buffer:
rendered.extend(affiliation_buffer)
affiliation_buffer.clear()
if page_index == 1 and deferred_meta:
rendered.extend(deferred_meta)
deferred_meta.clear()
first_page_meta_done = True
if references_heading_seen and (not references_emitted) and title.lower() != "references":
if reference_blocks:
append_reference_section(rendered, reference_blocks, reference_continuations)
references_emitted = True
references_section_active = False
rendered.append(f"### {title}")
if title.lower() == "references":
references_heading_seen = True
references_section_active = True
references_insert_index = len(rendered)
continue
if label == "abstract":
if page_index == 1 and affiliation_buffer:
rendered.extend(affiliation_buffer)
affiliation_buffer.clear()
if page_index == 1 and deferred_meta:
rendered.extend(deferred_meta)
deferred_meta.clear()
first_page_meta_done = True
rendered.append(clean_block_text(content))
continue
if label == "reference_content":
text = clean_block_text(content)
if text and (not is_reference_tail_noise_line(text)):
if is_reference_lead_block(text):
reference_blocks.append(block)
else:
reference_continuations.append(block)
continue
if label == "text":
text = clean_block_text(content)
if (
not text
or is_subfigure_label(text)
or is_frontmatter_noise_line(text)
or is_reference_tail_noise_line(text)
or is_embedded_figure_text_block(block, blocks, page_width=ocr_width, page_height=ocr_height)
):
continue
if references_section_active and raw_reference_blocks and bbox[1] >= first_reference_y - 10:
reference_continuations.append(block)
continue
if is_numbered_figure_caption(text):
linked_media = figure_caption_map.get(block.get("block_id"), []) or table_caption_map.get(
block.get("block_id"), []
)
if linked_media and page_image:
rendered_caption_media_ids.update(item.get("block_id") for item in linked_media)
union_bbox = [
min(item["block_bbox"][0] for item in linked_media),
min(item["block_bbox"][1] for item in linked_media),
max(item["block_bbox"][2] for item in linked_media),
max(item["block_bbox"][3] for item in linked_media),
]
asset_kind = "figure" if block.get("block_id") in figure_caption_map else "table"
asset_path = (
images_dir
/ "blocks"
/ f"page_{page_index:03d}_{asset_kind}_{union_bbox[0]}_{union_bbox[1]}_{union_bbox[2]}_{union_bbox[3]}.jpg"
)
if crop_block_asset(page_image, union_bbox, asset_path):
rendered.append(_image_embed_for_obsidian(asset_path.relative_to(vault)))
if asset_kind == "table":
for item in linked_media:
if item.get("block_label") == "table" and item.get("block_content"):
rendered.append(clean_block_text(item.get("block_content", "")))
break
rendered.append(text)
continue
if page_index == 1 and (not first_page_meta_done):
if handle_first_page_metadata_lines(text, rendered, affiliation_buffer, deferred_meta, footnotes):
continue
rendered.append(text)
continue
if label == "display_formula":
formula = clean_block_text(content)
formula = formula.strip()
if formula.startswith("$$") and formula.endswith("$$") and (len(formula) >= 4):
formula = formula[2:-2].strip()
rendered.append(f"$$\n{formula}\n$$")
continue
if label == "formula_number":
rendered.append(clean_block_text(content))
continue
if label in {"table", "image", "chart"}:
if block.get("block_id") in caption_linked_media_ids:
continue
if block.get("block_id") in rendered_caption_media_ids:
continue
if label in {"image", "chart"}:
cluster_id = cluster_index.get(block.get("block_id", -1))
if cluster_id is None or cluster_id in rendered_cluster_ids:
continue
rendered_cluster_ids.add(cluster_id)
cluster_blocks = clusters[cluster_id]
bbox = [
min(item["block_bbox"][0] for item in cluster_blocks),
min(item["block_bbox"][1] for item in cluster_blocks),
max(item["block_bbox"][2] for item in cluster_blocks),
max(item["block_bbox"][3] for item in cluster_blocks),
]
asset_name = f"page_{page_index:03d}_figure_{bbox[0]}_{bbox[1]}_{bbox[2]}_{bbox[3]}.jpg"
else:
asset_name = f"page_{page_index:03d}_{label}_{bbox[0]}_{bbox[1]}_{bbox[2]}_{bbox[3]}.jpg"
asset_path = images_dir / "blocks" / asset_name
if page_image and crop_block_asset(page_image, bbox, asset_path):
rendered.append(_image_embed_for_obsidian(asset_path.relative_to(vault)))
if label == "table" and content:
rendered.append(clean_block_text(content))
continue
if label == "figure_title":
caption_text = clean_block_text(content)
if is_subfigure_label(caption_text):
continue
if not is_formal_figure_legend(caption_text):
continue
linked_media = figure_caption_map.get(block.get("block_id"), []) or table_caption_map.get(
block.get("block_id"), []
)
if linked_media and page_image:
rendered_caption_media_ids.update(item.get("block_id") for item in linked_media)
union_bbox = [
min(item["block_bbox"][0] for item in linked_media),
min(item["block_bbox"][1] for item in linked_media),
max(item["block_bbox"][2] for item in linked_media),
max(item["block_bbox"][3] for item in linked_media),
]
asset_kind = "figure" if block.get("block_id") in figure_caption_map else "table"
asset_path = (
images_dir
/ "blocks"
/ f"page_{page_index:03d}_{asset_kind}_{union_bbox[0]}_{union_bbox[1]}_{union_bbox[2]}_{union_bbox[3]}.jpg"
)
if crop_block_asset(page_image, union_bbox, asset_path):
rendered.append(_image_embed_for_obsidian(asset_path.relative_to(vault)))
if asset_kind == "table":
for item in linked_media:
if item.get("block_label") == "table" and item.get("block_content"):
rendered.append(clean_block_text(item.get("block_content", "")))
break
rendered.append(caption_text)
continue
if label == "vision_footnote":
if is_formal_figure_legend(content):
rendered.append(clean_block_text(content))
continue
if is_embedded_vision_footnote_block(block, blocks, page_width=ocr_width, page_height=ocr_height):
continue
entries = parse_vision_footnote_entries(content)
marker_to_id = {}
for marker, body in entries:
footnote_counter += 1
footnote_id = f"p{page_index}-fn{footnote_counter}"
marker_to_id[marker] = footnote_id
footnotes.append(f"[^{footnote_id}]: {body or clean_block_text(content)}")
if rendered and marker_to_id:
last = rendered[-1]
if last.startswith("
"):
rendered[-1] = attach_table_footnotes(last, marker_to_id)
else:
for marker, footnote_id in marker_to_id.items():
rendered[-1] = attach_footnote_reference(rendered[-1], marker, footnote_id)
continue
if footnotes:
rendered.append("")
rendered.extend(footnotes)
rendered = dedupe_page_media_lines(rendered)
if reference_blocks and not references_emitted:
reference_lines = build_reference_section_lines(reference_blocks, reference_continuations)
if references_insert_index is not None:
rendered[references_insert_index:references_insert_index] = reference_lines
else:
rendered.extend(reference_lines)
return [part for part in rendered if part]
def postprocess_ocr_result(vault: Path, key: str, all_results: list[dict]) -> tuple[int, str, str, str]:
paths = pipeline_paths(vault)
ocr_root = paths["ocr"] / key
json_dir = ocr_root / "json"
images_dir = ocr_root / "images"
page_cache_dir = ocr_root / "pages"
meta_path = ocr_root / "meta.json"
json_dir.mkdir(parents=True, exist_ok=True)
images_dir.mkdir(parents=True, exist_ok=True)
page_cache_dir.mkdir(parents=True, exist_ok=True)
page_num = 0
meta = read_json(meta_path) if meta_path.exists() else {}
source_pdf = Path(meta.get("source_pdf", "")) if meta.get("source_pdf") else None
pdf_doc = None
try:
if source_pdf and source_pdf.exists():
pdf_doc = fitz.open(str(source_pdf))
for page_payload in all_results:
for res in page_payload.get("layoutParsingResults", []):
page_num += 1
render_page_blocks(vault, page_num, res, images_dir, page_cache_dir, pdf_doc=pdf_doc)
finally:
if pdf_doc is not None:
pdf_doc.close()
write_json(json_dir / "result.json", all_results)
# --- Phase 1: raw metadata and source metadata ---
artifacts = artifact_paths_for_key(vault, key)
# raw_meta.json
source_pdf_path = Path(meta.get("source_pdf", "")) if meta.get("source_pdf") else None
pdf_fingerprint = compute_pdf_fingerprint(source_pdf_path) if source_pdf_path and source_pdf_path.exists() else "unknown"
result_hash = compute_json_hash(all_results)
now_iso = datetime.now(timezone.utc).isoformat()
raw_meta = {
"ocr_provider": meta.get("ocr_provider", "PaddleOCR"),
"ocr_model": meta.get("ocr_model", "PaddleOCR"),
"ocr_raw_schema_version": "1.0.0",
"created_at": now_iso,
"updated_at": now_iso,
"pdf_fingerprint": pdf_fingerprint,
"result_hash": result_hash,
}
write_json(artifacts.raw_meta, raw_meta)
# source_metadata.json
source_meta = {
"zotero_key": meta.get("zotero_key", key),
"source_pdf": meta.get("source_pdf", ""),
"title": meta.get("title", ""),
"authors": meta.get("authors", []),
"year": meta.get("year", 0),
"journal": meta.get("journal", ""),
"doi": meta.get("doi", ""),
"bbt_citekey": meta.get("bbt_citekey", ""),
"source": "zotero_bbt",
}
write_json(artifacts.source_metadata, source_meta)
# Emit canonical raw blocks
artifacts.blocks_raw.parent.mkdir(parents=True, exist_ok=True)
artifacts.blocks_structured.parent.mkdir(parents=True, exist_ok=True)
all_raw_blocks = build_raw_blocks_for_result_lines(key, all_results)
write_raw_blocks_jsonl(artifacts.blocks_raw, all_raw_blocks)
structured = build_structured_blocks(all_raw_blocks)
write_structured_blocks_jsonl(artifacts.blocks_structured, structured)
# Write role-level span profiles
from paperforge.worker.ocr_profiles import write_role_span_profiles
write_role_span_profiles(structured, artifacts.blocks_structured.parent)
# --- Phase 2: metadata resolution ---
metadata_dir = ocr_root / "metadata"
metadata_dir.mkdir(parents=True, exist_ok=True)
frontmatter_candidates = extract_frontmatter_candidates(artifacts.blocks_structured)
resolved = resolve_metadata(source_meta, frontmatter_candidates)
write_resolved_metadata(metadata_dir / "resolved_metadata.json", resolved)
# --- Phase 2: figure inventory ---
figure_inventory = build_figure_inventory(structured)
write_figure_inventory(
artifacts.blocks_structured.parent / "figure_inventory.json",
figure_inventory,
)
# --- Phase 2: table inventory ---
table_inventory = build_table_inventory(structured)
write_table_inventory(
artifacts.blocks_structured.parent / "table_inventory.json",
table_inventory,
)
# --- Phase 2: object artifacts ---
from paperforge.worker.ocr_objects import extract_and_write_objects
ocr_asset_root = ocr_root / "assets"
ocr_render_root = ocr_root / "render"
page_dimensions_by_page: dict[int, tuple[int, int]] = {}
for block in structured:
page = int(block.get("page", 0) or 0)
width = int(block.get("page_width", 0) or 0)
height = int(block.get("page_height", 0) or 0)
if page and width and height and page not in page_dimensions_by_page:
page_dimensions_by_page[page] = (width, height)
extract_and_write_objects(
pdf_path=source_pdf_path,
figure_inventory=figure_inventory,
table_inventory=table_inventory,
asset_root=ocr_asset_root,
render_root=ocr_render_root,
page_dimensions_by_page=page_dimensions_by_page,
)
# --- Phase 3: structured renderer ---
from paperforge.worker.ocr_render import render_fulltext_markdown, write_render_outputs
markdown = render_fulltext_markdown(
structured_blocks=structured,
resolved_metadata=resolved,
figure_inventory=figure_inventory,
table_inventory=table_inventory,
page_count=page_num,
)
write_render_outputs(
render_root=ocr_root / "render",
compat_fulltext=ocr_root / "fulltext.md",
markdown=markdown,
)
# --- Phase 3: OCR health report ---
from paperforge.worker.ocr_health import build_ocr_health, write_ocr_health
health_report = build_ocr_health(
page_count=page_num,
raw_blocks_count=len(all_raw_blocks),
structured_blocks=structured,
figure_inventory=figure_inventory,
table_inventory=table_inventory,
)
write_ocr_health(ocr_root / "health", health_report)
meta["ocr_health_overall"] = health_report["overall"]
# --- Phase 5: role-based OCR index ---
from paperforge.worker.ocr_index import build_role_indexes, write_role_index
role_indexes = build_role_indexes(
structured_blocks=structured,
resolved_metadata=resolved,
)
write_role_index(ocr_root / "index", role_indexes)
# Update meta.json with version payloads
ocr_model = meta.get("ocr_model", meta.get("ocr_provider", "PaddleOCR"))
version_payload = build_version_payload(
pdf_fingerprint=pdf_fingerprint,
result_json_hash=result_hash,
ocr_model=ocr_model,
)
meta["raw_version"] = version_payload["raw_version"]
meta["derived_version"] = version_payload["derived_version"]
from paperforge.worker.ocr_versions import classify_version_state, expected_derived_payload, expected_raw_payload
state = classify_version_state(
meta=meta,
expected_raw=expected_raw_payload(ocr_model=meta.get("raw_version", {}).get("ocr_model", "unknown")),
expected_derived=expected_derived_payload(),
)
meta["raw_upgradable"] = state["raw_upgradable"]
meta["derived_stale"] = state["derived_stale"]
meta["version_state_updated_at"] = __import__("datetime").datetime.now().isoformat()
try:
meta["legacy_images_path"] = str(images_dir.relative_to(vault)).replace("\\", "/")
except ValueError:
meta["legacy_images_path"] = str(images_dir)
meta["path_map"] = {"structured_truth": "assets/", "legacy_compat": "images/"}
write_json(meta_path, meta)
fulltext_path = ocr_root / "fulltext.md"
markdown_dir = ocr_root / "markdown"
if markdown_dir.exists():
shutil.rmtree(markdown_dir)
markdown_path = str(fulltext_path.relative_to(paths["vault"])).replace("\\", "/") if page_num else ""
json_path = str((json_dir / "result.json").relative_to(paths["vault"])).replace("\\", "/")
fulltext_md_path = str(fulltext_path.resolve())
return (page_num, markdown_path, json_path, fulltext_md_path)
def _workspace_fulltext_candidates(lit_root: Path, zotero_key: str) -> list[Path]:
if not lit_root.exists():
return []
matches: list[Path] = []
for candidate in lit_root.rglob(f"{zotero_key} - *"):
if candidate.is_dir():
fulltext_path = candidate / "fulltext.md"
if fulltext_path.exists():
matches.append(fulltext_path)
return matches
def _rewrite_note_fields(text: str, **fields: object) -> str:
updated = text
for field, value in fields.items():
updated = _sync.update_frontmatter_field(updated, field, value)
return updated
def ocr_redo_papers(vault: Path, dry_run: bool = False, verbose: bool = False, no_progress: bool = False) -> int:
"""Scan for papers with ocr_redo: true, reset and immediately rerun OCR.
Spec 2.6: paperforge ocr redo [--dry-run]
"""
from paperforge.adapters.obsidian_frontmatter import (
extract_preserved_ocr_redo,
)
paths = pipeline_paths(vault)
ocr_root = paths.get("ocr")
lit_root = paths.get("literature")
if not lit_root or not lit_root.exists():
if verbose:
logger.info("No literature directory found, nothing to redo")
return 0
redo_entry_map: dict[str, tuple[Path, str]] = {}
for note_file in sorted(lit_root.rglob("*.md")):
if note_file.name in ("fulltext.md", "deep-reading.md", "discussion.md"):
continue
try:
text = note_file.read_text(encoding="utf-8")
except Exception:
continue
if not extract_preserved_ocr_redo(text):
continue
key_match = re.search(r"^zotero_key:\s*(.+)$", text, re.MULTILINE)
if not key_match:
continue
zotero_key = key_match.group(1).strip().strip('"').strip("'")
if not zotero_key or not re.match(r"^[A-Za-z0-9]{8}$", zotero_key):
if verbose:
logger.warning("Skipping invalid zotero_key: %s", zotero_key)
continue
existing = redo_entry_map.get(zotero_key)
if existing is None or note_file.name == f"{zotero_key}.md":
redo_entry_map[zotero_key] = (note_file, text)
redo_entries = [(key, note_file, text) for key, (note_file, text) in sorted(redo_entry_map.items())]
if not redo_entries:
print("No papers with ocr_redo: true found", flush=True)
return 0
print(f"Found {len(redo_entries)} paper(s) with ocr_redo: true:", flush=True)
for zotero_key, note_file, _text in redo_entries:
print(f" {zotero_key} ({note_file.parent.name})", flush=True)
if dry_run:
print("\nDry-run mode: no changes made. Run without --dry-run to rerun OCR.", flush=True)
return 0
reset_keys = []
for zotero_key, note_file, text in redo_entries:
for workspace_fulltext in _workspace_fulltext_candidates(lit_root, zotero_key):
workspace_fulltext.unlink(missing_ok=True)
if verbose:
logger.info("Deleted workspace fulltext for %s", zotero_key)
ocr_dir = ocr_root / zotero_key if ocr_root else None
if ocr_dir and ocr_dir.exists():
shutil.rmtree(ocr_dir)
if verbose:
logger.info("Deleted OCR directory for %s", zotero_key)
text = _rewrite_note_fields(
text,
do_ocr=True,
ocr_status="pending",
fulltext_md_path="",
ocr_redo=True,
)
note_file.write_text(text, encoding="utf-8")
reset_keys.append(zotero_key)
print(f"Reset {len(reset_keys)} paper(s): {', '.join(reset_keys)}", flush=True)
print("Launching OCR immediately for reset papers...", flush=True)
exit_code = run_ocr(vault, verbose=verbose, no_progress=no_progress, selected_keys=set(reset_keys))
success_keys: list[str] = []
failed_keys: list[str] = []
for zotero_key, note_file, _original_text in redo_entries:
current_text = note_file.read_text(encoding="utf-8")
meta_path = ocr_root / zotero_key / "meta.json" if ocr_root else None
meta = read_json(meta_path) if meta_path and meta_path.exists() else {}
status, _error = validate_ocr_meta(paths, meta) if meta else ("pending", "")
current_text = _rewrite_note_fields(current_text, ocr_status=status)
if status == "done":
current_text = _rewrite_note_fields(current_text, ocr_redo=False)
success_keys.append(zotero_key)
else:
current_text = _rewrite_note_fields(current_text, ocr_redo=True)
failed_keys.append(zotero_key)
note_file.write_text(current_text, encoding="utf-8")
refresh_index_entry(vault, zotero_key)
if success_keys:
print(f"Redo OCR done={len(success_keys)}: {', '.join(success_keys)}", flush=True)
if failed_keys:
print(f"Redo OCR pending/failed={len(failed_keys)}: {', '.join(failed_keys)}", flush=True)
return exit_code
def run_ocr(vault: Path, verbose: bool = False, no_progress: bool = False, selected_keys: set[str] | None = None) -> int:
from paperforge.pdf_resolver import resolve_pdf_path
paths = pipeline_paths(vault)
cleanup_blocked_ocr_dirs(paths)
# Zombie reset: reset stale processing jobs to pending
zombie_timeout = int(os.environ.get("PAPERFORGE_ZOMBIE_TIMEOUT_MINUTES", "30"))
ocr_root = paths.get("ocr")
if ocr_root and ocr_root.exists():
for meta_dir in ocr_root.iterdir():
meta_path = meta_dir / "meta.json"
if not meta_path.exists():
continue
try:
meta = read_json(meta_path)
except Exception:
continue
zkey = meta.get("zotero_key", meta_dir.name)
zstatus = str(meta.get("ocr_status", "") or "").strip().lower()
if zstatus in {"queued", "running"}:
zstarted = meta.get("ocr_started_at", "")
if zstarted:
try:
started_dt = datetime.fromisoformat(zstarted)
if (datetime.now(timezone.utc) - started_dt).total_seconds() > zombie_timeout * 60:
meta["ocr_status"] = "pending"
meta["retry_count"] = int(meta.get("retry_count", 0)) + 1
write_json(meta_path, meta)
logger.warning(
"Zombie reset %s: was %s, started at %s, reset to pending", zkey, zstatus, zstarted
)
except Exception:
pass
control_actions = load_control_actions(paths)
target_keys = {key for key, action in control_actions.items() if action.get("do_ocr", False)}
if selected_keys is not None:
target_keys &= set(selected_keys)
target_rows = []
for export_path in sorted(paths["exports"].glob("*.json")):
for item in load_export_rows(export_path):
if item["key"] not in target_keys:
continue
pdf_attachments = [a for a in item.get("attachments", []) if a.get("contentType") == "application/pdf"]
target_rows.append(
{
"zotero_key": item["key"],
"has_pdf": bool(pdf_attachments),
"pdf_path": pdf_attachments[0]["path"] if pdf_attachments else "",
}
)
for row in target_rows:
key = row["zotero_key"]
meta = ensure_ocr_meta(vault, row)
current = str(meta.get("ocr_status", "") or "").strip().lower()
if current == "error":
meta["ocr_status"] = "pending"
meta["ocr_job_id"] = ""
meta["ocr_started_at"] = ""
meta["ocr_finished_at"] = ""
meta["retry_count"] = 0
write_json(paths["ocr"] / key / "meta.json", meta)
elif current == "nopdf":
meta["ocr_status"] = "pending"
meta["error"] = ""
meta["retry_count"] = 0
write_json(paths["ocr"] / key / "meta.json", meta)
status, _error = validate_ocr_meta(paths, meta)
if status == "done_incomplete":
meta["ocr_status"] = "pending"
meta["ocr_job_id"] = ""
meta["ocr_started_at"] = ""
meta["ocr_finished_at"] = ""
meta["error"] = _error
meta["retry_count"] = 0
write_json(paths["ocr"] / key / "meta.json", meta)
ocr_queue = sync_ocr_queue(paths, target_rows)
print(f"DEBUG: target_keys count={len(target_keys)}, target_rows count={len(target_rows)}, queue count={len(ocr_queue)}", flush=True)
if ocr_queue:
print(f"DEBUG: first item={ocr_queue[0].get('zotero_key')} status={ocr_queue[0].get('queue_status')}", flush=True)
max_items_raw = os.environ.get("PADDLEOCR_MAX_ITEMS", "").strip()
max_items = 3
if max_items_raw:
try:
max_items = max(1, int(max_items_raw))
except ValueError:
max_items = 3
token = os.environ.get("PADDLEOCR_API_TOKEN", "").strip()
if not token:
token = os.environ.get("PADDLEOCR_API_TOKEN_USER", "").strip()
if not token:
try:
import winreg
with winreg.OpenKey(winreg.HKEY_CURRENT_USER, "Environment") as env_key:
token = str(winreg.QueryValueEx(env_key, "PADDLEOCR_API_TOKEN")[0]).strip()
except Exception:
token = ""
if not token:
# Fallback: parse vault-root .env
token = _read_dotenv(vault, "PADDLEOCR_API_TOKEN")
job_url = os.environ.get("PADDLEOCR_JOB_URL", "https://paddleocr.aistudio-app.com/api/v2/ocr/jobs").strip()
model = os.environ.get("PADDLEOCR_MODEL", "PaddleOCR-VL-1.6").strip()
optional_payload = {"useDocOrientationClassify": False, "useDocUnwarping": False, "useChartRecognition": False}
changed = 0
active_submitted = 0
queue_changed = False
_submitted: set[str] = set() # keys newly uploaded in this run
def _do_poll(job_id: str, token_val: str) -> requests.Response:
resp = requests.get(f"{job_url}/{job_id}", headers={"Authorization": f"bearer {token_val}"}, timeout=60)
resp.raise_for_status()
return resp
for queue_row in progress_bar(ocr_queue, desc="Processing OCR", disable=no_progress):
key = queue_row["zotero_key"]
meta = ensure_ocr_meta(vault, queue_row)
status = str(meta.get("ocr_status", "pending") or "pending").strip().lower()
queue_row["queue_status"] = status
if status == "done":
queue_changed = True
continue
if status in {"queued", "running"} and meta.get("ocr_job_id"):
active_submitted += 1
if not token:
continue
meta_path_poll = paths["ocr"] / key / "meta.json"
try:
response = retry_with_meta(_do_poll, meta_path_poll, meta["ocr_job_id"], token)
except Exception as e:
meta["ocr_status"] = "pending"
meta["error"] = str(e)
meta["last_error"] = str(e)
meta["retry_count"] = int(meta.get("retry_count", 0)) + 1
queue_row["queue_status"] = "pending"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
active_submitted = max(0, active_submitted - 1)
continue
try:
payload = response.json()["data"]
state = payload["state"]
except (json.JSONDecodeError, KeyError) as e:
meta["ocr_status"] = "pending"
meta["error"] = f"API schema mismatch during polling: {e}"
meta["last_error"] = meta["error"]
meta["raw_response"] = response.text[:1000]
meta["retry_count"] = int(meta.get("retry_count", 0)) + 1
queue_row["queue_status"] = "pending"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
active_submitted = max(0, active_submitted - 1)
continue
if state in {"pending", "running"}:
meta["ocr_status"] = state
queue_row["queue_status"] = state
meta["error"] = ""
elif state == "done":
try:
result_url = payload["resultUrl"]["jsonUrl"]
result_response = requests.get(result_url, timeout=120)
result_response.raise_for_status()
lines = [line.strip() for line in result_response.text.splitlines() if line.strip()]
all_results = []
for line in lines:
page_payload = json.loads(line)["result"]
all_results.append(page_payload)
page_num, markdown_path, json_path, fulltext_md_path = postprocess_ocr_result(
vault, key, all_results
)
except Exception as e:
meta["ocr_status"] = "pending"
meta["error"] = str(e)
meta["retry_count"] = int(meta.get("retry_count", 0)) + 1
queue_row["queue_status"] = "pending"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
active_submitted = max(0, active_submitted - 1)
continue
meta["ocr_status"] = "done"
# Per D-01: auto_analyze_after_ocr opt-in workflow streamlining
cfg_path = vault / "paperforge.json"
if cfg_path.exists():
try:
pf_cfg = read_json(cfg_path)
if pf_cfg.get("auto_analyze_after_ocr", False):
note_glob = list(paths["literature"].rglob(f"{key}.md"))
if not note_glob:
note_glob = list(paths["literature"].rglob(f"{key} - *.md"))
if note_glob:
note_path = max(note_glob, key=lambda p: len(p.parents))
text = note_path.read_text(encoding="utf-8")
text = re.sub(
r"^analyze:.*$",
"analyze: true",
text,
count=1,
flags=re.MULTILINE,
)
note_path.write_text(text, encoding="utf-8")
except Exception:
logger.warning("auto_analyze_after_ocr: failed for %s", key, exc_info=True)
meta["ocr_finished_at"] = datetime.now(timezone.utc).isoformat()
meta["page_count"] = page_num
meta["markdown_path"] = markdown_path
meta["json_path"] = json_path
meta["fulltext_md_path"] = fulltext_md_path
meta["error"] = ""
queue_row["queue_status"] = "done"
queue_changed = True
active_submitted = max(0, active_submitted - 1)
else:
meta["ocr_status"] = "error"
meta["error"] = payload.get("errorMsg", "Unknown OCR failure")
meta["library_record"] = key
queue_row["queue_status"] = "error"
active_submitted = max(0, active_submitted - 1)
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
# Upload pending items in batches until none remain (processes all do_ocr items, not just max_items)
def _do_upload(token_val: str, pdf_path: Path) -> requests.Response:
with open(pdf_path, "rb") as file_handle:
resp = requests.post(
job_url,
headers={"Authorization": f"bearer {token_val}"},
data={"model": model, "optionalPayload": json.dumps(optional_payload)},
files={"file": file_handle},
timeout=120,
)
resp.raise_for_status()
return resp
# Combined upload + poll loop: process all items in batches up to max_items concurrency
import time as _time
poll_interval = int(os.environ.get("PAPERFORGE_POLL_INTERVAL", "15"))
max_poll_cycles = int(os.environ.get("PAPERFORGE_POLL_MAX_CYCLES", "60"))
_completed_count = 0
_failed_count = 0
_token_warned = False
for _cycle in range(max_poll_cycles):
remaining = [r for r in ocr_queue if r.get("queue_status", "") not in ("done", "nopdf", "blocked")]
if not remaining:
break
available_slots = max(0, max_items - active_submitted)
if available_slots > 0:
upload_items = [r for r in remaining if r.get("queue_status", "") not in ("queued", "running")][
:available_slots
]
for queue_row in upload_items:
key = queue_row["zotero_key"]
meta = ensure_ocr_meta(vault, queue_row)
_sanitized_temp = None
status = str(meta.get("ocr_status", "pending") or "pending").strip().lower()
if status in {"done", "queued", "running"}:
continue
resolved_pdf = resolve_pdf_path(
queue_row.get("pdf_path", ""),
queue_row.get("has_pdf", False),
vault,
paths.get("zotero_dir") if "zotero_dir" in paths else None,
)
if not resolved_pdf:
meta["ocr_status"] = "nopdf"
queue_row["queue_status"] = "nopdf"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
print(f"OCR: {key} skipped (PDF not found)", flush=True)
continue
if not token:
meta["ocr_status"] = "blocked"
queue_row["queue_status"] = "blocked"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
if not _token_warned:
print("OCR: no API token configured — set PADDLEOCR_API_TOKEN in .env", flush=True)
_token_warned = True
continue
upload_pdf = resolved_pdf
if meta.get("needs_sanitize"):
try:
import tempfile
doc = fitz.open(str(resolved_pdf))
_sanitized_temp = Path(tempfile.mktemp(suffix=".pdf"))
doc.save(str(_sanitized_temp), garbage=4, deflate=True, clean=True)
doc.close()
upload_pdf = _sanitized_temp
meta["needs_sanitize"] = False
except Exception:
pass
if int(meta.get("retry_count", 0)) >= 3:
meta["ocr_status"] = "error"
meta["error"] = meta.get("error", "") or "Upload failed after 3 retries"
queue_row["queue_status"] = "error"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
continue
print(f"OCR: {key} uploading to PaddleOCR...", flush=True)
try:
response = retry_with_meta(_do_upload, paths["ocr"] / key / "meta.json", token, upload_pdf)
meta["ocr_job_id"] = response.json()["data"]["jobId"]
except Exception as e:
if _sanitized_temp is not None and _sanitized_temp.exists():
_sanitized_temp.unlink(missing_ok=True)
import requests as _requests
if isinstance(e, _requests.exceptions.HTTPError):
_status = getattr(getattr(e, "response", None), "status_code", 0)
if _status == 401:
meta["ocr_status"] = "blocked"
meta["error"] = "PaddleOCR token invalid"
meta["retry_count"] = 3
queue_row["queue_status"] = "blocked"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
print(f"OCR: {key} blocked (invalid API token)", flush=True)
continue
if isinstance(e, FileNotFoundError):
meta["ocr_status"] = "nopdf"
meta["error"] = "PDF not found"
queue_row["queue_status"] = "nopdf"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
continue
retry_count = int(meta.get("retry_count", 0)) + 1
meta["retry_count"] = retry_count
meta["error"] = str(e)
meta["last_error"] = str(e)
meta["ocr_status"] = "pending"
queue_row["queue_status"] = "pending"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
print(f"OCR: {key} upload failed (retry {retry_count}/3): {e}", flush=True)
continue
meta["ocr_status"] = "queued"
meta["ocr_started_at"] = datetime.now(timezone.utc).isoformat()
meta["error"] = ""
queue_row["queue_status"] = "queued"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
active_submitted += 1
print(f"OCR: {key} queued (job {meta['ocr_job_id']})", flush=True)
if _sanitized_temp is not None and _sanitized_temp.exists():
_sanitized_temp.unlink(missing_ok=True)
poll_items = [r for r in remaining if r.get("queue_status") in ("queued", "running")]
for queue_row in poll_items:
key = queue_row["zotero_key"]
meta = ensure_ocr_meta(vault, queue_row)
job_id = meta.get("ocr_job_id", "")
if not job_id:
continue
try:
response = retry_with_meta(_do_poll, paths["ocr"] / key / "meta.json", job_id, token)
payload = response.json()["data"]
state = payload["state"]
except Exception:
continue
if state == "done":
try:
result_url = payload["resultUrl"]["jsonUrl"]
result_response = requests.get(result_url, timeout=120)
if result_response.status_code == 404:
meta["ocr_status"] = "pending"
meta["ocr_job_id"] = ""
meta["needs_sanitize"] = True
meta["error"] = "Result object not found on provider (404)"
meta["retry_count"] = 0
queue_row["queue_status"] = "pending"
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
active_submitted = max(0, active_submitted - 1)
print(f"OCR: {key} result expired, will retry", flush=True)
continue
result_response.raise_for_status()
lines = [l.strip() for l in result_response.text.splitlines() if l.strip()]
results = [json.loads(l)["result"] for l in lines]
page_num, md_path, json_path, fulltext_md_path = postprocess_ocr_result(vault, key, results)
except Exception:
continue
meta["ocr_status"] = "done"
meta["ocr_finished_at"] = datetime.now(timezone.utc).isoformat()
meta["page_count"] = page_num
meta["markdown_path"] = md_path
meta["json_path"] = json_path
meta["fulltext_md_path"] = fulltext_md_path
meta["error"] = ""
queue_row["queue_status"] = "done"
queue_changed = True
active_submitted = max(0, active_submitted - 1)
_completed_count += 1
print(f"OCR: {key} completed ({page_num} pages)", flush=True)
elif state in ("error", "failed"):
meta["error"] = payload.get("errorMsg", "Unknown OCR failure")
meta["ocr_status"] = "queued"
queue_row["queue_status"] = "queued"
_failed_count += 1
print(f"OCR: {key} failed: {meta['error']}", flush=True)
else:
meta["ocr_status"] = state
queue_row["queue_status"] = state
write_json(paths["ocr"] / key / "meta.json", meta)
changed += 1
if any(r.get("queue_status") in ("queued", "running") for r in ocr_queue):
_time.sleep(poll_interval)
# Collect completed OCR keys for incremental index refresh (before filtering)
_done_ocr_keys = (
[r.get("zotero_key", "") for r in ocr_queue if r.get("queue_status") == "done"] if queue_changed else []
)
if queue_changed:
ocr_queue = [row for row in ocr_queue if str(row.get("queue_status", "")).lower() != "done"]
write_ocr_queue(paths, ocr_queue)
# Determine exit code: 0 = all settled, 1 = some items still pending
final_remaining = [r for r in ocr_queue if r.get("queue_status", "") not in ("done", "nopdf", "blocked", "error")]
pending_keys = [r["zotero_key"] for r in final_remaining]
final_statuses = {r["zotero_key"]: r.get("queue_status", "?") for r in ocr_queue}
done_count = sum(1 for s in final_statuses.values() if s == "done")
failed_count = sum(1 for s in final_statuses.values() if s in ("error", "blocked"))
pending_count = len(final_remaining)
summary_parts = []
if done_count:
summary_parts.append(f"done={done_count}")
if failed_count:
summary_parts.append(f"failed={failed_count}")
if pending_count:
summary_parts.append(f"pending={pending_count} ({', '.join(pending_keys)})")
print(f"OCR: {' '.join(summary_parts) if summary_parts else 'no items processed'}", flush=True)
if pending_keys:
print("OCR: re-run to continue polling incomplete items", flush=True)
try:
_sync.run_selection_sync(vault)
if _done_ocr_keys:
done_keys = [k for k in _done_ocr_keys if k]
for ocr_key in done_keys:
refresh_index_entry(vault, ocr_key)
if verbose:
print(f"ocr: refreshed {len(done_keys)} index entries incrementally")
else:
_sync.run_index_refresh(vault)
except ImportError:
_sync.run_index_refresh(vault)
except Exception as e:
logger.error("Post-OCR index refresh failed: %s", e)
print(f"ocr: updated {changed} records")
return 1 if pending_keys else 0