mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-10-05 13:01:12 +01:00
fix(downloads): reconcile downloads a restart interrupted (#1423)
This commit is contained in:
@@ -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``.
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user