From 0bc054e7b00646c6e731b11b06a8ad24502432f0 Mon Sep 17 00:00:00 2001 From: LLLin000 <809867916@qq.com> Date: Fri, 10 Jul 2026 23:55:21 +0800 Subject: [PATCH] fix: cooperative stop with state settlement and build-loop cancellation embed stop now sets status=stopping, waits for the build PID to exit (8s timeout), force-kills as last resort, then settles to idle. Build loop checks the cancellation flag (read_vector_build_state) between papers in the main loop and drains in-flight work before cleanly exiting without marking 'completed'. Fixes pre-existing UnboundLocalError from a local import shadowing the global mark_vector_build_state symbol in the run() function. --- paperforge/commands/embed.py | 55 +++++--- tests/integration/test_stop_idle_resume.py | 154 +++++++++++++++++++++ 2 files changed, 188 insertions(+), 21 deletions(-) create mode 100644 tests/integration/test_stop_idle_resume.py diff --git a/paperforge/commands/embed.py b/paperforge/commands/embed.py index 6711eb76..6de40077 100644 --- a/paperforge/commands/embed.py +++ b/paperforge/commands/embed.py @@ -159,32 +159,34 @@ def run(args: argparse.Namespace) -> int: if sub == "stop": state = read_vector_build_state(vault) pid = state.get("pid", 0) - if pid and state["status"] == "running": - import signal - - try: - os.kill(pid, signal.SIGTERM) - except Exception as exc: - result = PFResult( - ok=False, - command="embed stop", - version=PF_VERSION, - error=PFError(code=ErrorCode.INTERNAL_ERROR, message=f"Failed to stop embed build: {exc}"), - data={"state": "running", "pid": pid}, - ) - if args.json: - print(result.to_json()) - else: - print(result.error.message, file=sys.stderr) - return 1 + _st = state.get("status", "") + if pid and _st in ("running", "stopping"): mark_vector_build_state(vault, status="stopping", message="Stop requested") - result = PFResult(ok=True, command="embed stop", version=PF_VERSION, data={"state": "stopping", "pid": pid}) + # Wait for build process to notice the flag and exit (8s timeout) + import time as _time + _deadline = _time.time() + 8.0 + while _time.time() < _deadline: + if not _pid_alive(pid): + break + _time.sleep(0.2) + # Force-kill if still alive after timeout + if _pid_alive(pid): + import signal + try: + os.kill(pid, signal.SIGTERM) + except Exception: + pass + # Settle to idle, preserving progress + _current = read_vector_build_state(vault).get("current", state.get("current", 0)) + mark_vector_build_state(vault, status="idle", current=_current, pid=0, message="") + result = PFResult(ok=True, command="embed stop", version=PF_VERSION, data={"state": "stopped"}) else: result = PFResult(ok=True, command="embed stop", version=PF_VERSION, data={"state": "idle"}) if args.json: print(result.to_json()) else: - print("Stop requested." if result.data["state"] == "stopping" else "No active build.") + msg = "Build stopped." if result.data["state"] == "stopped" else "No active build." + print(msg) return 0 @@ -287,7 +289,6 @@ def run(args: argparse.Namespace) -> int: if stale: msg = "Previous build appears stale (crashed?). Recovering and rebuilding from scratch." print(msg) - from paperforge.embedding.build_state import mark_vector_build_state mark_vector_build_state(vault, status="idle", current=0, pid=0) resume = False @@ -397,6 +398,11 @@ def run(args: argparse.Namespace) -> int: if not key: continue + # ponytail: check cancellation flag between papers + if read_vector_build_state(vault).get("status") == "stopping": + logger.info("Build cancelled at paper %s", key) + break + has_body = _has_body_units_in_db(vault, key) has_object = _has_object_units_in_db(vault, key) @@ -581,6 +587,13 @@ def run(args: argparse.Namespace) -> int: print(result.to_json() if args.json else result.error.message, file=sys.stderr if not args.json else sys.stdout) return 1 + + # Check if we stopped or were cancelled — exit cleanly without marking completed + if read_vector_build_state(vault).get("status") == "stopping": + logger.info("Build stopped, exiting cleanly") + print("EMBED_DONE", flush=True) + return 0 + mark_vector_build_state( vault, status="completed", diff --git a/tests/integration/test_stop_idle_resume.py b/tests/integration/test_stop_idle_resume.py new file mode 100644 index 00000000..d49f4cd2 --- /dev/null +++ b/tests/integration/test_stop_idle_resume.py @@ -0,0 +1,154 @@ +"""Integration tests for the build lifecycle stop → idle → resume cycle. + +Verifies that embed stop --json marks build as stopping, settles to idle, +and that resume can continue a build from progress. +""" + +from __future__ import annotations + +import json +import subprocess +from pathlib import Path + +import pytest + +PYTHON = Path(r"D:\L\OB\Literature-hub\.venv\Scripts\python.exe") +REPO_ROOT = Path(__file__).resolve().parent.parent.parent + + +def _build_vault(tmp_path: Path) -> Path: + """Build a minimal vault with paperforge.json.""" + vault = tmp_path / "vault" + vault.mkdir(parents=True) + (vault / "paperforge.json").write_text( + '{"vault_config":{"system_dir":"System"}}', encoding="utf-8" + ) + (vault / "System" / "PaperForge").mkdir(parents=True) + return vault + + +def _run_embed_stop(vault: Path) -> dict: + """Run embed stop --json and return parsed data.""" + result = subprocess.run( + [ + str(PYTHON), + "-m", + "paperforge", + "--vault", + str(vault), + "embed", + "stop", + "--json", + ], + capture_output=True, + text=True, + cwd=str(REPO_ROOT), + timeout=30, + ) + parsed = json.loads(result.stdout) + return parsed + + +def _run_embed_status(vault: Path) -> dict: + """Run embed status --json and return parsed data.""" + result = subprocess.run( + [ + str(PYTHON), + "-m", + "paperforge", + "--vault", + str(vault), + "embed", + "status", + "--json", + ], + capture_output=True, + text=True, + cwd=str(REPO_ROOT), + timeout=30, + ) + assert result.returncode == 0 + parsed = json.loads(result.stdout) + assert parsed["ok"] is True + return parsed["data"] + + +def _mark_running(vault: Path, pid: int, current: int = 5, total: int = 100) -> None: + """Set build_state to 'running' with the given PID.""" + from paperforge.embedding.build_state import mark_vector_build_state + + mark_vector_build_state( + vault, + status="running", + pid=pid, + current=current, + total=total, + paper_id="TESTKEY", + model="test-model", + mode="api", + ) + + +def _read_build_state(vault: Path) -> dict: + """Read build_state table.""" + from paperforge.embedding.build_state import read_vector_build_state + + return read_vector_build_state(vault) + + +class TestStopIdleResume: + """Stop → idle → resume lifecycle.""" + + def test_stop_returns_stopped_and_idle(self, tmp_path: Path): + """embed stop --json returns {state: stopped} when a build is running, and settles to idle.""" + vault = _build_vault(tmp_path) + + # Set build to running with a dead PID (not 0 — 0 is falsy in Python) + _mark_running(vault, pid=99999, current=5, total=100) + + # Stop + result = _run_embed_stop(vault) + assert result["ok"] is True, f"Stop returned error: {result}" + assert result["data"]["state"] == "stopped", f"Expected stopped, got {result['data']['state']}" + + # Verify build state settled to idle + bs = _read_build_state(vault) + assert bs["status"] == "idle", f"Expected idle, got {bs}" + assert bs["pid"] == 0 + + def test_stop_when_idle_stays_idle(self, tmp_path: Path): + """embed stop --json on idle build returns {state: idle}.""" + vault = _build_vault(tmp_path) + + result = _run_embed_stop(vault) + assert result["ok"] is True + assert result["data"]["state"] == "idle" + + bs = _read_build_state(vault) + assert bs["status"] == "idle" + + def test_resume_after_stop_reads_progress(self, tmp_path: Path): + """After stop, status shows the idle state with progress, and resume mode continues from there.""" + vault = _build_vault(tmp_path) + + _mark_running(vault, pid=99998, current=42, total=200) + + # Stop + _run_embed_stop(vault) + + # Status should report progress and idle state + status = _run_embed_status(vault) + bs = status.get("build_state", {}) + assert bs.get("status") == "idle" + assert bs.get("current", 0) == 42 + assert bs.get("total", 0) == 200 + + # Resume mode: verify that a subsequent --resume build picks up from idle + # (we can't do a full build without a real vault, so we check that + # the resume path isn't blocked by the idle state) + from paperforge.embedding.build_state import read_vector_build_state + + bs_read = read_vector_build_state(vault) + assert bs_read["status"] == "idle" + # A resume build call would check this and start fresh or continue + # The state is correctly idle which allows resume