revert: changes to playlist processor

This commit is contained in:
arabcoders 2026-01-27 10:41:17 +03:00
parent da9f082978
commit dae141a6c2
3 changed files with 74 additions and 31 deletions

View file

@ -207,6 +207,9 @@ class Config(metaclass=Singleton):
live_premiere_buffer: int = 5 live_premiere_buffer: int = 5
"""The buffer time in minutes to add to video duration to wait before starting premiere download.""" """The buffer time in minutes to add to video duration to wait before starting premiere download."""
playlist_items_concurrency: int = 4
"""The number of concurrent playlist items to be processed at same time."""
auto_clear_history_days: int = 0 auto_clear_history_days: int = 0
"""Number of days after which completed download history is automatically cleared. 0 to disable.""" """Number of days after which completed download history is automatically cleared. 0 to disable."""
@ -265,6 +268,7 @@ class Config(metaclass=Singleton):
"max_workers_per_extractor", "max_workers_per_extractor",
"extract_info_timeout", "extract_info_timeout",
"debugpy_port", "debugpy_port",
"playlist_items_concurrency",
"download_path_depth", "download_path_depth",
"download_info_expires", "download_info_expires",
"auto_clear_history_days", "auto_clear_history_days",

View file

@ -1,10 +1,13 @@
"""Playlist processing.""" """Playlist processing."""
import asyncio
import logging import logging
from typing import TYPE_CHECKING, Any from typing import TYPE_CHECKING, Any
from app.library.Utils import merge_dict, ytdlp_reject from app.library.Utils import merge_dict, ytdlp_reject
from .utils import handle_task_exception
if TYPE_CHECKING: if TYPE_CHECKING:
from app.library.ItemDTO import Item from app.library.ItemDTO import Item
@ -56,54 +59,87 @@ async def process_playlist(
"n_entries": len(entries), "n_entries": len(entries),
} }
async def process_item(i: int, etr: dict) -> dict[str, str]: async def playlist_processor(i: int, etr: dict):
"""Process a single playlist item.""" acquired = False
item_name: str = ( try:
f"'{entry.get('title')}: {i}/{playlist_keys['n_entries']}' - '{etr.get('id')}: {etr.get('title')}'" item_name: str = (
) f"'{entry.get('title')}: {i}/{playlist_keys['n_entries']}' - '{etr.get('id')}: {etr.get('title')}'"
LOG.info(f"Processing '{item_name}'.") )
LOG.debug(f"Waiting to acquire lock for {item_name}")
await queue.processors.acquire()
acquired = True
LOG.debug(f"Acquired lock for {item_name}")
_status, _msg = ytdlp_reject(entry=etr, yt_params=yt_params) LOG.info(f"Processing '{item_name}'.")
if not _status:
return {"status": "error", "msg": _msg}
extras: dict[str, Any] = { _status, _msg = ytdlp_reject(entry=etr, yt_params=yt_params)
**playlist_keys, if not _status:
"playlist_index": i, return {"status": "error", "msg": _msg}
"playlist_index_number": i,
"playlist_autonumber": i,
}
for property in ("id", "title", "uploader", "uploader_id"): extras: dict[str, Any] = {
if property in entry: **playlist_keys,
extras[f"playlist_{property}"] = entry.get(property) "playlist_index": i,
"playlist_index_number": i,
"playlist_autonumber": i,
}
extractor_key = entry.get("ie_key") or entry.get("extractor_key") or entry.get("extractor") or "" for property in ("id", "title", "uploader", "uploader_id"):
if "thumbnail" not in etr and "youtube" in str(extractor_key).lower(): if property in entry:
extras["thumbnail"] = "https://img.youtube.com/vi/{id}/maxresdefault.jpg".format(**etr) extras[f"playlist_{property}"] = entry.get(property)
newItem: Item = item.new_with(url=etr.get("url") or etr.get("webpage_url"), extras=extras) extractor_key = entry.get("ie_key") or entry.get("extractor_key") or entry.get("extractor") or ""
if "thumbnail" not in etr and "youtube" in str(extractor_key).lower():
extras["thumbnail"] = "https://img.youtube.com/vi/{id}/maxresdefault.jpg".format(**etr)
if ("video" == etr.get("_type") and etr.get("url")) or ( newItem: Item = item.new_with(url=etr.get("url") or etr.get("webpage_url"), extras=extras)
"formats" in etr and isinstance(etr["formats"], list) and len(etr["formats"]) > 0
):
dct = merge_dict(merge_dict({"_type": "video"}, etr), entry)
dct.pop("entries", None)
return await queue.add(item=newItem, entry=dct, already=already)
return await queue.add(item=newItem, already=already) if ("video" == etr.get("_type") and etr.get("url")) or (
"formats" in etr and isinstance(etr["formats"], list) and len(etr["formats"]) > 0
):
dct = merge_dict(merge_dict({"_type": "video"}, etr), entry)
dct.pop("entries", None)
return await queue.add(item=newItem, entry=dct, already=already)
return await queue.add(item=newItem, already=already)
finally:
if acquired:
queue.processors.release()
max_downloads: int = -1 max_downloads: int = -1
ytdlp_opts: dict[str, Any] = item.get_ytdlp_opts().get_all() ytdlp_opts: dict[str, Any] = item.get_ytdlp_opts().get_all()
if ytdlp_opts.get("max_downloads") and isinstance(ytdlp_opts.get("max_downloads"), int): if ytdlp_opts.get("max_downloads") and isinstance(ytdlp_opts.get("max_downloads"), int):
max_downloads: int = ytdlp_opts.get("max_downloads") max_downloads: int = ytdlp_opts.get("max_downloads")
results: list[dict[str, str]] = [] batch_size: int = max(50, int(queue.config.playlist_items_concurrency) * 10)
async def run_batch(batch: list[tuple[int, dict]]) -> list[dict]:
tasks: list[asyncio.Task] = []
for i, etr in batch:
task = asyncio.create_task(
playlist_processor(i, etr),
name=f"playlist_processor_{etr.get('id')}_{i}",
)
task.add_done_callback(lambda t: handle_task_exception(t, LOG))
tasks.append(task)
batch_results: list[dict] = await asyncio.gather(*tasks)
await asyncio.sleep(0)
return batch_results
results: list[dict] = []
batch: list[tuple[int, dict]] = []
for i, etr in enumerate(entries, start=1): for i, etr in enumerate(entries, start=1):
if max_downloads > 0 and i > max_downloads: if max_downloads > 0 and i > max_downloads:
break break
results.append(await process_item(i, etr)) batch.append((i, etr))
if len(batch) >= batch_size:
results.extend(await run_batch(batch))
batch.clear()
if batch:
results.extend(await run_batch(batch))
log_msg: str = f"Playlist '{playlist_name}' processing completed with '{len(results)}' entries." log_msg: str = f"Playlist '{playlist_name}' processing completed with '{len(results)}' entries."
if max_downloads > 0 and len(entries) > max_downloads: if max_downloads > 0 and len(entries) > max_downloads:

View file

@ -1,3 +1,4 @@
import asyncio
import functools import functools
import glob import glob
import logging import logging
@ -41,6 +42,8 @@ class DownloadQueue(metaclass=Singleton):
"DataStore for the completed downloads." "DataStore for the completed downloads."
self.queue = DataStore(type=StoreType.QUEUE, connection=SqliteStore.get_instance()) self.queue = DataStore(type=StoreType.QUEUE, connection=SqliteStore.get_instance())
"DataStore for the download queue." "DataStore for the download queue."
self.processors = asyncio.Semaphore(self.config.playlist_items_concurrency)
"Semaphore to limit the number of concurrent processors."
self.pool = PoolManager(queue=self, config=self.config) self.pool = PoolManager(queue=self, config=self.config)
"Pool manager for coordinating download execution." "Pool manager for coordinating download execution."