diff --git a/shelfmark/core/activity_routes.py b/shelfmark/core/activity_routes.py index 5a1262a7..28b27cfc 100644 --- a/shelfmark/core/activity_routes.py +++ b/shelfmark/core/activity_routes.py @@ -2,11 +2,16 @@ from __future__ import annotations -from datetime import datetime, timezone from typing import Any, Callable, NamedTuple from flask import Flask, jsonify, request, session +from shelfmark.core.activity_view_state_service import ( + ADMIN_VIEWER_SCOPE, + NOAUTH_VIEWER_SCOPE, + ActivityViewStateService, + user_viewer_scope, +) from shelfmark.core.download_history_service import ACTIVE_DOWNLOAD_STATUS, DownloadHistoryService, VALID_TERMINAL_STATUSES from shelfmark.core.logger import setup_logger from shelfmark.core.models import ACTIVE_QUEUE_STATUSES, QueueStatus, TERMINAL_QUEUE_STATUSES @@ -14,9 +19,7 @@ from shelfmark.core.request_validation import RequestStatus from shelfmark.core.request_helpers import ( emit_ws_event, extract_release_source_id, - normalize_optional_text, normalize_positive_int, - now_utc_iso, populate_request_usernames, ) from shelfmark.core.user_db import UserDB @@ -24,16 +27,6 @@ from shelfmark.core.user_db import UserDB logger = setup_logger(__name__) -def _parse_timestamp(value: Any) -> float: - if not isinstance(value, str) or not value.strip(): - return 0.0 - normalized = value.strip().replace("Z", "+00:00") - try: - return datetime.fromisoformat(normalized).timestamp() - except ValueError: - return 0.0 - - def _require_authenticated(resolve_auth_mode: Callable[[], str]): auth_mode = resolve_auth_mode() if auth_mode == "none": @@ -116,6 +109,7 @@ class _ActorContext(NamedTuple): is_no_auth: bool is_admin: bool owner_scope: int | None + viewer_scope: str def _resolve_activity_actor( @@ -128,27 +122,35 @@ def _resolve_activity_actor( Returns (actor, error_response). On success actor is non-None. """ if resolve_auth_mode() == "none": - return _ActorContext(db_user_id=None, is_no_auth=True, is_admin=True, owner_scope=None), None + return _ActorContext( + db_user_id=None, + is_no_auth=True, + is_admin=True, + owner_scope=None, + viewer_scope=NOAUTH_VIEWER_SCOPE, + ), None db_user_id, db_gate = _resolve_db_user_id(user_db=user_db) if db_user_id is None: return None, db_gate is_admin = bool(session.get("is_admin")) + viewer_scope = ADMIN_VIEWER_SCOPE if is_admin else user_viewer_scope(db_user_id) return _ActorContext( db_user_id=db_user_id, is_no_auth=False, is_admin=is_admin, owner_scope=None if is_admin else db_user_id, + viewer_scope=viewer_scope, ), None -def _activity_ws_room(*, is_no_auth: bool, actor_db_user_id: int | None) -> str: +def _activity_ws_room(actor: _ActorContext) -> str: """Resolve the WebSocket room for activity events.""" - if is_no_auth: + if actor.is_no_auth or actor.is_admin: return "admins" - if actor_db_user_id is not None: - return f"user_{actor_db_user_id}" + if actor.db_user_id is not None: + return f"user_{actor.db_user_id}" return "admins" @@ -162,6 +164,19 @@ def _check_item_ownership(actor: _ActorContext, row: dict[str, Any]) -> Any | No return None +def _check_terminal_download(row: dict[str, Any]) -> Any | None: + final_status = str(row.get("final_status") or "").strip().lower() + if final_status not in VALID_TERMINAL_STATUSES: + return jsonify({"error": "Only terminal downloads can be dismissed"}), 409 + return None + + +def _check_terminal_request(row: dict[str, Any]) -> Any | None: + if _request_terminal_status(row) is None: + return jsonify({"error": "Only terminal requests can be dismissed"}), 409 + return None + + def _list_visible_requests(user_db: UserDB, *, is_admin: bool, db_user_id: int | None) -> list[dict[str, Any]]: if is_admin: request_rows = user_db.list_requests() @@ -233,29 +248,9 @@ def _build_download_status_from_db( download_payload["status_message"] = None status[final_status][task_id] = download_payload - # Include any queue items that don't have a DB row yet (race condition safety). - # Only include active items — terminal items without a DB row are orphans - # (e.g. admin cleared history) and should not be shown. - for task_id, (bucket_key, queue_payload) in queue_index.items(): - if bucket_key in ACTIVE_QUEUE_STATUSES and task_id not in status.get(bucket_key, {}): - status[bucket_key][task_id] = queue_payload - return status -def _collect_active_download_task_ids(status: dict[str, dict[str, Any]]) -> set[str]: - active_task_ids: set[str] = set() - for bucket_key in ACTIVE_QUEUE_STATUSES: - bucket = status.get(bucket_key) - if not isinstance(bucket, dict): - continue - for task_id in bucket.keys(): - normalized_task_id = str(task_id).strip() - if normalized_task_id: - active_task_ids.add(normalized_task_id) - return active_task_ids - - def _request_terminal_status(row: dict[str, Any]) -> str | None: request_status = row.get("status") if request_status == RequestStatus.PENDING: @@ -300,7 +295,11 @@ def _minimal_request_snapshot(request_row: dict[str, Any], request_id: int) -> d return {"kind": "request", "request": minimal_request} -def _request_history_entry(request_row: dict[str, Any]) -> dict[str, Any] | None: +def _request_history_entry( + request_row: dict[str, Any], + *, + dismissed_at: str | None, +) -> dict[str, Any] | None: request_id = normalize_positive_int(request_row.get("id")) if request_id is None: return None @@ -311,7 +310,7 @@ def _request_history_entry(request_row: dict[str, Any]) -> dict[str, Any] | None "user_id": request_row.get("user_id"), "item_type": "request", "item_key": item_key, - "dismissed_at": request_row.get("dismissed_at"), + "dismissed_at": dismissed_at, "snapshot": _minimal_request_snapshot(request_row, request_id), "origin": "request", "final_status": final_status, @@ -321,29 +320,13 @@ def _request_history_entry(request_row: dict[str, Any]) -> dict[str, Any] | None } -def _dedupe_dismissed_entries(entries: list[dict[str, str]]) -> list[dict[str, str]]: - seen: set[tuple[str, str]] = set() - result: list[dict[str, str]] = [] - for entry in entries: - item_type = str(entry.get("item_type") or "").strip().lower() - item_key = str(entry.get("item_key") or "").strip() - if item_type not in {"download", "request"} or not item_key: - continue - marker = (item_type, item_key) - if marker in seen: - continue - seen.add(marker) - result.append({"item_type": item_type, "item_key": item_key}) - return result - - def register_activity_routes( app: Flask, user_db: UserDB, *, + activity_view_state_service: ActivityViewStateService, download_history_service: DownloadHistoryService, resolve_auth_mode: Callable[[], str], - resolve_status_scope: Callable[[], tuple[bool, int | None, bool]], queue_status: Callable[..., dict[str, dict[str, Any]]], sync_request_delivery_states: Callable[..., list[dict[str, Any]]], emit_request_updates: Callable[[list[dict[str, Any]]], None], @@ -357,93 +340,65 @@ def register_activity_routes( if auth_gate is not None: return auth_gate - is_admin, db_user_id, can_access_status = resolve_status_scope() - if not can_access_status: - return ( - jsonify( - { - "error": "User identity unavailable for activity workflow", - "code": "user_identity_unavailable", - } - ), - 403, - ) + actor, actor_error = _resolve_activity_actor( + user_db=user_db, + resolve_auth_mode=resolve_auth_mode, + ) + if actor_error is not None: + return actor_error - owner_user_scope = None if is_admin else db_user_id - - live_queue = queue_status(user_id=owner_user_scope) - - try: - db_rows = download_history_service.get_undismissed( - user_id=owner_user_scope, - limit=200, - ) - except Exception as exc: - logger.warning("Failed to load undismissed download rows: %s", exc) - db_rows = [] + hidden_rows = activity_view_state_service.list_hidden(viewer_scope=actor.viewer_scope) + hidden_item_keys = {str(row.get("item_key") or "").strip() for row in hidden_rows} + dismissed_entries = [ + { + "item_type": str(row.get("item_type") or "").strip().lower(), + "item_key": str(row.get("item_key") or "").strip(), + } + for row in hidden_rows + if str(row.get("item_type") or "").strip().lower() in {"download", "request"} + and str(row.get("item_key") or "").strip() + ] + live_queue = queue_status(user_id=actor.owner_scope) + db_rows = download_history_service.list_recent( + user_id=actor.owner_scope, + limit=200, + ) + visible_db_rows = [ + row + for row in db_rows + if f"download:{str(row.get('task_id') or '').strip()}" not in hidden_item_keys + ] status = _build_download_status_from_db( - db_rows=db_rows, + db_rows=visible_db_rows, queue_status=live_queue, ) updated_requests = sync_request_delivery_states( user_db, queue_status=status, - user_id=owner_user_scope, + user_id=actor.owner_scope, ) emit_request_updates(updated_requests) - request_rows = _list_visible_requests(user_db, is_admin=is_admin, db_user_id=db_user_id) - - dismissed: list[dict[str, str]] = [] - dismissed_task_ids: list[str] = [] - try: - dismissed_task_ids = download_history_service.get_dismissed_keys( - user_id=owner_user_scope, - ) - except Exception as exc: - logger.warning("Failed to load dismissed download keys: %s", exc) - - # Only clear stale dismissals when active downloads overlap dismissed keys. - active_task_ids = _collect_active_download_task_ids(status) - stale_dismissed = active_task_ids & set(dismissed_task_ids) if active_task_ids else set() - if stale_dismissed: - try: - download_history_service.clear_dismissals_for_active( - task_ids=stale_dismissed, - user_id=owner_user_scope, - ) - dismissed_task_ids = [tid for tid in dismissed_task_ids if tid not in stale_dismissed] - except Exception as exc: - logger.warning("Failed to clear stale download dismissals for active tasks: %s", exc) - - dismissed.extend( - {"item_type": "download", "item_key": f"download:{task_id}"} - for task_id in dismissed_task_ids + request_rows = _list_visible_requests( + user_db, + is_admin=actor.is_admin, + db_user_id=actor.db_user_id, ) - - # Keep request dismissal state on the request rows directly. - try: - dismissed_request_rows = user_db.list_dismissed_requests(user_id=owner_user_scope) - for request_row in dismissed_request_rows: - request_id = normalize_positive_int(request_row.get("id")) - if request_id is None: - continue - dismissed.append({"item_type": "request", "item_key": f"request:{request_id}"}) - except Exception as exc: - logger.warning("Failed to load dismissed request keys: %s", exc) - - if not is_admin and db_user_id is None: - # In auth mode, if we can't identify a non-admin viewer, don't show dismissals. - dismissed = [] - else: - dismissed = _dedupe_dismissed_entries(dismissed) + visible_request_rows: list[dict[str, Any]] = [] + for row in request_rows: + request_id = normalize_positive_int(row.get("id")) + if request_id is None: + continue + if f"request:{request_id}" in hidden_item_keys: + continue + visible_request_rows.append(row) return jsonify( { "status": status, - "requests": request_rows, - "dismissed": dismissed, + "requests": visible_request_rows, + "dismissed": dismissed_entries, } ) @@ -476,18 +431,21 @@ def register_activity_routes( existing = download_history_service.get_by_task_id(task_id) if existing is None: - # Row already gone (e.g. admin cleared history) — treat as success - dismissal_item = {"item_type": "download", "item_key": f"download:{task_id}"} - else: - ownership_gate = _check_item_ownership(actor, existing) - if ownership_gate is not None: - return ownership_gate + return jsonify({"error": "Download not found"}), 404 - download_history_service.dismiss( - task_id=task_id, - user_id=actor.owner_scope, - ) - dismissal_item = {"item_type": "download", "item_key": f"download:{task_id}"} + ownership_gate = _check_item_ownership(actor, existing) + if ownership_gate is not None: + return ownership_gate + terminal_gate = _check_terminal_download(existing) + if terminal_gate is not None: + return terminal_gate + + activity_view_state_service.dismiss( + viewer_scope=actor.viewer_scope, + item_type="download", + item_key=f"download:{task_id}", + ) + dismissal_item = {"item_type": "download", "item_key": f"download:{task_id}"} elif item_type == "request": request_id = normalize_positive_int(_parse_item_key(item_key, "request")) @@ -501,13 +459,20 @@ def register_activity_routes( ownership_gate = _check_item_ownership(actor, request_row) if ownership_gate is not None: return ownership_gate + terminal_gate = _check_terminal_request(request_row) + if terminal_gate is not None: + return terminal_gate - user_db.update_request(request_id, dismissed_at=now_utc_iso()) + activity_view_state_service.dismiss( + viewer_scope=actor.viewer_scope, + item_type="request", + item_key=f"request:{request_id}", + ) dismissal_item = {"item_type": "request", "item_key": f"request:{request_id}"} else: return jsonify({"error": "item_type must be one of: download, request"}), 400 - room = _activity_ws_room(is_no_auth=actor.is_no_auth, actor_db_user_id=actor.db_user_id) + room = _activity_ws_room(actor) emit_ws_event( ws_manager, event_name="activity_update", @@ -541,8 +506,8 @@ def register_activity_routes( if not isinstance(items, list): return jsonify({"error": "items must be an array"}), 400 - download_task_ids: list[str] = [] - request_ids: list[int] = [] + dismissal_items: list[dict[str, str]] = [] + missing_item_keys: list[str] = [] for item in items: if not isinstance(item, dict): @@ -557,11 +522,15 @@ def register_activity_routes( return jsonify({"error": "download item_key must be in the format download:"}), 400 existing = download_history_service.get_by_task_id(task_id) if existing is None: + missing_item_keys.append(f"download:{task_id}") continue ownership_gate = _check_item_ownership(actor, existing) if ownership_gate is not None: return ownership_gate - download_task_ids.append(task_id) + terminal_gate = _check_terminal_download(existing) + if terminal_gate is not None: + return terminal_gate + dismissal_items.append({"item_type": "download", "item_key": f"download:{task_id}"}) continue if item_type == "request": @@ -570,28 +539,36 @@ def register_activity_routes( return jsonify({"error": "request item_key must be in the format request:"}), 400 request_row = user_db.get_request(request_id) if request_row is None: + missing_item_keys.append(f"request:{request_id}") continue ownership_gate = _check_item_ownership(actor, request_row) if ownership_gate is not None: return ownership_gate - request_ids.append(request_id) + terminal_gate = _check_terminal_request(request_row) + if terminal_gate is not None: + return terminal_gate + dismissal_items.append({"item_type": "request", "item_key": f"request:{request_id}"}) continue return jsonify({"error": "item_type must be one of: download, request"}), 400 - dismissed_download_count = download_history_service.dismiss_many( - task_ids=download_task_ids, - user_id=actor.owner_scope, + if missing_item_keys: + return ( + jsonify( + { + "error": "One or more activity items were not found", + "missing_item_keys": missing_item_keys, + } + ), + 404, + ) + + dismissed_count = activity_view_state_service.dismiss_many( + viewer_scope=actor.viewer_scope, + items=dismissal_items, ) - dismissed_request_count = user_db.dismiss_requests_batch( - request_ids=request_ids, - dismissed_at=now_utc_iso(), - ) - - dismissed_count = dismissed_download_count + dismissed_request_count - - room = _activity_ws_room(is_no_auth=actor.is_no_auth, actor_db_user_id=actor.db_user_id) + room = _activity_ws_room(actor) emit_ws_event( ws_manager, event_name="activity_update", @@ -628,30 +605,70 @@ def register_activity_routes( if offset < 0: return jsonify({"error": "offset must be a non-negative integer"}), 400 - # Fetch enough from each source to fill the requested page after merging. - merge_limit = offset + limit - download_history_rows = download_history_service.get_history( - user_id=actor.owner_scope, - limit=merge_limit, - offset=0, + history_rows = activity_view_state_service.list_history( + viewer_scope=actor.viewer_scope, + limit=limit, + offset=offset, ) - dismissed_request_rows = user_db.list_dismissed_requests(user_id=actor.owner_scope, limit=merge_limit) - request_history_rows = [ - entry - for entry in (_request_history_entry(row) for row in dismissed_request_rows) - if entry is not None - ] + payload: list[dict[str, Any]] = [] - combined = [*download_history_rows, *request_history_rows] - combined.sort( - key=lambda row: ( - _parse_timestamp(row.get("dismissed_at")), - str(row.get("id") or ""), - ), - reverse=True, - ) - paged = combined[offset:offset + limit] - return jsonify(paged) + for history_row in history_rows: + item_type = str(history_row.get("item_type") or "").strip().lower() + item_key = str(history_row.get("item_key") or "").strip() + dismissed_at = history_row.get("dismissed_at") + + if not isinstance(dismissed_at, str) or not dismissed_at.strip(): + raise RuntimeError(f"Activity history state missing dismissed_at for {item_key}") + + if item_type == "download": + task_id = _parse_item_key(item_key, "download") + if task_id is None: + raise RuntimeError(f"Invalid activity history item_key: {item_key}") + + download_row = download_history_service.get_by_task_id(task_id) + if download_row is None: + raise RuntimeError(f"Download history row not found for {item_key}") + + if not actor.is_admin: + owner_user_id = normalize_positive_int(download_row.get("user_id")) + if owner_user_id != actor.db_user_id: + raise RuntimeError(f"Viewer state out of scope for {item_key}") + + payload.append( + DownloadHistoryService.to_history_row( + download_row, + dismissed_at=dismissed_at, + ) + ) + continue + + if item_type == "request": + request_id = normalize_positive_int(_parse_item_key(item_key, "request")) + if request_id is None: + raise RuntimeError(f"Invalid activity history item_key: {item_key}") + + request_row = user_db.get_request(request_id) + if request_row is None: + raise RuntimeError(f"Request row not found for {item_key}") + + if not actor.is_admin: + owner_user_id = normalize_positive_int(request_row.get("user_id")) + if owner_user_id != actor.db_user_id: + raise RuntimeError(f"Viewer state out of scope for {item_key}") + + populate_request_usernames([request_row], user_db) + entry = _request_history_entry( + request_row, + dismissed_at=dismissed_at, + ) + if entry is None: + raise RuntimeError(f"Failed to build request history entry for {item_key}") + payload.append(entry) + continue + + raise RuntimeError(f"Unknown activity history item_type: {item_type}") + + return jsonify(payload) @app.route("/api/activity/history", methods=["DELETE"]) def api_activity_history_clear(): @@ -666,18 +683,18 @@ def register_activity_routes( if actor_error is not None: return actor_error - deleted_downloads = download_history_service.clear_dismissed(user_id=actor.owner_scope) - deleted_requests = user_db.delete_dismissed_requests(user_id=actor.owner_scope) - deleted_count = deleted_downloads + deleted_requests + cleared_count = activity_view_state_service.clear_history( + viewer_scope=actor.viewer_scope, + ) - room = _activity_ws_room(is_no_auth=actor.is_no_auth, actor_db_user_id=actor.db_user_id) + room = _activity_ws_room(actor) emit_ws_event( ws_manager, event_name="activity_update", room=room, payload={ "kind": "history_cleared", - "count": deleted_count, + "count": cleared_count, }, ) - return jsonify({"status": "cleared", "deleted_count": deleted_count}) + return jsonify({"status": "cleared", "cleared_count": cleared_count}) diff --git a/shelfmark/core/activity_view_state_service.py b/shelfmark/core/activity_view_state_service.py new file mode 100644 index 00000000..4539e53f --- /dev/null +++ b/shelfmark/core/activity_view_state_service.py @@ -0,0 +1,310 @@ +"""Persistence helpers for per-viewer activity visibility state.""" + +from __future__ import annotations + +import sqlite3 +import threading +from typing import Any + +from shelfmark.core.request_helpers import now_utc_iso + + +VALID_ACTIVITY_ITEM_TYPES = frozenset({"download", "request"}) +ADMIN_VIEWER_SCOPE = "admin:shared" +NOAUTH_VIEWER_SCOPE = "noauth:shared" +USER_VIEWER_SCOPE_PREFIX = "user:" + + +def user_viewer_scope(user_id: int) -> str: + if not isinstance(user_id, int) or user_id < 1: + raise ValueError("user_id must be a positive integer") + return f"{USER_VIEWER_SCOPE_PREFIX}{user_id}" + + +def normalize_viewer_scope(viewer_scope: Any) -> str: + if not isinstance(viewer_scope, str) or not viewer_scope.strip(): + raise ValueError("viewer_scope must be a non-empty string") + + normalized = viewer_scope.strip() + if normalized in {ADMIN_VIEWER_SCOPE, NOAUTH_VIEWER_SCOPE}: + return normalized + + if not normalized.startswith(USER_VIEWER_SCOPE_PREFIX): + raise ValueError( + "viewer_scope must be one of: admin:shared, noauth:shared, or user:" + ) + + raw_user_id = normalized[len(USER_VIEWER_SCOPE_PREFIX):].strip() + try: + parsed_user_id = int(raw_user_id) + except (TypeError, ValueError) as exc: + raise ValueError("viewer_scope user id must be a positive integer") from exc + + return user_viewer_scope(parsed_user_id) + + +def _normalize_item_type(item_type: Any) -> str: + if not isinstance(item_type, str) or not item_type.strip(): + raise ValueError("item_type must be a non-empty string") + normalized = item_type.strip().lower() + if normalized not in VALID_ACTIVITY_ITEM_TYPES: + raise ValueError("item_type must be one of: download, request") + return normalized + + +def _normalize_item_key(item_key: Any, *, item_type: str) -> str: + if not isinstance(item_key, str) or not item_key.strip(): + raise ValueError("item_key must be a non-empty string") + + normalized = item_key.strip() + expected_prefix = f"{item_type}:" + if not normalized.startswith(expected_prefix): + raise ValueError(f"item_key must be in the format {expected_prefix}") + if not normalized.split(":", 1)[1].strip(): + raise ValueError(f"item_key must be in the format {expected_prefix}") + return normalized + + +class ActivityViewStateService: + """Service for per-viewer activity dismissal and history visibility.""" + + def __init__(self, db_path: str): + self._db_path = db_path + self._lock = threading.Lock() + + def _connect(self) -> sqlite3.Connection: + conn = sqlite3.connect(self._db_path) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA foreign_keys = ON") + return conn + + def list_hidden( + self, + *, + viewer_scope: str, + limit: int | None = None, + ) -> list[dict[str, Any]]: + normalized_scope = normalize_viewer_scope(viewer_scope) + normalized_limit = None if limit is None else max(1, int(limit)) + query = """ + SELECT item_type, item_key, dismissed_at, cleared_at + FROM activity_view_state + WHERE viewer_scope = ? + AND dismissed_at IS NOT NULL + ORDER BY COALESCE(cleared_at, dismissed_at) DESC, id DESC + """ + params: list[Any] = [normalized_scope] + if normalized_limit is not None: + query += "\nLIMIT ?" + params.append(normalized_limit) + + conn = self._connect() + try: + rows = conn.execute(query, params).fetchall() + return [dict(row) for row in rows] + finally: + conn.close() + + def list_history( + self, + *, + viewer_scope: str, + limit: int = 50, + offset: int = 0, + ) -> list[dict[str, Any]]: + normalized_scope = normalize_viewer_scope(viewer_scope) + normalized_limit = max(1, min(int(limit), 5000)) + normalized_offset = max(0, int(offset)) + + conn = self._connect() + try: + rows = conn.execute( + """ + SELECT item_type, item_key, dismissed_at + FROM activity_view_state + WHERE viewer_scope = ? + AND dismissed_at IS NOT NULL + AND cleared_at IS NULL + ORDER BY dismissed_at DESC, id DESC + LIMIT ? OFFSET ? + """, + (normalized_scope, normalized_limit, normalized_offset), + ).fetchall() + return [dict(row) for row in rows] + finally: + conn.close() + + def dismiss( + self, + *, + viewer_scope: str, + item_type: str, + item_key: str, + ) -> int: + normalized_scope = normalize_viewer_scope(viewer_scope) + normalized_type = _normalize_item_type(item_type) + normalized_key = _normalize_item_key(item_key, item_type=normalized_type) + dismissed_at = now_utc_iso() + + with self._lock: + conn = self._connect() + try: + cursor = conn.execute( + """ + INSERT INTO activity_view_state ( + viewer_scope, + item_type, + item_key, + dismissed_at, + cleared_at + ) + VALUES (?, ?, ?, ?, NULL) + ON CONFLICT(viewer_scope, item_type, item_key) DO UPDATE SET + dismissed_at = excluded.dismissed_at, + cleared_at = NULL + """, + (normalized_scope, normalized_type, normalized_key, dismissed_at), + ) + conn.commit() + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + return max(rowcount, 0) + finally: + conn.close() + + def dismiss_many( + self, + *, + viewer_scope: str, + items: list[dict[str, str]], + ) -> int: + normalized_scope = normalize_viewer_scope(viewer_scope) + if not items: + return 0 + + seen: set[tuple[str, str]] = set() + normalized_items: list[tuple[str, str]] = [] + for item in items: + normalized_type = _normalize_item_type(item.get("item_type")) + normalized_key = _normalize_item_key(item.get("item_key"), item_type=normalized_type) + marker = (normalized_type, normalized_key) + if marker in seen: + continue + seen.add(marker) + normalized_items.append(marker) + + if not normalized_items: + return 0 + + dismissed_at = now_utc_iso() + with self._lock: + conn = self._connect() + try: + total = 0 + for normalized_type, normalized_key in normalized_items: + cursor = conn.execute( + """ + INSERT INTO activity_view_state ( + viewer_scope, + item_type, + item_key, + dismissed_at, + cleared_at + ) + VALUES (?, ?, ?, ?, NULL) + ON CONFLICT(viewer_scope, item_type, item_key) DO UPDATE SET + dismissed_at = excluded.dismissed_at, + cleared_at = NULL + """, + (normalized_scope, normalized_type, normalized_key, dismissed_at), + ) + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + total += max(rowcount, 0) + conn.commit() + return total + finally: + conn.close() + + def clear_history(self, *, viewer_scope: str) -> int: + normalized_scope = normalize_viewer_scope(viewer_scope) + cleared_at = now_utc_iso() + + with self._lock: + conn = self._connect() + try: + cursor = conn.execute( + """ + UPDATE activity_view_state + SET cleared_at = ? + WHERE viewer_scope = ? + AND dismissed_at IS NOT NULL + AND cleared_at IS NULL + """, + (cleared_at, normalized_scope), + ) + conn.commit() + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + return max(rowcount, 0) + finally: + conn.close() + + def clear_item_for_all_viewers(self, *, item_type: str, item_key: str) -> int: + normalized_type = _normalize_item_type(item_type) + normalized_key = _normalize_item_key(item_key, item_type=normalized_type) + + with self._lock: + conn = self._connect() + try: + cursor = conn.execute( + """ + DELETE FROM activity_view_state + WHERE item_type = ? AND item_key = ? + """, + (normalized_type, normalized_key), + ) + conn.commit() + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + return max(rowcount, 0) + finally: + conn.close() + + def delete_viewer_scope(self, *, viewer_scope: str) -> int: + normalized_scope = normalize_viewer_scope(viewer_scope) + + with self._lock: + conn = self._connect() + try: + cursor = conn.execute( + "DELETE FROM activity_view_state WHERE viewer_scope = ?", + (normalized_scope,), + ) + conn.commit() + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + return max(rowcount, 0) + finally: + conn.close() + + def delete_items(self, *, item_type: str, item_keys: list[str]) -> int: + normalized_type = _normalize_item_type(item_type) + normalized_keys = [ + _normalize_item_key(item_key, item_type=normalized_type) + for item_key in item_keys + ] + if not normalized_keys: + return 0 + + placeholders = ",".join("?" for _ in normalized_keys) + with self._lock: + conn = self._connect() + try: + cursor = conn.execute( + f""" + DELETE FROM activity_view_state + WHERE item_type = ? AND item_key IN ({placeholders}) + """, + (normalized_type, *normalized_keys), + ) + conn.commit() + rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 + return max(rowcount, 0) + finally: + conn.close() diff --git a/shelfmark/core/download_history_service.py b/shelfmark/core/download_history_service.py index 93020fa4..d639a394 100644 --- a/shelfmark/core/download_history_service.py +++ b/shelfmark/core/download_history_service.py @@ -1,4 +1,4 @@ -"""Persistence helpers for flat download terminal history.""" +"""Persistence helpers for canonical download activity rows.""" from __future__ import annotations @@ -61,20 +61,8 @@ def _normalize_limit(value: Any, *, default: int, minimum: int, maximum: int) -> return parsed -def _normalize_offset(value: Any, *, default: int) -> int: - if value is None: - return default - try: - parsed = int(value) - except (TypeError, ValueError) as exc: - raise ValueError("offset must be an integer") from exc - if parsed < 0: - return 0 - return parsed - - class DownloadHistoryService: - """Service for persisted terminal download history and dismissals.""" + """Service for persisted canonical download activity rows.""" def __init__(self, db_path: str): self._db_path = db_path @@ -132,7 +120,7 @@ class DownloadHistoryService: return None @classmethod - def _to_history_row(cls, row: dict[str, Any]) -> dict[str, Any]: + def to_history_row(cls, row: dict[str, Any], *, dismissed_at: str) -> dict[str, Any]: task_id = str(row.get("task_id") or "").strip() item_key = cls._to_item_key(task_id) download_payload = cls.to_download_payload(row) @@ -144,7 +132,7 @@ class DownloadHistoryService: "user_id": row.get("user_id"), "item_type": "download", "item_key": item_key, - "dismissed_at": row.get("dismissed_at"), + "dismissed_at": dismissed_at, "snapshot": { "kind": "download", "download": download_payload, @@ -287,7 +275,7 @@ class DownloadHistoryService: finally: conn.close() - def get_undismissed( + def list_recent( self, *, user_id: int | None, @@ -295,10 +283,10 @@ class DownloadHistoryService: ) -> list[dict[str, Any]]: normalized_user_id = normalize_optional_positive_int(user_id, "user_id") normalized_limit = _normalize_limit(limit, default=200, minimum=1, maximum=1000) - query = "SELECT * FROM download_history WHERE dismissed_at IS NULL" + query = "SELECT * FROM download_history" params: list[Any] = [] if normalized_user_id is not None: - query += " AND user_id = ?" + query += " WHERE user_id = ?" params.append(normalized_user_id) query += " ORDER BY terminal_at DESC, id DESC LIMIT ?" params.append(normalized_limit) @@ -309,133 +297,3 @@ class DownloadHistoryService: return [dict(row) for row in rows] finally: conn.close() - - def get_dismissed_keys(self, *, user_id: int | None, limit: int = 5000) -> list[str]: - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - normalized_limit = _normalize_limit(limit, default=5000, minimum=1, maximum=10000) - query = "SELECT task_id FROM download_history WHERE dismissed_at IS NOT NULL" - params: list[Any] = [] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - query += " ORDER BY dismissed_at DESC, id DESC LIMIT ?" - params.append(normalized_limit) - - conn = self._connect() - try: - rows = conn.execute(query, params).fetchall() - keys: list[str] = [] - for row in rows: - task_id = normalize_optional_text(row["task_id"]) - if task_id is not None: - keys.append(task_id) - return keys - finally: - conn.close() - - def dismiss(self, *, task_id: str, user_id: int | None) -> int: - normalized_task_id = _normalize_task_id(task_id) - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - query = "UPDATE download_history SET dismissed_at = ? WHERE task_id = ?" - params: list[Any] = [now_utc_iso(), normalized_task_id] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() - - def dismiss_many(self, *, task_ids: list[str], user_id: int | None) -> int: - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - normalized_task_ids = [_normalize_task_id(task_id) for task_id in task_ids] - if not normalized_task_ids: - return 0 - - placeholders = ",".join("?" for _ in normalized_task_ids) - query = f"UPDATE download_history SET dismissed_at = ? WHERE task_id IN ({placeholders})" - params: list[Any] = [now_utc_iso(), *normalized_task_ids] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() - - def get_history(self, *, user_id: int | None, limit: int, offset: int) -> list[dict[str, Any]]: - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - normalized_limit = _normalize_limit(limit, default=50, minimum=1, maximum=5000) - normalized_offset = _normalize_offset(offset, default=0) - - query = "SELECT * FROM download_history WHERE dismissed_at IS NOT NULL" - params: list[Any] = [] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - query += " ORDER BY dismissed_at DESC, id DESC LIMIT ? OFFSET ?" - params.extend([normalized_limit, normalized_offset]) - - conn = self._connect() - try: - rows = conn.execute(query, params).fetchall() - payload: list[dict[str, Any]] = [] - for row in rows: - row_dict = dict(row) - payload.append(self._to_history_row(row_dict)) - return payload - finally: - conn.close() - - def clear_dismissed(self, *, user_id: int | None) -> int: - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - query = "DELETE FROM download_history WHERE dismissed_at IS NOT NULL" - params: list[Any] = [] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() - - def clear_dismissals_for_active(self, *, task_ids: set[str], user_id: int | None) -> int: - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - normalized_task_ids = [_normalize_task_id(task_id) for task_id in task_ids] - if not normalized_task_ids: - return 0 - - placeholders = ",".join("?" for _ in normalized_task_ids) - query = f"UPDATE download_history SET dismissed_at = NULL WHERE task_id IN ({placeholders})" - params: list[Any] = [*normalized_task_ids] - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() diff --git a/shelfmark/core/queue.py b/shelfmark/core/queue.py index c2725f02..8fd5f11b 100644 --- a/shelfmark/core/queue.py +++ b/shelfmark/core/queue.py @@ -8,8 +8,11 @@ from threading import Lock, Event from typing import Dict, List, Optional, Tuple, Any, Callable from shelfmark.core.config import config as app_config +from shelfmark.core.logger import setup_logger from shelfmark.core.models import QueueStatus, QueueItem, DownloadTask, TERMINAL_QUEUE_STATUSES +logger = setup_logger(__name__) + class BookQueue: """Thread-safe download queue manager with priority support and cancellation.""" @@ -55,8 +58,8 @@ class BookQueue: if hook is not None: try: hook(task_id, task) - except Exception: - pass + except Exception as exc: + logger.warning("Queue hook failed while adding task %s: %s", task_id, exc) return True def get_next(self) -> Optional[Tuple[str, Event]]: @@ -294,8 +297,8 @@ class BookQueue: if hook is not None and hook_task is not None: try: hook(task_id, hook_task) - except Exception: - pass + except Exception as exc: + logger.warning("Queue hook failed while requeueing task %s: %s", task_id, exc) return True def reorder_queue(self, task_priorities: Dict[str, int]) -> bool: diff --git a/shelfmark/core/request_helpers.py b/shelfmark/core/request_helpers.py index 9cb54dba..c75a8a87 100644 --- a/shelfmark/core/request_helpers.py +++ b/shelfmark/core/request_helpers.py @@ -65,15 +65,9 @@ def coerce_int(value: Any, default: int) -> int: def normalize_optional_text(value: Any) -> str | None: - """Return a trimmed string or None for empty/missing input. - - Non-string values are coerced via ``str()`` before stripping; - ``None`` short-circuits to ``None``. - """ - if value is None: - return None + """Return a trimmed string or None for empty/non-string input.""" if not isinstance(value, str): - value = str(value) + return None normalized = value.strip() return normalized or None diff --git a/shelfmark/core/request_routes.py b/shelfmark/core/request_routes.py index 4a97dc8a..8714890a 100644 --- a/shelfmark/core/request_routes.py +++ b/shelfmark/core/request_routes.py @@ -123,15 +123,27 @@ def _resolve_title_from_book_data(book_data: Any) -> str: return "Unknown title" +def _normalize_optional_source_id(value: Any) -> str | None: + """Normalize source identifiers while allowing integer provider ids.""" + if isinstance(value, bool) or value is None: + return None + if isinstance(value, int): + value = str(value) + return normalize_optional_text(value) + + def _build_direct_release_data_from_book_data( *, book_data: dict[str, Any], content_type: str, ) -> dict[str, Any]: """Build release-level payload fields for direct-download requests.""" + source_id = _normalize_optional_source_id(book_data.get("provider_id")) or _normalize_optional_source_id( + book_data.get("id") + ) payload: dict[str, Any] = { "source": "direct_download", - "source_id": book_data.get("provider_id") or book_data.get("id"), + "source_id": source_id, "title": book_data.get("title"), "author": book_data.get("author"), "year": book_data.get("year"), @@ -169,8 +181,11 @@ def _normalize_direct_request_payload( if normalized_release_data.get("content_type") is None: normalized_release_data["content_type"] = content_type - if normalize_optional_text(normalized_release_data.get("source_id")) is None and isinstance(book_data, dict): - fallback_source_id = normalize_optional_text(book_data.get("provider_id")) or normalize_optional_text( + normalized_source_id = _normalize_optional_source_id(normalized_release_data.get("source_id")) + if normalized_source_id is not None: + normalized_release_data["source_id"] = normalized_source_id + elif isinstance(book_data, dict): + fallback_source_id = _normalize_optional_source_id(book_data.get("provider_id")) or _normalize_optional_source_id( book_data.get("id") ) if fallback_source_id is not None: diff --git a/shelfmark/core/user_db.py b/shelfmark/core/user_db.py index 3edf6c56..6b6ee9ed 100644 --- a/shelfmark/core/user_db.py +++ b/shelfmark/core/user_db.py @@ -7,6 +7,7 @@ import threading from typing import Any, Dict, List, Optional from shelfmark.core.auth_modes import AUTH_SOURCE_BUILTIN, AUTH_SOURCE_SET +from shelfmark.core.activity_view_state_service import user_viewer_scope from shelfmark.core.logger import setup_logger from shelfmark.core.request_helpers import normalize_optional_positive_int from shelfmark.core.models import QueueStatus @@ -57,8 +58,7 @@ CREATE TABLE IF NOT EXISTS download_requests ( reviewed_by INTEGER REFERENCES users(id), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, reviewed_at TIMESTAMP, - delivery_updated_at TIMESTAMP, - dismissed_at TIMESTAMP + delivery_updated_at TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_download_requests_user_status_created_at @@ -86,19 +86,32 @@ CREATE TABLE IF NOT EXISTS download_history ( status_message TEXT, download_path TEXT, queued_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, - terminal_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, - dismissed_at TIMESTAMP + terminal_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_download_history_user_status ON download_history (user_id, final_status, terminal_at DESC); -CREATE INDEX IF NOT EXISTS idx_download_history_dismissed -ON download_history (dismissed_at) WHERE dismissed_at IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_download_history_recent +ON download_history (user_id, terminal_at DESC, id DESC); -CREATE INDEX IF NOT EXISTS idx_download_history_undismissed -ON download_history (user_id, terminal_at DESC, id DESC) -WHERE dismissed_at IS NULL; +CREATE TABLE IF NOT EXISTS activity_view_state ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + viewer_scope TEXT NOT NULL, + item_type TEXT NOT NULL, + item_key TEXT NOT NULL, + dismissed_at TIMESTAMP, + cleared_at TIMESTAMP, + UNIQUE(viewer_scope, item_type, item_key) +); + +CREATE INDEX IF NOT EXISTS idx_activity_view_state_history +ON activity_view_state (viewer_scope, dismissed_at DESC, id DESC) +WHERE dismissed_at IS NOT NULL AND cleared_at IS NULL; + +CREATE INDEX IF NOT EXISTS idx_activity_view_state_hidden +ON activity_view_state (viewer_scope, item_type, item_key) +WHERE dismissed_at IS NOT NULL; """ @@ -168,7 +181,6 @@ class UserDB: conn.executescript(_CREATE_TABLES_SQL) self._migrate_auth_source_column(conn) self._migrate_request_delivery_columns(conn) - self._migrate_download_requests_dismissed_at(conn) self._migrate_download_history_queued_at(conn) conn.commit() # WAL mode must be changed outside an open transaction. @@ -224,17 +236,6 @@ class UserDB: """ ) - def _migrate_download_requests_dismissed_at(self, conn: sqlite3.Connection) -> None: - """Ensure download_requests.dismissed_at exists for request dismissal history.""" - columns = conn.execute("PRAGMA table_info(download_requests)").fetchall() - column_names = {str(col["name"]) for col in columns} - if "dismissed_at" not in column_names: - conn.execute("ALTER TABLE download_requests ADD COLUMN dismissed_at TIMESTAMP") - conn.execute( - "CREATE INDEX IF NOT EXISTS idx_download_requests_dismissed " - "ON download_requests (dismissed_at) WHERE dismissed_at IS NOT NULL" - ) - def _migrate_download_history_queued_at(self, conn: sqlite3.Connection) -> None: """Ensure download_history.queued_at exists for queue-time recording.""" columns = conn.execute("PRAGMA table_info(download_history)").fetchall() @@ -349,6 +350,25 @@ class UserDB: with self._lock: conn = self._connect() try: + request_rows = conn.execute( + "SELECT id FROM download_requests WHERE user_id = ?", + (user_id,), + ).fetchall() + request_item_keys = [f"request:{row['id']}" for row in request_rows] + if request_item_keys: + placeholders = ",".join("?" for _ in request_item_keys) + conn.execute( + f""" + DELETE FROM activity_view_state + WHERE item_type = 'request' + AND item_key IN ({placeholders}) + """, + request_item_keys, + ) + conn.execute( + "DELETE FROM activity_view_state WHERE viewer_scope = ?", + (user_viewer_scope(user_id),), + ) conn.execute("UPDATE download_requests SET reviewed_by = NULL WHERE reviewed_by = ?", (user_id,)) conn.execute("DELETE FROM users WHERE id = ?", (user_id,)) conn.commit() @@ -600,7 +620,6 @@ class UserDB: "delivery_state", "delivery_updated_at", "last_failure_reason", - "dismissed_at", } def update_request( @@ -660,11 +679,6 @@ class UserDB: if delivery_updated_at is not None and not isinstance(delivery_updated_at, str): raise ValueError("delivery_updated_at must be a string when provided") - if "dismissed_at" in updates: - dismissed_at = updates["dismissed_at"] - if dismissed_at is not None and not isinstance(dismissed_at, str): - raise ValueError("dismissed_at must be a string when provided") - if "content_type" in updates and not updates["content_type"]: raise ValueError("content_type is required") @@ -785,69 +799,3 @@ class UserDB: return int(row["count"]) if row else 0 finally: conn.close() - - def list_dismissed_requests(self, *, user_id: int | None, limit: int | None = None) -> List[Dict[str, Any]]: - """List dismissed requests, optionally scoped by owner user_id.""" - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - params: list[Any] = [] - query = "SELECT * FROM download_requests WHERE dismissed_at IS NOT NULL" - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - query += " ORDER BY dismissed_at DESC, id DESC" - if limit is not None and limit > 0: - query += " LIMIT ?" - params.append(limit) - - conn = self._connect() - try: - rows = conn.execute(query, params).fetchall() - results: List[Dict[str, Any]] = [] - for row in rows: - parsed = self._parse_request_row(row) - if parsed is not None: - results.append(parsed) - return results - finally: - conn.close() - - def dismiss_requests_batch(self, *, request_ids: list[int], dismissed_at: str) -> int: - """Set dismissed_at on multiple requests in a single UPDATE.""" - if not request_ids: - return 0 - placeholders = ",".join("?" for _ in request_ids) - query = f"UPDATE download_requests SET dismissed_at = ? WHERE id IN ({placeholders})" - params: list[Any] = [dismissed_at, *request_ids] - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() - - def delete_dismissed_requests(self, *, user_id: int | None) -> int: - """Delete dismissed terminal requests, optionally scoped by owner user_id.""" - normalized_user_id = normalize_optional_positive_int(user_id, "user_id") - params: list[Any] = [] - query = """ - DELETE FROM download_requests - WHERE dismissed_at IS NOT NULL - AND status IN ('fulfilled', 'rejected', 'cancelled') - """ - if normalized_user_id is not None: - query += " AND user_id = ?" - params.append(normalized_user_id) - - with self._lock: - conn = self._connect() - try: - cursor = conn.execute(query, params) - conn.commit() - rowcount = int(cursor.rowcount) if cursor.rowcount is not None else 0 - return max(rowcount, 0) - finally: - conn.close() diff --git a/shelfmark/main.py b/shelfmark/main.py index d20e8e21..1d880ec3 100644 --- a/shelfmark/main.py +++ b/shelfmark/main.py @@ -50,9 +50,11 @@ from shelfmark.core.requests_service import ( reopen_failed_request, sync_delivery_states_from_queue_status, ) +from shelfmark.core.activity_view_state_service import ActivityViewStateService from shelfmark.core.download_history_service import DownloadHistoryService from shelfmark.core.notifications import NotificationContext, NotificationEvent, notify_admin, notify_user from shelfmark.core.request_helpers import ( + emit_ws_event, coerce_bool, load_users_request_policy_settings, normalize_optional_text, @@ -129,10 +131,12 @@ from shelfmark.core.user_db import UserDB _user_db_path = _os.path.join(_os.environ.get("CONFIG_DIR", "/config"), "users.db") user_db: UserDB | None = None download_history_service: DownloadHistoryService | None = None +activity_view_state_service: ActivityViewStateService | None = None try: user_db = UserDB(_user_db_path) user_db.initialize() download_history_service = DownloadHistoryService(_user_db_path) + activity_view_state_service = ActivityViewStateService(_user_db_path) import shelfmark.config.users_settings as _ # noqa: F401 - registers users tab from shelfmark.core.oidc_routes import register_oidc_routes from shelfmark.core.admin_routes import register_admin_routes @@ -148,6 +152,7 @@ except (sqlite3.OperationalError, OSError) as e: ) user_db = None download_history_service = None + activity_view_state_service = None # Start download coordinator backend.start() @@ -408,13 +413,13 @@ if user_db is not None: queue_release=lambda *args, **kwargs: backend.queue_release(*args, **kwargs), ws_manager=ws_manager, ) - if download_history_service is not None: + if download_history_service is not None and activity_view_state_service is not None: register_activity_routes( app, user_db, + activity_view_state_service=activity_view_state_service, download_history_service=download_history_service, resolve_auth_mode=lambda: get_auth_mode(), - resolve_status_scope=lambda: _resolve_status_scope(), queue_status=lambda user_id=None: backend.queue_status(user_id=user_id), sync_request_delivery_states=sync_delivery_states_from_queue_status, emit_request_updates=lambda rows: _emit_request_update_events(rows), @@ -1164,6 +1169,24 @@ def _notify_admin_for_terminal_download_status(*, task_id: str, status: QueueSta ) +def _emit_activity_update_for_task(*, payload: dict[str, Any], task: Any) -> None: + owner_user_id = normalize_positive_int(getattr(task, "user_id", None)) + emit_ws_event( + ws_manager, + event_name="activity_update", + room="admins", + payload=payload, + ) + if owner_user_id is None: + return + emit_ws_event( + ws_manager, + event_name="activity_update", + room=f"user_{owner_user_id}", + payload=payload, + ) + + def _record_download_queued(task_id: str, task: Any) -> None: """Persist initial download record when a task enters the queue.""" if download_history_service is None: @@ -1197,6 +1220,32 @@ def _record_download_queued(task_id: str, task: Any) -> None: ) except Exception as exc: logger.warning("Failed to record download at queue time for task %s: %s", task_id, exc) + return + + if activity_view_state_service is None: + return + + try: + cleared_view_state = 0 + cleared_view_state += activity_view_state_service.clear_item_for_all_viewers( + item_type="download", + item_key=f"download:{task_id}", + ) + if request_id is not None: + cleared_view_state += activity_view_state_service.clear_item_for_all_viewers( + item_type="request", + item_key=f"request:{request_id}", + ) + if cleared_view_state > 0: + _emit_activity_update_for_task( + task=task, + payload={ + "kind": "activity_reset", + "task_id": task_id, + }, + ) + except Exception as exc: + logger.warning("Failed to reset activity viewer state for task %s: %s", task_id, exc) def _record_download_terminal_snapshot(task_id: str, status: QueueStatus, task: Any) -> None: @@ -1206,6 +1255,7 @@ def _record_download_terminal_snapshot(task_id: str, status: QueueStatus, task: if final_status is None: return + finalized_download = False if download_history_service is not None: try: download_history_service.finalize_download( @@ -1214,9 +1264,20 @@ def _record_download_terminal_snapshot(task_id: str, status: QueueStatus, task: status_message=normalize_optional_text(getattr(task, "status_message", None)), download_path=normalize_optional_text(getattr(task, "download_path", None)), ) + finalized_download = True except Exception as exc: logger.warning("Failed to finalize download history for task %s: %s", task_id, exc) + if finalized_download: + _emit_activity_update_for_task( + task=task, + payload={ + "kind": "download_terminal", + "task_id": task_id, + "status": final_status, + }, + ) + if user_db is None or status != QueueStatus.ERROR: return @@ -1237,6 +1298,11 @@ def _record_download_terminal_snapshot(task_id: str, status: QueueStatus, task: failure_reason=fallback_reason, ) if reopened_request is not None: + if activity_view_state_service is not None: + activity_view_state_service.clear_item_for_all_viewers( + item_type="request", + item_key=f"request:{request_id}", + ) _emit_request_update_events([reopened_request]) except Exception as exc: logger.warning( diff --git a/src/frontend/src/App.tsx b/src/frontend/src/App.tsx index 0b19dbff..a06b67af 100644 --- a/src/frontend/src/App.tsx +++ b/src/frontend/src/App.tsx @@ -135,17 +135,6 @@ type PendingOnBehalfDownload = actingAsUser: ActingAsUserSelection; }; -const mergeTerminalBucket = ( - persistedBucket: Record | undefined, - realtimeBucket: Record | undefined -): Record | undefined => { - const merged = { - ...(persistedBucket || {}), - ...(realtimeBucket || {}), - }; - return Object.keys(merged).length > 0 ? merged : undefined; -}; - function App() { const { toasts, showToast, removeToast } = useToast(); const { socket } = useSocket(); @@ -263,10 +252,12 @@ function App() { requestItems, dismissedActivityKeys, historyItems, + activityHistoryLoaded, pendingRequestCount, isActivitySnapshotLoading, activityHistoryLoading, activityHistoryHasMore, + prefetchActivityHistory, refreshActivitySnapshot, resetActivity, handleActivityTabChange, @@ -319,9 +310,9 @@ function App() { }; }, [currentStatus, dismissedDownloadTaskIds]); - // Use real-time buckets for active work and merge persisted terminal buckets - // so completed/errored entries survive restarts. Filter out dismissed items - // so the sidebar counts stay consistent with the activity panel. + // Use real-time buckets for active work and persisted activity snapshot + // buckets for terminal history. Filter out dismissed items so the sidebar + // counts stay consistent with the activity panel. const activitySidebarStatus = useMemo(() => { const filterDismissed = ( bucket: Record | undefined @@ -338,9 +329,9 @@ function App() { resolving: currentStatus.resolving, locating: currentStatus.locating, downloading: currentStatus.downloading, - complete: filterDismissed(mergeTerminalBucket(activityStatus.complete, currentStatus.complete)), - error: filterDismissed(mergeTerminalBucket(activityStatus.error, currentStatus.error)), - cancelled: filterDismissed(mergeTerminalBucket(activityStatus.cancelled, currentStatus.cancelled)), + complete: filterDismissed(activityStatus.complete), + error: filterDismissed(activityStatus.error), + cancelled: filterDismissed(activityStatus.cancelled), }; }, [activityStatus, currentStatus, dismissedDownloadTaskIds]); @@ -430,6 +421,13 @@ function App() { const [sidebarPinnedOpen, setSidebarPinnedOpen] = useState(false); const [headerHeight, setHeaderHeight] = useState(0); const headerObserverRef = useRef(null); + useEffect(() => { + if (!downloadsSidebarOpen) { + return; + } + prefetchActivityHistory(); + }, [downloadsSidebarOpen, prefetchActivityHistory]); + const headerRef = useCallback((el: HTMLDivElement | null) => { if (headerObserverRef.current) { headerObserverRef.current.disconnect(); @@ -1696,6 +1694,7 @@ function App() { requestItems={requestItems} dismissedItemKeys={dismissedActivityKeys} historyItems={historyItems} + historyLoaded={activityHistoryLoaded} historyHasMore={activityHistoryHasMore} historyLoading={activityHistoryLoading} onHistoryLoadMore={handleActivityHistoryLoadMore} diff --git a/src/frontend/src/components/ReleaseModal.tsx b/src/frontend/src/components/ReleaseModal.tsx index f26df459..58e8d2a3 100644 --- a/src/frontend/src/components/ReleaseModal.tsx +++ b/src/frontend/src/components/ReleaseModal.tsx @@ -402,17 +402,20 @@ function ShimmerBlock({ className }: { className: string }) { } // Loading skeleton for releases - matches ReleaseRow layout +// Renders enough rows to fill the container, fading out at the bottom via a gradient mask function ReleaseSkeleton() { + // Render enough rows to cover tall viewports; overflow is hidden by the mask + const rows = 8; return ( -
- {[1, 2, 3, 4, 5].map((i) => ( -
+
+ {Array.from({ length: rows }, (_, i) => ( +
{/* Thumbnail skeleton */} diff --git a/src/frontend/src/components/activity/ActivitySidebar.tsx b/src/frontend/src/components/activity/ActivitySidebar.tsx index fb438cf5..9814ebb8 100644 --- a/src/frontend/src/components/activity/ActivitySidebar.tsx +++ b/src/frontend/src/components/activity/ActivitySidebar.tsx @@ -1,4 +1,4 @@ -import { useEffect, useMemo, useRef, useState, type WheelEvent } from 'react'; +import { useCallback, useEffect, useMemo, useRef, useState, type WheelEvent } from 'react'; import { RequestRecord, StatusData } from '../../types'; import { downloadToActivityItem, DownloadStatusKey } from './activityMappers'; import { ActivityItem } from './activityTypes'; @@ -17,6 +17,7 @@ interface ActivitySidebarProps { requestItems: ActivityItem[]; dismissedItemKeys?: string[]; historyItems?: ActivityItem[]; + historyLoaded?: boolean; historyHasMore?: boolean; historyLoading?: boolean; onHistoryLoadMore?: () => void; @@ -247,6 +248,7 @@ export const ActivitySidebar = ({ requestItems, dismissedItemKeys = [], historyItems = [], + historyLoaded = false, historyHasMore = false, historyLoading = false, onHistoryLoadMore, @@ -274,6 +276,10 @@ export const ActivitySidebar = ({ () => new Set(dismissedItemKeys), [dismissedItemKeys] ); + const handleTabChange = useCallback((nextTab: ActivityTabKey) => { + setActiveTab(nextTab); + onActiveTabChange?.(nextTab); + }, [onActiveTabChange]); useEffect(() => { const mediaQuery = window.matchMedia('(min-width: 1024px)'); @@ -292,13 +298,9 @@ export const ActivitySidebar = ({ useEffect(() => { if (!showRequestsTab && activeTab === 'requests') { - setActiveTab('all'); + handleTabChange('all'); } - }, [showRequestsTab, activeTab]); - - useEffect(() => { - onActiveTabChange?.(activeTab); - }, [activeTab, onActiveTabChange]); + }, [showRequestsTab, activeTab, handleTabChange]); useEffect(() => { if (activeTab === 'downloads') { @@ -452,9 +454,10 @@ export const ActivitySidebar = ({ } return requestStatus === 'fulfilled' && item.kind === 'request'; }) - : activeTab === 'history' - ? historyItems - : mergedDownloadItems; + : activeTab === 'history' + ? historyItems + : mergedDownloadItems; + const isHistoryInitialLoad = activeTab === 'history' && !historyLoaded; const availableUsers = useMemo(() => { const userMap = new Map(); @@ -702,12 +705,12 @@ export const ActivitySidebar = ({ )} )} -