refactor: fix slowness in deleting big number of items
Co-authored-by: Copilot <copilot@github.com>
This commit is contained in:
parent
7d79e0e8b0
commit
30a0c30d94
15 changed files with 631 additions and 108 deletions
31
API.md
31
API.md
|
|
@ -133,6 +133,7 @@ This document describes the available endpoints and their usage. All endpoints r
|
||||||
- [`item_updated`](#item_updated)
|
- [`item_updated`](#item_updated)
|
||||||
- [`item_cancelled`](#item_cancelled)
|
- [`item_cancelled`](#item_cancelled)
|
||||||
- [`item_deleted`](#item_deleted)
|
- [`item_deleted`](#item_deleted)
|
||||||
|
- [`item_bulk_deleted`](#item_bulk_deleted)
|
||||||
- [`item_moved`](#item_moved)
|
- [`item_moved`](#item_moved)
|
||||||
- [`item_status`](#item_status)
|
- [`item_status`](#item_status)
|
||||||
- [`paused`](#paused)
|
- [`paused`](#paused)
|
||||||
|
|
@ -382,7 +383,7 @@ or an error:
|
||||||
**Body Parameters**:
|
**Body Parameters**:
|
||||||
- `type` (string, required): Type of items - `"queue"` or `"done"`
|
- `type` (string, required): Type of items - `"queue"` or `"done"`
|
||||||
- `ids` (array, optional): List of specific item IDs to delete. If provided, `status` filter is ignored
|
- `ids` (array, optional): List of specific item IDs to delete. If provided, `status` filter is ignored
|
||||||
- `status` (string, optional): Filter by status (e.g., `"finished"`, `"!finished"`). Required if `ids` not provided
|
- `status` (string, optional): Filter by status (e.g., `"finished"`, `"!finished"`, `"finished,skip"`, `"!finished,!skip"`). Required if `ids` not provided
|
||||||
- `remove_file` (boolean, optional): Whether to delete files from disk. Default: `true`.
|
- `remove_file` (boolean, optional): Whether to delete files from disk. Default: `true`.
|
||||||
|
|
||||||
> [!NOTE]
|
> [!NOTE]
|
||||||
|
|
@ -417,6 +418,15 @@ or an error:
|
||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
|
**Delete all completed and skipped items in one request:**
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"type": "done",
|
||||||
|
"status": "finished,skip",
|
||||||
|
"remove_file": false
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
**Response**:
|
**Response**:
|
||||||
```json
|
```json
|
||||||
{
|
{
|
||||||
|
|
@ -449,6 +459,7 @@ or an error:
|
||||||
**Notes**:
|
**Notes**:
|
||||||
- When using filter mode, all matching items will be deleted.
|
- When using filter mode, all matching items will be deleted.
|
||||||
- Filter mode with `{ "status": "!finished" }` is useful for cleaning up failed/pending downloads.
|
- Filter mode with `{ "status": "!finished" }` is useful for cleaning up failed/pending downloads.
|
||||||
|
- `status` also accepts comma-separated include filters (`finished,skip`) and comma-separated exclude filters (`!finished,!skip`).
|
||||||
- Filter mode returns a `deleted` count indicating how many items were removed.
|
- Filter mode returns a `deleted` count indicating how many items were removed.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
@ -3177,6 +3188,24 @@ Emitted when a download item is deleted from the queue or history.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
##### `item_bulk_deleted`
|
||||||
|
|
||||||
|
Emitted when multiple items are cleared in bulk operation.
|
||||||
|
|
||||||
|
**Event**:
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"event": "item_bulk_deleted",
|
||||||
|
"data": {
|
||||||
|
"count": 10001,
|
||||||
|
"status": "finished,skip",
|
||||||
|
"ids": ["id1", "id2", "..."] // optional, included when clear is driven by explicit ids
|
||||||
|
}
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
##### `item_moved`
|
##### `item_moved`
|
||||||
|
|
||||||
Emitted when a download item is moved between queue and history.
|
Emitted when a download item is moved between queue and history.
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,7 @@ class NotificationEvents:
|
||||||
ITEM_COMPLETED: str = Events.ITEM_COMPLETED
|
ITEM_COMPLETED: str = Events.ITEM_COMPLETED
|
||||||
ITEM_CANCELLED: str = Events.ITEM_CANCELLED
|
ITEM_CANCELLED: str = Events.ITEM_CANCELLED
|
||||||
ITEM_DELETED: str = Events.ITEM_DELETED
|
ITEM_DELETED: str = Events.ITEM_DELETED
|
||||||
|
ITEM_BULK_DELETED: str = Events.ITEM_BULK_DELETED
|
||||||
ITEM_PAUSED: str = Events.ITEM_PAUSED
|
ITEM_PAUSED: str = Events.ITEM_PAUSED
|
||||||
ITEM_RESUMED: str = Events.ITEM_RESUMED
|
ITEM_RESUMED: str = Events.ITEM_RESUMED
|
||||||
ITEM_MOVED: str = Events.ITEM_MOVED
|
ITEM_MOVED: str = Events.ITEM_MOVED
|
||||||
|
|
|
||||||
|
|
@ -120,6 +120,48 @@ class DataStore:
|
||||||
|
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
async def get_many_by_ids(self, ids: Iterable[str]) -> list[tuple[str, Download]]:
|
||||||
|
ids_list = list(ids)
|
||||||
|
if not ids_list:
|
||||||
|
return []
|
||||||
|
|
||||||
|
items: list[tuple[str, Download]] = []
|
||||||
|
missing_ids: list[str] = []
|
||||||
|
|
||||||
|
for item_id in ids_list:
|
||||||
|
cached = self._dict.get(item_id)
|
||||||
|
if cached:
|
||||||
|
items.append((item_id, cached))
|
||||||
|
continue
|
||||||
|
missing_ids.append(item_id)
|
||||||
|
|
||||||
|
if StoreType.HISTORY == self._type and missing_ids:
|
||||||
|
loaded = await self._connection.get_many_by_ids(str(self._type), missing_ids)
|
||||||
|
for item_id, item in loaded:
|
||||||
|
self._dict[item_id] = Download(info=item)
|
||||||
|
|
||||||
|
items.extend((item_id, download) for item_id in ids_list if (download := self._dict.get(item_id)))
|
||||||
|
|
||||||
|
seen: set[str] = set()
|
||||||
|
ordered: list[tuple[str, Download]] = []
|
||||||
|
for item_id, download in items:
|
||||||
|
if item_id in seen:
|
||||||
|
continue
|
||||||
|
seen.add(item_id)
|
||||||
|
ordered.append((item_id, download))
|
||||||
|
return ordered
|
||||||
|
|
||||||
|
async def get_many_by_status(self, status_filter: str) -> list[tuple[str, Download]]:
|
||||||
|
if StoreType.HISTORY != self._type:
|
||||||
|
return []
|
||||||
|
|
||||||
|
items = await self._connection.get_many_by_status(str(self._type), status_filter)
|
||||||
|
downloads: list[tuple[str, Download]] = []
|
||||||
|
for item_id, item in items:
|
||||||
|
self._dict[item_id] = Download(info=item)
|
||||||
|
downloads.append((item_id, self._dict[item_id]))
|
||||||
|
return downloads
|
||||||
|
|
||||||
def items(self):
|
def items(self):
|
||||||
return self._dict.items()
|
return self._dict.items()
|
||||||
|
|
||||||
|
|
@ -175,11 +217,41 @@ class DataStore:
|
||||||
return [(item_id, Download(info=item)) for item_id, item in items], total_items, current_page, total_pages
|
return [(item_id, Download(info=item)) for item_id, item in items], total_items, current_page, total_pages
|
||||||
|
|
||||||
async def bulk_delete(self, ids: Iterable[str]) -> int:
|
async def bulk_delete(self, ids: Iterable[str]) -> int:
|
||||||
deleted = await self._connection.bulk_delete(str(self._type), ids)
|
ids_list = list(ids)
|
||||||
for _id in ids:
|
deleted = await self._connection.bulk_delete(str(self._type), ids_list)
|
||||||
|
for _id in ids_list:
|
||||||
self._dict.pop(_id, None)
|
self._dict.pop(_id, None)
|
||||||
return deleted
|
return deleted
|
||||||
|
|
||||||
|
async def bulk_delete_by_status(self, status_filter: str) -> int:
|
||||||
|
deleted = await self._connection.bulk_delete_by_status(str(self._type), status_filter)
|
||||||
|
if deleted > 0:
|
||||||
|
self._drop_cached_by_status(status_filter)
|
||||||
|
return deleted
|
||||||
|
|
||||||
|
def _drop_cached_by_status(self, status_filter: str) -> None:
|
||||||
|
raw_statuses = [entry.strip() for entry in status_filter.split(",") if entry.strip()]
|
||||||
|
if not raw_statuses:
|
||||||
|
return
|
||||||
|
|
||||||
|
if all(entry.startswith("!") for entry in raw_statuses):
|
||||||
|
excluded = {entry[1:].strip() for entry in raw_statuses if entry[1:].strip()}
|
||||||
|
if not excluded:
|
||||||
|
return
|
||||||
|
|
||||||
|
for item_id, download in list(self._dict.items()):
|
||||||
|
if download.info and download.info.status not in excluded:
|
||||||
|
self._dict.pop(item_id, None)
|
||||||
|
return
|
||||||
|
|
||||||
|
included = {entry for entry in raw_statuses if not entry.startswith("!")}
|
||||||
|
if not included:
|
||||||
|
return
|
||||||
|
|
||||||
|
for item_id, download in list(self._dict.items()):
|
||||||
|
if download.info and download.info.status in included:
|
||||||
|
self._dict.pop(item_id, None)
|
||||||
|
|
||||||
async def test(self) -> bool:
|
async def test(self) -> bool:
|
||||||
await self._connection.count(str(self._type))
|
await self._connection.count(str(self._type))
|
||||||
return True
|
return True
|
||||||
|
|
|
||||||
|
|
@ -37,6 +37,7 @@ class Events:
|
||||||
ITEM_COMPLETED: str = "item_completed"
|
ITEM_COMPLETED: str = "item_completed"
|
||||||
ITEM_CANCELLED: str = "item_cancelled"
|
ITEM_CANCELLED: str = "item_cancelled"
|
||||||
ITEM_DELETED: str = "item_deleted"
|
ITEM_DELETED: str = "item_deleted"
|
||||||
|
ITEM_BULK_DELETED: str = "item_bulk_deleted"
|
||||||
ITEM_PAUSED: str = "item_paused"
|
ITEM_PAUSED: str = "item_paused"
|
||||||
ITEM_RESUMED: str = "item_resumed"
|
ITEM_RESUMED: str = "item_resumed"
|
||||||
ITEM_MOVED: str = "item_moved"
|
ITEM_MOVED: str = "item_moved"
|
||||||
|
|
@ -87,6 +88,7 @@ class Events:
|
||||||
Events.ITEM_UPDATED,
|
Events.ITEM_UPDATED,
|
||||||
Events.ITEM_CANCELLED,
|
Events.ITEM_CANCELLED,
|
||||||
Events.ITEM_DELETED,
|
Events.ITEM_DELETED,
|
||||||
|
Events.ITEM_BULK_DELETED,
|
||||||
Events.ITEM_MOVED,
|
Events.ITEM_MOVED,
|
||||||
Events.ITEM_STATUS,
|
Events.ITEM_STATUS,
|
||||||
Events.PAUSED,
|
Events.PAUSED,
|
||||||
|
|
|
||||||
|
|
@ -372,6 +372,103 @@ class DownloadQueue(metaclass=Singleton):
|
||||||
|
|
||||||
return status
|
return status
|
||||||
|
|
||||||
|
async def clear_bulk(self, ids: list[str], remove_file: bool = False) -> dict[str, int | str]:
|
||||||
|
if not ids:
|
||||||
|
return {"deleted": 0}
|
||||||
|
|
||||||
|
items = await self.done.get_many_by_ids(ids)
|
||||||
|
if not items:
|
||||||
|
return {"deleted": 0}
|
||||||
|
|
||||||
|
if self.config.remove_files is not True:
|
||||||
|
remove_file = False
|
||||||
|
|
||||||
|
removed_files = 0
|
||||||
|
deleted_ids: list[str] = []
|
||||||
|
deleted_titles: list[str] = []
|
||||||
|
|
||||||
|
for item_id, item in items:
|
||||||
|
item_ref: str = f"{item_id=} {item.info.id=} {item.info.title=}"
|
||||||
|
filename: str = ""
|
||||||
|
|
||||||
|
LOG.debug(f"{remove_file=} {item_ref} - Removing local files: {item.info.status=}")
|
||||||
|
|
||||||
|
if remove_file and "finished" == item.info.status and item.info.filename:
|
||||||
|
filename = str(item.info.filename)
|
||||||
|
if item.info.folder:
|
||||||
|
filename = f"{item.info.folder}/{item.info.filename}"
|
||||||
|
|
||||||
|
try:
|
||||||
|
rf = Path(
|
||||||
|
calc_download_path(
|
||||||
|
base_path=Path(self.config.download_path),
|
||||||
|
folder=filename,
|
||||||
|
create_path=False,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
if rf.is_file() and rf.exists():
|
||||||
|
if rf.stem and rf.suffix:
|
||||||
|
for file_ref in rf.parent.glob(f"{glob.escape(rf.stem)}.*"):
|
||||||
|
if file_ref.is_file() and file_ref.exists() and not file_ref.name.startswith("."):
|
||||||
|
removed_files += 1
|
||||||
|
LOG.debug(f"Removing '{item_ref}' local file '{file_ref.name}'.")
|
||||||
|
file_ref.unlink(missing_ok=True)
|
||||||
|
else:
|
||||||
|
LOG.debug(f"Removing '{item_ref}' local file '{rf.name}'.")
|
||||||
|
rf.unlink(missing_ok=True)
|
||||||
|
removed_files += 1
|
||||||
|
else:
|
||||||
|
LOG.warning(f"Failed to remove '{item_ref}' local file '{filename}'. File not found.")
|
||||||
|
except Exception as e:
|
||||||
|
LOG.error(f"Unable to remove '{item_ref}' local file '{filename}'. {e!s}")
|
||||||
|
|
||||||
|
deleted_ids.append(item_id)
|
||||||
|
deleted_titles.append(item.info.title or item.info.id or item_id)
|
||||||
|
|
||||||
|
deleted_count = await self.done.bulk_delete(deleted_ids)
|
||||||
|
if deleted_count < 1:
|
||||||
|
return {"deleted": 0}
|
||||||
|
|
||||||
|
title = "History Removed" if removed_files > 0 else "History Cleared"
|
||||||
|
message = f"Removed {deleted_count} item{'s' if deleted_count != 1 else ''} from history."
|
||||||
|
if removed_files > 0:
|
||||||
|
message += f" Also removed {removed_files} local file{'s' if removed_files != 1 else ''}."
|
||||||
|
|
||||||
|
self._notify.emit(
|
||||||
|
Events.ITEM_BULK_DELETED,
|
||||||
|
data={"ids": deleted_ids, "count": deleted_count},
|
||||||
|
title=title,
|
||||||
|
message=message,
|
||||||
|
)
|
||||||
|
|
||||||
|
summary = ", ".join(deleted_titles[:5])
|
||||||
|
if deleted_count > 5:
|
||||||
|
summary += ", ..."
|
||||||
|
LOG.info(f"Bulk cleared {deleted_count} history item(s): {summary}")
|
||||||
|
|
||||||
|
return {"deleted": deleted_count}
|
||||||
|
|
||||||
|
async def clear_by_status(self, status_filter: str, remove_file: bool = False) -> dict[str, int | str]:
|
||||||
|
if self.config.remove_files is not True:
|
||||||
|
remove_file = False
|
||||||
|
|
||||||
|
if not remove_file:
|
||||||
|
deleted_count = await self.done.bulk_delete_by_status(status_filter)
|
||||||
|
if deleted_count < 1:
|
||||||
|
return {"deleted": 0}
|
||||||
|
|
||||||
|
self._notify.emit(
|
||||||
|
Events.ITEM_BULK_DELETED,
|
||||||
|
data={"count": deleted_count, "status": status_filter},
|
||||||
|
title="History Cleared",
|
||||||
|
message=f"Cleared {deleted_count} item{'s' if deleted_count != 1 else ''} from history.",
|
||||||
|
)
|
||||||
|
LOG.info(f"Bulk cleared {deleted_count} history item(s) by status '{status_filter}'.")
|
||||||
|
return {"deleted": deleted_count}
|
||||||
|
|
||||||
|
items = await self.done.get_many_by_status(status_filter)
|
||||||
|
return await self.clear_bulk([item_id for item_id, _ in items], remove_file=remove_file)
|
||||||
|
|
||||||
async def get(self, mode: str = "all") -> dict[str, list[dict[str, ItemDTO]]]:
|
async def get(self, mode: str = "all") -> dict[str, list[dict[str, ItemDTO]]]:
|
||||||
"""
|
"""
|
||||||
Get the download queue and the download history.
|
Get the download queue and the download history.
|
||||||
|
|
|
||||||
|
|
@ -7,6 +7,7 @@ from collections.abc import Iterable
|
||||||
from dataclasses import fields
|
from dataclasses import fields
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
from email.utils import formatdate
|
from email.utils import formatdate
|
||||||
|
from urllib.parse import quote_plus
|
||||||
|
|
||||||
from aiohttp import web
|
from aiohttp import web
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
|
|
@ -25,6 +26,51 @@ LOG: logging.Logger = logging.getLogger(__name__)
|
||||||
ITEM_DTO_FIELDS: set[str] = {f.name for f in fields(ItemDTO)}
|
ITEM_DTO_FIELDS: set[str] = {f.name for f in fields(ItemDTO)}
|
||||||
|
|
||||||
|
|
||||||
|
def _memory_db_url(db_path: str) -> str:
|
||||||
|
if db_path == ":memory:":
|
||||||
|
return "sqlite+aiosqlite:///file::memory:?cache=shared&uri=true"
|
||||||
|
|
||||||
|
memory_name = db_path[len(":memory:") :].lstrip(":") or "default"
|
||||||
|
quoted_name = quote_plus(memory_name)
|
||||||
|
return f"sqlite+aiosqlite:///file:{quoted_name}?mode=memory&cache=shared&uri=true"
|
||||||
|
|
||||||
|
|
||||||
|
def _decode_row(row: dict) -> ItemDTO:
|
||||||
|
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
||||||
|
data = json.loads(row["data"])
|
||||||
|
data.pop("_id", None)
|
||||||
|
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
||||||
|
item._id = row["id"]
|
||||||
|
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
||||||
|
return item
|
||||||
|
|
||||||
|
|
||||||
|
def _build_status_filter_clause(status_filter: str | None) -> tuple[str | None, dict[str, str]]:
|
||||||
|
if not status_filter:
|
||||||
|
return None, {}
|
||||||
|
|
||||||
|
raw_statuses = [entry.strip() for entry in status_filter.split(",") if entry.strip()]
|
||||||
|
if not raw_statuses:
|
||||||
|
return None, {}
|
||||||
|
|
||||||
|
if all(entry.startswith("!") for entry in raw_statuses):
|
||||||
|
statuses = [entry[1:].strip() for entry in raw_statuses if entry[1:].strip()]
|
||||||
|
if not statuses:
|
||||||
|
return None, {}
|
||||||
|
|
||||||
|
placeholders = ", ".join(f":status_{index}" for index in range(len(statuses)))
|
||||||
|
params = {f"status_{index}": status for index, status in enumerate(statuses)}
|
||||||
|
return f"json_extract(data, '$.status') NOT IN ({placeholders})", params
|
||||||
|
|
||||||
|
statuses = [entry for entry in raw_statuses if not entry.startswith("!")]
|
||||||
|
if not statuses:
|
||||||
|
return None, {}
|
||||||
|
|
||||||
|
placeholders = ", ".join(f":status_{index}" for index in range(len(statuses)))
|
||||||
|
params = {f"status_{index}": status for index, status in enumerate(statuses)}
|
||||||
|
return f"json_extract(data, '$.status') IN ({placeholders})", params
|
||||||
|
|
||||||
|
|
||||||
class Terminator:
|
class Terminator:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
@ -111,16 +157,7 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
)
|
)
|
||||||
rows = result.mappings().all()
|
rows = result.mappings().all()
|
||||||
|
|
||||||
items: list[tuple[str, ItemDTO]] = []
|
return [(row["id"], _decode_row(row)) for row in rows]
|
||||||
for row in rows:
|
|
||||||
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
|
||||||
data = json.loads(row["data"])
|
|
||||||
data.pop("_id", None)
|
|
||||||
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
|
||||||
item._id = row["id"]
|
|
||||||
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
|
||||||
items.append((row["id"], item))
|
|
||||||
return items
|
|
||||||
|
|
||||||
async def exists(self, type_value: str, key: str | None = None, url: str | None = None) -> bool:
|
async def exists(self, type_value: str, key: str | None = None, url: str | None = None) -> bool:
|
||||||
return await self.get(type_value, key=key, url=url) is not None
|
return await self.get(type_value, key=key, url=url) is not None
|
||||||
|
|
@ -154,13 +191,7 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
if not row:
|
if not row:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
return _decode_row(row)
|
||||||
data = json.loads(row["data"])
|
|
||||||
data.pop("_id", None)
|
|
||||||
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
|
||||||
item._id = row["id"]
|
|
||||||
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
|
||||||
return item
|
|
||||||
|
|
||||||
async def get_by_id(self, type_value: str, id: str) -> ItemDTO | None:
|
async def get_by_id(self, type_value: str, id: str) -> ItemDTO | None:
|
||||||
await self.get_connection()
|
await self.get_connection()
|
||||||
|
|
@ -173,13 +204,43 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
if not row:
|
if not row:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
return _decode_row(row)
|
||||||
data = json.loads(row["data"])
|
|
||||||
data.pop("_id", None)
|
async def get_many_by_ids(self, type_value: str, ids: Iterable[str]) -> list[tuple[str, ItemDTO]]:
|
||||||
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
ids_list = list(ids)
|
||||||
item._id = row["id"]
|
if not ids_list:
|
||||||
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
return []
|
||||||
return item
|
|
||||||
|
await self.get_connection()
|
||||||
|
placeholders = ", ".join(f":id_{index}" for index in range(len(ids_list)))
|
||||||
|
params: dict[str, str] = {"type_value": type_value}
|
||||||
|
params.update({f"id_{index}": item_id for index, item_id in enumerate(ids_list)})
|
||||||
|
|
||||||
|
result = await self._conn.execute(
|
||||||
|
text(
|
||||||
|
f'SELECT "id", "data", "created_at" FROM "history" WHERE "type" = :type_value AND "id" IN ({placeholders})' # noqa: S608
|
||||||
|
),
|
||||||
|
params,
|
||||||
|
)
|
||||||
|
rows = result.mappings().all()
|
||||||
|
order = {item_id: index for index, item_id in enumerate(ids_list)}
|
||||||
|
rows.sort(key=lambda row: order.get(row["id"], len(ids_list)))
|
||||||
|
return [(row["id"], _decode_row(row)) for row in rows]
|
||||||
|
|
||||||
|
async def get_many_by_status(self, type_value: str, status_filter: str) -> list[tuple[str, ItemDTO]]:
|
||||||
|
await self.get_connection()
|
||||||
|
params: dict[str, str] = {"type_value": type_value}
|
||||||
|
where_clauses = ['"type" = :type_value']
|
||||||
|
|
||||||
|
status_clause, status_params = _build_status_filter_clause(status_filter)
|
||||||
|
if status_clause:
|
||||||
|
where_clauses.append(status_clause)
|
||||||
|
params.update(status_params)
|
||||||
|
|
||||||
|
query = f'SELECT "id", "data", "created_at" FROM "history" WHERE {" AND ".join(where_clauses)} ORDER BY "created_at" DESC' # noqa: S608
|
||||||
|
result = await self._conn.execute(text(query), params)
|
||||||
|
rows = result.mappings().all()
|
||||||
|
return [(row["id"], _decode_row(row)) for row in rows]
|
||||||
|
|
||||||
async def get_item(self, type_value: str, **kwargs) -> ItemDTO | None:
|
async def get_item(self, type_value: str, **kwargs) -> ItemDTO | None:
|
||||||
"""
|
"""
|
||||||
|
|
@ -265,12 +326,7 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
if not row:
|
if not row:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
item = _decode_row(row)
|
||||||
data = json.loads(row["data"])
|
|
||||||
data.pop("_id", None)
|
|
||||||
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
|
||||||
item._id = row["id"]
|
|
||||||
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
|
||||||
|
|
||||||
return item if any(matches_condition(k, v, item.__dict__) for k, v in kwargs.items()) else None
|
return item if any(matches_condition(k, v, item.__dict__) for k, v in kwargs.items()) else None
|
||||||
|
|
||||||
|
|
@ -279,14 +335,10 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
where_clauses: list[str] = ['"type" = :type_value']
|
where_clauses: list[str] = ['"type" = :type_value']
|
||||||
params: dict[str, str] = {"type_value": type_value}
|
params: dict[str, str] = {"type_value": type_value}
|
||||||
|
|
||||||
if status_filter:
|
status_clause, status_params = _build_status_filter_clause(status_filter)
|
||||||
if status_filter.startswith("!"):
|
if status_clause:
|
||||||
status_value = status_filter[1:]
|
where_clauses.append(status_clause)
|
||||||
where_clauses.append("json_extract(data, '$.status') != :status")
|
params.update(status_params)
|
||||||
params["status"] = status_value
|
|
||||||
else:
|
|
||||||
where_clauses.append("json_extract(data, '$.status') = :status")
|
|
||||||
params["status"] = status_filter
|
|
||||||
|
|
||||||
where_clause: str = " AND ".join(where_clauses)
|
where_clause: str = " AND ".join(where_clauses)
|
||||||
count_query: str = f'SELECT COUNT(*) as count FROM "history" WHERE {where_clause}' # noqa: S608
|
count_query: str = f'SELECT COUNT(*) as count FROM "history" WHERE {where_clause}' # noqa: S608
|
||||||
|
|
@ -308,14 +360,10 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
where_clauses: list[str] = ['"type" = :type_value']
|
where_clauses: list[str] = ['"type" = :type_value']
|
||||||
params: dict[str, str | int] = {"type_value": type_value}
|
params: dict[str, str | int] = {"type_value": type_value}
|
||||||
|
|
||||||
if status_filter:
|
status_clause, status_params = _build_status_filter_clause(status_filter)
|
||||||
if status_filter.startswith("!"):
|
if status_clause:
|
||||||
status_value = status_filter[1:]
|
where_clauses.append(status_clause)
|
||||||
where_clauses.append("json_extract(data, '$.status') != :status")
|
params.update(status_params)
|
||||||
params["status"] = status_value
|
|
||||||
else:
|
|
||||||
where_clauses.append("json_extract(data, '$.status') = :status")
|
|
||||||
params["status"] = status_filter
|
|
||||||
|
|
||||||
where_clause: str = " AND ".join(where_clauses)
|
where_clause: str = " AND ".join(where_clauses)
|
||||||
count_query: str = f'SELECT COUNT(*) as count FROM "history" WHERE {where_clause}' # noqa: S608
|
count_query: str = f'SELECT COUNT(*) as count FROM "history" WHERE {where_clause}' # noqa: S608
|
||||||
|
|
@ -338,15 +386,7 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
result = await self._conn.execute(text(query), params)
|
result = await self._conn.execute(text(query), params)
|
||||||
rows = result.mappings().all()
|
rows = result.mappings().all()
|
||||||
|
|
||||||
items: list[tuple[str, ItemDTO]] = []
|
items = [(row["id"], _decode_row(row)) for row in rows]
|
||||||
for row in rows:
|
|
||||||
row_date = datetime.strptime(row["created_at"], "%Y-%m-%d %H:%M:%S") # noqa: DTZ007
|
|
||||||
data = json.loads(row["data"])
|
|
||||||
data.pop("_id", None)
|
|
||||||
item = init_class(ItemDTO, data, ITEM_DTO_FIELDS)
|
|
||||||
item._id = row["id"]
|
|
||||||
item.datetime = formatdate(row_date.replace(tzinfo=UTC).timestamp())
|
|
||||||
items.append((row["id"], item))
|
|
||||||
|
|
||||||
return items, total_items, page, total_pages
|
return items, total_items, page, total_pages
|
||||||
|
|
||||||
|
|
@ -378,6 +418,23 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
await self._conn.commit()
|
await self._conn.commit()
|
||||||
return result.rowcount if result else 0
|
return result.rowcount if result else 0
|
||||||
|
|
||||||
|
async def bulk_delete_by_status(self, type_value: str, status_filter: str) -> int:
|
||||||
|
await self.get_connection()
|
||||||
|
params: dict[str, str] = {"type_value": type_value}
|
||||||
|
where_clauses = ['"type" = :type_value']
|
||||||
|
|
||||||
|
status_clause, status_params = _build_status_filter_clause(status_filter)
|
||||||
|
if status_clause:
|
||||||
|
where_clauses.append(status_clause)
|
||||||
|
params.update(status_params)
|
||||||
|
|
||||||
|
result = await self._conn.execute(
|
||||||
|
text(f'DELETE FROM "history" WHERE {" AND ".join(where_clauses)}'), # noqa: S608
|
||||||
|
params,
|
||||||
|
)
|
||||||
|
await self._conn.commit()
|
||||||
|
return result.rowcount if result else 0
|
||||||
|
|
||||||
async def enqueue_upsert(self, type_value: str, item: ItemDTO) -> None:
|
async def enqueue_upsert(self, type_value: str, item: ItemDTO) -> None:
|
||||||
await self._enqueue(_Op("upsert", type_value, item, None, None))
|
await self._enqueue(_Op("upsert", type_value, item, None, None))
|
||||||
|
|
||||||
|
|
@ -553,7 +610,7 @@ class SqliteStore(metaclass=ThreadSafe):
|
||||||
from app.main import ROOT_PATH
|
from app.main import ROOT_PATH
|
||||||
|
|
||||||
if self._db_path.startswith(":memory"):
|
if self._db_path.startswith(":memory"):
|
||||||
db_url: str = "sqlite+aiosqlite:///file::memory:?cache=shared&uri=true"
|
db_url: str = _memory_db_url(self._db_path)
|
||||||
else:
|
else:
|
||||||
os.makedirs(os.path.dirname(self._db_path) or ".", exist_ok=True)
|
os.makedirs(os.path.dirname(self._db_path) or ".", exist_ok=True)
|
||||||
db_url: str = f"sqlite+aiosqlite:///{self._db_path}"
|
db_url: str = f"sqlite+aiosqlite:///{self._db_path}"
|
||||||
|
|
|
||||||
|
|
@ -177,9 +177,15 @@ async def items_delete(request: Request, queue: DownloadQueue, encoder: Encoder)
|
||||||
status=web.HTTPBadRequest.status_code,
|
status=web.HTTPBadRequest.status_code,
|
||||||
)
|
)
|
||||||
|
|
||||||
ds = queue.queue if storeType == StoreType.QUEUE else queue.done
|
|
||||||
|
|
||||||
if ids:
|
if ids:
|
||||||
|
if storeType == StoreType.HISTORY:
|
||||||
|
result = await queue.clear_bulk(ids, remove_file=remove_file)
|
||||||
|
return web.json_response(
|
||||||
|
data={"items": {}, "deleted": int(result.get("deleted", 0))},
|
||||||
|
status=web.HTTPOk.status_code,
|
||||||
|
dumps=encoder.encode,
|
||||||
|
)
|
||||||
|
|
||||||
return web.json_response(
|
return web.json_response(
|
||||||
data={
|
data={
|
||||||
"items": await (
|
"items": await (
|
||||||
|
|
@ -198,6 +204,23 @@ async def items_delete(request: Request, queue: DownloadQueue, encoder: Encoder)
|
||||||
status=web.HTTPBadRequest.status_code,
|
status=web.HTTPBadRequest.status_code,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
if storeType == StoreType.HISTORY:
|
||||||
|
result = await queue.clear_by_status(status_filter, remove_file=remove_file)
|
||||||
|
deleted = int(result.get("deleted", 0))
|
||||||
|
if deleted < 1:
|
||||||
|
return web.json_response(
|
||||||
|
data={"items": {}, "deleted": 0},
|
||||||
|
status=web.HTTPOk.status_code,
|
||||||
|
dumps=encoder.encode,
|
||||||
|
)
|
||||||
|
|
||||||
|
return web.json_response(
|
||||||
|
data={"items": {}, "deleted": deleted},
|
||||||
|
status=web.HTTPOk.status_code,
|
||||||
|
dumps=encoder.encode,
|
||||||
|
)
|
||||||
|
|
||||||
|
ds = queue.queue
|
||||||
items_to_delete = []
|
items_to_delete = []
|
||||||
page = 1
|
page = 1
|
||||||
per_page = 1000
|
per_page = 1000
|
||||||
|
|
@ -223,11 +246,7 @@ async def items_delete(request: Request, queue: DownloadQueue, encoder: Encoder)
|
||||||
|
|
||||||
return web.json_response(
|
return web.json_response(
|
||||||
data={
|
data={
|
||||||
"items": await (
|
"items": await queue.cancel(items_to_delete),
|
||||||
queue.cancel(items_to_delete)
|
|
||||||
if storeType == StoreType.QUEUE
|
|
||||||
else queue.clear(items_to_delete, remove_file=remove_file)
|
|
||||||
),
|
|
||||||
"deleted": len(items_to_delete),
|
"deleted": len(items_to_delete),
|
||||||
},
|
},
|
||||||
status=web.HTTPOk.status_code,
|
status=web.HTTPOk.status_code,
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,8 @@ from collections import OrderedDict
|
||||||
from dataclasses import asdict
|
from dataclasses import asdict
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
from email.utils import formatdate
|
from email.utils import formatdate
|
||||||
|
from unittest.mock import AsyncMock, Mock
|
||||||
|
from uuid import uuid4
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
|
|
@ -28,9 +30,9 @@ async def reset_sqlite_store() -> None:
|
||||||
|
|
||||||
|
|
||||||
async def make_db(data: int = 0) -> SqliteStore:
|
async def make_db(data: int = 0) -> SqliteStore:
|
||||||
"""Create a temporary database with test data."""
|
"""Create a named in-memory database with test data."""
|
||||||
await reset_sqlite_store()
|
await reset_sqlite_store()
|
||||||
ins = SqliteStore.get_instance(db_path=":memory:")
|
ins = SqliteStore.get_instance(db_path=f":memory:test-datastore-{uuid4().hex}")
|
||||||
await ins.get_connection()
|
await ins.get_connection()
|
||||||
|
|
||||||
base_time = datetime.now(UTC)
|
base_time = datetime.now(UTC)
|
||||||
|
|
@ -216,6 +218,31 @@ class TestDataStore:
|
||||||
ok = await store.test()
|
ok = await store.test()
|
||||||
assert ok is True
|
assert ok is True
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_bulk_delete_by_status_drops_matching_cached_items(self) -> None:
|
||||||
|
connection = Mock()
|
||||||
|
connection.bulk_delete_by_status = AsyncMock(return_value=2)
|
||||||
|
store = DataStore(StoreType.HISTORY, connection)
|
||||||
|
|
||||||
|
finished = make_item(id="finished")
|
||||||
|
finished.status = "finished"
|
||||||
|
skipped = make_item(id="skipped")
|
||||||
|
skipped.status = "skip"
|
||||||
|
pending = make_item(id="pending")
|
||||||
|
pending.status = "pending"
|
||||||
|
|
||||||
|
store._dict[finished._id] = StubDownload(info=finished)
|
||||||
|
store._dict[skipped._id] = StubDownload(info=skipped)
|
||||||
|
store._dict[pending._id] = StubDownload(info=pending)
|
||||||
|
|
||||||
|
deleted = await store.bulk_delete_by_status("finished,skip")
|
||||||
|
|
||||||
|
assert deleted == 2
|
||||||
|
connection.bulk_delete_by_status.assert_awaited_once_with("done", "finished,skip")
|
||||||
|
assert finished._id not in store._dict
|
||||||
|
assert skipped._id not in store._dict
|
||||||
|
assert pending._id in store._dict
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_get_item_returns_none_when_no_kwargs(self) -> None:
|
async def test_get_item_returns_none_when_no_kwargs(self) -> None:
|
||||||
"""Test that get_item returns None when no kwargs provided."""
|
"""Test that get_item returns None when no kwargs provided."""
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
import json
|
import json
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
from uuid import uuid4
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
import pytest_asyncio
|
import pytest_asyncio
|
||||||
|
|
@ -25,9 +26,10 @@ async def reset_sqlite_store() -> None:
|
||||||
|
|
||||||
|
|
||||||
async def make_db(data: int = 100) -> SqliteStore:
|
async def make_db(data: int = 100) -> SqliteStore:
|
||||||
"""Create a temporary database with test data."""
|
"""Create a named in-memory database with test data."""
|
||||||
await reset_sqlite_store()
|
await reset_sqlite_store()
|
||||||
ins = SqliteStore.get_instance(db_path=":memory:")
|
db_path = f":memory:test-datastore-pagination-{uuid4().hex}"
|
||||||
|
ins = SqliteStore.get_instance(db_path=db_path)
|
||||||
await ins.get_connection()
|
await ins.get_connection()
|
||||||
|
|
||||||
base_time = datetime.now(UTC)
|
base_time = datetime.now(UTC)
|
||||||
|
|
@ -235,6 +237,13 @@ class TestDataStorePagination:
|
||||||
"folder": "/downloads",
|
"folder": "/downloads",
|
||||||
"status": "downloading",
|
"status": "downloading",
|
||||||
}
|
}
|
||||||
|
item_data_skip = {
|
||||||
|
"url": "https://example.com/skip",
|
||||||
|
"title": "Skipped Video",
|
||||||
|
"id": "skip-video",
|
||||||
|
"folder": "/downloads",
|
||||||
|
"status": "skip",
|
||||||
|
}
|
||||||
|
|
||||||
db = await make_db(data=100)
|
db = await make_db(data=100)
|
||||||
try:
|
try:
|
||||||
|
|
@ -258,6 +267,16 @@ class TestDataStorePagination:
|
||||||
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
await db.execute_raw(
|
||||||
|
'INSERT INTO "history" ("id", "type", "url", "data", "created_at") VALUES (?, ?, ?, ?, ?)',
|
||||||
|
(
|
||||||
|
"skip-id",
|
||||||
|
str(StoreType.HISTORY),
|
||||||
|
item_data_skip["url"],
|
||||||
|
json.dumps(item_data_skip),
|
||||||
|
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
||||||
|
),
|
||||||
|
)
|
||||||
datastore = DataStore(type=StoreType.HISTORY, connection=db)
|
datastore = DataStore(type=StoreType.HISTORY, connection=db)
|
||||||
|
|
||||||
# Filter for finished items only
|
# Filter for finished items only
|
||||||
|
|
@ -276,6 +295,14 @@ class TestDataStorePagination:
|
||||||
assert len(items) == 1
|
assert len(items) == 1
|
||||||
assert total == 1
|
assert total == 1
|
||||||
assert items[0][1].info.status == "pending"
|
assert items[0][1].info.status == "pending"
|
||||||
|
|
||||||
|
items, total, _page, _total_pages = await datastore.get_items_paginated(
|
||||||
|
page=1, per_page=200, status_filter="finished,skip"
|
||||||
|
)
|
||||||
|
|
||||||
|
assert total == 101
|
||||||
|
assert len(items) == 101
|
||||||
|
assert {item.info.status for _, item in items} == {"finished", "skip"}
|
||||||
finally:
|
finally:
|
||||||
await db.close()
|
await db.close()
|
||||||
|
|
||||||
|
|
@ -296,6 +323,13 @@ class TestDataStorePagination:
|
||||||
"folder": "/downloads",
|
"folder": "/downloads",
|
||||||
"status": "error",
|
"status": "error",
|
||||||
}
|
}
|
||||||
|
item_data_skip = {
|
||||||
|
"url": "https://example.com/skip2",
|
||||||
|
"title": "Skipped Video 2",
|
||||||
|
"id": "skip-video-2",
|
||||||
|
"folder": "/downloads",
|
||||||
|
"status": "skip",
|
||||||
|
}
|
||||||
db = await make_db(data=0)
|
db = await make_db(data=0)
|
||||||
try:
|
try:
|
||||||
await db.execute_raw(
|
await db.execute_raw(
|
||||||
|
|
@ -318,6 +352,16 @@ class TestDataStorePagination:
|
||||||
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
await db.execute_raw(
|
||||||
|
'INSERT INTO "history" ("id", "type", "url", "data", "created_at") VALUES (?, ?, ?, ?, ?)',
|
||||||
|
(
|
||||||
|
"skip-id-2",
|
||||||
|
str(StoreType.HISTORY),
|
||||||
|
item_data_skip["url"],
|
||||||
|
json.dumps(item_data_skip),
|
||||||
|
datetime.now(UTC).strftime("%Y-%m-%d %H:%M:%S"),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
datastore = DataStore(type=StoreType.HISTORY, connection=db)
|
datastore = DataStore(type=StoreType.HISTORY, connection=db)
|
||||||
|
|
||||||
|
|
@ -326,14 +370,22 @@ class TestDataStorePagination:
|
||||||
page=1, per_page=50, status_filter="!finished"
|
page=1, per_page=50, status_filter="!finished"
|
||||||
)
|
)
|
||||||
|
|
||||||
assert total == 2 # Only 2 non-finished items
|
assert total == 3
|
||||||
assert len(items) == 2
|
assert len(items) == 3
|
||||||
for _item_id, item in items:
|
for _item_id, item in items:
|
||||||
assert item.info.status != "finished"
|
assert item.info.status != "finished"
|
||||||
|
|
||||||
# Verify we have pending and error
|
# Verify we have pending, error, and skip
|
||||||
statuses = {item.info.status for _, item in items}
|
statuses = {item.info.status for _, item in items}
|
||||||
assert statuses == {"pending", "error"}
|
assert statuses == {"pending", "error", "skip"}
|
||||||
|
|
||||||
|
items, total, _page, _total_pages = await datastore.get_items_paginated(
|
||||||
|
page=1, per_page=50, status_filter="!finished,!skip"
|
||||||
|
)
|
||||||
|
|
||||||
|
assert total == 2
|
||||||
|
assert len(items) == 2
|
||||||
|
assert {item.info.status for _, item in items} == {"pending", "error"}
|
||||||
finally:
|
finally:
|
||||||
await db.close()
|
await db.close()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1362,6 +1362,78 @@ class TestQueueManager:
|
||||||
done_store._connection.flush.assert_awaited_once()
|
done_store._connection.flush.assert_awaited_once()
|
||||||
assert status[item.info._id] == "ok", "Clear should still report success after flushing deletes"
|
assert status[item.info._id] == "ok", "Clear should still report success after flushing deletes"
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_clear_bulk_uses_bulk_delete_and_aggregate_notification(self) -> None:
|
||||||
|
queue_manager = object.__new__(DownloadQueue)
|
||||||
|
queue_manager.config = Mock(remove_files=False, download_path="/tmp")
|
||||||
|
queue_manager._notify = Mock()
|
||||||
|
|
||||||
|
item_one = Mock()
|
||||||
|
item_one.info = make_item(id="done-id-1", title="Finished clip 1")
|
||||||
|
item_one.info._id = "done-id-1"
|
||||||
|
item_one.info.status = "finished"
|
||||||
|
|
||||||
|
item_two = Mock()
|
||||||
|
item_two.info = make_item(id="done-id-2", title="Finished clip 2")
|
||||||
|
item_two.info._id = "done-id-2"
|
||||||
|
item_two.info.status = "finished"
|
||||||
|
|
||||||
|
done_store = Mock()
|
||||||
|
done_store.get_many_by_ids = AsyncMock(return_value=[("done-id-1", item_one), ("done-id-2", item_two)])
|
||||||
|
done_store.bulk_delete = AsyncMock(return_value=2)
|
||||||
|
queue_manager.done = done_store
|
||||||
|
|
||||||
|
result = await DownloadQueue.clear_bulk(queue_manager, ["done-id-1", "done-id-2"], remove_file=False)
|
||||||
|
|
||||||
|
assert result == {"deleted": 2}
|
||||||
|
done_store.get_many_by_ids.assert_awaited_once_with(["done-id-1", "done-id-2"])
|
||||||
|
done_store.bulk_delete.assert_awaited_once_with(["done-id-1", "done-id-2"])
|
||||||
|
queue_manager._notify.emit.assert_called_once()
|
||||||
|
assert queue_manager._notify.emit.call_args.args[0] == Events.ITEM_BULK_DELETED
|
||||||
|
assert queue_manager._notify.emit.call_args.kwargs["data"] == {"ids": ["done-id-1", "done-id-2"], "count": 2}
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_clear_by_status_uses_status_fetch_before_bulk_delete(self) -> None:
|
||||||
|
queue_manager = object.__new__(DownloadQueue)
|
||||||
|
queue_manager.config = Mock(remove_files=False, download_path="/tmp")
|
||||||
|
queue_manager._notify = Mock()
|
||||||
|
|
||||||
|
done_store = Mock()
|
||||||
|
done_store.bulk_delete_by_status = AsyncMock(return_value=1)
|
||||||
|
done_store.get_many_by_status = AsyncMock()
|
||||||
|
queue_manager.done = done_store
|
||||||
|
|
||||||
|
result = await DownloadQueue.clear_by_status(queue_manager, "finished", remove_file=False)
|
||||||
|
|
||||||
|
assert result == {"deleted": 1}
|
||||||
|
done_store.bulk_delete_by_status.assert_awaited_once_with("finished")
|
||||||
|
done_store.get_many_by_status.assert_not_called()
|
||||||
|
queue_manager._notify.emit.assert_called_once()
|
||||||
|
assert queue_manager._notify.emit.call_args.args[0] == Events.ITEM_BULK_DELETED
|
||||||
|
assert queue_manager._notify.emit.call_args.kwargs["data"] == {"count": 1, "status": "finished"}
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_clear_by_status_with_file_removal_fetches_matching_items(self) -> None:
|
||||||
|
queue_manager = object.__new__(DownloadQueue)
|
||||||
|
queue_manager.config = Mock(remove_files=True, download_path="/tmp")
|
||||||
|
queue_manager._notify = Mock()
|
||||||
|
|
||||||
|
item = Mock()
|
||||||
|
item.info = make_item(id="done-id", title="Finished clip")
|
||||||
|
item.info._id = "done-id"
|
||||||
|
item.info.status = "finished"
|
||||||
|
|
||||||
|
done_store = Mock()
|
||||||
|
done_store.get_many_by_status = AsyncMock(return_value=[("done-id", item)])
|
||||||
|
queue_manager.done = done_store
|
||||||
|
queue_manager.clear_bulk = AsyncMock(return_value={"deleted": 1})
|
||||||
|
|
||||||
|
result = await DownloadQueue.clear_by_status(queue_manager, "finished", remove_file=True)
|
||||||
|
|
||||||
|
assert result == {"deleted": 1}
|
||||||
|
done_store.get_many_by_status.assert_awaited_once_with("finished")
|
||||||
|
queue_manager.clear_bulk.assert_awaited_once_with(["done-id"], remove_file=True)
|
||||||
|
|
||||||
|
|
||||||
class TestPoolManager:
|
class TestPoolManager:
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
|
|
|
||||||
|
|
@ -10,31 +10,6 @@ from app.library.Events import Event, EventBus, EventListener, Events
|
||||||
class TestEvents:
|
class TestEvents:
|
||||||
"""Test the Events constants class."""
|
"""Test the Events constants class."""
|
||||||
|
|
||||||
def test_events_constants_exist(self):
|
|
||||||
"""Test that all expected event constants exist."""
|
|
||||||
# Basic lifecycle events
|
|
||||||
assert Events.STARTUP == "startup"
|
|
||||||
assert Events.LOADED == "loaded"
|
|
||||||
assert Events.STARTED == "started"
|
|
||||||
assert Events.SHUTDOWN == "shutdown"
|
|
||||||
|
|
||||||
# Connection events
|
|
||||||
assert Events.CONNECTED == "connected"
|
|
||||||
assert Events.CONFIG_UPDATE == "config_update"
|
|
||||||
|
|
||||||
# Log events
|
|
||||||
assert Events.LOG_INFO == "log_info"
|
|
||||||
assert Events.LOG_WARNING == "log_warning"
|
|
||||||
assert Events.LOG_ERROR == "log_error"
|
|
||||||
assert Events.LOG_SUCCESS == "log_success"
|
|
||||||
|
|
||||||
# Item events
|
|
||||||
assert Events.ITEM_ADDED == "item_added"
|
|
||||||
assert Events.ITEM_UPDATED == "item_updated"
|
|
||||||
assert Events.ITEM_COMPLETED == "item_completed"
|
|
||||||
assert Events.ITEM_CANCELLED == "item_cancelled"
|
|
||||||
assert Events.ITEM_DELETED == "item_deleted"
|
|
||||||
|
|
||||||
def test_events_get_all(self):
|
def test_events_get_all(self):
|
||||||
"""Test Events.get_all() method returns all constants."""
|
"""Test Events.get_all() method returns all constants."""
|
||||||
all_events = Events.get_all()
|
all_events = Events.get_all()
|
||||||
|
|
@ -59,6 +34,7 @@ class TestEvents:
|
||||||
"item_completed",
|
"item_completed",
|
||||||
"item_cancelled",
|
"item_cancelled",
|
||||||
"item_deleted",
|
"item_deleted",
|
||||||
|
"item_bulk_deleted",
|
||||||
"item_paused",
|
"item_paused",
|
||||||
"item_resumed",
|
"item_resumed",
|
||||||
"item_moved",
|
"item_moved",
|
||||||
|
|
@ -94,6 +70,7 @@ class TestEvents:
|
||||||
Events.LOG_SUCCESS,
|
Events.LOG_SUCCESS,
|
||||||
Events.ITEM_ADDED,
|
Events.ITEM_ADDED,
|
||||||
Events.ITEM_UPDATED,
|
Events.ITEM_UPDATED,
|
||||||
|
Events.ITEM_BULK_DELETED,
|
||||||
]
|
]
|
||||||
|
|
||||||
for expected in expected_frontend:
|
for expected in expected_frontend:
|
||||||
|
|
|
||||||
49
app/tests/test_history_routes.py
Normal file
49
app/tests/test_history_routes.py
Normal file
|
|
@ -0,0 +1,49 @@
|
||||||
|
import json
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import AsyncMock, Mock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from app.library.encoder import Encoder
|
||||||
|
from app.library.DataStore import StoreType
|
||||||
|
from app.routes.api.history import items_delete
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeRequest:
|
||||||
|
def __init__(self, *, payload: dict[str, Any] | None = None) -> None:
|
||||||
|
self._payload = payload or {}
|
||||||
|
self.query: dict[str, str] = {}
|
||||||
|
self.match_info: dict[str, str] = {}
|
||||||
|
|
||||||
|
async def json(self) -> dict[str, Any]:
|
||||||
|
return self._payload
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_items_delete_uses_bulk_history_status_clear() -> None:
|
||||||
|
request = _FakeRequest(payload={"type": StoreType.HISTORY.value, "status": "finished,skip", "remove_file": False})
|
||||||
|
queue = Mock()
|
||||||
|
queue.clear_by_status = AsyncMock(return_value={"deleted": 12})
|
||||||
|
encoder = Encoder()
|
||||||
|
|
||||||
|
response = await items_delete(request, queue, encoder)
|
||||||
|
|
||||||
|
assert response.status == 200
|
||||||
|
queue.clear_by_status.assert_awaited_once_with("finished,skip", remove_file=False)
|
||||||
|
body = json.loads(response.body.decode("utf-8"))
|
||||||
|
assert body == {"items": {}, "deleted": 12}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_items_delete_uses_bulk_history_id_clear() -> None:
|
||||||
|
request = _FakeRequest(payload={"type": StoreType.HISTORY.value, "ids": ["a", "b"], "remove_file": False})
|
||||||
|
queue = Mock()
|
||||||
|
queue.clear_bulk = AsyncMock(return_value={"deleted": 2})
|
||||||
|
encoder = Encoder()
|
||||||
|
|
||||||
|
response = await items_delete(request, queue, encoder)
|
||||||
|
|
||||||
|
assert response.status == 200
|
||||||
|
queue.clear_bulk.assert_awaited_once_with(["a", "b"], remove_file=False)
|
||||||
|
body = json.loads(response.body.decode("utf-8"))
|
||||||
|
assert body == {"items": {}, "deleted": 2}
|
||||||
|
|
@ -1,4 +1,7 @@
|
||||||
from datetime import UTC, datetime, timedelta
|
from datetime import UTC, datetime, timedelta
|
||||||
|
import os
|
||||||
|
from tempfile import mkstemp
|
||||||
|
from unittest.mock import AsyncMock, patch
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
|
|
@ -9,9 +12,11 @@ from app.library.sqlite_store import SqliteStore
|
||||||
|
|
||||||
|
|
||||||
async def make_store() -> SqliteStore:
|
async def make_store() -> SqliteStore:
|
||||||
"""Create an isolated in-memory SqliteStore instance with schema."""
|
"""Create an isolated temporary SqliteStore instance with schema."""
|
||||||
SqliteStore._reset_singleton()
|
SqliteStore._reset_singleton()
|
||||||
store = SqliteStore.get_instance(db_path=":memory:")
|
fd, db_path = mkstemp(prefix="ytptube-sqlite-store-", suffix=".db")
|
||||||
|
os.close(fd)
|
||||||
|
store = SqliteStore.get_instance(db_path=db_path)
|
||||||
await store.get_connection()
|
await store.get_connection()
|
||||||
assert store._engine is not None, "Engine should be initialized after _ensure_conn"
|
assert store._engine is not None, "Engine should be initialized after _ensure_conn"
|
||||||
return store
|
return store
|
||||||
|
|
@ -76,6 +81,31 @@ async def test_sqlalchemy_engine_disposed_on_close() -> None:
|
||||||
assert store._sessionmaker is None, "Sessionmaker should be None after close"
|
assert store._sessionmaker is None, "Sessionmaker should be None after close"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_named_in_memory_databases_are_isolated() -> None:
|
||||||
|
SqliteStore._reset_singleton()
|
||||||
|
first = SqliteStore.get_instance(db_path=":memory:named-a")
|
||||||
|
await first.get_connection()
|
||||||
|
await first.execute_raw(
|
||||||
|
'INSERT INTO "history" ("id", "type", "url", "data", "created_at") VALUES (?, ?, ?, ?, ?)',
|
||||||
|
(
|
||||||
|
"first",
|
||||||
|
"queue",
|
||||||
|
"https://example.com/a",
|
||||||
|
'{"id":"a","title":"A","url":"https://example.com/a","folder":"/downloads","status":"finished"}',
|
||||||
|
"2024-01-01 00:00:00",
|
||||||
|
),
|
||||||
|
)
|
||||||
|
await first.close()
|
||||||
|
|
||||||
|
SqliteStore._reset_singleton()
|
||||||
|
second = SqliteStore.get_instance(db_path=":memory:named-b")
|
||||||
|
await second.get_connection()
|
||||||
|
rows = await second.fetch_raw('SELECT "id" FROM "history" WHERE "id" = ?', ("first",))
|
||||||
|
assert rows == []
|
||||||
|
await second.close()
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_enqueue_upsert_and_fetch_saved():
|
async def test_enqueue_upsert_and_fetch_saved():
|
||||||
store = await make_store()
|
store = await make_store()
|
||||||
|
|
@ -131,14 +161,17 @@ async def test_count_with_status_filters():
|
||||||
store = await make_store()
|
store = await make_store()
|
||||||
finished_items = [make_item(i, status="finished") for i in range(2)]
|
finished_items = [make_item(i, status="finished") for i in range(2)]
|
||||||
pending_items = [make_item(i + 10, status="pending") for i in range(3)]
|
pending_items = [make_item(i + 10, status="pending") for i in range(3)]
|
||||||
|
skipped_items = [make_item(i + 20, status="skip") for i in range(2)]
|
||||||
|
|
||||||
for itm in finished_items + pending_items:
|
for itm in finished_items + pending_items + skipped_items:
|
||||||
await store.enqueue_upsert("history", itm)
|
await store.enqueue_upsert("history", itm)
|
||||||
await store.flush()
|
await store.flush()
|
||||||
|
|
||||||
assert await store.count("history", status_filter="finished") == 2
|
assert await store.count("history", status_filter="finished") == 2
|
||||||
assert await store.count("history", status_filter="pending") == 3
|
assert await store.count("history", status_filter="pending") == 3
|
||||||
assert await store.count("history", status_filter="!finished") == 3
|
assert await store.count("history", status_filter="finished,skip") == 4
|
||||||
|
assert await store.count("history", status_filter="!finished") == 5
|
||||||
|
assert await store.count("history", status_filter="!finished,!skip") == 3
|
||||||
await store.close()
|
await store.close()
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -227,6 +260,44 @@ async def test_enqueue_bulk_delete_returns_count_and_bulk_path():
|
||||||
await store.close()
|
await store.close()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_get_many_by_ids_and_status_and_bulk_delete_by_status():
|
||||||
|
store = await make_store()
|
||||||
|
finished = [make_item(i, status="finished") for i in range(2)]
|
||||||
|
pending = [make_item(i + 10, status="pending") for i in range(2)]
|
||||||
|
skipped = [make_item(i + 20, status="skip") for i in range(2)]
|
||||||
|
|
||||||
|
for item in finished + pending + skipped:
|
||||||
|
await store.enqueue_upsert("history", item)
|
||||||
|
await store.flush()
|
||||||
|
|
||||||
|
by_ids = await store.get_many_by_ids("history", [pending[1]._id, finished[0]._id])
|
||||||
|
assert [item_id for item_id, _ in by_ids] == [pending[1]._id, finished[0]._id]
|
||||||
|
|
||||||
|
by_status = await store.get_many_by_status("history", "finished")
|
||||||
|
assert {item_id for item_id, _ in by_status} == {finished[0]._id, finished[1]._id}
|
||||||
|
|
||||||
|
completed = await store.get_many_by_status("history", "finished,skip")
|
||||||
|
assert {item_id for item_id, _ in completed} == {
|
||||||
|
finished[0]._id,
|
||||||
|
finished[1]._id,
|
||||||
|
skipped[0]._id,
|
||||||
|
skipped[1]._id,
|
||||||
|
}
|
||||||
|
|
||||||
|
deleted = await store.bulk_delete_by_status("history", "finished")
|
||||||
|
assert deleted == 2
|
||||||
|
|
||||||
|
remaining = await store.fetch_saved("history")
|
||||||
|
assert {item_id for item_id, _ in remaining} == {
|
||||||
|
pending[0]._id,
|
||||||
|
pending[1]._id,
|
||||||
|
skipped[0]._id,
|
||||||
|
skipped[1]._id,
|
||||||
|
}
|
||||||
|
await store.close()
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_paginate_out_of_range_returns_last_page():
|
async def test_paginate_out_of_range_returns_last_page():
|
||||||
store = await make_store()
|
store = await make_store()
|
||||||
|
|
|
||||||
|
|
@ -1131,16 +1131,13 @@ const clearCompleted = async (): Promise<void> => {
|
||||||
|
|
||||||
selectedElms.value = [];
|
selectedElms.value = [];
|
||||||
|
|
||||||
await Promise.all([
|
await deleteHistoryItems({ status: 'finished,skip', removeFile: false });
|
||||||
deleteHistoryItems({ status: 'finished', removeFile: false }),
|
|
||||||
deleteHistoryItems({ status: 'skip', removeFile: false }),
|
|
||||||
]);
|
|
||||||
|
|
||||||
await reloadHistory({ order: 'DESC', perPage: config.app.default_pagination });
|
await reloadHistory({ order: 'DESC', perPage: config.app.default_pagination });
|
||||||
};
|
};
|
||||||
|
|
||||||
const clearIncomplete = async (): Promise<void> => {
|
const clearIncomplete = async (): Promise<void> => {
|
||||||
if (false === (await box.confirm('Clear all in-complete downloads?'))) {
|
if (false === (await box.confirm('Clear all incomplete downloads?'))) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
1
ui/app/types/sockets.d.ts
vendored
1
ui/app/types/sockets.d.ts
vendored
|
|
@ -25,6 +25,7 @@ export type WSEP = {
|
||||||
item_updated: EventPayload<StoreItem>;
|
item_updated: EventPayload<StoreItem>;
|
||||||
item_cancelled: EventPayload<StoreItem>;
|
item_cancelled: EventPayload<StoreItem>;
|
||||||
item_deleted: EventPayload<StoreItem>;
|
item_deleted: EventPayload<StoreItem>;
|
||||||
|
item_bulk_deleted: EventPayload<{ count: number; status?: string; ids?: string[] }>;
|
||||||
item_moved: EventPayload<{ to: 'queue' | 'history'; item: StoreItem }>;
|
item_moved: EventPayload<{ to: 'queue' | 'history'; item: StoreItem }>;
|
||||||
item_status: EventPayload<{ status?: string; msg?: string; preset?: string }>;
|
item_status: EventPayload<{ status?: string; msg?: string; preset?: string }>;
|
||||||
paused: EventPayload<{ paused?: boolean }>;
|
paused: EventPayload<{ paused?: boolean }>;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue