From 355839fb57927884603b107cde341789ac225de3 Mon Sep 17 00:00:00 2001 From: Oded Falik Date: Wed, 23 Sep 2026 11:48:23 -0700 Subject: [PATCH] Keep PDF processing outside MCP request loop --- CHANGELOG.md | 9 + README.md | 15 +- paper_intelligence/processing_worker.py | 26 +++ paper_intelligence/server.py | 6 +- paper_intelligence/tools/convert.py | 10 +- paper_intelligence/tools/search.py | 227 ++++++++++++++++---- paper_intelligence/utils/chromadb_client.py | 12 +- tests/test_integration.py | 7 +- tests/test_reliability.py | 143 ++++++++++++ 9 files changed, 399 insertions(+), 56 deletions(-) create mode 100644 paper_intelligence/processing_worker.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 3d83a63..a5423a3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] +### Changed +- First-use processing now runs in a detached subprocess so Marker/Torch startup cannot block the MCP request loop; status responses include an ETA and retry interval. +- Hybrid collection searches remain bounded: grep covers the collection, while semantic fanout above 20 papers returns an explicit partial status instead of risking a client timeout. +- Scoped semantic searches reuse one embedding model across paper databases. + +### Fixed +- Direct PDF lookup and status checks no longer import Torch just to derive an output path. +- Direct PDF lookup never scans sibling directories, and explicitly requested library scans skip unreadable children. + ## [0.5.2] - 2026-08-28 ### Fixed diff --git a/README.md b/README.md index 892fed9..1c90827 100644 --- a/README.md +++ b/README.md @@ -134,11 +134,16 @@ many MCP clients. The first call returns promptly: } ``` -Processing continues after that response. Poll `get_paper_info` using `paper_dir`; when -it reports `status: "ready"`, retry the original search. Already-processed grep searches -do not initialize the semantic model. RAG and hybrid searches initialize it when needed. -Search no longer performs an unconditional remote-library sync, which previously allowed -a local query to block for up to five minutes. +Processing runs in a detached process and continues after that response without blocking +status calls. Poll `get_paper_info` using `paper_dir`; when it reports `status: "ready"`, +retry the original search. Already-processed grep searches do not initialize the semantic +model. RAG and hybrid searches initialize one shared model per call when needed. + +To keep collection-wide calls below normal MCP deadlines, semantic search accepts at most +20 paper directories per call. A broader `hybrid` call returns bounded grep results with +`status: "partial"` and an explicit `semantic_search.status: "skipped_source_limit"`; +select up to 20 papers for semantic retrieval. Search no longer performs an unconditional +remote-library sync, which previously allowed a local query to block for up to five minutes. ### `get_paper_info` diff --git a/paper_intelligence/processing_worker.py b/paper_intelligence/processing_worker.py new file mode 100644 index 0000000..83cc489 --- /dev/null +++ b/paper_intelligence/processing_worker.py @@ -0,0 +1,26 @@ +"""Detached worker for first-use PDF conversion and indexing.""" + +import argparse +from pathlib import Path + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--source", required=True) + parser.add_argument("--paper-dir", required=True) + parser.add_argument("--use-llm", action="store_true") + args = parser.parse_args() + + # Import only in the child. Importing the processing stack in the MCP server can + # make the request loop unresponsive while Marker, Torch, and embedding models load. + from .tools.search import _run_processing_job + + _run_processing_job( + Path(args.source), + Path(args.paper_dir), + use_llm=args.use_llm, + ) + + +if __name__ == "__main__": + main() diff --git a/paper_intelligence/server.py b/paper_intelligence/server.py index dbe1a76..f2cd9e7 100644 --- a/paper_intelligence/server.py +++ b/paper_intelligence/server.py @@ -17,7 +17,8 @@ "Pass PDF paths directly to search. First use queues 1-3 minute background " "processing and returns status='processing'; call get_paper_info on the returned " "paper_dir until status='ready', then retry search. Searches inspect only the " - "sources explicitly supplied." + "sources explicitly supplied. Semantic search is limited to 20 papers per call; " + "broader hybrid searches return collection-wide grep results and a partial status." ), ) @@ -49,6 +50,9 @@ def search( include_context: Include surrounding lines in results use_llm: Use LLM for better PDF conversion (slower) + Semantic retrieval is limited to 20 ready papers per call. Broader hybrid + searches return bounded grep results plus an explicit partial status. + Returns: Search results with content, location, and relevance scores """ diff --git a/paper_intelligence/tools/convert.py b/paper_intelligence/tools/convert.py index e511cb6..f3c656d 100644 --- a/paper_intelligence/tools/convert.py +++ b/paper_intelligence/tools/convert.py @@ -5,10 +5,14 @@ from pathlib import Path from typing import Optional -# Set device preference for Apple Silicon MPS -if not os.environ.get("TORCH_DEVICE"): + +def _set_device_preference() -> None: + """Choose an accelerator only inside the detached conversion worker.""" + if os.environ.get("TORCH_DEVICE"): + return try: import torch + if torch.backends.mps.is_available(): os.environ["TORCH_DEVICE"] = "mps" elif torch.cuda.is_available(): @@ -87,6 +91,8 @@ def convert_pdf( - message: Status message - images_dir: Path to extracted images (if any) """ + _set_device_preference() + from marker.converters.pdf import PdfConverter from marker.models import create_model_dict from marker.output import text_from_rendered diff --git a/paper_intelligence/tools/search.py b/paper_intelligence/tools/search.py index 05365ee..4e1a94a 100644 --- a/paper_intelligence/tools/search.py +++ b/paper_intelligence/tools/search.py @@ -8,8 +8,9 @@ import os import re import sqlite3 +import subprocess +import sys import threading -from concurrent.futures import Future, ThreadPoolExecutor from datetime import datetime, timezone from pathlib import Path from typing import Literal, Optional @@ -18,8 +19,10 @@ from ..utils.markdown_parser import MarkdownParser PROCESSING_STATE_FILENAME = ".paper-intelligence-processing.json" -_PROCESSING_EXECUTOR = ThreadPoolExecutor(max_workers=1, thread_name_prefix="paper-processing") -_PROCESSING_JOBS: dict[str, Future] = {} +PROCESSING_RETRY_AFTER_SECONDS = 30 +PROCESSING_ETA = "usually 1-3 minutes" +MAX_SEMANTIC_PAPERS = 20 +_PROCESSING_JOBS: dict[str, subprocess.Popen] = {} _PROCESSING_LOCK = threading.Lock() @@ -134,7 +137,9 @@ def _read_processing_state(paper_dir: Path) -> Optional[dict]: def _write_processing_state(paper_dir: Path, state: dict) -> None: paper_dir.mkdir(parents=True, exist_ok=True) state_path = _processing_state_path(paper_dir) - temporary_path = state_path.with_suffix(f"{state_path.suffix}.tmp") + temporary_path = state_path.with_name( + f".{state_path.name}.{os.getpid()}.{threading.get_ident()}.tmp" + ) temporary_path.write_text(json.dumps(state, indent=2), encoding="utf-8") temporary_path.replace(state_path) @@ -154,7 +159,9 @@ def _run_processing_job(source: Path, paper_dir: Path, use_llm: bool) -> None: "paper_dir": str(paper_dir), "started_at": started_at, "worker_pid": os.getpid(), - "message": "PDF processing is running in the background.", + "estimated_completion": PROCESSING_ETA, + "retry_after_seconds": PROCESSING_RETRY_AFTER_SECONDS, + "message": "PDF processing is running in a detached background process.", }) try: @@ -186,38 +193,113 @@ def _run_processing_job(source: Path, paper_dir: Path, use_llm: bool) -> None: }) -def _schedule_processing(source: Path, paper_dir: Path, use_llm: bool) -> dict: - """Start or resume one background processing job and return observable status.""" - key = str(paper_dir) - with _PROCESSING_LOCK: - future = _PROCESSING_JOBS.get(key) - if future is None or future.done(): - queued_at = datetime.now(timezone.utc).isoformat() - _write_processing_state(paper_dir, { - "status": "queued", - "source": str(source), - "paper_dir": key, - "queued_at": queued_at, - "message": "PDF processing is queued and will continue after this call returns.", - }) - _PROCESSING_JOBS[key] = _PROCESSING_EXECUTOR.submit( - _run_processing_job, source, paper_dir, use_llm - ) +def _pid_is_running(pid: object) -> bool: + try: + os.kill(int(pid), 0) + except (OSError, TypeError, ValueError): + return False + return True + + +def _processing_is_active(key: str, state: Optional[dict]) -> bool: + process = _PROCESSING_JOBS.get(key) + if process is not None and process.poll() is None: + return True + return bool( + state + and state.get("status") in {"queued", "processing"} + and _pid_is_running(state.get("worker_pid")) + ) + + +def _start_processing_process(source: Path, paper_dir: Path, use_llm: bool) -> subprocess.Popen: + """Launch heavy PDF work outside the MCP process. + + Marker and Torch initialization can monopolize an in-process worker long enough + that even the original tool response and status calls miss the MCP deadline. + A detached interpreter keeps the MCP request loop responsive and lets processing + survive a client timeout or server reconnect. + """ + command = [ + sys.executable, + "-m", + "paper_intelligence.processing_worker", + "--source", + str(source), + "--paper-dir", + str(paper_dir), + ] + if use_llm: + command.append("--use-llm") + return subprocess.Popen( + command, + stdin=subprocess.DEVNULL, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + close_fds=True, + start_new_session=True, + ) - state = _read_processing_state(paper_dir) or {} + +def _status_response(source: Path, paper_dir: Path, state: dict) -> dict: + key = str(paper_dir) return { **state, "status": state.get("status", "queued"), "source": str(source), "paper_dir": key, - "retry_after_seconds": 30, + "estimated_completion": PROCESSING_ETA, + "retry_after_seconds": PROCESSING_RETRY_AFTER_SECONDS, "next_step": ( - f"Call get_paper_info with paper_dir={key!r}; when status is 'ready', " + f"Call get_paper_info with paper_dir={key!r} in " + f"{PROCESSING_RETRY_AFTER_SECONDS} seconds; when status is 'ready', " "retry search with the original source." ), } +def _schedule_processing(source: Path, paper_dir: Path, use_llm: bool) -> dict: + """Start or resume one detached processing job and return immediately.""" + key = str(paper_dir) + with _PROCESSING_LOCK: + state = _read_processing_state(paper_dir) + if _processing_is_active(key, state): + return _status_response(source, paper_dir, state or {}) + + queued_state = { + "status": "queued", + "source": str(source), + "paper_dir": key, + "queued_at": datetime.now(timezone.utc).isoformat(), + "estimated_completion": PROCESSING_ETA, + "retry_after_seconds": PROCESSING_RETRY_AFTER_SECONDS, + "message": ( + "PDF processing has started in a detached background process and " + "will continue after this call returns." + ), + } + _write_processing_state(paper_dir, queued_state) + try: + process = _start_processing_process(source, paper_dir, use_llm) + except OSError as exc: + failed_state = { + **queued_state, + "status": "failed", + "message": f"Could not start background processing: {exc}", + } + _write_processing_state(paper_dir, failed_state) + return _status_response(source, paper_dir, failed_state) + _PROCESSING_JOBS[key] = process + queued_state["worker_pid"] = process.pid + persisted_state = _read_processing_state(paper_dir) + if not persisted_state or persisted_state.get("status") == "queued": + _write_processing_state(paper_dir, queued_state) + else: + # The child may already have advanced the state to "processing". + queued_state = persisted_state + return _status_response(source, paper_dir, queued_state) + + def _find_paper_dirs( search_paths: list[str], auto_process: bool = True, @@ -263,10 +345,23 @@ def _find_paper_dirs( # Directory discovery is limited to the directory explicitly supplied by the # caller and its immediate paper children. if path.is_dir(): - candidates = [path] if (path / "paper.md").exists() else [ - subdir for subdir in path.iterdir() - if subdir.is_dir() and (subdir / "paper.md").exists() - ] + if (path / "paper.md").exists(): + candidates = [path] + else: + candidates = [] + try: + children = path.iterdir() + for subdir in children: + try: + if subdir.is_dir() and (subdir / "paper.md").exists(): + candidates.append(subdir) + except OSError: + # An explicitly selected library may contain unreadable + # siblings; one sibling must not abort discovery of others. + continue + except OSError: + pass + for candidate in candidates: if _is_fully_processed(candidate) or not auto_process: paper_dirs.append(candidate) @@ -283,6 +378,7 @@ def grep_search( paper_dirs: list[Path], case_sensitive: bool = False, regex: bool = False, + max_results: Optional[int] = None, ) -> list[dict]: """Perform grep-style text search across paper markdown files. @@ -359,6 +455,8 @@ def grep_search( "header_context": header_context, "match_type": "grep", }) + if max_results is not None and len(results) >= max_results: + return results except Exception: continue @@ -366,6 +464,12 @@ def grep_search( return results +def _get_rag_client_class(): + from ..utils.chromadb_client import RAGClient + + return RAGClient + + def rag_search( query: str, paper_dirs: list[Path], @@ -381,9 +485,10 @@ def rag_search( Returns: List of matches with content, score, and metadata """ - from ..utils.chromadb_client import RAGClient + RAGClient = _get_rag_client_class() results = [] + embed_model = None for paper_dir in paper_dirs: chroma_dir = paper_dir / "chroma" @@ -392,7 +497,11 @@ def rag_search( continue try: - rag_client = RAGClient(persist_directory=chroma_dir) + rag_client = RAGClient( + persist_directory=chroma_dir, + embed_model=embed_model, + ) + embed_model = rag_client.embed_model raw_results = rag_client.query("paper", query, top_k) for r in raw_results: @@ -487,16 +596,30 @@ def search( paper_dirs=dirs, case_sensitive=case_sensitive, regex=regex, + max_results=top_k * 2 if mode == "hybrid" else top_k, ) results.extend(grep_results) + semantic_status = None if mode in ("rag", "hybrid"): - rag_results = rag_search( - query=query, - paper_dirs=dirs, - top_k=top_k, - ) - results.extend(rag_results) + if len(dirs) > MAX_SEMANTIC_PAPERS: + semantic_status = { + "status": "skipped_source_limit", + "papers_requested": len(dirs), + "max_papers": MAX_SEMANTIC_PAPERS, + "message": ( + "Semantic search was skipped to keep this MCP call bounded. " + f"Select at most {MAX_SEMANTIC_PAPERS} paper directories for " + "RAG, or use grep mode for the whole collection." + ), + } + else: + rag_results = rag_search( + query=query, + paper_dirs=dirs, + top_k=top_k, + ) + results.extend(rag_results) # Deduplicate by content similarity if mode == "hybrid": @@ -530,13 +653,21 @@ def search( "mode": mode, "success": True, } - if processing_jobs: + if semantic_status: + result["semantic_search"] = semantic_status + if processing_jobs or semantic_status: result["status"] = "partial" - result["processing"] = processing_jobs - result["message"] = ( - "Searched ready sources; remaining first-use processing continues in " - "the background." - ) + if processing_jobs: + result["processing"] = processing_jobs + messages = [] + if processing_jobs: + messages.append( + "Searched ready sources; remaining first-use processing continues " + "in the background." + ) + if semantic_status: + messages.append(semantic_status["message"]) + result["message"] = " ".join(messages) else: result["status"] = "ready" return result @@ -627,6 +758,16 @@ def get_paper_info(paper_dir: str) -> dict: if processing and processing.get("status") in {"queued", "processing", "failed"}: info["processing"] = processing info["status"] = processing["status"] + if ( + info["status"] in {"queued", "processing"} + and processing.get("worker_pid") + and not _processing_is_active(str(path), processing) + ): + info["status"] = "stalled" + info["message"] = ( + "The background worker stopped before processing completed. Retry " + "search to start a new worker." + ) elif ( info["has_markdown"] and info["has_index"] diff --git a/paper_intelligence/utils/chromadb_client.py b/paper_intelligence/utils/chromadb_client.py index 1652577..9ddff1f 100644 --- a/paper_intelligence/utils/chromadb_client.py +++ b/paper_intelligence/utils/chromadb_client.py @@ -46,6 +46,7 @@ def __init__( model_name: Optional[str] = None, chunk_size: int = DEFAULT_CHUNK_SIZE, chunk_overlap: int = DEFAULT_CHUNK_OVERLAP, + embed_model=None, ): """Initialize RAG client with persistent ChromaDB storage. @@ -54,6 +55,7 @@ def __init__( model_name: HuggingFace embedding model name chunk_size: Size of text chunks for embedding chunk_overlap: Overlap between chunks + embed_model: Existing embedding model to reuse across paper databases """ self.persist_directory = Path(persist_directory) self.persist_directory.mkdir(parents=True, exist_ok=True) @@ -64,10 +66,12 @@ def __init__( self.device = get_device() # Initialize embedding model with GPU support - self.embed_model = HuggingFaceEmbedding( - model_name=self.model_name, - device=self.device, - ) + self.embed_model = embed_model + if self.embed_model is None: + self.embed_model = HuggingFaceEmbedding( + model_name=self.model_name, + device=self.device, + ) # Set global settings (no LLM needed for embedding/retrieval only) Settings.embed_model = self.embed_model diff --git a/tests/test_integration.py b/tests/test_integration.py index 5817d86..1ffc212 100644 --- a/tests/test_integration.py +++ b/tests/test_integration.py @@ -8,6 +8,7 @@ """ import shutil +import time from pathlib import Path import pytest @@ -391,13 +392,16 @@ def test_search_auto_processes_pdf(self, temp_output_dir): pdf_copy = temp_output_dir / "test_auto.pdf" shutil.copy(SAMPLE_PDF, pdf_copy) + started_at = time.monotonic() result = search( query="the", sources=[str(pdf_copy)], mode="grep", top_k=5, ) + first_response_seconds = time.monotonic() - started_at + assert first_response_seconds < 5 assert result["success"] assert result["status"] == "processing" assert result["processing"][0]["retry_after_seconds"] == 30 @@ -405,7 +409,8 @@ def test_search_auto_processes_pdf(self, temp_output_dir): # Processing continues independently of the first MCP response. paper_dir = temp_output_dir / "test_auto" from paper_intelligence.tools.search import _PROCESSING_JOBS - _PROCESSING_JOBS[str(paper_dir)].result(timeout=300) + process = _PROCESSING_JOBS[str(paper_dir)] + assert process.wait(timeout=300) == 0 completed = search( query="the", diff --git a/tests/test_reliability.py b/tests/test_reliability.py index b0b6282..69e5815 100644 --- a/tests/test_reliability.py +++ b/tests/test_reliability.py @@ -25,6 +25,14 @@ def test_direct_pdf_does_not_discover_sibling_temp_content(tmp_path, monkeypatch unrelated.mkdir() (unrelated / "metadata.json").write_text("unrelated", encoding="utf-8") + original_iterdir = Path.iterdir + + def reject_parent_scan(path): + if path == tmp_path: + raise AssertionError("a direct PDF must not scan its parent") + return original_iterdir(path) + + monkeypatch.setattr(Path, "iterdir", reject_parent_scan) scheduled = [] def fake_schedule(job_source, paper_dir, use_llm): @@ -68,6 +76,62 @@ def test_first_use_search_returns_actionable_background_status(tmp_path, monkeyp assert "get_paper_info" in result["processing"][0]["next_step"] +def test_first_use_launches_detached_process_and_returns_status(tmp_path, monkeypatch): + search_module = importlib.import_module("paper_intelligence.tools.search") + source = tmp_path / "paper.pdf" + source.write_bytes(b"%PDF-1.4\n") + paper_dir = tmp_path / "paper" + launched = [] + + class FakeProcess: + pid = 4321 + + def poll(self): + return None + + def fake_start(job_source, job_dir, use_llm): + launched.append((job_source, job_dir, use_llm)) + return FakeProcess() + + search_module._PROCESSING_JOBS.clear() + monkeypatch.setattr(search_module, "_start_processing_process", fake_start) + + result = search_module._schedule_processing(source, paper_dir, use_llm=False) + + assert launched == [(source, paper_dir, False)] + assert result["status"] == "queued" + assert result["worker_pid"] == 4321 + assert result["estimated_completion"] == "usually 1-3 minutes" + assert result["retry_after_seconds"] == 30 + assert "get_paper_info" in result["next_step"] + + +def test_processing_subprocess_is_detached(tmp_path, monkeypatch): + search_module = importlib.import_module("paper_intelligence.tools.search") + captured = {} + + class FakeProcess: + pid = 1234 + + def fake_popen(command, **kwargs): + captured["command"] = command + captured["kwargs"] = kwargs + return FakeProcess() + + monkeypatch.setattr(search_module.subprocess, "Popen", fake_popen) + process = search_module._start_processing_process( + tmp_path / "paper.pdf", tmp_path / "paper", use_llm=True + ) + + assert process.pid == 1234 + assert captured["command"][1:3] == ["-m", "paper_intelligence.processing_worker"] + assert captured["command"][-1] == "--use-llm" + assert captured["kwargs"]["start_new_session"] is True + assert captured["kwargs"]["stdin"] is subprocess.DEVNULL + assert captured["kwargs"]["stdout"] is subprocess.DEVNULL + assert captured["kwargs"]["stderr"] is subprocess.DEVNULL + + def test_processing_job_persists_completion_status(tmp_path, monkeypatch): search_module = importlib.import_module("paper_intelligence.tools.search") @@ -100,6 +164,27 @@ def test_status_module_import_does_not_load_embedding_stack(): subprocess.run([sys.executable, "-c", code], check=True, timeout=10) +def test_direct_pdf_discovery_does_not_import_torch(): + code = """ +import pathlib +import tempfile +import sys +import importlib + +search_module = importlib.import_module("paper_intelligence.tools.search") +with tempfile.TemporaryDirectory() as directory: + source = pathlib.Path(directory) / "paper.pdf" + source.write_bytes(b"%PDF-1.4\\n") + search_module._schedule_processing = lambda *_args: {"status": "queued"} + ready, processing = search_module._find_paper_dirs([str(source)]) + assert ready == [] + assert processing == [{"status": "queued"}] + search_module.get_paper_info(str(source)) +assert "torch" not in sys.modules +""" + subprocess.run([sys.executable, "-c", code], check=True, timeout=10) + + def test_get_paper_info_counts_chunks_without_embedding_model(tmp_path): from paper_intelligence.tools.search import get_paper_info @@ -120,6 +205,64 @@ def test_get_paper_info_counts_chunks_without_embedding_model(tmp_path): assert result["chunk_count"] == 2 +def test_collection_hybrid_search_skips_unbounded_semantic_fanout(monkeypatch): + search_module = importlib.import_module("paper_intelligence.tools.search") + paper_dirs = [ + Path(f"/papers/paper-{index}") + for index in range(search_module.MAX_SEMANTIC_PAPERS + 1) + ] + + monkeypatch.setattr( + search_module, + "_find_paper_dirs", + lambda *_args, **_kwargs: (paper_dirs, []), + ) + monkeypatch.setattr( + search_module, + "grep_search", + lambda *_args, **_kwargs: [{"content": "bounded", "match_type": "grep"}], + ) + + def rag_must_not_run(*_args, **_kwargs): + raise AssertionError("collection-wide semantic fanout must be rejected promptly") + + monkeypatch.setattr(search_module, "rag_search", rag_must_not_run) + + result = search_module.search("bounded", ["/papers"], mode="hybrid") + + assert result["success"] is True + assert result["status"] == "partial" + assert result["num_results"] == 1 + assert result["semantic_search"]["status"] == "skipped_source_limit" + assert result["semantic_search"]["papers_requested"] == len(paper_dirs) + + +def test_rag_search_reuses_one_embedding_model_across_papers(tmp_path, monkeypatch): + search_module = importlib.import_module("paper_intelligence.tools.search") + first = tmp_path / "first" + second = tmp_path / "second" + (first / "chroma").mkdir(parents=True) + (second / "chroma").mkdir(parents=True) + seen_models = [] + shared_model = object() + + class FakeRAGClient: + def __init__(self, persist_directory, embed_model=None): + seen_models.append(embed_model) + self.embed_model = embed_model or shared_model + + def query(self, collection_name, query, top_k): + return [] + + monkeypatch.setattr( + search_module, "_get_rag_client_class", lambda: FakeRAGClient + ) + + search_module.rag_search("query", [first, second]) + + assert seen_models == [None, shared_model] + + def test_mcp_search_does_not_sync_whole_library(tmp_path, monkeypatch): """Already-indexed local searches must not enter the five-minute rclone path.""" import paper_intelligence.server as server