Channel features: subscribers, new-videos badges, activity-driven sync
- Store subscriber count from YouTube statistics and show it on channel page - Sync 50 videos per channel with playlistItems pagination support - Show per-channel and per-category new-videos counters (2-day window) - Replace hourly videos sync with activity trigger (2h idle) and incremental backfill (hard cap 200 per channel) - Clicking the sidebar new-videos count filters the category feed to recent videos only (new_only) - Update agent-team docs: deploy after green checks
This commit is contained in:
parent
fde9a439df
commit
e10df8dcbd
27 changed files with 1420 additions and 107 deletions
|
|
@ -21,19 +21,9 @@ def _run_subscriptions_sync() -> None:
|
|||
db.close()
|
||||
|
||||
|
||||
def _run_videos_sync() -> None:
|
||||
db = SessionLocal()
|
||||
try:
|
||||
sync.sync_videos(db)
|
||||
except sync.SyncInProgress:
|
||||
logger.info("Scheduled videos sync skipped: already running")
|
||||
except Exception:
|
||||
logger.exception("Scheduled videos sync failed")
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def create_scheduler() -> BackgroundScheduler:
|
||||
# Videos sync is not scheduled anymore: it runs on user activity via
|
||||
# services.sync_trigger.maybe_trigger_videos_sync (see require_session).
|
||||
scheduler = BackgroundScheduler(timezone="UTC")
|
||||
scheduler.add_job(
|
||||
_run_subscriptions_sync,
|
||||
|
|
@ -41,10 +31,4 @@ def create_scheduler() -> BackgroundScheduler:
|
|||
hours=settings.subscriptions_sync_interval_hours,
|
||||
id="subscriptions_sync",
|
||||
)
|
||||
scheduler.add_job(
|
||||
_run_videos_sync,
|
||||
"interval",
|
||||
minutes=settings.videos_sync_interval_minutes,
|
||||
id="videos_sync",
|
||||
)
|
||||
return scheduler
|
||||
|
|
|
|||
|
|
@ -52,6 +52,10 @@ def get_videos_sync_status(db: Session) -> dict:
|
|||
return _get_status(db, VIDEOS_SYNC_STATUS_KEY, _videos_lock)
|
||||
|
||||
|
||||
def is_videos_sync_running() -> bool:
|
||||
return _videos_lock.locked()
|
||||
|
||||
|
||||
def sync_subscriptions(db: Session) -> dict:
|
||||
if not _subscriptions_lock.acquire(blocking=False):
|
||||
raise SyncInProgress("Subscriptions sync already in progress")
|
||||
|
|
@ -108,10 +112,13 @@ def sync_subscriptions(db: Session) -> dict:
|
|||
subscribed_ids = [c.youtube_channel_id for c in existing.values() if c.subscribed]
|
||||
try:
|
||||
uploads = youtube_client.fetch_uploads_playlists(credentials, subscribed_ids)
|
||||
for channel_id, uploads_playlist_id in uploads.items():
|
||||
for channel_id, data in uploads.items():
|
||||
channel = existing.get(channel_id)
|
||||
if channel is not None:
|
||||
channel.uploads_playlist_id = uploads_playlist_id
|
||||
if data["uploads_playlist_id"] is not None:
|
||||
channel.uploads_playlist_id = data["uploads_playlist_id"]
|
||||
if data["subscriber_count"] is not None:
|
||||
channel.subscriber_count = data["subscriber_count"]
|
||||
db.commit()
|
||||
except Exception:
|
||||
logger.exception("Failed to fetch uploads playlists during subscriptions sync")
|
||||
|
|
@ -166,21 +173,38 @@ def sync_videos(db: Session) -> dict:
|
|||
)
|
||||
channel_by_youtube_id = {c.youtube_channel_id: c for c in channels}
|
||||
|
||||
# One query for all known video ids (grouped by channel) so there is
|
||||
# no N+1 at the DB level; the YouTube API is still queried per channel.
|
||||
known_ids_by_channel: dict[int, set[str]] = {}
|
||||
for youtube_video_id, channel_id in db.query(Video.youtube_video_id, Video.channel_id).all():
|
||||
known_ids_by_channel.setdefault(channel_id, set()).add(youtube_video_id)
|
||||
|
||||
candidate_video_ids: set[str] = set()
|
||||
for channel in channels:
|
||||
try:
|
||||
video_ids = youtube_client.fetch_playlist_video_ids(
|
||||
credentials, channel.uploads_playlist_id, settings.videos_per_channel_sync
|
||||
# videos_backfill_cap is a hard history-depth limit per
|
||||
# channel: only the newest N playlist items are ever
|
||||
# considered; anything older is not backfilled by design.
|
||||
# The early stop on consecutive known ids saves playlistItems
|
||||
# pages on repeated syncs inside that window.
|
||||
new_ids = youtube_client.fetch_playlist_video_ids_incremental(
|
||||
credentials,
|
||||
channel.uploads_playlist_id,
|
||||
known_ids_by_channel.get(channel.id, set()),
|
||||
settings.videos_backfill_cap,
|
||||
settings.videos_known_stop_threshold,
|
||||
)
|
||||
candidate_video_ids.update(video_ids)
|
||||
candidate_video_ids.update(new_ids)
|
||||
except Exception:
|
||||
logger.exception("Failed to fetch playlist items for channel %s", channel.youtube_channel_id)
|
||||
|
||||
# Details are fetched only for ids that are not in the DB yet (the
|
||||
# incremental fetch above already filters out known ids), so every
|
||||
# returned item is a new video to insert. Existing videos' metadata is
|
||||
# deliberately not refreshed by the videos sync.
|
||||
details = youtube_client.fetch_videos_details(credentials, list(candidate_video_ids))
|
||||
|
||||
existing = {v.youtube_video_id: v for v in db.query(Video).all()}
|
||||
added = 0
|
||||
updated = 0
|
||||
skipped = 0
|
||||
|
||||
for item in details:
|
||||
|
|
@ -197,9 +221,8 @@ def sync_videos(db: Session) -> dict:
|
|||
duration_seconds = parse_iso8601_duration(item["duration_iso8601"])
|
||||
youtube_url = settings.youtube_watch_url_template.format(video_id=item["youtube_video_id"])
|
||||
|
||||
video = existing.get(item["youtube_video_id"])
|
||||
if video is None:
|
||||
video = Video(
|
||||
db.add(
|
||||
Video(
|
||||
youtube_video_id=item["youtube_video_id"],
|
||||
channel_id=channel.id,
|
||||
title=item["title"],
|
||||
|
|
@ -209,18 +232,8 @@ def sync_videos(db: Session) -> dict:
|
|||
duration_seconds=duration_seconds,
|
||||
youtube_url=youtube_url,
|
||||
)
|
||||
db.add(video)
|
||||
existing[item["youtube_video_id"]] = video
|
||||
added += 1
|
||||
else:
|
||||
video.channel_id = channel.id
|
||||
video.title = item["title"]
|
||||
video.description = item["description"]
|
||||
video.thumbnail_url = item["thumbnail_url"]
|
||||
video.published_at = published_at
|
||||
video.duration_seconds = duration_seconds
|
||||
video.youtube_url = youtube_url
|
||||
updated += 1
|
||||
)
|
||||
added += 1
|
||||
|
||||
db.commit()
|
||||
|
||||
|
|
@ -230,12 +243,14 @@ def sync_videos(db: Session) -> dict:
|
|||
"finished_at": _now_iso(),
|
||||
"error": None,
|
||||
"videos_added": added,
|
||||
"videos_updated": updated,
|
||||
# Kept for API compatibility: always 0, because the videos sync
|
||||
# never refreshes metadata of existing videos.
|
||||
"videos_updated": 0,
|
||||
"videos_skipped": skipped,
|
||||
"channels_checked": len(channels),
|
||||
}
|
||||
_save_status(db, VIDEOS_SYNC_STATUS_KEY, result)
|
||||
logger.info("Videos sync completed: added=%d updated=%d skipped=%d", added, updated, skipped)
|
||||
logger.info("Videos sync completed: added=%d skipped=%d", added, skipped)
|
||||
return result
|
||||
|
||||
except Exception as exc:
|
||||
|
|
|
|||
91
backend/app/services/sync_trigger.py
Normal file
91
backend/app/services/sync_trigger.py
Normal file
|
|
@ -0,0 +1,91 @@
|
|||
"""Activity-triggered videos sync.
|
||||
|
||||
The hourly scheduled videos sync is gone; instead, every authenticated request
|
||||
(require_session) calls maybe_trigger_videos_sync(). The check must stay cheap:
|
||||
if the sync is not running and enough idle time has passed since the last
|
||||
completed sync, a background thread is started which calls sync.sync_videos
|
||||
with its own DB session.
|
||||
"""
|
||||
|
||||
import json
|
||||
import logging
|
||||
import threading
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from app.config import settings
|
||||
from app.db import SessionLocal
|
||||
from app.services import sync
|
||||
from app.services.state import get_setting
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def is_videos_sync_due(finished_at_iso: str | None, now: datetime, idle_hours: int) -> bool:
|
||||
"""Pure predicate: should an automatic videos sync start now?
|
||||
|
||||
True when the last completed sync's finished_at is missing/unparseable or
|
||||
older than `idle_hours`."""
|
||||
if not finished_at_iso:
|
||||
return True
|
||||
if not isinstance(finished_at_iso, str):
|
||||
return True
|
||||
try:
|
||||
finished_at = datetime.fromisoformat(finished_at_iso.replace("Z", "+00:00"))
|
||||
except (ValueError, TypeError):
|
||||
return True
|
||||
if finished_at.tzinfo is None:
|
||||
finished_at = finished_at.replace(tzinfo=timezone.utc)
|
||||
if now.tzinfo is None:
|
||||
now = now.replace(tzinfo=timezone.utc)
|
||||
return now - finished_at >= timedelta(hours=idle_hours)
|
||||
|
||||
|
||||
def _extract_finished_at(raw: str | None) -> str | None:
|
||||
if not raw:
|
||||
return None
|
||||
try:
|
||||
payload = json.loads(raw)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
if not isinstance(payload, dict):
|
||||
return None
|
||||
return payload.get("finished_at")
|
||||
|
||||
|
||||
def _run_videos_sync() -> None:
|
||||
db = SessionLocal()
|
||||
try:
|
||||
sync.sync_videos(db)
|
||||
except sync.SyncInProgress:
|
||||
# Lost the race against a manual sync or another trigger: fine, the
|
||||
# other run will record the status.
|
||||
logger.info("Triggered videos sync skipped: already running")
|
||||
except Exception:
|
||||
logger.exception("Triggered videos sync failed")
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def maybe_trigger_videos_sync() -> None:
|
||||
"""Cheap per-request check; never raises. Starts a background videos sync
|
||||
when due, or does nothing when a sync is already running."""
|
||||
try:
|
||||
if sync.is_videos_sync_running():
|
||||
return
|
||||
|
||||
db = SessionLocal()
|
||||
try:
|
||||
raw = get_setting(db, sync.VIDEOS_SYNC_STATUS_KEY)
|
||||
finished_at = _extract_finished_at(raw)
|
||||
if not is_videos_sync_due(finished_at, datetime.now(timezone.utc), settings.videos_sync_idle_hours):
|
||||
return
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
logger.info(
|
||||
"Videos sync due (last finished: %s), starting in background thread", finished_at
|
||||
)
|
||||
threading.Thread(target=_run_videos_sync, daemon=True).start()
|
||||
except Exception:
|
||||
# Never break the request because of the trigger itself.
|
||||
logger.exception("Failed to evaluate videos sync trigger")
|
||||
|
|
@ -9,6 +9,12 @@ 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
|
||||
|
|
@ -102,26 +108,108 @@ def fetch_subscriptions(credentials: Credentials) -> list[dict]:
|
|||
|
||||
|
||||
def fetch_playlist_video_ids(credentials: Credentials, playlist_id: str, max_results: int) -> list[str]:
|
||||
with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client:
|
||||
params = {
|
||||
"part": "contentDetails",
|
||||
"playlistId": playlist_id,
|
||||
"maxResults": min(max_results, 50),
|
||||
}
|
||||
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()
|
||||
"""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
|
||||
|
||||
video_ids = []
|
||||
for item in data.get("items", []):
|
||||
video_id = item.get("contentDetails", {}).get("videoId")
|
||||
if video_id:
|
||||
video_ids.append(video_id)
|
||||
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] = []
|
||||
|
||||
|
|
@ -171,14 +259,24 @@ def unsubscribe(credentials: Credentials, youtube_subscription_id: str) -> None:
|
|||
_raise_for_status(response)
|
||||
|
||||
|
||||
def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> dict[str, str]:
|
||||
result: dict[str, str] = {}
|
||||
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",
|
||||
"part": "snippet,contentDetails,statistics",
|
||||
"id": ",".join(batch),
|
||||
"maxResults": BATCH_SIZE,
|
||||
}
|
||||
|
|
@ -191,7 +289,10 @@ def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) ->
|
|||
uploads = (
|
||||
item.get("contentDetails", {}).get("relatedPlaylists", {}).get("uploads")
|
||||
)
|
||||
if channel_id and uploads:
|
||||
result[channel_id] = uploads
|
||||
if channel_id:
|
||||
result[channel_id] = {
|
||||
"uploads_playlist_id": uploads or None,
|
||||
"subscriber_count": _parse_subscriber_count(item.get("statistics", {})),
|
||||
}
|
||||
|
||||
return result
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue