feat: 增加视频清单搜索,并在开发期缓存剧目清单
This commit is contained in:
@@ -4,6 +4,7 @@ from server.services.covers import covers
|
||||
from server.services.device import DeviceProfile
|
||||
from server.services.http import http
|
||||
from server.services.rapt import API_BASE
|
||||
from server.services import catalog_cache
|
||||
from server.services.sessions import sessions
|
||||
|
||||
CLASSIFY_MAX_PAGES = 50
|
||||
@@ -39,7 +40,27 @@ class CatalogService:
|
||||
raise RuntimeError("没有获取到分类")
|
||||
return rows
|
||||
|
||||
async def stream(self, token: str, raw_device: dict | None):
|
||||
async def stream(self, token: str, raw_device: dict | None, refresh: bool = False):
|
||||
if catalog_cache.enabled():
|
||||
if refresh:
|
||||
catalog_cache.clear()
|
||||
else:
|
||||
cached = catalog_cache.read()
|
||||
if cached:
|
||||
for event in cached:
|
||||
yield event
|
||||
return
|
||||
collected: list[dict] = []
|
||||
failed = False
|
||||
async for event in self._fetch(token, raw_device):
|
||||
if event.get("failed"):
|
||||
failed = True
|
||||
collected.append(event)
|
||||
yield event
|
||||
if catalog_cache.enabled() and not failed and collected and collected[-1].get("type") == "done":
|
||||
catalog_cache.write(collected)
|
||||
|
||||
async def _fetch(self, token: str, raw_device: dict | None):
|
||||
token = token.strip()
|
||||
if not token:
|
||||
raise RuntimeError("请先登录")
|
||||
@@ -79,7 +100,10 @@ class CatalogService:
|
||||
for event in self._iter_paged(key, name, token, device, path, method, body, max_pages, note):
|
||||
loop.call_soon_threadsafe(queue.put_nowait, event)
|
||||
except Exception as exc:
|
||||
loop.call_soon_threadsafe(queue.put_nowait, self._event(key, name, [], str(exc) or "获取清单失败", True))
|
||||
loop.call_soon_threadsafe(
|
||||
queue.put_nowait,
|
||||
self._event(key, name, [], str(exc) or "获取清单失败", done=True, failed=True),
|
||||
)
|
||||
finally:
|
||||
loop.call_soon_threadsafe(queue.put_nowait, None)
|
||||
|
||||
@@ -101,8 +125,11 @@ class CatalogService:
|
||||
task.cancel()
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
def _event(self, key: str, name: str, items: list[dict], note: str = "", done: bool = False) -> dict:
|
||||
return {"type": "page", "key": key, "name": name, "note": note, "items": items, "done": done}
|
||||
def _event(self, key: str, name: str, items: list[dict], note: str = "", done: bool = False, failed: bool = False) -> dict:
|
||||
event = {"type": "page", "key": key, "name": name, "note": note, "items": items, "done": done}
|
||||
if failed:
|
||||
event["failed"] = True
|
||||
return event
|
||||
|
||||
def _drama(self, item: dict) -> dict | None:
|
||||
vid = item.get("id")
|
||||
|
||||
Reference in New Issue
Block a user