From 119d3374e3bc123982966ff49cbcfa912f96718a Mon Sep 17 00:00:00 2001 From: CaliBrain Date: Fri, 2 Oct 2026 22:35:09 -0400 Subject: [PATCH] fix(downloads): reconcile downloads a restart interrupted (#1423) --- shelfmark/core/download_history_service.py | 13 ++++- shelfmark/core/startup_reconcile.py | 61 +++++++++++----------- tests/core/test_startup_reconcile.py | 58 +++++++++++++++++--- 3 files changed, 93 insertions(+), 39 deletions(-) diff --git a/shelfmark/core/download_history_service.py b/shelfmark/core/download_history_service.py index c09acbab..e43998d7 100644 --- a/shelfmark/core/download_history_service.py +++ b/shelfmark/core/download_history_service.py @@ -406,7 +406,7 @@ class DownloadHistoryService: return self._row_to_dict(row) finally: conn.close() - + def list_active(self) -> list[dict[str, Any]]: """Return every row still marked active, oldest first.""" conn = self._connect() @@ -423,6 +423,17 @@ class DownloadHistoryService: finally: conn.close() + def list_request_ids(self) -> set[int]: + """Return the id of every request that has at least one download row.""" + conn = self._connect() + try: + rows = conn.execute( + "SELECT DISTINCT request_id FROM download_history WHERE request_id IS NOT NULL" + ).fetchall() + return {int(row["request_id"]) for row in rows} + finally: + conn.close() + def list_request_task_ids(self, *, user_id: int) -> list[str]: """Return the task ids of every request-linked download owned by ``user_id``. diff --git a/shelfmark/core/startup_reconcile.py b/shelfmark/core/startup_reconcile.py index 591a6edc..1ce984af 100644 --- a/shelfmark/core/startup_reconcile.py +++ b/shelfmark/core/startup_reconcile.py @@ -30,10 +30,12 @@ logger = setup_logger(__name__) # it relabels a stale row that has none. INTERRUPTED_MESSAGE = "Interrupted" +# The failure reason a reopened request is given. +REOPEN_REASON = "Interrupted by a restart" + # Only what was queued this recently is put back for another attempt. Anything older was -# orphaned by an earlier restart nobody reconciled: it is closed out here, but reopening a -# backlog all at once (or filling an admin's approval queue with old requests) is not -# something a startup should do. +# orphaned by an earlier restart nobody reconciled, and reopening a backlog all at once (or +# filling an admin's approval queue with old requests) is not something a startup should do. REOPEN_WINDOW = timedelta(hours=24) @@ -69,25 +71,33 @@ def reconcile_interrupted_downloads( ) -> dict[str, int]: """Close out downloads the last process left active, and reopen the recent ones. - Every history row still "active" is finalised as an error with the message "Interrupted", - which keeps its retry payload, so the manual retry it already offers keeps working. A - fulfilled request whose download was queued within ``reopen_window`` goes back to pending - through ``UserDB.reopen_failed_request``, upstream's own path for a failed request, so it - is approved again or, where automatic downloads are on, tried again. A request whose own - delivery state says it arrived is left alone, and so is its row. + A history row still "active" with no request is finalised as an error with the message + "Interrupted", which keeps its retry payload, so the manual retry it already offers keeps + working. + + A request-linked row is closed out only together with its request. Upstream retries a + request-linked error from its request, not from the row, so finalising the row alone + would take away the retry the activity API offers for it today, and the next activity + read would then reopen the request whatever its age. So when the download was queued + within ``reopen_window``, the request goes back to pending through + ``UserDB.reopen_failed_request``, upstream's own path for a failed request, and the row + is finalised. Every other request-linked row is left as it is, still offering its retry. Returns how many rows were closed out and how many requests were reopened. """ now = now or datetime.now(UTC) marked = 0 reopened = 0 - seen_requests: set[int] = set() for row in history_service.list_active(): request_id = row.get("request_id") - request = user_db.get_request(int(request_id)) if request_id else None - if request is not None and request.get("delivery_state") == QueueStatus.COMPLETE: - continue # the request says it arrived, so do not call it interrupted + if request_id: + if not _is_recent(row.get("queued_at"), now=now, window=reopen_window): + continue + # Refused when the request is gone, no longer fulfilled, or says it arrived. + if user_db.reopen_failed_request(int(request_id), failure_reason=REOPEN_REASON) is None: + continue + reopened += 1 history_service.finalize_download( task_id=str(row["task_id"]), @@ -96,30 +106,19 @@ def reconcile_interrupted_downloads( ) marked += 1 - if request is None: - continue - seen_requests.add(int(request["id"])) - if request.get("status") != RequestStatus.FULFILLED: - continue - if _is_recent(row.get("queued_at"), now=now, window=reopen_window) and ( - user_db.reopen_failed_request( - int(request["id"]), failure_reason="Interrupted by a restart" - ) - is not None - ): - reopened += 1 - # A request can be "queued" with no history row at all, for instance when the row was - # never written. After a restart nothing is queued, so it is stranded the same way. + # never written. After a restart nothing is queued, so it is stranded the same way. A + # request that has a row is not one of these: delivery_state is only synced while + # something polls the queue, so a download that finished with nobody watching still says + # "queued", and the activity API settles it from the row on its next read. + requests_with_history = history_service.list_request_ids() for request in user_db.list_requests(status=RequestStatus.FULFILLED): - if int(request["id"]) in seen_requests: + if int(request["id"]) in requests_with_history: continue if request.get("delivery_state") != QueueStatus.QUEUED: continue if _is_recent(_request_stamp(request), now=now, window=reopen_window) and ( - user_db.reopen_failed_request( - int(request["id"]), failure_reason="Interrupted by a restart" - ) + user_db.reopen_failed_request(int(request["id"]), failure_reason=REOPEN_REASON) is not None ): reopened += 1 diff --git a/tests/core/test_startup_reconcile.py b/tests/core/test_startup_reconcile.py index aa173a64..87a0a030 100644 --- a/tests/core/test_startup_reconcile.py +++ b/tests/core/test_startup_reconcile.py @@ -12,7 +12,9 @@ from datetime import UTC, datetime, timedelta import pytest from shelfmark.core import startup_reconcile +from shelfmark.core.activity_routes import _build_download_status_from_db from shelfmark.core.download_history_service import DownloadHistoryService +from shelfmark.core.requests_service import sync_delivery_states_from_queue_status from shelfmark.core.user_db import UserDB NOW = datetime(2026, 9, 30, 12, 0, tzinfo=UTC) @@ -90,6 +92,14 @@ def _run(user_db, history, **kwargs): return startup_reconcile.reconcile_interrupted_downloads(user_db, history, now=NOW, **kwargs) +def _read_activity(user_db, history): + """Sync request delivery states the way an admin's activity snapshot does.""" + status = _build_download_status_from_db( + db_rows=history.list_recent(user_id=None), queue_status={} + ) + sync_delivery_states_from_queue_status(user_db, queue_status=status) + + def test_a_recently_interrupted_download_is_closed_out_and_its_request_reopened(db): user_db, history, path = db user = user_db.create_user(username="reader", role="user") @@ -108,19 +118,30 @@ def test_a_recently_interrupted_download_is_closed_out_and_its_request_reopened( assert "restart" in stored["last_failure_reason"] -def test_an_older_orphan_is_closed_out_but_its_request_is_left_for_a_slower_retry(db): +def test_an_older_request_linked_orphan_is_left_active_and_keeps_its_retry(db): user_db, history, path = db user = user_db.create_user(username="reader", role="user") request = _request(user_db, user["id"], hours_ago=72) _download( - history, path, task_id="t1", user_id=user["id"], request_id=request["id"], hours_ago=72 + history, + path, + task_id="t1", + user_id=user["id"], + request_id=request["id"], + hours_ago=72, + retry_payload={"source": "audiobookbay", "source_id": "rel-1"}, ) - assert _run(user_db, history) == {"marked": 1, "reopened": 0} + assert _run(user_db, history) == {"marked": 0, "reopened": 0} - assert history.get_by_task_id("t1")["final_status"] == "error" - stored = user_db.get_request(request["id"]) - assert (stored["status"], stored["delivery_state"]) == ("fulfilled", "queued") + row = history.get_by_task_id("t1") + assert row["final_status"] == "active" + assert history.is_retry_available(row) is True + + # Closing the row out as an error would end its retry, and the activity read would then + # reopen the request regardless of the window. + _read_activity(user_db, history) + assert user_db.get_request(request["id"])["status"] == "fulfilled" def test_a_download_with_no_request_is_closed_out_and_keeps_its_manual_retry(db): @@ -174,14 +195,37 @@ def test_a_recent_request_stuck_queued_with_no_history_is_reopened(db): assert user_db.get_request(old["id"])["status"] == "fulfilled" +def test_a_request_whose_download_finished_is_not_reopened(db): + """delivery_state only moves while the queue is polled, so it can still say "queued".""" + user_db, history, path = db + user = user_db.create_user(username="reader", role="user") + request = _request(user_db, user["id"], hours_ago=2) + _download( + history, + path, + task_id="t1", + user_id=user["id"], + request_id=request["id"], + hours_ago=2, + final_status="complete", + ) + + assert _run(user_db, history) == {"marked": 0, "reopened": 0} + + stored = user_db.get_request(request["id"]) + assert (stored["status"], stored["delivery_state"]) == ("fulfilled", "queued") + assert stored["release_data"] is not None + + def test_a_request_that_is_not_fulfilled_keeps_its_status(db): user_db, history, path = db user = user_db.create_user(username="reader", role="user") request = _request(user_db, user["id"], status="pending", state="none") _download(history, path, task_id="t1", user_id=user["id"], request_id=request["id"]) - assert _run(user_db, history) == {"marked": 1, "reopened": 0} + assert _run(user_db, history) == {"marked": 0, "reopened": 0} + assert history.get_by_task_id("t1")["final_status"] == "active" assert user_db.get_request(request["id"])["status"] == "pending"