import logging import httpx from google.oauth2.credentials import Credentials from app.config import settings logger = logging.getLogger(__name__) BATCH_SIZE = 50 # Hard safety cap for playlistItems pagination: at most this many pages # (BATCH_SIZE items each, i.e. 500 ids) per playlist, so the pageToken loop # can never run away. Raise both the cap and the caller's max_results # together if more is ever needed. MAX_PLAYLIST_PAGES = 10 class YouTubeQuotaExceeded(Exception): pass class YouTubeAPIError(Exception): pass class YouTubeInsufficientScope(Exception): """Raised when Google rejects a call because the stored token was granted under an older, narrower scope (e.g. readonly tokens issued before the unsubscribe feature needed write access) -- the fix is reconnecting.""" pass def _headers(credentials: Credentials) -> dict: return {"Authorization": f"Bearer {credentials.token}"} def _raise_for_status(response: httpx.Response) -> None: if response.status_code in (200, 204): return try: payload = response.json() reason = payload.get("error", {}).get("errors", [{}])[0].get("reason", "") message = payload.get("error", {}).get("message", response.text) except Exception: reason = "" message = response.text if response.status_code == 403 and reason in ("quotaExceeded", "dailyLimitExceeded", "rateLimitExceeded"): raise YouTubeQuotaExceeded(message) if response.status_code in (401, 403) and ( reason == "insufficientPermissions" or "insufficient authentication scopes" in message.lower() ): raise YouTubeInsufficientScope(message) logger.error("YouTube API error %s: %s", response.status_code, message) raise YouTubeAPIError(f"{response.status_code}: {message}") def fetch_subscriptions(credentials: Credentials) -> list[dict]: subscriptions: list[dict] = [] page_token: str | None = None with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: while True: params = { "part": "snippet,contentDetails", "mine": "true", "maxResults": 50, } if page_token: params["pageToken"] = page_token response = client.get(f"{settings.youtube_api_base_url}/subscriptions", params=params, headers=_headers(credentials)) _raise_for_status(response) data = response.json() for item in data.get("items", []): snippet = item.get("snippet", {}) resource_id = snippet.get("resourceId", {}) channel_id = resource_id.get("channelId") if not channel_id: continue thumbnails = snippet.get("thumbnails", {}) thumbnail = ( thumbnails.get("high") or thumbnails.get("medium") or thumbnails.get("default") or {} ).get("url") subscriptions.append( { "youtube_channel_id": channel_id, # The subscription resource's own id -- distinct from # the channel id, required to later call # subscriptions.delete (unsubscribe). "youtube_subscription_id": item.get("id"), "title": snippet.get("title", ""), "description": snippet.get("description", ""), "thumbnail_url": thumbnail, } ) page_token = data.get("nextPageToken") if not page_token: break return subscriptions def fetch_playlist_video_ids(credentials: Credentials, playlist_id: str, max_results: int) -> list[str]: """Low-level primitive: fetch up to `max_results` playlist item ids, newest first, with no knowledge of what is already synced. The videos sync uses fetch_playlist_video_ids_incremental instead; this stays as the plain paginated helper (kept for tests and any future non-incremental callers).""" video_ids: list[str] = [] page_token: str | None = None pages_fetched = 0 with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: while pages_fetched < MAX_PLAYLIST_PAGES and len(video_ids) < max_results: params = { "part": "contentDetails", "playlistId": playlist_id, "maxResults": min(max_results - len(video_ids), BATCH_SIZE), } if page_token: params["pageToken"] = page_token response = client.get(f"{settings.youtube_api_base_url}/playlistItems", params=params, headers=_headers(credentials)) if response.status_code == 404: return [] _raise_for_status(response) data = response.json() for item in data.get("items", []): video_id = item.get("contentDetails", {}).get("videoId") if video_id: video_ids.append(video_id) pages_fetched += 1 page_token = data.get("nextPageToken") if not page_token: break return video_ids def fetch_playlist_video_ids_incremental( credentials: Credentials, playlist_id: str, known_ids: set[str], max_results: int, stop_threshold: int, ) -> list[str]: """Fetch up to `max_results` previously-unknown video ids from a playlist, stopping pagination early once `stop_threshold` consecutive ids that are already in `known_ids` are encountered. Uploads playlists are ordered newest-first, so a long run of known ids means we reached history that was already synced and there is nothing new further down. `max_results` is a hard history-depth limit: at most the newest `max_results` videos of the channel are ever considered, and everything older than that is intentionally not backfilled (we do not mirror full channel history). The early stop on known ids only saves pages on repeated syncs within that window. Returns only the unknown ids, in playlist order.""" new_ids: list[str] = [] seen_new: set[str] = set() consecutive_known = 0 page_token: str | None = None pages_fetched = 0 with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: while pages_fetched < MAX_PLAYLIST_PAGES and len(new_ids) < max_results: params = { "part": "contentDetails", "playlistId": playlist_id, "maxResults": BATCH_SIZE, } if page_token: params["pageToken"] = page_token response = client.get(f"{settings.youtube_api_base_url}/playlistItems", params=params, headers=_headers(credentials)) if response.status_code == 404: return new_ids _raise_for_status(response) data = response.json() for item in data.get("items", []): video_id = item.get("contentDetails", {}).get("videoId") if not video_id: continue if video_id in known_ids or video_id in seen_new: consecutive_known += 1 if consecutive_known >= stop_threshold: return new_ids else: seen_new.add(video_id) new_ids.append(video_id) consecutive_known = 0 if len(new_ids) >= max_results: return new_ids pages_fetched += 1 page_token = data.get("nextPageToken") if not page_token: break return new_ids def fetch_videos_details(credentials: Credentials, video_ids: list[str]) -> list[dict]: results: list[dict] = [] with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: for i in range(0, len(video_ids), BATCH_SIZE): batch = video_ids[i : i + BATCH_SIZE] params = { "part": "snippet,contentDetails,status", "id": ",".join(batch), "maxResults": BATCH_SIZE, } response = client.get(f"{settings.youtube_api_base_url}/videos", params=params, headers=_headers(credentials)) _raise_for_status(response) data = response.json() for item in data.get("items", []): snippet = item.get("snippet", {}) thumbnails = snippet.get("thumbnails", {}) thumbnail = ( thumbnails.get("high") or thumbnails.get("medium") or thumbnails.get("default") or {} ).get("url") results.append( { "youtube_video_id": item.get("id"), "youtube_channel_id": snippet.get("channelId"), "title": snippet.get("title", ""), "description": snippet.get("description", ""), "thumbnail_url": thumbnail, "published_at": snippet.get("publishedAt"), "duration_iso8601": item.get("contentDetails", {}).get("duration"), } ) return results def unsubscribe(credentials: Credentials, youtube_subscription_id: str) -> None: with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: response = client.delete( f"{settings.youtube_api_base_url}/subscriptions", params={"id": youtube_subscription_id}, headers=_headers(credentials), ) if response.status_code == 404: # Already gone (unsubscribed elsewhere, or stale id) -- treat as success. return _raise_for_status(response) def _parse_subscriber_count(statistics: dict) -> int | None: raw = statistics.get("subscriberCount") if raw is None: return None try: return int(raw) except (TypeError, ValueError): return None def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> dict[str, dict]: result: dict[str, dict] = {} with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: for i in range(0, len(channel_ids), BATCH_SIZE): batch = channel_ids[i : i + BATCH_SIZE] params = { "part": "snippet,contentDetails,statistics", "id": ",".join(batch), "maxResults": BATCH_SIZE, } response = client.get(f"{settings.youtube_api_base_url}/channels", params=params, headers=_headers(credentials)) _raise_for_status(response) data = response.json() for item in data.get("items", []): channel_id = item.get("id") uploads = ( item.get("contentDetails", {}).get("relatedPlaylists", {}).get("uploads") ) if channel_id: result[channel_id] = { "uploads_playlist_id": uploads or None, "subscriber_count": _parse_subscriber_count(item.get("statistics", {})), } return result