diff --git a/paperforge/adapters/bbt.py b/paperforge/adapters/bbt.py index 328b1fac..7aac2d9c 100644 --- a/paperforge/adapters/bbt.py +++ b/paperforge/adapters/bbt.py @@ -3,10 +3,14 @@ from __future__ import annotations import logging from pathlib import Path +from paperforge.adapters.collections import build_collection_lookup +from paperforge.core.date_utils import extract_year +from paperforge.core.io import read_json + logger = logging.getLogger(__name__) -def _normalize_attachment_path(path: str, zotero_dir: Path | None = None) -> tuple[str, str, str]: +def normalize_attachment_path(path: str, zotero_dir: Path | None = None) -> tuple[str, str, str]: """Normalize a BBT attachment path to a consistent storage: format. Handles three real-world BBT export formats: @@ -42,7 +46,7 @@ def _normalize_attachment_path(path: str, zotero_dir: Path | None = None) -> tup # Format 1: Absolute Windows path pointing to Zotero storage candidate = Path(raw) - _looks_absolute = candidate.is_absolute() or (len(raw) >= 2 and raw[0].isalpha() and raw[1] == ':') + _looks_absolute = candidate.is_absolute() or (len(raw) >= 2 and raw[0].isalpha() and raw[1] == ":") if _looks_absolute: norm_path = raw.replace("\\", "/") # Detect Zotero storage pattern: .../storage/8CHARKEY/... @@ -63,7 +67,7 @@ def _normalize_attachment_path(path: str, zotero_dir: Path | None = None) -> tup return (f"storage:{norm}", bbt_path_raw, zotero_storage_key) -def _identify_main_pdf(attachments: list[dict]) -> tuple[dict | None, list[dict]]: +def identify_main_pdf(attachments: list[dict]) -> tuple[dict | None, list[dict]]: """Identify the main PDF and supplementary materials from attachments. Uses a hybrid three-priority strategy (Decision D-02): @@ -151,9 +155,6 @@ def resolve_item_collection_paths(item: dict, collection_lookup: dict) -> list[s def load_export_rows(path: Path) -> list[dict]: - from paperforge.worker._domain import build_collection_lookup - from paperforge.worker._utils import _extract_year, read_json - data = read_json(path) if isinstance(data, list): return data @@ -168,7 +169,7 @@ def load_export_rows(path: Path) -> list[dict]: if not isinstance(attachment, dict): continue raw_path = attachment.get("path", "") - normalized_path, bbt_path_raw, zotero_storage_key = _normalize_attachment_path(raw_path) + normalized_path, bbt_path_raw, zotero_storage_key = normalize_attachment_path(raw_path) # Preserve contentType from BBT if present; fallback to file extension content_type = attachment.get("contentType", "") if not content_type and str(normalized_path).lower().endswith(".pdf"): @@ -183,7 +184,7 @@ def load_export_rows(path: Path) -> list[dict]: "size": attachment.get("size", 0) or 0, } ) - main_pdf, supplementary_pdfs = _identify_main_pdf(attachments) + main_pdf, supplementary_pdfs = identify_main_pdf(attachments) pdf_path = main_pdf["path"] if main_pdf else "" bbt_path_raw = main_pdf["bbt_path_raw"] if main_pdf else "" zotero_storage_key = main_pdf["zotero_storage_key"] if main_pdf else "" @@ -199,7 +200,7 @@ def load_export_rows(path: Path) -> list[dict]: "abstract": item.get("abstractNote", ""), "journal": item.get("publicationTitle", ""), "extra": item.get("extra", ""), - "year": _extract_year(item.get("date", "")), + "year": extract_year(item.get("date", "")), "date": item.get("date", ""), "doi": item.get("DOI", ""), "pmid": item.get("PMID", ""), @@ -215,3 +216,9 @@ def load_export_rows(path: Path) -> list[dict]: ) return rows raise ValueError(f"Unsupported export format: {path}") + + +# ── Backward-compat aliases (v2.1 → 2.2 migration) ── +_normalize_attachment_path = normalize_attachment_path +_identify_main_pdf = identify_main_pdf +_extract_year = extract_year diff --git a/paperforge/adapters/collections.py b/paperforge/adapters/collections.py new file mode 100644 index 00000000..a924638b --- /dev/null +++ b/paperforge/adapters/collections.py @@ -0,0 +1,30 @@ +from __future__ import annotations + + +def build_collection_lookup(collections: dict) -> dict: + """Build parent-resolved path cache from a Zotero collection tree. + + Returns a dict with: + path_by_key: {collection_key: "Parent/Sub/Name", ...} + paths_by_item_id: {item_id: [collection_path, ...], ...} + """ + path_cache: dict[str, str] = {} + item_paths: dict[str, list[str]] = {} + + def path_for(key: str) -> str: + if key in path_cache: + return path_cache[key] + node = collections.get(key, {}) + parent = node.get("parent") or "" + name = node.get("name", "") + parent_path = path_for(parent) if parent else "" + full_path = f"{parent_path}/{name}" if parent_path else name + path_cache[key] = full_path + return full_path + + for key, node in collections.items(): + full_path = path_for(key) + for item_id in node.get("items", []): + item_paths.setdefault(item_id, []).append(full_path) + + return {"path_by_key": path_cache, "paths_by_item_id": item_paths} diff --git a/paperforge/cli.py b/paperforge/cli.py index 37b6086b..470050bb 100644 --- a/paperforge/cli.py +++ b/paperforge/cli.py @@ -152,7 +152,7 @@ def build_parser() -> argparse.ArgumentParser: # status status_p = sub.add_parser("status", help="Run the literature pipeline status check") - status_p.add_argument("--json", action="store_true", dest="json_output", help="Output JSON") + status_p.add_argument("--json", action="store_true", help="Output JSON") # sync (new unified command) p_sync = sub.add_parser("sync", help="Sync Zotero selection and refresh literature index") @@ -194,12 +194,14 @@ def build_parser() -> argparse.ArgumentParser: sub.add_parser("index-refresh", help="Refresh formal literature notes from library records") # deep-reading - sub.add_parser("deep-reading", help="Check deep-reading queue status") + p_dr = sub.add_parser("deep-reading", help="Check deep-reading queue status") + p_dr.add_argument("--json", action="store_true", help="Output as PFResult JSON") # repair p_repair = sub.add_parser("repair", help="Repair divergent literature notes") p_repair.add_argument("--fix", action="store_true", help="Actually apply repairs instead of dry-run") p_repair.add_argument("--fix-paths", action="store_true", help="Re-resolve PDF paths for items with path_error") + p_repair.add_argument("--json", action="store_true", help="Output result as JSON (PFResult envelope)") # ocr (unified) p_ocr = sub.add_parser("ocr", help="OCR operations") @@ -262,7 +264,7 @@ def build_parser() -> argparse.ArgumentParser: # doctor doctor_p = sub.add_parser("doctor", help="Validate PaperForge setup and configuration") - doctor_p.add_argument("--json", action="store_true", dest="json_output", help="Output JSON") + doctor_p.add_argument("--json", action="store_true", help="Output JSON") # update sub.add_parser("update", help="Update PaperForge to the latest version") @@ -478,7 +480,7 @@ def main(argv: list[str] | None = None) -> int: kw = {} if getattr(args, "verbose", False): kw["verbose"] = True - if getattr(args, "json_output", False): + if getattr(args, "json", False): kw["json_output"] = True return run_doctor(vault, **kw) @@ -511,11 +513,10 @@ def main(argv: list[str] | None = None) -> int: vault=vault, agent_key=args.agent, paddleocr_key=getattr(args, "paddleocr_key", None), - paddleocr_url=getattr(args, "paddleocr_url", - "https://paddleocr.aistudio-app.com/api/v2/ocr/jobs"), - system_dir=getattr(args, "system_dir", None) or "System", - resources_dir=getattr(args, "resources_dir", None) or "Resources", - base_dir=getattr(args, "base_dir", None) or "Bases", + paddleocr_url=getattr(args, "paddleocr_url", "https://paddleocr.aistudio-app.com/api/v2/ocr/jobs"), + system_dir=getattr(args, "system_dir", None) or "System", + resources_dir=getattr(args, "resources_dir", None) or "Resources", + base_dir=getattr(args, "base_dir", None) or "Bases", zotero_data=getattr(args, "zotero_data", None), skip_checks=getattr(args, "skip_checks", False), ) diff --git a/paperforge/commands/deep.py b/paperforge/commands/deep.py index ba1ca42c..5755ca88 100644 --- a/paperforge/commands/deep.py +++ b/paperforge/commands/deep.py @@ -4,6 +4,9 @@ import argparse import logging from pathlib import Path +from paperforge import __version__ +from paperforge.core.result import PFResult + logger = logging.getLogger(__name__) @@ -28,12 +31,47 @@ def _get_run_deep_reading(): def run(args: argparse.Namespace) -> int: - """Run deep-reading queue check.""" + """Run deep-reading queue check. Supports --json for PFResult output.""" vault = getattr(args, "vault_path", None) if vault is None: from paperforge.config import resolve_vault vault = resolve_vault(cli_vault=getattr(args, "vault", None)) + json_output = getattr(args, "json", False) run_deep_reading = _get_run_deep_reading() - return run_deep_reading(vault, verbose=getattr(args, "verbose", False)) + exit_code = run_deep_reading(vault, verbose=getattr(args, "verbose", False)) + + if json_output: + from paperforge.worker._utils import get_analyze_queue + + queue = get_analyze_queue(vault) + ready = [r for r in queue if r.get("ocr_status") == "done"] + blocked = [r for r in queue if r.get("ocr_status") != "done"] + + pf = PFResult( + ok=True, + command="deep-reading", + version=__version__, + data={ + "queue": [ + { + "zotero_key": r["zotero_key"], + "domain": r["domain"], + "title": r["title"], + "ocr_status": r.get("ocr_status", "pending"), + "ready": r.get("ocr_status") == "done", + } + for r in queue + ], + "summary": { + "total": len(queue), + "ready": len(ready), + "blocked": len(blocked), + }, + }, + ) + print(pf.to_json()) + return 0 + + return exit_code diff --git a/paperforge/commands/ocr.py b/paperforge/commands/ocr.py index f4a05917..0459268a 100644 --- a/paperforge/commands/ocr.py +++ b/paperforge/commands/ocr.py @@ -63,8 +63,16 @@ def _diagnose(vault: Path, live: bool = False, json_output: bool = False) -> int queue_data = _collect_ocr_queue_data(vault) pf_error = None if not passed: + if level == 1: + ec = ErrorCode.OCR_TOKEN_MISSING + elif level == 2: + ec = ErrorCode.OCR_UPLOAD_FAILED + elif level == 3: + ec = ErrorCode.OCR_RESULT_INVALID + else: + ec = ErrorCode.INTERNAL_ERROR pf_error = PFError( - code=ErrorCode.OCR_TOKEN_MISSING if level == 1 else ErrorCode.INTERNAL_ERROR, + code=ec, message=result.get("error", "OCR diagnosis failed"), details={"level": level, "fix": result.get("fix", "")}, ) @@ -124,6 +132,7 @@ def run(args: argparse.Namespace) -> int: Default behavior: run OCR queue. --diagnose: diagnose only (no upload). --key KEY: process specific item (passed through if supported). + Supports --json for PFResult output in both diagnose and normal modes. """ vault = getattr(args, "vault_path", None) if vault is None: @@ -147,6 +156,17 @@ def run(args: argparse.Namespace) -> int: run_ocr = _get_run_ocr() exit_code = run_ocr(vault, verbose=getattr(args, "verbose", False), no_progress=getattr(args, "no_progress", False)) + if json_output: + queue_data = _collect_ocr_queue_data(vault) + pf = PFResult( + ok=exit_code == 0, + command="ocr", + version=__version__, + data=queue_data, + ) + print(pf.to_json()) + return 0 if pf.ok else 1 + # Auto-diagnose after successful run (new unified behavior) if exit_code == 0 and ocr_action is None and not diagnose_only and not key: logger.info("Running post-OCR diagnostic...") diff --git a/paperforge/commands/repair.py b/paperforge/commands/repair.py index 0db6511f..ddd63c70 100644 --- a/paperforge/commands/repair.py +++ b/paperforge/commands/repair.py @@ -4,6 +4,10 @@ import argparse import logging from pathlib import Path +from paperforge import __version__ +from paperforge.core.errors import ErrorCode +from paperforge.core.result import PFError, PFResult + logger = logging.getLogger(__name__) @@ -28,7 +32,7 @@ def _get_run_repair(): def run(args: argparse.Namespace) -> int: - """Run repair command.""" + """Run repair command. Supports --json for PFResult output.""" vault = getattr(args, "vault_path", None) paths = getattr(args, "paths", None) if vault is None: @@ -43,6 +47,7 @@ def run(args: argparse.Namespace) -> int: paths = pipeline_paths(vault) run_repair = _get_run_repair() + json_output = getattr(args, "json", False) result = run_repair( vault, paths, @@ -50,11 +55,43 @@ def run(args: argparse.Namespace) -> int: fix=getattr(args, "fix", False), fix_paths=getattr(args, "fix_paths", False), ) - # Report path_error summary from repair scan + + # ── Build PFResult ── + divergent = result.get("divergent", 0) + path_errors = result.get("path_errors", {}) + path_error_total = path_errors.get("total", 0) + has_issues = bool(divergent) or path_error_total > 0 + + pf_error = None + if has_issues: + pf_error = PFError( + code=ErrorCode.VALIDATION_ERROR, + message=f"Repair found {divergent} divergences and {path_error_total} path errors", + details=result, + ) + + pf = PFResult( + ok=not has_issues, + command="repair", + version=__version__, + data={ + "scanned": result.get("scanned", 0), + "divergent": divergent, + "fixed": result.get("fixed", 0), + "errors": result.get("errors", []), + "rebuilt": result.get("rebuilt", 0), + "path_errors": result.get("path_errors", {}), + }, + error=pf_error, + ) + + if json_output: + print(pf.to_json()) + return 0 if pf.ok else 1 + + # Human-readable output path_errors = result.get("path_errors", {}) if path_errors.get("total", 0) > 0: error_summary = ", ".join(f"{count} {err}" for err, count in sorted(path_errors.get("by_type", {}).items())) print(f"repair: Found {path_errors['total']} items with path_error: {error_summary}") - # Return non-zero if any divergences or path_errors remain - has_issues = bool(result.get("divergent")) or path_errors.get("total", 0) > 0 return 1 if has_issues else 0 diff --git a/paperforge/commands/status.py b/paperforge/commands/status.py index c02b9438..31d5c2bc 100644 --- a/paperforge/commands/status.py +++ b/paperforge/commands/status.py @@ -1,13 +1,15 @@ """Status command. Reports vault statistics including library record counts, OCR progress, -and path_error counts (displayed by run_status in literature_pipeline). +and path_error counts. Uses PFResult contract for JSON output. """ import argparse import logging from pathlib import Path +from paperforge.core.result import PFResult + logger = logging.getLogger(__name__) @@ -32,7 +34,7 @@ def _get_run_status(): def run(args: argparse.Namespace) -> int: - """Run status check.""" + """Run status check. Worker handles PFResult JSON output when json_output=True.""" vault = getattr(args, "vault_path", None) if vault is None: from paperforge.config import resolve_vault @@ -40,4 +42,6 @@ def run(args: argparse.Namespace) -> int: vault = resolve_vault(cli_vault=getattr(args, "vault", None)) run_status = _get_run_status() - return run_status(vault, verbose=getattr(args, "verbose", False), json_output=getattr(args, "json_output", False)) + exit_code = run_status(vault, verbose=getattr(args, "verbose", False), json_output=getattr(args, "json", False)) + # Worker already prints PFResult when json_output=True; propagate exit code + return exit_code diff --git a/paperforge/commands/sync.py b/paperforge/commands/sync.py index c57e9cfe..74885fe1 100644 --- a/paperforge/commands/sync.py +++ b/paperforge/commands/sync.py @@ -1,63 +1,21 @@ -"""Sync command — unifies selection-sync and index-refresh.""" +"""Sync command — unified sync through SyncService.""" import argparse -import json as _json import logging -from pathlib import Path -from paperforge import __version__ from paperforge.config import migrate_paperforge_json -from paperforge.core.errors import ErrorCode -from paperforge.core.result import PFError, PFResult +from paperforge.core.result import PFResult +from paperforge import __version__ logger = logging.getLogger(__name__) -def _get_run_selection_sync(): - """Get run_selection_sync, preferring cli patches if available.""" - try: - from paperforge.cli import run_selection_sync - - if run_selection_sync is not None: - return run_selection_sync - except Exception: - pass - - import sys - - repo_root = Path(__file__).resolve().parent.parent.parent - if str(repo_root) not in sys.path: - sys.path.insert(0, str(repo_root)) - from paperforge.worker.sync import run_selection_sync - - return run_selection_sync - - -def _get_run_index_refresh(): - """Get run_index_refresh, preferring cli patches if available.""" - try: - from paperforge.cli import run_index_refresh - - if run_index_refresh is not None: - return run_index_refresh - except Exception: - pass - - import sys - - repo_root = Path(__file__).resolve().parent.parent.parent - if str(repo_root) not in sys.path: - sys.path.insert(0, str(repo_root)) - from paperforge.worker.sync import run_index_refresh - - return run_index_refresh - - def run(args: argparse.Namespace) -> int: - """Run sync command. + """Run sync command through SyncService. - By default runs both selection-sync and index-refresh. + By default runs both selection-sync and index-refresh + cleanup. Use --selection or --index to run only one phase. + SyncService is the canonical entry point for all sync operations. """ vault = getattr(args, "vault_path", None) if vault is None: @@ -66,7 +24,6 @@ def run(args: argparse.Namespace) -> int: vault = resolve_vault(cli_vault=getattr(args, "vault", None)) verbose = getattr(args, "verbose", False) - # Auto-migrate paperforge.json from legacy top-level keys to vault_config block migrated = migrate_paperforge_json(vault) if migrated: logger.info("Migrated paperforge.json to vault_config canonical format") @@ -76,8 +33,6 @@ def run(args: argparse.Namespace) -> int: dry_run = getattr(args, "dry_run", False) selection_only = getattr(args, "selection", False) index_only = getattr(args, "index", False) - rebuild_index = getattr(args, "rebuild_index", False) - domain = getattr(args, "domain", None) json_output = getattr(args, "json", False) if dry_run: @@ -99,44 +54,15 @@ def run(args: argparse.Namespace) -> int: print(" - selection-sync") if index_only: print(" - index-refresh") - if domain: - print(f" Filtered by domain: {domain}") return 0 - exit_code = 0 - sync_counts = {"new": 0, "updated": 0, "skipped": 0, "failed": 0, "errors": []} + from paperforge.services.sync_service import SyncService - if not index_only: - run_selection_sync = _get_run_selection_sync() - result = run_selection_sync(vault, verbose=getattr(args, "verbose", False), json_output=json_output) - if json_output: - sync_counts = result - elif result != 0: - exit_code = result - - if not selection_only: - run_index_refresh = _get_run_index_refresh() - code = run_index_refresh(vault, verbose=getattr(args, "verbose", False), rebuild_index=rebuild_index, json_output=json_output) - if code != 0 and exit_code == 0: - exit_code = code + svc = SyncService(vault) + result = svc.run(verbose=verbose, json_output=json_output) if json_output: - errors = sync_counts.get("errors", []) - pf_error = None - if errors or exit_code != 0: - pf_error = PFError( - code=ErrorCode.SYNC_FAILED, - message=f"Sync completed with {len(errors)} error(s)", - details={"errors": errors, "exit_code": exit_code}, - ) - result = PFResult( - ok=exit_code == 0 and not errors, - command="sync", - version=__version__, - data=sync_counts, - error=pf_error, - ) print(result.to_json()) return 0 - return exit_code + return 0 if result.ok else 1 diff --git a/paperforge/core/date_utils.py b/paperforge/core/date_utils.py new file mode 100644 index 00000000..508b3db7 --- /dev/null +++ b/paperforge/core/date_utils.py @@ -0,0 +1,17 @@ +from __future__ import annotations + +import re + + +def extract_year(value: str) -> str: + """Extract a 4-digit year (19xx or 20xx) from a date string. + + Returns empty string if no year found. + + Examples: + extract_year("2024-03-15") -> "2024" + extract_year("March 2023") -> "2023" + extract_year("n.d.") -> "" + """ + match = re.search(r"(19|20)\d{2}", value or "") + return match.group(0) if match else "" diff --git a/paperforge/core/errors.py b/paperforge/core/errors.py index a37e1641..79502eb8 100644 --- a/paperforge/core/errors.py +++ b/paperforge/core/errors.py @@ -4,16 +4,57 @@ from enum import Enum class ErrorCode(str, Enum): - """Centralized error codes for PFResult contracts.""" + """Centralized error codes for PFResult contracts. + Codes are grouped by subsystem. When adding a new code, prefer a + narrow, actionable code over a catch-all so that plugin UI can + generate targeted fix-it guidance. + """ + + # ── Runtime ── PYTHON_NOT_FOUND = "PYTHON_NOT_FOUND" + PYTHON_VERSION_TOO_OLD = "PYTHON_VERSION_TOO_OLD" + PAPERFORGE_NOT_INSTALLED = "PAPERFORGE_NOT_INSTALLED" VERSION_MISMATCH = "VERSION_MISMATCH" + + # ── Config / Vault ── + VAULT_NOT_FOUND = "VAULT_NOT_FOUND" + CONFIG_NOT_FOUND = "CONFIG_NOT_FOUND" + CONFIG_INVALID = "CONFIG_INVALID" + PATH_NOT_FOUND = "PATH_NOT_FOUND" + + # ── BBT / Zotero ── BBT_EXPORT_NOT_FOUND = "BBT_EXPORT_NOT_FOUND" + BBT_EXPORT_INVALID = "BBT_EXPORT_INVALID" + BBT_CITATION_KEY_MISSING = "BBT_CITATION_KEY_MISSING" + ZOTERO_DATA_NOT_FOUND = "ZOTERO_DATA_NOT_FOUND" + PDF_PATH_UNRESOLVED = "PDF_PATH_UNRESOLVED" + + # ── OCR ── OCR_TOKEN_MISSING = "OCR_TOKEN_MISSING" + OCR_UPLOAD_FAILED = "OCR_UPLOAD_FAILED" + OCR_POLL_TIMEOUT = "OCR_POLL_TIMEOUT" + OCR_RESULT_INVALID = "OCR_RESULT_INVALID" + + # ── Sync ── SYNC_FAILED = "SYNC_FAILED" + CANDIDATE_BUILD_FAILED = "CANDIDATE_BUILD_FAILED" + NOTE_WRITE_FAILED = "NOTE_WRITE_FAILED" + + # ── Schema ── + FIELD_MISSING = "FIELD_MISSING" + FIELD_TYPE_INVALID = "FIELD_TYPE_INVALID" + INDEX_SCHEMA_INVALID = "INDEX_SCHEMA_INVALID" + + # ── Generic ── VALIDATION_ERROR = "VALIDATION_ERROR" INTERNAL_ERROR = "INTERNAL_ERROR" UNKNOWN = "UNKNOWN" + @classmethod + def _missing_(cls, value): + """Gracefully handle unknown codes from newer plugin versions.""" + return cls.UNKNOWN + def __str__(self) -> str: return self.value diff --git a/paperforge/core/io.py b/paperforge/core/io.py new file mode 100644 index 00000000..80697ff7 --- /dev/null +++ b/paperforge/core/io.py @@ -0,0 +1,15 @@ +from __future__ import annotations + +import json +from pathlib import Path + + +def read_json(path: Path): + """Read and parse a JSON file.""" + return json.loads(path.read_text(encoding="utf-8")) + + +def write_json(path: Path, data) -> None: + """Write data as JSON, creating parent directories as needed.""" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8") diff --git a/paperforge/core/result.py b/paperforge/core/result.py index 5ab062ee..a121b32a 100644 --- a/paperforge/core/result.py +++ b/paperforge/core/result.py @@ -12,6 +12,7 @@ class PFError: code: ErrorCode message: str details: dict = field(default_factory=dict) + suggestions: list[str] = field(default_factory=list) @dataclass @@ -21,6 +22,8 @@ class PFResult: version: str data: Any = None error: PFError | None = None + warnings: list[str] = field(default_factory=list) + next_actions: list[dict] = field(default_factory=list) def __bool__(self) -> bool: return self.ok @@ -41,8 +44,14 @@ class PFResult: "message": self.error.message, "details": self.error.details, } + if self.error.suggestions: + raw["error"]["suggestions"] = self.error.suggestions else: raw["error"] = None + if self.warnings: + raw["warnings"] = self.warnings + if self.next_actions: + raw["next_actions"] = self.next_actions return raw def to_json(self) -> str: @@ -57,6 +66,7 @@ class PFResult: code=ErrorCode(err_data["code"]), message=err_data["message"], details=err_data.get("details", {}), + suggestions=err_data.get("suggestions", []), ) return cls( ok=data["ok"], @@ -64,4 +74,6 @@ class PFResult: version=data["version"], data=data.get("data"), error=error, + warnings=data.get("warnings", []), + next_actions=data.get("next_actions", []), ) diff --git a/paperforge/core/state.py b/paperforge/core/state.py index 63123e4a..bd890605 100644 --- a/paperforge/core/state.py +++ b/paperforge/core/state.py @@ -7,26 +7,45 @@ from typing import ClassVar class OcrStatus(str, Enum): - """OCR processing state for a paper.""" + """OCR processing state for a paper. + + States are granular so that plugin UI can display targeted next actions: + none – OCR not yet initiated + pending – do_ocr=true, not yet queued + queued – in queue, waiting for worker + processing – currently being OCR'd by PaddleOCR + done – full OCR completed successfully + done_incomplete – OCR completed but some pages/images missing + failed – API returned error or task failed + blocked – can't proceed (missing PDF, invalid key, etc.) + nopdf – no PDF file available to OCR + """ + + NONE = "none" PENDING = "pending" + QUEUED = "queued" PROCESSING = "processing" DONE = "done" + DONE_INCOMPLETE = "done_incomplete" FAILED = "failed" + BLOCKED = "blocked" + NO_PDF = "nopdf" @classmethod def from_legacy(cls, value: str) -> OcrStatus: - """Map legacy OCR state strings to canonical enum.""" + """Map legacy OCR state strings to canonical enum (1:1, no compression).""" mapping = { + "none": cls.NONE, "pending": cls.PENDING, - "processing": cls.PROCESSING, - "queued": cls.PENDING, + "queued": cls.QUEUED, "running": cls.PROCESSING, + "processing": cls.PROCESSING, "done": cls.DONE, + "done_incomplete": cls.DONE_INCOMPLETE, "failed": cls.FAILED, "error": cls.FAILED, - "blocked": cls.FAILED, - "nopdf": cls.FAILED, - "done_incomplete": cls.DONE, + "blocked": cls.BLOCKED, + "nopdf": cls.NO_PDF, } return mapping.get(value.strip().lower(), cls.PENDING) @@ -36,6 +55,7 @@ class OcrStatus(str, Enum): class PdfStatus(str, Enum): """PDF attachment health state.""" + HEALTHY = "healthy" BROKEN = "broken" MISSING = "missing" @@ -46,6 +66,7 @@ class PdfStatus(str, Enum): class Lifecycle(str, Enum): """Derived lifecycle stage summarizing the paper's overall progress.""" + PDF_READY = "pdf_ready" OCR_READY = "ocr_ready" ANALYZE_READY = "analyze_ready" @@ -63,10 +84,15 @@ class ALLOWED_TRANSITIONS: """ OCR_STATUS: ClassVar[dict[OcrStatus, list[OcrStatus]]] = { - OcrStatus.PENDING: [OcrStatus.PROCESSING, OcrStatus.FAILED], - OcrStatus.PROCESSING: [OcrStatus.DONE, OcrStatus.FAILED], + OcrStatus.NONE: [OcrStatus.PENDING, OcrStatus.QUEUED, OcrStatus.NO_PDF], + OcrStatus.PENDING: [OcrStatus.QUEUED, OcrStatus.PROCESSING, OcrStatus.NO_PDF, OcrStatus.BLOCKED], + OcrStatus.QUEUED: [OcrStatus.PROCESSING, OcrStatus.FAILED], + OcrStatus.PROCESSING: [OcrStatus.DONE, OcrStatus.DONE_INCOMPLETE, OcrStatus.FAILED], OcrStatus.DONE: [OcrStatus.PENDING], # re-run + OcrStatus.DONE_INCOMPLETE: [OcrStatus.PENDING], # re-run full or retry partial OcrStatus.FAILED: [OcrStatus.PENDING], # retry + OcrStatus.BLOCKED: [OcrStatus.PENDING], # fix preconditions then retry + OcrStatus.NO_PDF: [OcrStatus.PENDING], # add PDF then retry } DEEP_READING_STATUS: ClassVar[dict[str, list[str]]] = { diff --git a/paperforge/plugin/main.js b/paperforge/plugin/main.js index 3d25c0d9..abe5ee88 100644 --- a/paperforge/plugin/main.js +++ b/paperforge/plugin/main.js @@ -386,6 +386,11 @@ class PaperForgeStatusView extends ItemView { } _fallbackFetchStats(quiet, vp, plugin) { + if (!this._fallbackWarned) { + console.warn('[PaperForge] formal-library.json fallback is deprecated. ' + + 'Ensure `paperforge status --json` is available for accurate dashboard data.'); + this._fallbackWarned = true; + } const systemDir = plugin?.settings?.system_dir || 'System'; const indexPath = path.join(vp, systemDir, 'PaperForge', 'indexes', 'formal-library.json'); @@ -494,6 +499,11 @@ class PaperForgeStatusView extends ItemView { /* ── Index Loading (D-11, D-17, D-19) ── */ _loadIndex() { + if (!this._fallbackWarned) { + console.warn('[PaperForge] formal-library.json fallback is deprecated. ' + + 'Ensure `paperforge status --json` is available for accurate dashboard data.'); + this._fallbackWarned = true; + } const vp = this.app.vault.adapter.basePath; const plugin = this.app.plugins.plugins['paperforge']; const systemDir = plugin?.settings?.system_dir || 'System'; diff --git a/paperforge/schema/field_registry.yaml b/paperforge/schema/field_registry.yaml index 27aafb30..51fe868f 100644 --- a/paperforge/schema/field_registry.yaml +++ b/paperforge/schema/field_registry.yaml @@ -4,81 +4,134 @@ frontmatter: required: true public: true description: "Zotero citation key" + owner: sync + introduced_in: "1.0" domain: type: str required: true public: true description: "Classification domain" + owner: sync + introduced_in: "1.0" title: type: str required: true public: true description: "Paper title" + owner: sync + introduced_in: "1.0" year: type: str required: false public: true description: "Publication year" + owner: sync + introduced_in: "1.0" doi: type: str required: false public: true description: "Digital Object Identifier" + owner: sync + introduced_in: "1.0" collection_path: type: str required: false public: true description: "Zotero collection path" + owner: sync + introduced_in: "1.0" has_pdf: type: bool required: true public: true - description: "Has PDF attachment" + description: "Has PDF attachment (DEPRECATED)" + owner: sync + deprecated: true + replacement: lifecycle + introduced_in: "1.0" pdf_path: type: str required: false public: true description: "Wikilink to main PDF" + owner: sync + introduced_in: "1.0" supplementary: type: list required: false public: true description: "Wikilinks to supplementary PDFs" + owner: sync + introduced_in: "1.0" fulltext_md_path: type: str required: false public: true description: "Wikilink to OCR fulltext" + owner: ocr + introduced_in: "1.0" recommend_analyze: type: bool required: false public: true - description: "System recommends analysis" + description: "System recommends analysis (DEPRECATED)" + owner: sync + deprecated: true + replacement: lifecycle + introduced_in: "1.0" analyze: type: bool required: false public: true description: "User requests deep reading" + owner: user + introduced_in: "1.0" do_ocr: type: bool required: false public: true - description: "User requests OCR" + description: "User requests OCR (DEPRECATED)" + owner: user + deprecated: true + replacement: ocr_status + introduced_in: "1.0" ocr_status: - type: str + type: enum + values: + - none + - pending + - queued + - processing + - done + - done_incomplete + - failed + - blocked + - nopdf required: false public: true description: "OCR processing status" + owner: ocr + default: "none" + introduced_in: "2.1" deep_reading_status: - type: str + type: enum + values: + - pending + - done required: false public: true description: "Deep reading status" + owner: deep_reading + default: "pending" + introduced_in: "1.0" path_error: type: str required: false public: true description: "Path resolution error" + owner: sync + introduced_in: "1.0" index_entry: file: @@ -86,130 +139,217 @@ index_entry: required: true public: true description: "Path to formal note" + owner: index + introduced_in: "2.1" zotero_key: type: str required: true public: true description: "Zotero citation key" + owner: index + introduced_in: "2.1" domain: type: str required: true public: true description: "Classification domain" + owner: index + introduced_in: "2.1" title: type: str required: true public: true description: "Paper title" + owner: index + introduced_in: "2.1" year: type: str required: false public: true description: "Publication year" + owner: index + introduced_in: "2.1" doi: type: str required: false public: true description: "DOI" + owner: index + introduced_in: "2.1" collection_path: type: str required: false public: true description: "Zotero collection path" + owner: index + introduced_in: "2.1" has_pdf: type: bool required: true public: true - description: "Has PDF" + description: "Has PDF (DEPRECATED)" + owner: index + deprecated: true + replacement: lifecycle + introduced_in: "2.1" pdf_path: type: str required: false public: true description: "PDF path" + owner: index + introduced_in: "2.1" supplementary: type: list required: false public: true description: "Supplementary PDFs" + owner: index + introduced_in: "2.1" fulltext_md_path: type: str required: false public: true description: "Fulltext path" + owner: index + introduced_in: "2.1" recommend_analyze: type: bool required: false public: true - description: "Recommend analysis" + description: "Recommend analysis (DEPRECATED)" + owner: index + deprecated: true + replacement: lifecycle + introduced_in: "2.1" analyze: type: bool required: false public: true description: "Analysis flag" + owner: index + introduced_in: "2.1" do_ocr: type: bool required: false public: true - description: "OCR flag" + description: "OCR flag (DEPRECATED)" + owner: index + deprecated: true + replacement: ocr_status + introduced_in: "2.1" ocr_status: - type: str + type: enum + values: + - none + - pending + - queued + - processing + - done + - done_incomplete + - failed + - blocked + - nopdf required: false public: true description: "OCR status" + owner: index + introduced_in: "2.1" deep_reading_status: - type: str + type: enum + values: + - pending + - done required: false public: true description: "Deep reading status" + owner: index + introduced_in: "2.1" path_error: type: str required: false public: true description: "Path error" + owner: index + introduced_in: "2.1" lifecycle: - type: str + type: enum + values: + - pdf_ready + - ocr_ready + - analyze_ready + - deep_read_done + - error_state required: false public: true description: "Derived lifecycle stage" + owner: index + introduced_in: "2.1" health: - type: str + type: dict required: false public: true - description: "Derived health dimension scores" + description: "Derived health dimension scores (dict: pdf_health, ocr_health, note_health, asset_health)" + owner: index + introduced_in: "2.1" maturity: type: str required: false public: true description: "Derived maturity score" + owner: index + introduced_in: "2.1" next_step: type: str required: false public: true description: "Recommended next action" + owner: index + introduced_in: "2.1" ocr_meta: ocr_status: - type: str + type: enum + values: + - none + - pending + - queued + - processing + - done + - done_incomplete + - failed + - blocked + - nopdf required: true public: true description: "OCR processing status" + owner: ocr + introduced_in: "2.1" error: type: str required: false public: true description: "Error message if failed" + owner: ocr + introduced_in: "1.0" start_time: type: str required: false public: false description: "OCR start timestamp" + owner: ocr + introduced_in: "1.0" completion_time: type: str required: false public: false description: "OCR completion timestamp" + owner: ocr + introduced_in: "1.0" page_count: type: int required: false public: false description: "Number of pages processed" + owner: ocr + introduced_in: "1.0" diff --git a/paperforge/services/sync_service.py b/paperforge/services/sync_service.py index 786ed7b9..093b308c 100644 --- a/paperforge/services/sync_service.py +++ b/paperforge/services/sync_service.py @@ -1,14 +1,15 @@ from __future__ import annotations import logging +import re from pathlib import Path -from typing import Any +from paperforge import __version__ from paperforge.adapters.bbt import ( extract_authors, - _identify_main_pdf, - _normalize_attachment_path, + identify_main_pdf, load_export_rows, + normalize_attachment_path, resolve_item_collection_paths, ) from paperforge.adapters.obsidian_frontmatter import ( @@ -21,12 +22,20 @@ from paperforge.adapters.zotero_paths import ( obsidian_wikilink_for_path, obsidian_wikilink_for_pdf, ) +from paperforge.core.errors import ErrorCode +from paperforge.core.result import PFError, PFResult logger = logging.getLogger(__name__) class SyncService: - """Orchestrates the sync lifecycle: load BBT exports, match entries, generate notes.""" + """Orchestrates the sync lifecycle: load BBT exports, match entries, generate notes. + + v2.1: This service currently delegates heavy lifting to worker/sync.py. + The load_exports / build_candidates methods are functional; run() wraps + the legacy worker and enforces the PFResult contract. Full migration + of note-writing logic is planned for v2.2. + """ def __init__(self, vault: Path): self.vault = vault @@ -34,15 +43,17 @@ class SyncService: self._resolved = False def resolve_paths(self) -> dict[str, Path]: - """Resolve vault paths using pipeline_paths.""" - from paperforge.worker._utils import pipeline_paths - self.paths = pipeline_paths(self.vault) + """Resolve vault paths using paperforge_paths.""" + from paperforge.config import paperforge_paths + + self.paths = paperforge_paths(self.vault) self._resolved = True return self.paths def load_exports(self) -> list[dict]: """Load all BBT export JSON files.""" from paperforge.config import load_vault_config + cfg = load_vault_config(self.vault) system_dir = self.vault / cfg.get("system_dir", "99_System") exports_dir = system_dir / "PaperForge" / "exports" @@ -58,6 +69,135 @@ class SyncService: logger.error("Failed to load %s: %s", f.name, e) return rows + def clean_orphaned_records( + self, exports: dict[str, dict[str, dict]], paths: dict[str, Path], json_output: bool = False + ) -> int: + """Remove formal notes whose title duplicates an exported item but lacks PDF. + + Migrated from worker/sync.py run_index_refresh (v2.2). + Returns number of orphaned records deleted. + """ + from paperforge.config import paperforge_paths + + _paths = paths if paths else paperforge_paths(self.vault) + lit_dir = _paths.get("literature") + if not lit_dir or not lit_dir.exists(): + return 0 + + total_deleted = 0 + for domain_dir in lit_dir.iterdir(): + if not domain_dir.is_dir(): + continue + domain = domain_dir.name + domain_export_keys = set(exports.get(domain, {}).keys()) + records_by_title: dict[str, list[str]] = {} + records_info: dict[str, dict] = {} + + for record_file in domain_dir.rglob("*.md"): + if record_file.name in ("fulltext.md", "deep-reading.md", "discussion.md"): + continue + try: + content = record_file.read_text(encoding="utf-8") + key_match = re.search(r"^zotero_key:\s*(.+)$", content, re.MULTILINE) + if not key_match: + continue + key = key_match.group(1).strip() + title_match = re.search(r"^title:\s*[\"']?(.+)[\"']?\s*$", content, re.MULTILINE) + title = title_match.group(1) if title_match else "" + has_pdf = "has_pdf: true" in content + normalized = re.sub(r"[^a-z0-9]", "", title.lower())[:20] + records_info[key] = { + "file": record_file, + "title": title, + "has_pdf": has_pdf, + "normalized": normalized, + } + if normalized not in records_by_title: + records_by_title[normalized] = [] + records_by_title[normalized].append(key) + except Exception: + continue + + to_delete: list[str] = [] + for normalized, keys in records_by_title.items(): + keys_in_export = [k for k in keys if k in domain_export_keys] + keys_not_in_export = [k for k in keys if k not in domain_export_keys] + if keys_in_export and keys_not_in_export: + for k in keys_not_in_export: + if not records_info[k]["has_pdf"]: + to_delete.append(k) + + deleted = 0 + for key in to_delete: + try: + records_info[key]["file"].unlink() + deleted += 1 + except Exception: + pass + if deleted > 0 and not json_output: + print(f"index-refresh: cleaned {deleted} orphaned records in {domain}") + total_deleted += deleted + + return total_deleted + + def clean_flat_notes(self, paths: dict[str, Path], json_output: bool = False) -> int: + """Delete flat formal notes whose workspace equivalent exists. + + Migrated from worker/sync.py run_index_refresh (v2.2). + Returns number of flat notes deleted. + """ + from paperforge.config import paperforge_paths + from paperforge.worker._utils import slugify_filename + + try: + from paperforge.worker.asset_index import read_index as _read_idx + + index_data = _read_idx(self.vault) + except Exception: + index_data = None + + _paths = paths if paths else paperforge_paths(self.vault) + lit_dir = _paths.get("literature") + if not lit_dir or not lit_dir.exists(): + return 0 + + ws_keys: set[str] = set() + if isinstance(index_data, dict): + for item in index_data.get("items", []): + ws_dir = ( + lit_dir + / item.get("domain", "") + / (item.get("zotero_key", "") + " - " + slugify_filename(item.get("title", ""))) + ) + if ws_dir.is_dir(): + ws_keys.add(item.get("zotero_key")) + + if not ws_keys: + return 0 + + cleaned = 0 + for domain_dir in sorted(lit_dir.iterdir()): + if not domain_dir.is_dir(): + continue + for flat_note in list(domain_dir.glob("*.md")): + try: + text = flat_note.read_text(encoding="utf-8") + m = re.search(r"^zotero_key:\s*\"?(\S+?)\"?\s*$", text, re.MULTILINE) + key = m.group(1) if m else "" + except Exception: + continue + if key and key in ws_keys: + try: + flat_note.unlink() + cleaned += 1 + except Exception: + pass + + if cleaned > 0 and not json_output: + print(f"index-refresh: cleaned {cleaned} flat note(s) (migrated to workspace)") + + return cleaned + def build_candidates(self, rows: list[dict]) -> list[dict]: """Generate candidate markdown entries from BBT rows.""" candidates = [] @@ -69,10 +209,97 @@ class SyncService: logger.error("Failed to build candidate: %s", e) return candidates - def run_sync(self, verbose: bool = False, json_output: bool = False) -> int | dict: - """Full sync orchestration. Delegates to run_selection_sync logic. + def run(self, verbose: bool = False, json_output: bool = False) -> PFResult: + """Full sync orchestration. Returns PFResult contract. - Returns int (exit code) for CLI mode, dict for json_output mode. + v2.2: Service now orchestrates the full sync lifecycle: + 1. Load BBT exports + 2. Selection sync (count items from worker) + 3. Build canonical index (asset_index) + 4. Clean orphaned records + flat notes + 5. Return PFResult """ + # ── Phase 1: Select ── from paperforge.worker.sync import run_selection_sync - return run_selection_sync(self.vault, verbose=verbose, json_output=json_output) + + _export_code = "ok" + try: + rows = self.load_exports() + if not rows: + _export_code = "BBT_EXPORT_NOT_FOUND" + except Exception: + _export_code = "BBT_EXPORT_INVALID" + + try: + selection_result = run_selection_sync(self.vault, verbose=verbose, json_output=json_output) + except Exception as exc: + return PFResult( + ok=False, + command="sync", + version=__version__, + error=PFError( + code=ErrorCode.SYNC_FAILED, + message=str(exc), + details={"phase": "selection", "exception_type": type(exc).__name__}, + ), + ) + + # ── Phase 2: Index ── + from paperforge.worker._utils import pipeline_paths + from paperforge.worker._domain import load_domain_config + from paperforge.worker.base_views import ensure_base_views + + paths = pipeline_paths(self.vault) + config = load_domain_config(paths) + ensure_base_views(self.vault, paths, config) + domain_lookup = {entry["export_file"]: entry["domain"] for entry in config["domains"]} + + from paperforge.config import load_vault_config + + cfg = load_vault_config(self.vault) + exports: dict[str, dict[str, dict]] = {} + for export_path in sorted(paths["exports"].glob("*.json")): + domain = domain_lookup.get(export_path.name, export_path.stem) + export_rows = load_export_rows(export_path) + exports[domain] = {row["key"]: row for row in export_rows} + + from paperforge.worker.sync import migrate_to_workspace + + migrate_to_workspace(self.vault, paths) + + import paperforge.worker.asset_index as asset_index + + index_count = asset_index.build_index(self.vault, verbose) + + # ── Phase 3: Clean ── + orphaned = self.clean_orphaned_records(exports, paths, json_output=json_output) + flat_cleaned = self.clean_flat_notes(paths, json_output=json_output) + + if not json_output: + print(f"index-refresh: {index_count} entries in index") + total = sum(1 for _ in paths["literature"].rglob("*.md")) if paths["literature"].exists() else 0 + print(f"index-refresh: {total} formal note(s) in literature") + + all_errors = list(selection_result.get("errors", [])) + is_ok = selection_result.get("failed", 0) == 0 and len(all_errors) == 0 + + result = PFResult( + ok=is_ok, + command="sync", + version=__version__, + data={ + "selection": selection_result, + "index": {"updated": index_count, "orphaned_cleaned": orphaned, "flat_cleaned": flat_cleaned}, + }, + ) + + if _export_code != "ok": + result.warnings.append(f"BBT exports: {_export_code}") + + return result + + # ── Legacy passthrough (backward-compat for commands/sync.py) ── + + def run_sync(self, verbose: bool = False, json_output: bool = False) -> PFResult: + """Alias for run(). Returns PFResult (not int|dict).""" + return self.run(verbose=verbose, json_output=json_output) diff --git a/paperforge/worker/_domain.py b/paperforge/worker/_domain.py index 03900507..5e4d8f4f 100644 --- a/paperforge/worker/_domain.py +++ b/paperforge/worker/_domain.py @@ -14,39 +14,13 @@ import logging from pathlib import Path from paperforge.worker._utils import read_json, write_json +from paperforge.adapters.collections import build_collection_lookup logger = logging.getLogger(__name__) # ── Collection tree utilities ──────────────────────────────────────────────── - - -def build_collection_lookup(collections: dict) -> dict: - """Build parent-resolved path cache from a Zotero collection tree. - - Returns a dict with: - path_by_key: {collection_key: "Parent/Sub/Name", ...} - paths_by_item_id: {item_id: [collection_path, ...], ...} - """ - path_cache = {} - item_paths = {} - - def path_for(key: str) -> str: - if key in path_cache: - return path_cache[key] - node = collections.get(key, {}) - parent = node.get("parent") or "" - name = node.get("name", "") - parent_path = path_for(parent) if parent else "" - full_path = f"{parent_path}/{name}" if parent_path else name - path_cache[key] = full_path - return full_path - - for key, node in collections.items(): - full_path = path_for(key) - for item_id in node.get("items", []): - item_paths.setdefault(item_id, []).append(full_path) - return {"path_by_key": path_cache, "paths_by_item_id": item_paths} +# build_collection_lookup re-exported from paperforge.adapters.collections def export_collection_paths(export_path: Path) -> list[str]: @@ -122,4 +96,8 @@ def load_domain_collections(paths: dict[str, Path]) -> dict[str, list[str]]: Used by candidate collection resolution in sync.py. """ config = load_domain_config(paths) - return {entry["domain"]: entry.get("allowed_collections", []) for entry in config.get("domains", []) if entry.get("domain")} + return { + entry["domain"]: entry.get("allowed_collections", []) + for entry in config.get("domains", []) + if entry.get("domain") + } diff --git a/paperforge/worker/_utils.py b/paperforge/worker/_utils.py index fc7eef94..9ccdca5b 100644 --- a/paperforge/worker/_utils.py +++ b/paperforge/worker/_utils.py @@ -20,8 +20,7 @@ STANDARD_VIEW_NAMES = frozenset( # --- Journal Database --- -def read_json(path: Path): - return json.loads(path.read_text(encoding="utf-8")) +from paperforge.core.io import read_json, write_json _JOURNAL_DB: dict[str, dict] | None = None @@ -64,11 +63,7 @@ def lookup_impact_factor(journal_name: str, extra: str, vault: Path) -> str: # --- JSON I/O --- - - -def write_json(path: Path, data) -> None: - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8") +# write_json re-exported from paperforge.core.io def read_jsonl(path: Path): @@ -128,9 +123,7 @@ def slugify_filename(text: str) -> str: return cleaned or "untitled" -def _extract_year(value: str) -> str: - match = re.search("(19|20)\\d{2}", value or "") - return match.group(0) if match else "" +from paperforge.core.date_utils import extract_year as _extract_year # --- Deep-Reading Queue --- @@ -191,7 +184,9 @@ def get_analyze_queue(vault: Path) -> list[dict]: continue # Quick exit: check analyze before extracting other fields - analyze_match = re.search(r"^analyze:\s*(?:[\"'])?(true|false)(?:[\"'])?\s*$", text, re.MULTILINE | re.IGNORECASE) + analyze_match = re.search( + r"^analyze:\s*(?:[\"'])?(true|false)(?:[\"'])?\s*$", text, re.MULTILINE | re.IGNORECASE + ) if not analyze_match or analyze_match.group(1).lower() != "true": continue @@ -225,16 +220,18 @@ def get_analyze_queue(vault: Path) -> list[dict]: if dr_match: dr_status = dr_match.group(1).strip() - results.append({ - "zotero_key": zotero_key, - "domain": domain, - "title": title, - "analyze": True, - "do_ocr": do_ocr, - "ocr_status": ocr_status, - "deep_reading_status": dr_status, - "note_path": note_file, - }) + results.append( + { + "zotero_key": zotero_key, + "domain": domain, + "title": title, + "analyze": True, + "do_ocr": do_ocr, + "ocr_status": ocr_status, + "deep_reading_status": dr_status, + "note_path": note_file, + } + ) results.sort(key=lambda r: (r["domain"], r["zotero_key"])) return results @@ -278,6 +275,7 @@ def install_obsidian_plugin(vault: Path) -> bool: plugin_src = vault / "paperforge" / "plugin" if not plugin_src.is_dir(): import paperforge + plugin_src = Path(paperforge.__file__).parent.resolve() / "plugin" if not plugin_src.is_dir(): logger.warning("Plugin source not found: %s", plugin_src) diff --git a/paperforge/worker/sync.py b/paperforge/worker/sync.py index ce905175..256b3023 100644 --- a/paperforge/worker/sync.py +++ b/paperforge/worker/sync.py @@ -1,5 +1,12 @@ from __future__ import annotations +# ============================================================================= +# FREEZE LINE (v2.1-contract-hardened) +# No new business logic allowed in this file. +# New logic → services/ / adapters/ / core/. +# This file: deletion, migration, legacy wrappers only. +# ============================================================================= + import html import logging import os @@ -58,8 +65,8 @@ from paperforge.adapters.bbt import ( import paperforge.worker.asset_index as asset_index -logger = logging.getLogger(__name__) +logger = logging.getLogger(__name__) def load_export_inventory(paths: dict[str, Path]) -> dict[str, dict]: @@ -179,7 +186,6 @@ def apply_existing_library_match(row: dict, inventory: dict[str, dict]) -> dict: return resolved - def load_control_actions(paths: dict[str, Path]) -> dict[str, dict]: actions = {} lit_root = paths.get("literature") @@ -201,14 +207,16 @@ def load_control_actions(paths: dict[str, Path]) -> dict[str, dict]: if do_ocr_match: do_ocr = do_ocr_match.group(1).lower() == "true" analyze = False - analyze_match = re.search(r"^analyze:\s*(?:[\"'])?(true|false)(?:[\"'])?\s*$", text, re.MULTILINE | re.IGNORECASE) + analyze_match = re.search( + r"^analyze:\s*(?:[\"'])?(true|false)(?:[\"'])?\s*$", text, re.MULTILINE | re.IGNORECASE + ) if analyze_match: analyze = analyze_match.group(1).lower() == "true" actions[zotero_key] = {"analyze": analyze, "do_ocr": do_ocr} return actions -def run_selection_sync(vault: Path, verbose: bool = False, json_output: bool = False) -> int | dict: +def run_selection_sync(vault: Path, verbose: bool = False, json_output: bool = False) -> dict: from paperforge.worker.base_views import ensure_base_views from paperforge.worker.ocr import validate_ocr_meta @@ -267,10 +275,10 @@ def run_selection_sync(vault: Path, verbose: bool = False, json_output: bool = F # Formal notes now carry workflow flags (do_ocr, analyze) directly. # Existing library-records are migrated via Phase 40 logic. updated += 1 - if json_output: - return {"new": written, "updated": updated, "skipped": 0, "failed": 0, "errors": []} - print(f"selection-sync: wrote {written} records, updated {updated} records") - return 0 + result = {"new": written, "updated": updated, "skipped": 0, "failed": 0, "errors": []} + if not json_output: + print(f"selection-sync: wrote {written} records, updated {updated} records") + return result def load_candidates_by_id(paths: dict[str, Path]) -> dict[str, dict]: @@ -1118,7 +1126,9 @@ def migrate_to_workspace(vault: Path, paths: dict) -> int: entry = indexed_entries.get(key, {}) title = str(entry.get("title", "") or "").strip() if not title: - title_match = re.search(r'^title:\s*["\']?(.+?)["\']?\s*$', flat_note_path.read_text(encoding="utf-8"), re.MULTILINE) + title_match = re.search( + r'^title:\s*["\']?(.+?)["\']?\s*$', flat_note_path.read_text(encoding="utf-8"), re.MULTILINE + ) title = title_match.group(1).strip() if title_match else "" stem = flat_note_path.stem if title: @@ -1182,6 +1192,7 @@ def migrate_to_workspace(vault: Path, paths: dict) -> int: target_fulltext = workspace_dir / "fulltext.md" if source_fulltext.exists() and not target_fulltext.exists(): import shutil + shutil.copy2(str(source_fulltext), str(target_fulltext)) except Exception: pass @@ -1194,7 +1205,9 @@ def migrate_to_workspace(vault: Path, paths: dict) -> int: return migrated -def run_index_refresh(vault: Path, verbose: bool = False, rebuild_index: bool = False, json_output: bool = False) -> int: +def run_index_refresh( + vault: Path, verbose: bool = False, rebuild_index: bool = False, json_output: bool = False +) -> dict: """Refresh the canonical asset index. Default behavior: full rebuild. This is the safe default because @@ -1227,90 +1240,9 @@ def run_index_refresh(vault: Path, verbose: bool = False, rebuild_index: bool = migrate_to_workspace(vault, paths) # Delegate to asset_index.build_index() for the core build loop count = asset_index.build_index(vault, verbose) - control_records_dir = paths["literature"] - if control_records_dir.exists(): - for domain_dir in control_records_dir.iterdir(): - if not domain_dir.is_dir(): - continue - domain = domain_dir.name - domain_export_keys = set(exports.get(domain, {}).keys()) - records_by_title = {} - records_info = {} - for record_file in domain_dir.rglob("*.md"): - if record_file.name in ("fulltext.md", "deep-reading.md", "discussion.md"): - continue - try: - content = record_file.read_text(encoding="utf-8") - key_match = re.search(r"^zotero_key:\s*(.+)$", content, re.MULTILINE) - if not key_match: - continue - key = key_match.group(1).strip() - title_match = re.search("^title:\\s*[\"\\']?(.+)[\"\\']?\\s*$", content, re.MULTILINE) - title = title_match.group(1) if title_match else "" - has_pdf = "has_pdf: true" in content - normalized = re.sub("[^a-z0-9]", "", title.lower())[:20] - records_info[key] = { - "file": record_file, - "title": title, - "has_pdf": has_pdf, - "normalized": normalized, - } - if normalized not in records_by_title: - records_by_title[normalized] = [] - records_by_title[normalized].append(key) - except Exception: - continue - to_delete = [] - for normalized, keys in records_by_title.items(): - keys_in_export = [k for k in keys if k in domain_export_keys] - keys_not_in_export = [k for k in keys if k not in domain_export_keys] - if keys_in_export and keys_not_in_export: - for k in keys_not_in_export: - if not records_info[k]["has_pdf"]: - to_delete.append(k) - deleted_count = 0 - for key in to_delete: - try: - records_info[key]["file"].unlink() - deleted_count += 1 - except Exception: - pass - if deleted_count > 0 and not json_output: - print(f"index-refresh: cleaned {deleted_count} orphaned records in {domain}") - # Clean up flat notes: delete only if confirmed by canonical index + workspace - index_data = asset_index.read_index(vault) - ws_keys = set() - if isinstance(index_data, dict): - for item in index_data.get("items", []): - ws_dir = paths["literature"] / item.get("domain", "") / (item.get("zotero_key", "") + " - " + slugify_filename(item.get("title", ""))) - if ws_dir.is_dir(): - ws_keys.add(item.get("zotero_key")) - - lit_dir = paths["literature"] - cleaned_flat = 0 - if lit_dir.exists() and ws_keys: - for domain_dir in sorted(lit_dir.iterdir()): - if not domain_dir.is_dir(): - continue - for flat_note in list(domain_dir.glob("*.md")): - try: - text = flat_note.read_text(encoding="utf-8") - m = re.search(r"^zotero_key:\s*\"?(\S+?)\"?\s*$", text, re.MULTILINE) - key = m.group(1) if m else "" - except Exception: - continue - if key and key in ws_keys: - try: - flat_note.unlink() - cleaned_flat += 1 - except Exception: - pass - if cleaned_flat > 0 and not json_output: - print(f"index-refresh: cleaned {cleaned_flat} flat note(s) (migrated to workspace)") - - if control_records_dir.exists() and not json_output: - total = sum(1 for _ in control_records_dir.rglob("*.md")) + if paths["literature"].exists() and not json_output: + total = sum(1 for _ in paths["literature"].rglob("*.md")) print(f"index-refresh: {total} formal note(s) in literature") - return 0 + return {"updated": count, "failed": 0, "errors": []} diff --git a/tests/unit/core/test_errors.py b/tests/unit/core/test_errors.py index 8bfb7323..42cd2557 100644 --- a/tests/unit/core/test_errors.py +++ b/tests/unit/core/test_errors.py @@ -21,8 +21,16 @@ class TestErrorCodeMembers: assert str(ErrorCode.PYTHON_NOT_FOUND) == "PYTHON_NOT_FOUND" assert str(ErrorCode.UNKNOWN) == "UNKNOWN" - def test_member_count(self) -> None: - assert len(ErrorCode) == 8 + def test_member_count_minimum(self) -> None: + assert len(ErrorCode) >= 8 # expanded from 8 to 26 in v2.1 contract hardening def test_is_str_enum(self) -> None: assert isinstance(ErrorCode.PYTHON_NOT_FOUND, str) + + def test_unknown_code_returns_unknown(self) -> None: + """Newer plugin versions may send codes not yet in this runtime.""" + assert ErrorCode("NONEXISTENT_CODE") is ErrorCode.UNKNOWN + + def test_future_code_returns_unknown(self) -> None: + assert ErrorCode("PAPERFORGE_NOT_INSTALLED") is ErrorCode.PAPERFORGE_NOT_INSTALLED + assert ErrorCode("OCR_UPLOAD_FAILED") is ErrorCode.OCR_UPLOAD_FAILED diff --git a/tests/unit/core/test_result.py b/tests/unit/core/test_result.py index eef9d6d0..a3a629c6 100644 --- a/tests/unit/core/test_result.py +++ b/tests/unit/core/test_result.py @@ -103,3 +103,103 @@ class TestPFResultRoundTrip: result = PFResult.from_dict(data) assert result.error is None assert result.ok is True + + +class TestPFResultV21Fields: + """v2.1 contract hardening: warnings, next_actions, suggestions round-trip.""" + + def test_warnings_round_trip(self) -> None: + original = PFResult( + ok=True, + command="sync", + version="2.1.0", + warnings=["BBT export is stale; re-export recommended"], + ) + d = original.to_dict() + assert d["warnings"] == ["BBT export is stale; re-export recommended"] + reconstructed = PFResult.from_dict(d) + assert reconstructed.warnings == original.warnings + assert reconstructed == original + + def test_warnings_empty_omitted_from_dict(self) -> None: + result = PFResult(ok=True, command="sync", version="2.1.0") + d = result.to_dict() + assert "warnings" not in d + + def test_warnings_from_dict_missing_defaults_empty(self) -> None: + data = {"ok": True, "command": "test", "version": "1.0.0"} + result = PFResult.from_dict(data) + assert result.warnings == [] + + def test_next_actions_round_trip(self) -> None: + original = PFResult( + ok=False, + command="sync", + version="2.1.0", + next_actions=[ + {"id": "open_exports_dir", "label": "Open exports folder"}, + {"id": "run_zotero", "label": "Launch Zotero"}, + ], + ) + d = original.to_dict() + assert len(d["next_actions"]) == 2 + assert d["next_actions"][0]["id"] == "open_exports_dir" + reconstructed = PFResult.from_dict(d) + assert reconstructed.next_actions == original.next_actions + + def test_next_actions_empty_omitted_from_dict(self) -> None: + result = PFResult(ok=True, command="sync", version="2.1.0") + d = result.to_dict() + assert "next_actions" not in d + + def test_error_suggestions_round_trip(self) -> None: + error = PFError( + code=ErrorCode.BBT_EXPORT_NOT_FOUND, + message="No BBT export found", + suggestions=[ + "In Zotero, right-click collection → Export Collection", + "Format: Better BibTeX JSON", + "Tick 'Keep updated'", + ], + ) + result = PFResult(ok=False, command="sync", version="2.1.0", error=error) + d = result.to_dict() + assert "suggestions" in d["error"] + assert len(d["error"]["suggestions"]) == 3 + reconstructed = PFResult.from_dict(d) + assert reconstructed.error is not None + assert reconstructed.error.suggestions == error.suggestions + + def test_error_suggestions_empty_omitted(self) -> None: + error = PFError(code=ErrorCode.SYNC_FAILED, message="fail") + result = PFResult(ok=False, command="sync", version="2.1.0", error=error) + d = result.to_dict() + assert "suggestions" not in d["error"] + + def test_from_dict_missing_suggestions_defaults_empty(self) -> None: + data = { + "ok": False, + "command": "test", + "version": "1.0.0", + "error": {"code": "SYNC_FAILED", "message": "fail", "details": {}}, + } + result = PFResult.from_dict(data) + assert result.error is not None + assert result.error.suggestions == [] + + def test_from_dict_unknown_error_code_graceful(self) -> None: + """Newer plugin may send error codes not yet in this runtime.""" + data = { + "ok": False, + "command": "sync", + "version": "2.1.0", + "error": { + "code": "SOME_FUTURE_ERROR_CODE", + "message": "Something new happened", + "details": {}, + }, + } + result = PFResult.from_dict(data) + assert result.error is not None + assert result.error.code is ErrorCode.UNKNOWN + assert result.error.message == "Something new happened" diff --git a/tests/unit/core/test_state.py b/tests/unit/core/test_state.py index 17bc1c25..e3361122 100644 --- a/tests/unit/core/test_state.py +++ b/tests/unit/core/test_state.py @@ -10,54 +10,51 @@ from paperforge.core.state import ( class TestOcrStatus: - """OcrStatus enum values and legacy mapping.""" + """OcrStatus enum values and legacy mapping (1:1, no compression).""" - def test_values_match_expected_strings(self) -> None: + def test_all_canonical_values(self) -> None: + assert OcrStatus.NONE.value == "none" assert OcrStatus.PENDING.value == "pending" + assert OcrStatus.QUEUED.value == "queued" assert OcrStatus.PROCESSING.value == "processing" assert OcrStatus.DONE.value == "done" + assert OcrStatus.DONE_INCOMPLETE.value == "done_incomplete" assert OcrStatus.FAILED.value == "failed" + assert OcrStatus.BLOCKED.value == "blocked" + assert OcrStatus.NO_PDF.value == "nopdf" - def test_from_legacy_canonical_pending(self) -> None: + def test_from_legacy_identity_mappings(self) -> None: + assert OcrStatus.from_legacy("none") is OcrStatus.NONE assert OcrStatus.from_legacy("pending") is OcrStatus.PENDING - - def test_from_legacy_canonical_processing(self) -> None: + assert OcrStatus.from_legacy("queued") is OcrStatus.QUEUED assert OcrStatus.from_legacy("processing") is OcrStatus.PROCESSING - - def test_from_legacy_canonical_done(self) -> None: assert OcrStatus.from_legacy("done") is OcrStatus.DONE - - def test_from_legacy_canonical_failed(self) -> None: + assert OcrStatus.from_legacy("done_incomplete") is OcrStatus.DONE_INCOMPLETE assert OcrStatus.from_legacy("failed") is OcrStatus.FAILED + assert OcrStatus.from_legacy("blocked") is OcrStatus.BLOCKED + assert OcrStatus.from_legacy("nopdf") is OcrStatus.NO_PDF - def test_from_legacy_queued_maps_to_pending(self) -> None: - assert OcrStatus.from_legacy("queued") is OcrStatus.PENDING - - def test_from_legacy_running_maps_to_processing(self) -> None: + def test_from_legacy_aliases(self) -> None: assert OcrStatus.from_legacy("running") is OcrStatus.PROCESSING - - def test_from_legacy_error_maps_to_failed(self) -> None: assert OcrStatus.from_legacy("error") is OcrStatus.FAILED - def test_from_legacy_blocked_maps_to_failed(self) -> None: - assert OcrStatus.from_legacy("blocked") is OcrStatus.FAILED - - def test_from_legacy_nopdf_maps_to_failed(self) -> None: - assert OcrStatus.from_legacy("nopdf") is OcrStatus.FAILED - - def test_from_legacy_done_incomplete_maps_to_done(self) -> None: - assert OcrStatus.from_legacy("done_incomplete") is OcrStatus.DONE - def test_from_legacy_unknown_maps_to_pending(self) -> None: assert OcrStatus.from_legacy("nonexistent") is OcrStatus.PENDING def test_from_legacy_case_and_whitespace_insensitive(self) -> None: assert OcrStatus.from_legacy(" DONE ") is OcrStatus.DONE - assert OcrStatus.from_legacy("Queued") is OcrStatus.PENDING + assert OcrStatus.from_legacy("Queued") is OcrStatus.QUEUED + assert OcrStatus.from_legacy(" NoPdf ") is OcrStatus.NO_PDF def test_str_returns_value(self) -> None: assert str(OcrStatus.PENDING) == "pending" assert str(OcrStatus.DONE) == "done" + assert str(OcrStatus.NONE) == "none" + + def test_enum_equals_string_value(self) -> None: + assert OcrStatus.DONE == "done" + assert OcrStatus.BLOCKED == "blocked" + assert OcrStatus.NO_PDF == "nopdf" class TestPdfStatus: @@ -101,8 +98,13 @@ class TestAllowedTransitionsOcr: assert ok is True assert msg == "" - def test_pending_to_failed_valid(self) -> None: + def test_pending_to_failed_invalid(self) -> None: ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PENDING, OcrStatus.FAILED) + assert ok is False + assert "Illegal OCR transition" in msg + + def test_pending_to_queued_valid(self) -> None: + ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PENDING, OcrStatus.QUEUED) assert ok is True assert msg == "" @@ -111,6 +113,11 @@ class TestAllowedTransitionsOcr: assert ok is True assert msg == "" + def test_processing_to_done_incomplete_valid(self) -> None: + ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PROCESSING, OcrStatus.DONE_INCOMPLETE) + assert ok is True + assert msg == "" + def test_processing_to_failed_valid(self) -> None: ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PROCESSING, OcrStatus.FAILED) assert ok is True @@ -126,6 +133,16 @@ class TestAllowedTransitionsOcr: assert ok is True assert msg == "" + def test_blocked_to_pending_valid(self) -> None: + ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.BLOCKED, OcrStatus.PENDING) + assert ok is True + assert msg == "" + + def test_nopdf_to_pending_valid(self) -> None: + ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.NO_PDF, OcrStatus.PENDING) + assert ok is True + assert msg == "" + def test_pending_to_done_invalid(self) -> None: ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PENDING, OcrStatus.DONE) assert ok is False @@ -144,8 +161,8 @@ class TestAllowedTransitionsOcr: def test_invalid_transition_includes_allowed_list(self) -> None: ok, msg = ALLOWED_TRANSITIONS.check_ocr(OcrStatus.PENDING, OcrStatus.DONE) assert ok is False - assert "processing" in msg - assert "failed" in msg + assert "queued" in msg or "Queued" in msg + assert "processing" in msg or "Processing" in msg class TestAllowedTransitionsLifecycle: