From 5ce877f186cbed75bb8fd3ec9c1f4a1f581639c2 Mon Sep 17 00:00:00 2001 From: yipkaio Date: Sun, 27 Sep 2026 17:57:45 +0800 Subject: [PATCH] Recover receipts interrupted by API restart --- app/database.py | 16 +++++++++++ app/main.py | 13 +++++++++ docs/user-guide.md | 2 +- tests/test_persistence_api.py | 53 +++++++++++++++++++++++++++++++++++ 4 files changed, 83 insertions(+), 1 deletion(-) diff --git a/app/database.py b/app/database.py index acb3944..48c7d47 100644 --- a/app/database.py +++ b/app/database.py @@ -123,6 +123,22 @@ def start(self, receipt_id: str, content_type: str, size: int, image_path: str, raise DuplicateReceiptError(duplicate[0]) raise + def recover_interrupted(self) -> int: + """Close uploads left in progress when the sole API worker was stopped. + + Call only during worker startup, before it can accept new uploads. Keep + the retained file and OCR text so the owner can inspect or reprocess it. + """ + with self.connect() as db: + db.execute("BEGIN IMMEDIATE") + cursor = db.execute( + "UPDATE receipts SET processing_status='FAILED', " + "error='Processing was interrupted. Reprocess from saved OCR or delete and upload again.', " + "updated_at=? WHERE processing_status='PROCESSING'", + (now(),), + ) + return cursor.rowcount + def probable_duplicates(self, receipt_id: str, extraction: dict) -> list[str]: """Return strict identity matches; vendor+amount alone are never enough.""" required = (extraction.get("receipt_number"), extraction.get("date"), diff --git a/app/main.py b/app/main.py index c0cd670..3d8e307 100644 --- a/app/main.py +++ b/app/main.py @@ -265,6 +265,15 @@ def has_expected_signature(content_type: str, prefix: bytes) -> bool: @asynccontextmanager async def lifespan(api): + # Production runs one API worker. After that worker stops, no previous + # upload can complete; release its retained records before accepting traffic. + recovered = await run_in_threadpool( + ReceiptStore(get_settings().database_path).recover_interrupted + ) + if recovered: + logging.getLogger(__name__).warning( + "Marked %d interrupted receipt upload(s) as failed", recovered + ) async def cleanup(): while True: try: @@ -1082,6 +1091,10 @@ async def process() -> ReceiptProcessed: result = await process() await run_in_threadpool(store.complete, result.model_dump(mode="json")) return result + except asyncio.CancelledError: + # A canceled HTTP task must not leave a permanent PROCESSING row. + await run_in_threadpool(store.fail, receipt_id, "Receipt processing was interrupted") + raise except HTTPException as exc: await run_in_threadpool(store.fail, receipt_id, str(exc.detail)) exc.headers = {**(exc.headers or {}), "X-Receipt-ID": receipt_id} diff --git a/docs/user-guide.md b/docs/user-guide.md index d076b78..911214e 100644 --- a/docs/user-guide.md +++ b/docs/user-guide.md @@ -20,7 +20,7 @@ The sign-in image is a current capture from the deployed page. The dashboard, hi *Upload a receipt* -Upload a JPEG, PNG or bounded PDF through the web workspace. The selected file appears beside the form, and its preview remains visible after processing. Select **Open processed receipt** to inspect extracted fields and any review reasons, or **Upload another** to clear the form for the next file. This local preview is a convenience; the saved receipt detail shows the retained original evidence. Telegram receipts reach the same authenticated backend through the OpenClaw relay. +Upload a JPEG, PNG or bounded PDF through the web workspace. The selected file appears beside the form, and its preview remains visible after processing. Select **Open processed receipt** to inspect extracted fields and any review reasons, or **Upload another** to clear the form for the next file. This local preview is a convenience; the saved receipt detail shows the retained original evidence. Telegram receipts reach the same authenticated backend through the OpenClaw relay. If an API restart interrupts an upload, refresh Receipt history after the service returns: the interrupted record becomes **Failed**. Open it to reprocess from saved OCR text (when available), or move it to Deleted receipts and upload the file again. An active upload cannot be deleted. ![Pending reviews](assets/screenshots/2026-09-23/03-pending-reviews.jpg) diff --git a/tests/test_persistence_api.py b/tests/test_persistence_api.py index 4682157..82ada60 100644 --- a/tests/test_persistence_api.py +++ b/tests/test_persistence_api.py @@ -31,6 +31,59 @@ def test_result_survives_new_application(monkeypatch, tmp_path): assert restarted.get("/receipts/" + original["receipt_id"]).status_code == 401 +def test_restart_releases_interrupted_upload_for_deletion_and_retry(monkeypatch, tmp_path): + client, _, _, _ = configured_client(monkeypatch, tmp_path) + store = ReceiptStore(Path(os.environ["DATABASE_PATH"])) + receipt_id = str(uuid4()) + store.start(receipt_id, "image/jpeg", 10, str(tmp_path / "stuck.jpg"), + "Office supplies", "stuck-content") + assert client.get(f"/receipts/{receipt_id}", headers=HEADERS).json()["processing_status"] == "PROCESSING" + with TestClient(create_app()) as restarted: + saved = restarted.get(f"/receipts/{receipt_id}", headers=HEADERS).json() + assert saved["processing_status"] == "FAILED" + assert "interrupted" in saved["error"].lower() + assert saved["content_sha256"] == "stuck-content" + assert restarted.post( + f"/receipts/{receipt_id}/lifecycle", headers=HEADERS, + json={"request_id": str(uuid4()), "action": "DELETE", + "expected_version": 0, "expected_record_version": 0, + "reviewer": "Test reviewer", "reason": "Abandoned upload from interrupted processing"}, + ).status_code == 200 + assert store.get(receipt_id)["lifecycle_state"] == "DELETED" + + +def test_restart_leaves_completed_and_failed_receipts_unchanged(monkeypatch, tmp_path): + configured_client(monkeypatch, tmp_path) + store = ReceiptStore(Path(os.environ["DATABASE_PATH"])) + failed_id = str(uuid4()) + store.start(failed_id, "image/jpeg", 10, "private/path", None) + store.fail(failed_id, "Existing failure") + with TestClient(create_app()) as restarted: + saved = restarted.get(f"/receipts/{failed_id}", headers=HEADERS).json() + assert saved["processing_status"] == "FAILED" + assert saved["error"] == "Existing failure" + + +def test_restart_retains_ocr_for_reprocessing(monkeypatch, tmp_path): + client, _, _, _ = configured_client(monkeypatch, tmp_path) + store = ReceiptStore(Path(os.environ["DATABASE_PATH"])) + receipt_id = str(uuid4()) + store.start(receipt_id, "image/jpeg", 10, "private/path", "Office supplies") + store.save_ocr(receipt_id, "Vendor and total from saved receipt", "tesseract", 0.94) + with TestClient(client.app) as restarted: + response = restarted.post( + f"/receipts/{receipt_id}/reprocess", headers=HEADERS, + json={"request_id": str(uuid4()), "expected_record_version": 0, + "expected_lifecycle_version": 0, "reviewer": "Test reviewer", + "reason": "Recover extraction after interrupted upload"}, + ) + assert response.status_code == 200, response.text + assert response.json()["status"] == "SUCCEEDED" + saved = store.get(receipt_id) + assert saved["processing_status"] == "REVIEW_QUEUE" + assert saved["ocr_text"] == "Vendor and total from saved receipt" + + def test_review_filter_and_pagination(monkeypatch, tmp_path): client, _, extractor, _ = configured_client(monkeypatch, tmp_path) client.post("/receipts/upload", headers=HEADERS, files=FILES)