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 "" def _resolve_paddleocr_token(vault: Path) -> str: """Resolve PaddleOCR API token from canonical sources. Priority: PADDLEOCR_API_TOKEN env → PADDLEOCR_API_TOKEN_USER env → Windows HKCU Environment → vault/.env or System/PaperForge/.env. Returns empty string when no token is found. """ 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: token = _read_dotenv(vault, "PADDLEOCR_API_TOKEN") return token 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_errors import ( OCRAPISchemaError, OCRArtifactIntegrityError, OCRNetworkError, OCRPDFResolveError, OCRPostprocessError, classify_ocr_error, ) 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.ocr_object_writeback import apply_object_writebacks from paperforge.worker.sync import ( load_control_actions, load_export_rows, ) logger = logging.getLogger(__name__) OCR_QUEUE_STATUSES = { "pending", "queued", "running", "done", "done_degraded", "retryable_error", "fatal_error", "nopdf", "blocked", } OCR_SETTLED_STATUSES = {"done", "done_degraded", "fatal_error", "blocked"} def _ocr_pipeline_v3_enabled() -> bool: """Check if the V3 pipeline is active. Now enabled by default — v3 normalizes role candidates before figure/table matching, then commits final roles after matching. Set OCR_PIPELINE_V3=0 (or any falsy value) to revert to legacy normalize-then-match order. """ val = os.environ.get("OCR_PIPELINE_V3", "").strip() if val: return val.lower() not in {"0", "false", "no", "off"} return True _LEGACY_HARD_DEGRADE_PREFIXES = ( "rendered text gaps", ) def _extract_hard_degraded_reasons(health_report: dict) -> list[str]: hard = health_report.get("hard_degraded_reasons") if hard is not None: return [str(r) for r in hard if str(r or "").strip()] legacy = health_report.get("degraded_reasons", []) or [] return [ str(r) for r in legacy if str(r or "").strip() and str(r).startswith(_LEGACY_HARD_DEGRADE_PREFIXES) ] def apply_ocr_error_state(row: dict, meta: dict, error_state: dict) -> None: status = error_state["status"] row["queue_status"] = status meta["ocr_status"] = status meta["error"] = error_state["last_error"] meta["error_type"] = error_state["error_type"] meta["error_stage"] = error_state["error_stage"] meta["retryable"] = error_state["retryable"] meta["last_error"] = error_state["last_error"] 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 not in ("done", "done_degraded"): 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 if title.lower() == "abstract" or re.match(r"^\d+\s", title): rendered.append(f"## {title}") else: 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) source_pdf_path = Path(meta.get("source_pdf", "")) if meta.get("source_pdf") else None if source_pdf_path and source_pdf_path.exists(): from paperforge.worker.ocr_pdf_spans import ( backfill_missing_text_from_pdf, backfill_span_metadata_from_pdf, ) # Attach PDF-derived font/geometry evidence first backfill_span_metadata_from_pdf(all_raw_blocks, source_pdf_path) # Recover empty OCR text using the same source PDF backfill_missing_text_from_pdf(all_raw_blocks, source_pdf_path) write_raw_blocks_jsonl(artifacts.blocks_raw, all_raw_blocks) normalize_mode = "seed_only" if _ocr_pipeline_v3_enabled() else "legacy" structured, doc_structure = build_structured_blocks( all_raw_blocks, source_metadata=source_meta, structure_output_dir=artifacts.blocks_structured.parent, normalize_mode=normalize_mode, ) if _ocr_pipeline_v3_enabled(): from paperforge.worker.ocr_pre_match_normalize import pre_match_normalize 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, page_blocks=all_raw_blocks, structured_blocks=structured, ) write_resolved_metadata(metadata_dir / "resolved_metadata.json", resolved) figure_inventory = build_figure_inventory(structured) # --- Author bio: residual figure inventory cleanup (Pass B, P1) --- from paperforge.worker.ocr_bio import ( residual_author_bio_pass, post_ref_bio_cleanup, prune_figure_inventory_after_bio, _resolve_ref_start_page, ) residual_author_bio_pass( figure_inventory, structured, include_ambiguous=False, include_weak_matched=False, ) # --- Author bio: post-ref bio cleanup (Pass C, P0) --- ref_start_page = _resolve_ref_start_page(structured) if ref_start_page is not None: post_ref_bio_cleanup(figure_inventory, structured, ref_start_page=ref_start_page) prune_figure_inventory_after_bio(figure_inventory) # --- Phase 2a: reader figure synthesis --- from paperforge.worker.ocr_figure_reader import synthesize_reader_figures reader_payload = synthesize_reader_figures(figure_inventory, structured_blocks=structured) reader_figures_dir = ocr_root / "structure" reader_figures_dir.mkdir(parents=True, exist_ok=True) write_json(reader_figures_dir / "reader_figures.json", reader_payload) # --- Phase 2: table inventory --- table_inventory = build_table_inventory(structured) # --- Phase 2: object writeback --- if _ocr_pipeline_v3_enabled(): from paperforge.worker.ocr_post_match_normalize import post_match_normalize structured, doc_structure = post_match_normalize( structured, figure_inventory, table_inventory, document_structure=doc_structure, source_frontmatter_anchors=getattr(doc_structure, "source_frontmatter_anchors", None), ) apply_object_writebacks( structured_blocks=structured, figure_inventory=figure_inventory, table_inventory=table_inventory, ) write_figure_inventory( artifacts.blocks_structured.parent / "figure_inventory.json", figure_inventory, ) write_table_inventory( artifacts.blocks_structured.parent / "table_inventory.json", table_inventory, ) # Re-persist structured blocks with writeback roles (table_html, figure_asset) # ponytail: writes entire list again; if throughput matters, write only changed blocks write_structured_blocks_jsonl(artifacts.blocks_structured, structured) # --- 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, structured_blocks=structured, ) # --- Phase 3: structured renderer --- from paperforge.worker.ocr_render import render_fulltext_markdown, write_render_outputs, RenderOutput rendered = render_fulltext_markdown( structured_blocks=structured, resolved_metadata=resolved, figure_inventory=figure_inventory, table_inventory=table_inventory, page_count=page_num, document_structure=doc_structure, reader_payload=reader_payload, return_events=True, ) # --- health report --- from paperforge.worker.ocr_health import ( build_ocr_health, build_ocr_raw_integrity_health, write_ocr_health, ) ocr_raw_integrity = build_ocr_raw_integrity_health(all_raw_blocks) 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, doc_structure=doc_structure, reader_payload=reader_payload, rendered_markdown=rendered.markdown, ) health_report["ocr_raw_integrity"] = ocr_raw_integrity write_ocr_health(ocr_root / "health", health_report) meta["ocr_health_overall"] = health_report["overall"] # Persist decision log from paperforge.worker.ocr_decisions import collect_decisions, write_decision_log write_decision_log(ocr_root / "health" / "decision_log.jsonl", collect_decisions(structured)) # --- Phase 4: render output --- meta = write_render_outputs( render_root=ocr_root / "render", user_fulltext=ocr_root / "fulltext.md", markdown=rendered.markdown, heading_events=rendered.heading_events, emitted_block_events=rendered.emitted_block_events, meta=meta, rebuild_increment=False, ) # --- 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) # --- Phase 6: structure tree --- from paperforge.retrieval.structure_tree import build_structure_tree, write_structure_tree structure_tree = build_structure_tree( heading_events=rendered.heading_events, emitted_block_events=rendered.emitted_block_events, structured_blocks=structured, ) write_structure_tree(ocr_root / "index", structure_tree) # 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 _find_notes_by_key(lit_root: Path | None) -> dict[str, Path]: """Scan literature directory and build {zotero_key: note_path} map.""" notes: dict[str, Path] = {} if not lit_root or not lit_root.exists(): return notes 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 key_match = re.search(r"^zotero_key:\s*(.+)$", text, re.MULTILINE) if not key_match: continue zkey = key_match.group(1).strip().strip('"').strip("'") if not zkey or not re.match(r"^[A-Za-z0-9]{8}$", zkey): continue existing = notes.get(zkey) if existing is None or note_file.name == f"{zkey}.md": notes[zkey] = note_file return notes def redo_papers_for_keys( vault: Path, keys: list[str], *, verbose: bool = False, progress_callback: Callable[[str], None] | None = None, stop_check: Callable[[], bool] | None = None, ) -> dict: """Redo specific papers: full artifact delete + OCR run + post-check. Processes papers sequentially so per-paper progress and cooperative stop are natural. Each paper goes through: 1. Delete derived artifacts and workspace fulltexts 2. Rewrite note as pending 3. Run OCR for this single paper 4. Check result, update note, refresh index Args: vault: Vault root. keys: Zotero citation keys to redo. verbose: Log additional details. progress_callback: Called with key after each paper's cycle completes. stop_check: Called before each paper; return True to stop. Returns: dict with success_keys, failed_keys, exit_code. """ from paperforge.worker._utils import pipeline_paths paths = pipeline_paths(vault) ocr_root = paths.get("ocr") lit_root = paths.get("literature") notes_by_key = _find_notes_by_key(lit_root) success_keys: list[str] = [] failed_keys: list[str] = [] _ocr_exit_code = 0 for key in keys: if stop_check is not None and stop_check(): break note_file = notes_by_key.get(key) if note_file is None: logger.warning("Redo: note not found for %s, counting as failed", key) failed_keys.append(key) _ocr_exit_code = 1 if progress_callback is not None: progress_callback(key) continue try: original_text = note_file.read_text(encoding="utf-8") except Exception: logger.warning("Redo: cannot read note for %s, counting as failed", key) failed_keys.append(key) if progress_callback is not None: progress_callback(key) _ocr_exit_code = 1 continue # ── Guard: nopdf papers must not be touched ── meta_path = ocr_root / key / "meta.json" if ocr_root else None _existing_meta = read_json(meta_path) if meta_path and meta_path.exists() else {} if str(_existing_meta.get("ocr_status", "") or "").strip().lower() == "nopdf": logger.warning("Redo: %s is nopdf — cannot redo (missing PDF), skipping", key) failed_keys.append(key) _ocr_exit_code = 1 if progress_callback is not None: progress_callback(key) continue # ── Phase 1: delete artifacts ── for wf in _workspace_fulltext_candidates(lit_root, key): wf.unlink(missing_ok=True) ocr_dir = ocr_root / key if ocr_root else None if ocr_dir and ocr_dir.exists(): shutil.rmtree(ocr_dir) # ── Phase 2: rewrite note as pending ── text = _rewrite_note_fields( original_text, do_ocr=True, ocr_status="pending", fulltext_md_path="", ocr_redo=True, ) note_file.write_text(text, encoding="utf-8") # ── Phase 3: run OCR for this paper ── _paper_exit = run_ocr(vault, verbose=verbose, no_progress=True, selected_keys={key}) if _paper_exit != 0: _ocr_exit_code = _paper_exit # ── Phase 4: post-OCR check ── meta_path = ocr_root / 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 = note_file.read_text(encoding="utf-8") 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(key) else: current_text = _rewrite_note_fields(current_text, ocr_redo=True) failed_keys.append(key) _ocr_exit_code = _ocr_exit_code or 1 note_file.write_text(current_text, encoding="utf-8") refresh_index_entry(vault, key) if progress_callback is not None: progress_callback(key) return { "success_keys": success_keys, "failed_keys": failed_keys, "exit_code": _ocr_exit_code, } 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 == "retryable_error": retry_count = int(meta.get("retry_count", 0)) + 1 meta["retry_count"] = retry_count if retry_count > 3: meta["ocr_status"] = "fatal_error" meta["error"] = meta.get("last_error", "Max retries exceeded") else: meta["ocr_status"] = "pending" meta["ocr_job_id"] = "" meta["error"] = meta.get("last_error", "") write_json(paths["ocr"] / key / "meta.json", meta) elif current == "fatal_error": pass 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 = _resolve_paddleocr_token(vault) 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: error_state = classify_ocr_error(OCRNetworkError(str(e)), stage="poll") apply_ocr_error_state(queue_row, meta, error_state) 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: error_state = classify_ocr_error( OCRAPISchemaError(f"API schema mismatch during polling: {e}"), stage="poll" ) apply_ocr_error_state(queue_row, meta, error_state) meta["raw_response"] = response.text[:1000] write_json(paths["ocr"] / key / "meta.json", meta) json_dir = paths["ocr"] / key / "json" json_dir.mkdir(parents=True, exist_ok=True) raw_error_dir = paths["ocr"] / key / "raw" raw_error_dir.mkdir(parents=True, exist_ok=True) write_json(raw_error_dir / "api_error_payload.json", {"raw_response": response.text[:5000]}) 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 (KeyboardInterrupt, SystemExit): raise except (OCRPDFResolveError, OCRArtifactIntegrityError, OCRPostprocessError) as exc: error_state = classify_ocr_error(exc, stage="postprocess") apply_ocr_error_state(queue_row, meta, error_state) write_json(paths["ocr"] / key / "meta.json", meta) changed += 1 active_submitted = max(0, active_submitted - 1) continue except Exception as exc: error_state = classify_ocr_error(OCRPostprocessError(str(exc)), stage="postprocess") apply_ocr_error_state(queue_row, meta, error_state) write_json(paths["ocr"] / key / "meta.json", meta) changed += 1 active_submitted = max(0, active_submitted - 1) continue health_path_poll = paths["ocr"] / key / "health" / "ocr_health.json" health_report_poll = read_json(health_path_poll) if health_path_poll.exists() else {} hard_reasons_poll = _extract_hard_degraded_reasons(health_report_poll) if hard_reasons_poll: meta["ocr_status"] = "done_degraded" meta["degraded_reason"] = "; ".join(hard_reasons_poll) else: meta["ocr_status"] = "done" meta.pop("degraded_reason", None) # 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 OCR_SETTLED_STATUSES] 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 (KeyboardInterrupt, SystemExit): raise except (OCRPDFResolveError, OCRArtifactIntegrityError, OCRPostprocessError) as exc: error_state = classify_ocr_error(exc, stage="postprocess") apply_ocr_error_state(queue_row, meta, error_state) write_json(paths["ocr"] / key / "meta.json", meta) changed += 1 active_submitted = max(0, active_submitted - 1) continue except Exception as exc: error_state = classify_ocr_error(OCRPostprocessError(str(exc)), stage="postprocess") apply_ocr_error_state(queue_row, meta, error_state) continue health_path_loop = paths["ocr"] / key / "health" / "ocr_health.json" health_report_loop = read_json(health_path_loop) if health_path_loop.exists() else {} hard_reasons_loop = _extract_hard_degraded_reasons(health_report_loop) if hard_reasons_loop: meta["ocr_status"] = "done_degraded" meta["degraded_reason"] = "; ".join(hard_reasons_loop) else: meta["ocr_status"] = "done" meta.pop("degraded_reason", None) 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", "done_degraded", "fatal_error", "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