lllin000_PaperForge/paperforge/worker/sync.py

1262 lines
54 KiB
Python

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
import re
import urllib.parse
from datetime import datetime, timezone
from pathlib import Path
from xml.etree import ElementTree as ET
import requests
import paperforge.worker.asset_index as asset_index
from paperforge.adapters.bbt import (
collection_fields,
load_export_rows,
)
from paperforge.adapters.obsidian_frontmatter import (
_legacy_control_flags,
_read_frontmatter_optional_bool_from_text,
canonicalize_decision,
compute_final_collection,
extract_preserved_deep_reading,
extract_preserved_ocr_redo,
extract_preserved_tags,
update_frontmatter_field,
)
from paperforge.adapters.zotero_paths import (
obsidian_wikilink_for_path,
obsidian_wikilink_for_pdf,
)
from paperforge.config import load_vault_config
from paperforge.worker._domain import load_domain_collections, load_domain_config
from paperforge.worker._utils import (
_extract_year,
lookup_impact_factor,
pipeline_paths,
read_json,
read_jsonl,
slugify_filename,
write_json,
write_jsonl,
yaml_block,
yaml_list,
yaml_quote,
)
logger = logging.getLogger(__name__)
def load_export_inventory(paths: dict[str, Path]) -> dict[str, dict]:
inventory = {"doi": {}, "pmid": {}, "title": {}}
for export_path in sorted(paths["exports"].glob("*.json")):
domain = export_path.stem
for item in load_export_rows(export_path):
record = {
"zotero_key": item.get("key", ""),
"domain": domain,
"title": item.get("title", ""),
"doi": item.get("doi", ""),
"pmid": item.get("pmid", ""),
"collections": item.get("collections", []),
}
doi = str(record.get("doi", "") or "").strip().lower()
pmid = str(record.get("pmid", "") or "").strip()
title = normalize_candidate_title(record.get("title", ""))
if doi and doi not in inventory["doi"]:
inventory["doi"][doi] = record
if pmid and pmid not in inventory["pmid"]:
inventory["pmid"][pmid] = record
if title and title not in inventory["title"]:
inventory["title"][title] = record
return inventory
def find_existing_library_match(row: dict, inventory: dict[str, dict]) -> dict | None:
doi = str(row.get("doi", "") or "").strip().lower()
if doi and doi in inventory["doi"]:
return inventory["doi"][doi]
pmid = str(row.get("pmid", "") or "").strip()
if pmid and pmid in inventory["pmid"]:
return inventory["pmid"][pmid]
title = normalize_candidate_title(row.get("title", ""))
if title and title in inventory["title"]:
return inventory["title"][title]
return None
def resolve_collection_choice(domain: str, raw_value: str, catalog: dict[str, list[str]]) -> dict[str, str]:
text = str(raw_value or "").strip()
if not text:
return {"resolved": "", "match": "", "input": ""}
allowed = [path for path in catalog.get(domain, []) if path]
if not allowed:
return {"resolved": "", "match": "no_catalog", "input": text}
lower_text = text.lower()
if "/" not in text:
leaf_matches = [path for path in allowed if path.split("/")[-1].strip().lower() == lower_text]
leaf_matches = sorted(set(leaf_matches))
if len(leaf_matches) > 1:
return {"resolved": "", "match": "ambiguous_leaf", "input": text}
exact_map = {path: path for path in allowed}
if text in exact_map:
return {"resolved": text, "match": "exact", "input": text}
lower_exact = {path.lower(): path for path in allowed}
if lower_text in lower_exact:
return {"resolved": lower_exact[lower_text], "match": "exact_ci", "input": text}
leaf_matches = [path for path in allowed if path.split("/")[-1].strip().lower() == lower_text]
if len(leaf_matches) == 1:
return {"resolved": leaf_matches[0], "match": "leaf", "input": text}
suffix_matches = [path for path in allowed if path.lower().endswith("/" + lower_text) or path.lower() == lower_text]
suffix_matches = sorted(set(suffix_matches))
if len(suffix_matches) == 1:
return {"resolved": suffix_matches[0], "match": "suffix", "input": text}
compact = re.sub("\\s+", "", lower_text)
compact_matches = []
for path in allowed:
path_compact = re.sub("\\s+", "", path.lower())
if path_compact.endswith("/" + compact) or path_compact == compact:
compact_matches.append(path)
compact_matches = sorted(set(compact_matches))
if len(compact_matches) == 1:
return {"resolved": compact_matches[0], "match": "compact_suffix", "input": text}
match = "ambiguous" if leaf_matches or suffix_matches or compact_matches else "unresolved"
return {"resolved": "", "match": match, "input": text}
def apply_candidate_collection_resolution(row: dict, catalog: dict[str, list[str]]) -> dict:
resolved = dict(row)
domain = str(resolved.get("domain", "") or "").strip()
recommended = resolve_collection_choice(domain, resolved.get("recommended_collection", ""), catalog)
user = resolve_collection_choice(domain, resolved.get("user_collection", ""), catalog)
resolved["recommended_collection"] = recommended.get("resolved", "")
resolved["user_collection_resolved"] = user.get("resolved", "")
if str(resolved.get("user_collection", "") or "").strip():
resolved["final_collection"] = user.get("resolved", "")
resolved["collection_resolution"] = (
f"user_{user.get('match', 'unresolved')}" if user.get("match") else "user_unresolved"
)
else:
resolved["final_collection"] = recommended.get("resolved", "")
resolved["collection_resolution"] = (
f"recommended_{recommended.get('match', 'unresolved')}"
if recommended.get("match")
else "recommended_unresolved"
)
return resolved
def apply_existing_library_match(row: dict, inventory: dict[str, dict]) -> dict:
resolved = dict(row)
match = find_existing_library_match(resolved, inventory)
if not match:
resolved["existing_zotero_key"] = ""
resolved["existing_collections"] = []
resolved["duplicate_hint"] = ""
return resolved
resolved["existing_zotero_key"] = str(match.get("zotero_key", "") or "").strip()
resolved["existing_collections"] = list(match.get("collections", []) or [])
collections_text = " | ".join(resolved["existing_collections"])
if collections_text:
resolved["duplicate_hint"] = f"已存在于 Zotero: {resolved['existing_zotero_key']} ({collections_text})"
else:
resolved["duplicate_hint"] = f"已存在于 Zotero: {resolved['existing_zotero_key']}"
return resolved
def load_control_actions(paths: dict[str, Path]) -> dict[str, dict]:
actions = {}
lit_root = paths.get("literature")
if not lit_root or not lit_root.exists():
return actions
for note_file in 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
zotero_key = key_match.group(1).strip().strip('"').strip("'")
do_ocr = False
do_ocr_match = re.search(r"^do_ocr:\s*(?:[\"'])?(true|false)(?:[\"'])?\s*$", text, re.MULTILINE | re.IGNORECASE)
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
)
if analyze_match:
analyze = analyze_match.group(1).lower() == "true"
ocr_redo = bool(extract_preserved_ocr_redo(text))
actions[zotero_key] = {"analyze": analyze, "do_ocr": do_ocr, "ocr_redo": ocr_redo}
return actions
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
paths = pipeline_paths(vault)
config = load_domain_config(paths)
ensure_base_views(vault, paths, config)
domain_lookup = {entry["export_file"]: entry["domain"] for entry in config["domains"]}
written = 0
updated = 0
for export_path in sorted(paths["exports"].glob("*.json")):
domain = domain_lookup.get(export_path.name, export_path.stem)
for item in load_export_rows(export_path):
pdf_attachments = [a for a in item.get("attachments", []) if a.get("contentType") == "application/pdf"]
has_pdf = bool(pdf_attachments)
raw_pdf_path = pdf_attachments[0].get("path", "") if pdf_attachments else ""
from paperforge.pdf_resolver import resolve_pdf_path
cfg = load_vault_config(vault)
zotero_dir = vault / cfg.get("system_dir", "System") / "Zotero"
resolved_pdf = resolve_pdf_path(raw_pdf_path, has_pdf, vault, zotero_dir)
collection_meta = collection_fields(item.get("collections", []))
meta_path = paths["ocr"] / item["key"] / "meta.json"
meta = read_json(meta_path) if meta_path.exists() else {}
validated_ocr_status, validated_error = validate_ocr_meta(paths, meta) if meta else ("pending", "")
if meta:
meta["ocr_status"] = validated_ocr_status
if validated_error:
meta["error"] = validated_error
write_json(meta_path, meta)
note_path = paths["literature"] / domain / f"{item['key']}.md"
note_text = note_path.read_text(encoding="utf-8") if note_path.exists() else ""
fulltext_md_path = obsidian_wikilink_for_path(
vault, meta.get("fulltext_md_path", "") or meta.get("markdown_path", "")
)
ocr_status = meta.get("ocr_status", "pending")
record_ocr_status = "nopdf" if not has_pdf or not resolved_pdf else ocr_status
creators = item.get("creators", [])
first_author = ""
for c in creators:
if c.get("creatorType") == "author":
first_author = f"{c.get('firstName', '')} {c.get('lastName', '')}".strip()
break
journal = item.get("journal", "")
extra = item.get("extra", "")
impact_factor = lookup_impact_factor(journal, extra, vault)
# Convert supplementary storage: paths to wikilinks
supplementary_wikilinks = []
for supp_path in item.get("supplementary", []):
if supp_path:
wikilink = obsidian_wikilink_for_pdf(supp_path, vault, zotero_dir)
if wikilink:
supplementary_wikilinks.append(wikilink)
pdf_wikilink = obsidian_wikilink_for_pdf(resolved_pdf, vault, zotero_dir) if resolved_pdf else ""
# Phase 37: library-records deprecated — skip creation.
# Formal notes now carry workflow flags (do_ocr, analyze) directly.
# Existing library-records are migrated via Phase 40 logic.
updated += 1
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]:
candidates = read_json(paths["candidates"])
return {row["candidate_id"]: row for row in candidates}
def save_candidates(paths: dict[str, Path], candidate_map: dict[str, dict]) -> None:
collection_catalog = load_domain_collections(paths)
export_inventory = load_export_inventory(paths)
rows = []
for row in candidate_map.values():
copy = dict(row)
copy["decision"] = canonicalize_decision(copy.get("decision", ""))
copy = apply_candidate_collection_resolution(copy, collection_catalog)
copy = apply_existing_library_match(copy, export_inventory)
copy["final_collection"] = compute_final_collection(copy)
rows.append(copy)
write_json(paths["candidates"], rows)
def writeback_command_for_candidate(row: dict) -> dict | None:
final_collection = str(row.get("final_collection", "") or "").strip()
if not final_collection:
return None
candidate_id = str(row.get("candidate_id", "") or "").strip()
if not candidate_id:
return None
command = {
"command_id": f"wb-native-{candidate_id}",
"status": "queued",
"source_candidate_id": candidate_id,
"target_domain": str(row.get("domain", "") or "").strip(),
"target_collection": final_collection,
"requested_at": datetime.now(timezone.utc).isoformat(),
}
existing_zotero_key = str(row.get("existing_zotero_key", "") or "").strip()
if existing_zotero_key:
command.update({"action": "attach_existing_item_to_collection", "existing_zotero_key": existing_zotero_key})
return command
doi = str(row.get("doi", "") or "").strip()
pmid = str(row.get("pmid", "") or "").strip()
if doi:
command.update(
{
"action": "create_item_from_identifier",
"identifier_type": "doi",
"identifier": doi,
"metadata_fallback": {
"title": row.get("title", ""),
"authors": row.get("authors", []),
"year": str(row.get("year", "") or ""),
"journal": row.get("journal", ""),
"doi": doi,
"pmid": pmid,
"abstractNote": row.get("abstract_short", ""),
},
}
)
return command
if pmid:
command.update(
{
"action": "create_item_from_identifier",
"identifier_type": "pmid",
"identifier": pmid,
"metadata_fallback": {
"title": row.get("title", ""),
"authors": row.get("authors", []),
"year": str(row.get("year", "") or ""),
"journal": row.get("journal", ""),
"doi": doi,
"pmid": pmid,
"abstractNote": row.get("abstract_short", ""),
},
}
)
return command
command.update(
{
"action": "create_item_from_metadata",
"metadata": {
"title": row.get("title", ""),
"authors": row.get("authors", []),
"year": str(row.get("year", "") or ""),
"journal": row.get("journal", ""),
"doi": doi,
"pmid": pmid,
"abstractNote": row.get("abstract_short", ""),
},
}
)
return command
def sync_writeback_queue(paths: dict[str, Path], candidate_map: dict[str, dict]) -> tuple[list[dict], int]:
existing_rows = read_jsonl(paths["queue"])
existing_by_candidate = {
str(row.get("source_candidate_id", "") or "").strip(): dict(row)
for row in existing_rows
if str(row.get("source_candidate_id", "") or "").strip()
}
queue_rows: list[dict] = []
queued_candidates: set[str] = set()
created = 0
for candidate_id, row in candidate_map.items():
decision = canonicalize_decision(row.get("decision", ""))
if decision != "纳入":
continue
if str(row.get("import_status", "") or "").strip() == "imported":
continue
candidate_copy = dict(row)
candidate_copy["final_collection"] = compute_final_collection(candidate_copy)
final_collection = str(candidate_copy.get("final_collection", "") or "").strip()
if not final_collection:
candidate_copy["import_status"] = "needs_collection_resolution"
candidate_map[candidate_id] = candidate_copy
continue
command = writeback_command_for_candidate(candidate_copy)
if not command:
candidate_copy["import_status"] = "blocked"
candidate_map[candidate_id] = candidate_copy
continue
existing = existing_by_candidate.get(candidate_id)
if existing and str(existing.get("status", "") or "").strip() in {"queued", "running", "processed"}:
merged = dict(existing)
merged["target_collection"] = command["target_collection"]
merged["target_domain"] = command.get("target_domain", merged.get("target_domain", ""))
if merged.get("status") != "processed":
merged["requested_at"] = command["requested_at"]
queue_rows.append(merged)
else:
queue_rows.append(command)
created += 1
candidate_copy["import_status"] = "queued_for_writeback"
candidate_map[candidate_id] = candidate_copy
queued_candidates.add(candidate_id)
for row in existing_rows:
candidate_id = str(row.get("source_candidate_id", "") or "").strip()
status = str(row.get("status", "") or "").strip()
if candidate_id in queued_candidates:
continue
if status == "processed":
queue_rows.append(row)
write_jsonl(paths["queue"], queue_rows)
return (queue_rows, created)
def load_bridge_config(paths: dict[str, Path]) -> dict:
config_path = paths["bridge_config"]
if not config_path.exists():
sample = read_json(paths["bridge_config_sample"])
write_json(config_path, sample)
return read_json(config_path)
def apply_writeback_log(paths: dict[str, Path], candidate_map: dict[str, dict]) -> int:
log_rows = read_jsonl(paths["log"])
changed = 0
latest_by_candidate: dict[str, dict] = {}
for row in log_rows:
candidate_id = str(row.get("source_candidate_id", "") or "").strip()
if candidate_id:
latest_by_candidate[candidate_id] = row
for candidate_id, log_row in latest_by_candidate.items():
candidate = dict(candidate_map.get(candidate_id, {}))
if not candidate:
continue
status = str(log_row.get("status", "") or "").strip()
if status == "success":
candidate["import_status"] = "imported"
candidate["zotero_key"] = str(log_row.get("zotero_key", "") or "").strip()
changed += 1
elif status == "error":
candidate["import_status"] = "writeback_error"
changed += 1
candidate_map[candidate_id] = candidate
return changed
def invoke_native_bridge(paths: dict[str, Path], max_commands: int = 5) -> dict:
config = load_bridge_config(paths)
base_url = str(config.get("server_base_url", "http://127.0.0.1:23119")).rstrip("/")
endpoint = str(config.get("process_endpoint", "/literaturePipeline/processQueue"))
payload = {
"queuePath": str(paths["queue"]).replace("\\", "/"),
"logPath": str(paths["log"]).replace("\\", "/"),
"configPath": str(paths["bridge_config"]).replace("\\", "/"),
"maxCommands": max_commands,
}
response = requests.post(f"{base_url}{endpoint}", json=payload, timeout=30)
response.raise_for_status()
result = response.json()
if not result.get("ok", False):
raise RuntimeError(result.get("error", "Native bridge returned non-ok result"))
return result
def normalize_candidate_title(text: str) -> str:
return re.sub("\\s+", " ", str(text or "").strip().lower())
def candidate_identity_keys(row: dict) -> dict[str, str]:
doi = str(row.get("doi", "") or "").strip().lower()
pmid = str(row.get("pmid", "") or "").strip()
title = normalize_candidate_title(row.get("title", ""))
return {"doi": doi, "pmid": pmid, "title": title}
def candidate_id_from_payload(row: dict) -> str:
source = str(row.get("source", "") or "candidate").strip().lower()
doi = re.sub("[^a-z0-9]+", "-", str(row.get("doi", "") or "").strip().lower()).strip("-")
pmid = re.sub("[^0-9]+", "", str(row.get("pmid", "") or "").strip())
fallback = re.sub("[^a-z0-9]+", "-", normalize_candidate_title(row.get("title", ""))).strip("-")[:80]
suffix = doi or pmid or fallback or datetime.now(timezone.utc).strftime("%Y%m%d%H%M%S")
return f"{source}-{suffix}"
def _normalize_candidate_value(value):
if value is None:
return ""
if isinstance(value, list):
return [str(item).strip() for item in value if str(item).strip()]
return str(value).strip()
def _authors_from_pubmed(value) -> list[str] | str:
if isinstance(value, list):
return [str(item).strip() for item in value if str(item).strip()]
return str(value or "").strip()
def _authors_from_openalex(value) -> list[str]:
authors = []
if isinstance(value, list):
for item in value:
if isinstance(item, dict):
author = item.get("author") or {}
name = author.get("display_name") or item.get("display_name") or ""
if name:
authors.append(str(name).strip())
elif str(item).strip():
authors.append(str(item).strip())
return authors
def _abstract_from_openalex(inverted_index) -> str:
if not isinstance(inverted_index, dict):
return ""
tokens: list[tuple[int, str]] = []
for word, positions in inverted_index.items():
if not isinstance(positions, list):
continue
for pos in positions:
try:
tokens.append((int(pos), str(word)))
except Exception:
continue
if not tokens:
return ""
tokens.sort(key=lambda item: item[0])
return " ".join((word for _, word in tokens))
def _authors_from_arxiv(value) -> list[str]:
authors = []
if isinstance(value, list):
for item in value:
if isinstance(item, dict):
name = item.get("name", "") or item.get("author", "")
if name:
authors.append(str(name).strip())
elif str(item).strip():
authors.append(str(item).strip())
return authors
def _authors_from_scholar(value) -> list[str]:
if isinstance(value, list):
return [str(item).strip() for item in value if str(item).strip()]
if isinstance(value, str):
parts = re.split("\\s*,\\s*|\\s+and\\s+", value)
return [part.strip() for part in parts if part.strip()]
return []
def adapt_pubmed_candidate(row: dict) -> dict:
payload = dict(row.get("payload") or {})
adapted = {
"candidate_id": row.get("candidate_id") or f"pubmed-{payload.get('pmid') or payload.get('PMID') or ''}",
"domain": row.get("domain", ""),
"title": payload.get("title", ""),
"authors": _authors_from_pubmed(payload.get("authors", [])),
"year": payload.get("year", ""),
"journal": payload.get("journal", ""),
"doi": payload.get("doi", ""),
"pmid": payload.get("pmid") or payload.get("PMID", ""),
"source": row.get("source", "") or "pubmed_search",
"requester_skill": row.get("requester_skill", ""),
"request_context": row.get("request_context", ""),
"abstract_short": payload.get("abstract", "") or payload.get("abstract_short", ""),
"decision": row.get("decision", ""),
"recommended_collection": row.get("recommended_collection", ""),
"recommend_confidence": row.get("recommend_confidence", ""),
"recommend_reason": row.get("recommend_reason", ""),
"user_collection": row.get("user_collection", ""),
"final_collection": row.get("final_collection", ""),
"duplicate_hint": row.get("duplicate_hint", ""),
"import_status": row.get("import_status", ""),
"note": row.get("note", ""),
"candidate_source_type": row.get("candidate_source_type", "") or "external_search",
"source_context": row.get("source_context", ""),
"task_relevance_reason": row.get("task_relevance_reason", ""),
}
return adapted
def adapt_openalex_candidate(row: dict) -> dict:
payload = dict(row.get("payload") or {})
primary_location = payload.get("primary_location") or {}
source_info = primary_location.get("source") or {}
openalex_id = str(payload.get("id", "") or "").rstrip("/").split("/")[-1]
adapted = {
"candidate_id": row.get("candidate_id") or f"openalex-{openalex_id}",
"domain": row.get("domain", ""),
"title": payload.get("display_name", "") or payload.get("title", ""),
"authors": _authors_from_openalex(payload.get("authorships", [])),
"year": payload.get("publication_year", "") or payload.get("year", ""),
"journal": source_info.get("display_name", "") or payload.get("journal", ""),
"doi": str(payload.get("doi", "") or "").replace("https://doi.org/", ""),
"pmid": "",
"source": row.get("source", "") or "openalex_search",
"requester_skill": row.get("requester_skill", ""),
"request_context": row.get("request_context", ""),
"abstract_short": _abstract_from_openalex(payload.get("abstract_inverted_index")),
"decision": row.get("decision", ""),
"recommended_collection": row.get("recommended_collection", ""),
"recommend_confidence": row.get("recommend_confidence", ""),
"recommend_reason": row.get("recommend_reason", ""),
"user_collection": row.get("user_collection", ""),
"final_collection": row.get("final_collection", ""),
"duplicate_hint": row.get("duplicate_hint", ""),
"import_status": row.get("import_status", ""),
"note": row.get("note", ""),
"candidate_source_type": row.get("candidate_source_type", "") or "external_search",
"source_context": row.get("source_context", ""),
"task_relevance_reason": row.get("task_relevance_reason", ""),
}
return adapted
def adapt_arxiv_candidate(row: dict) -> dict:
payload = dict(row.get("payload") or {})
arxiv_id = str(payload.get("id", "") or payload.get("entry_id", "")).rstrip("/").split("/")[-1]
adapted = {
"candidate_id": row.get("candidate_id") or f"arxiv-{arxiv_id}",
"domain": row.get("domain", ""),
"title": payload.get("title", ""),
"authors": _authors_from_arxiv(payload.get("authors", [])),
"year": str(payload.get("published", "") or "")[:4],
"journal": payload.get("journal_ref", "") or "arXiv",
"doi": payload.get("doi", ""),
"pmid": "",
"source": row.get("source", "") or "arxiv_search",
"requester_skill": row.get("requester_skill", ""),
"request_context": row.get("request_context", ""),
"abstract_short": payload.get("summary", "") or payload.get("abstract", ""),
"decision": row.get("decision", ""),
"recommended_collection": row.get("recommended_collection", ""),
"recommend_confidence": row.get("recommend_confidence", ""),
"recommend_reason": row.get("recommend_reason", ""),
"user_collection": row.get("user_collection", ""),
"final_collection": row.get("final_collection", ""),
"duplicate_hint": row.get("duplicate_hint", ""),
"import_status": row.get("import_status", ""),
"note": row.get("note", ""),
"candidate_source_type": row.get("candidate_source_type", "") or "external_search",
"source_context": row.get("source_context", ""),
"task_relevance_reason": row.get("task_relevance_reason", ""),
}
return adapted
def adapt_google_scholar_candidate(row: dict) -> dict:
payload = dict(row.get("payload") or {})
adapted = {
"candidate_id": row.get("candidate_id")
or f"google-scholar-{re.sub('[^a-z0-9]+', '-', normalize_candidate_title(payload.get('title', ''))).strip('-')[:80]}",
"domain": row.get("domain", ""),
"title": payload.get("title", ""),
"authors": _authors_from_scholar(payload.get("authors", [])),
"year": str(payload.get("year", "") or ""),
"journal": payload.get("journal", "") or payload.get("venue", ""),
"doi": payload.get("doi", ""),
"pmid": "",
"source": row.get("source", "") or "google_scholar_search",
"requester_skill": row.get("requester_skill", ""),
"request_context": row.get("request_context", ""),
"abstract_short": payload.get("abstract", "") or payload.get("snippet", ""),
"decision": row.get("decision", ""),
"recommended_collection": row.get("recommended_collection", ""),
"recommend_confidence": row.get("recommend_confidence", ""),
"recommend_reason": row.get("recommend_reason", ""),
"user_collection": row.get("user_collection", ""),
"final_collection": row.get("final_collection", ""),
"duplicate_hint": row.get("duplicate_hint", ""),
"import_status": row.get("import_status", ""),
"note": row.get("note", ""),
"candidate_source_type": row.get("candidate_source_type", "") or "external_search",
"source_context": row.get("source_context", ""),
"task_relevance_reason": row.get("task_relevance_reason", ""),
}
return adapted
def adapt_candidate_event(row: dict) -> dict:
adapter = str(row.get("adapter", "") or "").strip()
if adapter == "pubmed_search":
return adapt_pubmed_candidate(row)
if adapter == "openalex_search":
return adapt_openalex_candidate(row)
if adapter == "arxiv_search":
return adapt_arxiv_candidate(row)
if adapter == "google_scholar_search":
return adapt_google_scholar_candidate(row)
return dict(row)
def default_user_agent() -> str:
return "ResearchLiteraturePipeline/1.0 (+local-vault)"
def build_search_event(base_task: dict, source: str, payload: dict) -> dict:
return {
"adapter": source,
"payload": payload,
"source": source,
"requester_skill": base_task.get("requester_skill", ""),
"request_context": base_task.get("request_context", ""),
"domain": base_task.get("domain", ""),
"recommended_collection": base_task.get("recommended_collection", ""),
"recommend_confidence": base_task.get("recommend_confidence", ""),
"recommend_reason": base_task.get("recommend_reason", "") or f"{source} 检索命中",
"candidate_source_type": base_task.get("candidate_source_type", "") or "external_search",
"task_relevance_reason": base_task.get("task_relevance_reason", ""),
"note": base_task.get("note", ""),
}
def _pubmed_abstract_and_doi_map(xml_text: str) -> dict[str, dict]:
root = ET.fromstring(xml_text)
result = {}
for article in root.findall(".//PubmedArticle"):
pmid = (article.findtext(".//MedlineCitation/PMID") or "").strip()
if not pmid:
continue
abstract_parts = []
for node in article.findall(".//Abstract/AbstractText"):
label = (node.attrib.get("Label") or "").strip()
text = " ".join("".join(node.itertext()).split())
if not text:
continue
abstract_parts.append(f"{label}: {text}" if label else text)
doi = ""
for id_node in article.findall(".//PubmedData/ArticleIdList/ArticleId"):
if (id_node.attrib.get("IdType") or "").lower() == "doi":
doi = "".join(id_node.itertext()).strip()
if doi:
break
if not doi:
for id_node in article.findall(".//ELocationID"):
if (id_node.attrib.get("EIdType") or "").lower() == "doi":
doi = "".join(id_node.itertext()).strip()
if doi:
break
result[pmid] = {"abstract": "\n".join(abstract_parts).strip(), "doi": doi}
return result
def search_pubmed(task: dict, limit: int) -> tuple[list[dict], dict]:
query = str(task.get("query", "") or "").strip()
if not query:
return ([], {"count": 0, "ids": []})
base_url = "https://eutils.ncbi.nlm.nih.gov/entrez/eutils"
headers = {"User-Agent": os.environ.get("LIT_PIPELINE_USER_AGENT", default_user_agent())}
common_params = {}
email = os.environ.get("NCBI_EMAIL", "").strip()
api_key = os.environ.get("NCBI_API_KEY", "").strip()
if email:
common_params["email"] = email
if api_key:
common_params["api_key"] = api_key
esearch = requests.get(
f"{base_url}/esearch.fcgi",
params={
"db": "pubmed",
"retmode": "json",
"sort": "relevance",
"term": query,
"retmax": limit,
**common_params,
},
headers=headers,
timeout=60,
)
esearch.raise_for_status()
search_payload = esearch.json().get("esearchresult", {})
ids = [str(item).strip() for item in search_payload.get("idlist", []) if str(item).strip()]
if not ids:
return ([], {"count": 0, "ids": []})
id_text = ",".join(ids)
esummary = requests.get(
f"{base_url}/esummary.fcgi",
params={"db": "pubmed", "retmode": "json", "id": id_text, **common_params},
headers=headers,
timeout=60,
)
esummary.raise_for_status()
summary_payload = esummary.json().get("result", {})
efetch = requests.get(
f"{base_url}/efetch.fcgi",
params={"db": "pubmed", "retmode": "xml", "id": id_text, **common_params},
headers=headers,
timeout=90,
)
efetch.raise_for_status()
abstract_map = _pubmed_abstract_and_doi_map(efetch.text)
rows = []
for pmid in ids:
item = summary_payload.get(pmid, {}) or {}
title = html.unescape(str(item.get("title", "") or "")).strip()
if not title:
continue
article_ids = item.get("articleids", []) or []
doi = ""
for article_id in article_ids:
if str(article_id.get("idtype", "")).lower() == "doi":
doi = str(article_id.get("value", "")).strip()
if doi:
break
if not doi:
doi = abstract_map.get(pmid, {}).get("doi", "")
rows.append(
{
"pmid": pmid,
"title": title,
"authors": [author.get("name", "") for author in item.get("authors") or [] if author.get("name")],
"year": _extract_year(str(item.get("pubdate", "") or "")),
"journal": item.get("fulljournalname", "") or item.get("source", ""),
"doi": doi,
"abstract": abstract_map.get(pmid, {}).get("abstract", ""),
}
)
return (rows, {"count": len(rows), "ids": ids})
def search_openalex(task: dict, limit: int) -> tuple[list[dict], dict]:
query = str(task.get("query", "") or "").strip()
if not query:
return ([], {"count": 0})
headers = {"User-Agent": os.environ.get("LIT_PIPELINE_USER_AGENT", default_user_agent())}
params = {"search": query, "per-page": limit}
api_key = os.environ.get("OPENALEX_API_KEY", "").strip()
if api_key:
params["api_key"] = api_key
mailto = os.environ.get("OPENALEX_MAILTO", "").strip()
if mailto:
params["mailto"] = mailto
response = requests.get("https://api.openalex.org/works", params=params, headers=headers, timeout=60)
response.raise_for_status()
payload = response.json()
results = payload.get("results", []) or []
return (results, {"count": len(results), "meta": payload.get("meta", {})})
def search_arxiv(task: dict, limit: int) -> tuple[list[dict], dict]:
query = str(task.get("query", "") or "").strip()
if not query:
return ([], {"count": 0})
encoded_query = urllib.parse.quote(f"all:{query}")
url = f"https://export.arxiv.org/api/query?search_query={encoded_query}&start=0&max_results={limit}&sortBy=relevance&sortOrder=descending"
headers = {"User-Agent": os.environ.get("LIT_PIPELINE_USER_AGENT", default_user_agent())}
response = requests.get(url, headers=headers, timeout=60)
response.raise_for_status()
ns = {"atom": "http://www.w3.org/2005/Atom", "arxiv": "http://arxiv.org/schemas/atom"}
root = ET.fromstring(response.text)
entries = []
for entry in root.findall("atom:entry", ns):
entries.append(
{
"id": (entry.findtext("atom:id", default="", namespaces=ns) or "").strip(),
"title": " ".join((entry.findtext("atom:title", default="", namespaces=ns) or "").split()),
"summary": " ".join((entry.findtext("atom:summary", default="", namespaces=ns) or "").split()),
"published": (entry.findtext("atom:published", default="", namespaces=ns) or "").strip(),
"authors": [
{"name": (node.findtext("atom:name", default="", namespaces=ns) or "").strip()}
for node in entry.findall("atom:author", ns)
],
"doi": (entry.findtext("arxiv:doi", default="", namespaces=ns) or "").strip(),
"journal_ref": (entry.findtext("arxiv:journal_ref", default="", namespaces=ns) or "").strip(),
}
)
return (entries, {"count": len(entries)})
def _coerce_source_name(value: str) -> str:
text = str(value or "").strip().lower()
aliases = {
"pubmed": "pubmed_search",
"pubmed_search": "pubmed_search",
"openalex": "openalex_search",
"openalex_search": "openalex_search",
"arxiv": "arxiv_search",
"arxiv_search": "arxiv_search",
}
return aliases.get(text, text)
def run_search_command(vault: Path, args) -> int:
paths = pipeline_paths(vault)
paths["search_tasks"].mkdir(parents=True, exist_ok=True)
task_id = f"search-{datetime.now().strftime('%Y%m%d-%H%M%S')}"
sources = args.sources or ["pubmed_search", "openalex_search", "arxiv_search"]
task = {
"task_id": task_id,
"query": args.query,
"domain": args.domain,
"recommended_collection": args.recommended_collection or "",
"requester_skill": args.requester_skill or "",
"request_context": args.request_context or "",
"sources": sources,
"limit": args.limit,
"recommend_reason": args.recommend_reason or "围绕检索主题补充候选文献",
"candidate_source_type": "external_search",
}
task_path = paths["search_tasks"] / f"{task_id}.json"
write_json(task_path, task)
print(f"search: task written -> {task_path}")
code = run_search_sources(vault)
if code:
return code
if not args.skip_ingest:
return run_ingest_candidates(vault)
return 0
def normalize_candidate_payload(row: dict) -> dict:
normalized = {key: _normalize_candidate_value(value) for key, value in adapt_candidate_event(row).items()}
normalized["candidate_id"] = str(normalized.get("candidate_id", "") or "").strip() or candidate_id_from_payload(
normalized
)
normalized["title"] = str(normalized.get("title", "") or "").strip()
normalized["domain"] = str(normalized.get("domain", "") or "").strip()
normalized["source"] = str(normalized.get("source", "") or "").strip() or "candidate_ingest"
normalized["candidate_source_type"] = (
str(normalized.get("candidate_source_type", "") or "").strip() or normalized["source"]
)
normalized["decision"] = canonicalize_decision(normalized.get("decision", ""))
normalized["import_status"] = str(normalized.get("import_status", "") or "").strip() or "pending"
normalized["recommend_confidence"] = str(normalized.get("recommend_confidence", "") or "").strip() or "0"
normalized["status"] = "candidate"
return normalized
def merge_candidate_record(existing: dict | None, incoming: dict) -> dict:
merged = dict(existing or {})
preserve_if_existing = {"decision", "user_collection", "note", "import_status"}
for key, value in incoming.items():
if existing and key in preserve_if_existing:
current = merged.get(key, "")
if str(current).strip():
continue
merged[key] = value
merged["decision"] = canonicalize_decision(merged.get("decision", ""))
merged["final_collection"] = compute_final_collection(merged)
return merged
def resolve_existing_candidate(candidate_map: dict[str, dict], incoming: dict) -> tuple[str | None, dict | None]:
candidate_id = incoming.get("candidate_id", "")
if candidate_id and candidate_id in candidate_map:
return (candidate_id, candidate_map[candidate_id])
incoming_keys = candidate_identity_keys(incoming)
for existing_id, existing in candidate_map.items():
existing_keys = candidate_identity_keys(existing)
if incoming_keys["doi"] and incoming_keys["doi"] == existing_keys["doi"]:
return (existing_id, existing)
if incoming_keys["pmid"] and incoming_keys["pmid"] == existing_keys["pmid"]:
return (existing_id, existing)
if incoming_keys["title"] and incoming_keys["title"] == existing_keys["title"]:
return (existing_id, existing)
return (None, None)
def _harvest_csv_paths(paths: dict[str, Path]) -> list[Path]:
root = paths["harvest_root"]
if not root.exists():
return []
return sorted(root.rglob("*-05-reference-harvest-candidates.csv"))
def _normalize_harvest_value(value: str) -> str:
text = str(value or "").strip()
if text == "<inherit-from-task-context>":
return ""
return text
def _normalize_harvest_row(row: dict[str, str]) -> dict:
normalized = normalize_candidate_payload({key: _normalize_harvest_value(value) for key, value in row.items()})
normalized["source"] = normalized.get("source", "") or "reference_harvest"
normalized["candidate_source_type"] = normalized.get("candidate_source_type", "") or "reference_harvest"
return normalized
def _merge_harvest_candidate(existing: dict | None, incoming: dict) -> dict:
return merge_candidate_record(existing, incoming)
def next_key(domain: str, export_rows: list[dict]) -> str:
prefix = "ORTHO" if domain == "骨科" else "SPORT"
existing = [row.get("key", "") for row in export_rows]
max_num = 0
for key in existing:
if key.startswith(prefix):
suffix = key[len(prefix) :]
if suffix.isdigit():
max_num = max(max_num, int(suffix))
return f"{prefix}{max_num + 1:03d}"
def frontmatter_note(entry: dict, existing_text: str = "", preserved_tags: list[str] | None = None, preserved_ocr_redo: bool | None = None) -> str:
preserved_deep = extract_preserved_deep_reading(existing_text)
if preserved_tags is None and existing_text:
preserved_tags = extract_preserved_tags(existing_text)
if preserved_ocr_redo is None and existing_text:
preserved_ocr_redo = extract_preserved_ocr_redo(existing_text)
first_author = entry.get("first_author", "")
if not first_author:
authors = entry.get("authors", [])
first_author = authors[0] if authors else ""
lines = [
"---",
f"title: {yaml_quote(entry.get('title', ''))}",
f"aliases: [{yaml_quote(entry.get('title', ''))}, {yaml_quote(entry.get('citation_key', ''))}]",
f"year: {entry.get('year', '')}",
f"journal: {yaml_quote(entry.get('journal', ''))}",
f"first_author: {yaml_quote(first_author)}",
f"zotero_key: {yaml_quote(entry.get('zotero_key', ''))}",
f"citation_key: {yaml_quote(entry.get('citation_key', ''))}",
f"domain: {yaml_quote(entry.get('domain', ''))}",
f"doi: {yaml_quote(entry.get('doi', ''))}",
f"pmid: {yaml_quote(entry.get('pmid', ''))}",
]
lines.extend(yaml_list("collection_path", entry.get("collections", [])))
lines.extend(yaml_list("collection_tags", entry.get("collection_tags", [])))
lines.extend([
f"impact_factor: {yaml_quote(entry.get('impact_factor', ''))}",
])
lines.extend(yaml_block(entry.get("abstract", "")))
lines.extend(
[
f"has_pdf: {('true' if entry.get('has_pdf') else 'false')}",
f"do_ocr: {('true' if entry.get('do_ocr') else 'false')}",
f"analyze: {('true' if entry.get('analyze') else 'false')}",
f"ocr_status: {yaml_quote(entry.get('ocr_status', 'pending'))}",
f"deep_reading_status: {yaml_quote(entry.get('deep_reading_status', 'pending'))}",
f"pdf_path: {yaml_quote(entry.get('pdf_path', ''))}",
f"fulltext_md_path: {yaml_quote('[[{}]]'.format(entry['fulltext_path']) if entry.get('ocr_status') == 'done' and entry.get('fulltext_path') else '')}",
]
)
if preserved_ocr_redo:
lines.append("ocr_redo: true")
if preserved_tags is not None:
lines.append("tags:")
for tag in preserved_tags:
lines.append(f" - {tag}")
else:
lines.extend(
[
"tags:",
" - 文献阅读",
f" - {entry.get('domain', '')}",
]
)
lines.extend(
[
"---",
"",
f"# {entry.get('title', '')}",
"",
"## 📄 文献基本信息",
"",
f"- Zotero Key: `{entry.get('zotero_key', '')}`",
f"- Collection: {', '.join(entry.get('collections', [entry.get('collection_path', '')]))}",
f"- 作者:{', '.join(entry.get('authors', []))}",
"",
"## 摘要",
"",
entry.get("abstract", "") or "暂无摘要",
"",
]
)
if preserved_deep:
lines.extend(["", preserved_deep, ""])
return "\n".join(lines)
def analyze_selected_keys(paths: dict[str, Path]) -> set[str]:
return {key for key, row in load_control_actions(paths).items() if row.get("analyze")}
def migrate_to_workspace(vault: Path, paths: dict) -> int:
"""Migrate flat literature notes into paper workspace directories.
Copies each flat note at <literature_dir>/<domain>/<key> - <Title>.md into:
<literature_dir>/<domain>/<key> - <Title>/<key> - <Title>.md
Preserves any existing ## 🔍 精读 section inside the main note (single source of truth).
Creates:
<literature_dir>/<domain>/<key> - <Title>/ai/ (empty directory)
Idempotent: skips papers whose workspace directory already exists.
The original flat note is preserved (copy-not-move per D-12).
Returns: number of papers migrated (0 means all are already workspace).
"""
index = asset_index.read_index(vault)
items = index.get("items", []) if isinstance(index, dict) else []
indexed_entries: dict[str, dict] = {
str(entry.get("zotero_key", "") or "").strip(): entry
for entry in items
if str(entry.get("zotero_key", "") or "").strip()
}
flat_notes: list[tuple[Path, str, str]] = []
lit_root = paths.get("literature")
if lit_root and lit_root.exists():
for note_path in lit_root.rglob("*.md"):
if note_path.name in ("fulltext.md", "deep-reading.md", "discussion.md"):
continue
relative = note_path.relative_to(lit_root)
if len(relative.parts) != 2:
continue
domain = relative.parts[0]
try:
text = note_path.read_text(encoding="utf-8")
except Exception:
continue
key_match = re.search(r'^zotero_key:\s*"?(\S+?)"?\s*$', text, re.MULTILINE)
zotero_key = key_match.group(1) if key_match else ""
if not zotero_key:
stem = note_path.stem
if " - " in stem:
zotero_key = stem.split(" - ", 1)[0].strip()
else:
zotero_key = stem
if not zotero_key:
continue
flat_notes.append((note_path, domain, zotero_key))
if not indexed_entries and not flat_notes:
return 0
migrated = 0
for flat_note_path, domain, key in flat_notes:
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 = title_match.group(1).strip() if title_match else ""
stem = flat_note_path.stem
if title:
title_slug = slugify_filename(title)
elif " - " in stem:
title_slug = stem.split(" - ", 1)[1].strip()
else:
title_slug = stem.strip()
workspace_dir = paths["literature"] / domain / f"{key} - {title_slug}"
main_note_path = workspace_dir / f"{key}.md"
# Self-healing: old workspaces have {key} - {title}.md inside, rename to {key}.md
if not main_note_path.exists():
for old_candidate in workspace_dir.glob(f"{key} - *.md"):
old_candidate.rename(main_note_path)
break
legacy_flags = _legacy_control_flags(paths, key)
if workspace_dir.exists():
if main_note_path.exists() and any(value is not None for value in legacy_flags.values()):
try:
workspace_content = main_note_path.read_text(encoding="utf-8")
flat_content = flat_note_path.read_text(encoding="utf-8")
updated_content = workspace_content
for field in ("do_ocr", "analyze"):
legacy_value = legacy_flags[field]
if legacy_value is None:
continue
workspace_value = _read_frontmatter_optional_bool_from_text(workspace_content, field)
flat_value = _read_frontmatter_optional_bool_from_text(flat_content, field)
if workspace_value is None or workspace_value == flat_value:
updated_content = update_frontmatter_field(
updated_content,
field,
"true" if legacy_value else "false",
)
if updated_content != workspace_content:
main_note_path.write_text(updated_content, encoding="utf-8")
except Exception:
pass
continue
# Create workspace directory
workspace_dir.mkdir(parents=True, exist_ok=True)
# Read flat note content
content = flat_note_path.read_text(encoding="utf-8")
if legacy_flags["do_ocr"] is not None:
content = update_frontmatter_field(content, "do_ocr", "true" if legacy_flags["do_ocr"] else "false")
if legacy_flags["analyze"] is not None:
content = update_frontmatter_field(content, "analyze", "true" if legacy_flags["analyze"] else "false")
# Write main note to workspace (copy of flat note)
main_note_path.write_text(content, encoding="utf-8")
# Create ai/ directory
ai_dir = workspace_dir / "ai"
ai_dir.mkdir(exist_ok=True)
# WS-04: Bridge OCR fulltext to workspace if available
meta_path = paths.get("ocr", Path()) / key / "meta.json"
if meta_path.exists():
try:
meta = read_json(meta_path)
if meta.get("ocr_status") == "done":
source_fulltext = meta_path.parent / "fulltext.md"
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
migrated += 1
if migrated > 0:
print(f"migrate_to_workspace: migrated {migrated} paper(s) to workspace structure")
return migrated
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
selection-sync may affect many papers. Workers that modify individual
papers (ocr, deep-reading, repair) use asset_index.refresh_index_entry()
for incremental refresh by key.
Args:
vault: Path to the vault root.
verbose: If True, print detailed progress.
rebuild_index: If True, force full rebuild (default: True for sync).
json_output: If True, suppress human-readable print output.
"""
from paperforge.worker.base_views import ensure_base_views
paths = pipeline_paths(vault)
config = load_domain_config(paths)
ensure_base_views(vault, paths, config)
domain_lookup = {entry["export_file"]: entry["domain"] for entry in config["domains"]}
cfg = load_vault_config(vault)
zotero_dir = vault / cfg.get("system_dir", "System") / "Zotero"
exports = {}
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}
# Migrate flat notes to workspace before build_index (D-11: first-sync migration)
migrate_to_workspace(vault, paths)
# Delegate to asset_index.build_index() for the core build loop
count = asset_index.build_index(vault, verbose)
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 {"updated": count, "failed": 0, "errors": []}