import asyncio 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 RECOMMEND_MAX_PAGES = 80 FIXED_SECTIONS = ( {"key": "trending", "name": "Trending Now", "note": "", "loading": True}, {"key": "new_releases", "name": "New Releases", "note": "", "loading": True}, {"key": "recommend", "name": "Recommend", "note": "", "loading": True}, ) class CatalogService: def outlines(self, categories: list[dict]) -> list[dict]: rows = [ {"key": f"type-{item['id']}", "name": item["label"], "note": "", "loading": True} for item in categories ] rows.extend(dict(item) for item in FIXED_SECTIONS) return rows def _categories(self, token: str, device: DeviceProfile) -> list[dict]: data = http.request_data("GET", API_BASE + "/app/open/classifyv2", token, device, None) rows = [] if isinstance(data, list): for item in data: if not isinstance(item, dict): continue type_id = item.get("id") label = str(item.get("label") or "").strip() if type_id is None or not label: continue rows.append({"id": type_id, "label": label}) if not rows: raise RuntimeError("没有获取到分类") return rows 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("请先登录") device = sessions.bind(token, DeviceProfile.merge(raw_device)) categories = await asyncio.to_thread(self._categories, token, device) yield {"type": "outline", "sections": self.outlines(categories)} jobs = [ ( f"type-{item['id']}", item["label"], "/app/open/classify_video", "GET", lambda page, type_id=item["id"]: {"type": type_id, "page": page}, CLASSIFY_MAX_PAGES, "", ) for item in categories ] jobs.append(("trending", "Trending Now", "/app/open/orgin?class=trending_now", "GET", None, 1, "")) jobs.append(("new_releases", "New Releases", "/app/open/orgin?class=new_releases", "GET", None, 1, "")) jobs.append(( "recommend", "Recommend", "/app/open/recommend", "POST", lambda page: {"page": page, "limit": 10}, RECOMMEND_MAX_PAGES, "", )) loop = asyncio.get_running_loop() queue: asyncio.Queue = asyncio.Queue() async def run(job): key, name, path, method, body, max_pages, note = job def work(): try: 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 "获取清单失败", done=True, failed=True), ) finally: loop.call_soon_threadsafe(queue.put_nowait, None) await asyncio.to_thread(work) tasks = [asyncio.create_task(run(job)) for job in jobs] finished = 0 try: while finished < len(tasks): event = await queue.get() if event is None: finished += 1 continue yield event yield {"type": "done"} finally: for task in tasks: if not task.done(): task.cancel() await asyncio.gather(*tasks, return_exceptions=True) 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") title = item.get("title") if vid is None or not title: return None image = (item.get("image") or "").strip() return { "id": vid, "title": title, "image": covers.local_url(image) if image else "", } def _iter_paged(self, key: str, name: str, token: str, device: DeviceProfile, path: str, method: str, body, max_pages: int, note: str = ""): seen = set() page = 1 sent = False while page <= max_pages: query = body(page) if body and method == "GET" else None payload = body(page) if body and method != "GET" else None data = http.request_data(method, API_BASE + path, token, device, payload, query) batch = data if isinstance(data, list) else [] if not batch and isinstance(data, dict) and isinstance(data.get("list"), list): batch = data["list"] if not batch: break fresh = [] for item in batch: if not isinstance(item, dict): continue drama = self._drama(item) if drama is None or drama["id"] in seen: continue seen.add(drama["id"]) fresh.append(drama) if not fresh: break sent = True yield self._event(key, name, fresh) if body is None: break page += 1 yield self._event(key, name, [], note, True) catalog = CatalogService()