myYouTube/backend/app/services/sync.py
vrubelroman e10df8dcbd 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
2026-09-17 21:07:33 +00:00

286 lines
10 KiB
Python

import json
import logging
import threading
from datetime import datetime, timezone
from sqlalchemy.orm import Session
from app.config import settings
from app.core.duration import parse_iso8601_duration
from app.models.channel import Channel
from app.models.video import Video
from app.services import google_oauth, youtube_client
from app.services.state import get_setting, set_setting
logger = logging.getLogger(__name__)
SUBSCRIPTIONS_SYNC_STATUS_KEY = "sync_subscriptions_status"
VIDEOS_SYNC_STATUS_KEY = "sync_videos_status"
_subscriptions_lock = threading.Lock()
_videos_lock = threading.Lock()
class SyncInProgress(Exception):
pass
def _now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
def _parse_youtube_datetime(value: str) -> datetime:
return datetime.fromisoformat(value.replace("Z", "+00:00"))
def _get_status(db: Session, key: str, lock: threading.Lock) -> dict:
raw = get_setting(db, key)
status = json.loads(raw) if raw else {"status": "never_run"}
status["running"] = lock.locked()
return status
def _save_status(db: Session, key: str, status: dict) -> None:
set_setting(db, key, json.dumps(status))
def get_subscriptions_sync_status(db: Session) -> dict:
return _get_status(db, SUBSCRIPTIONS_SYNC_STATUS_KEY, _subscriptions_lock)
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")
started_at = _now_iso()
try:
_save_status(
db, SUBSCRIPTIONS_SYNC_STATUS_KEY,
{"status": "running", "started_at": started_at, "finished_at": None, "error": None},
)
credentials = google_oauth.get_credentials(db)
subscriptions = youtube_client.fetch_subscriptions(credentials)
seen_channel_ids = {sub["youtube_channel_id"] for sub in subscriptions}
existing = {c.youtube_channel_id: c for c in db.query(Channel).all()}
added = 0
updated = 0
now = datetime.now(timezone.utc)
for sub in subscriptions:
channel = existing.get(sub["youtube_channel_id"])
if channel is None:
channel = Channel(
youtube_channel_id=sub["youtube_channel_id"],
youtube_subscription_id=sub["youtube_subscription_id"],
title=sub["title"],
description=sub["description"],
thumbnail_url=sub["thumbnail_url"],
subscribed=True,
last_synced_at=now,
)
db.add(channel)
existing[sub["youtube_channel_id"]] = channel
added += 1
else:
channel.youtube_subscription_id = sub["youtube_subscription_id"]
channel.title = sub["title"]
channel.description = sub["description"]
channel.thumbnail_url = sub["thumbnail_url"]
channel.subscribed = True
channel.last_synced_at = now
updated += 1
unsubscribed = 0
for channel_id, channel in existing.items():
if channel_id not in seen_channel_ids and channel.subscribed:
channel.subscribed = False
unsubscribed += 1
db.commit()
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, data in uploads.items():
channel = existing.get(channel_id)
if channel is not None:
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")
result = {
"status": "completed",
"started_at": started_at,
"finished_at": _now_iso(),
"error": None,
"channels_added": added,
"channels_updated": updated,
"channels_unsubscribed": unsubscribed,
}
_save_status(db, SUBSCRIPTIONS_SYNC_STATUS_KEY, result)
logger.info(
"Subscriptions sync completed: added=%d updated=%d unsubscribed=%d", added, updated, unsubscribed
)
return result
except Exception as exc:
logger.exception("Subscriptions sync failed")
db.rollback()
result = {
"status": "failed",
"started_at": started_at,
"finished_at": _now_iso(),
"error": str(exc),
}
_save_status(db, SUBSCRIPTIONS_SYNC_STATUS_KEY, result)
raise
finally:
_subscriptions_lock.release()
def sync_videos(db: Session) -> dict:
if not _videos_lock.acquire(blocking=False):
raise SyncInProgress("Videos sync already in progress")
started_at = _now_iso()
try:
_save_status(
db, VIDEOS_SYNC_STATUS_KEY,
{"status": "running", "started_at": started_at, "finished_at": None, "error": None},
)
credentials = google_oauth.get_credentials(db)
channels = (
db.query(Channel)
.filter(Channel.subscribed.is_(True), Channel.uploads_playlist_id.isnot(None))
.all()
)
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:
# 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(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))
added = 0
skipped = 0
for item in details:
channel = channel_by_youtube_id.get(item["youtube_channel_id"])
if channel is None:
skipped += 1
continue
published_at = _parse_youtube_datetime(item["published_at"]) if item["published_at"] else None
if published_at is None:
skipped += 1
continue
duration_seconds = parse_iso8601_duration(item["duration_iso8601"])
youtube_url = settings.youtube_watch_url_template.format(video_id=item["youtube_video_id"])
db.add(
Video(
youtube_video_id=item["youtube_video_id"],
channel_id=channel.id,
title=item["title"],
description=item["description"],
thumbnail_url=item["thumbnail_url"],
published_at=published_at,
duration_seconds=duration_seconds,
youtube_url=youtube_url,
)
)
added += 1
db.commit()
result = {
"status": "completed",
"started_at": started_at,
"finished_at": _now_iso(),
"error": None,
"videos_added": added,
# 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 skipped=%d", added, skipped)
return result
except Exception as exc:
logger.exception("Videos sync failed")
db.rollback()
result = {
"status": "failed",
"started_at": started_at,
"finished_at": _now_iso(),
"error": str(exc),
}
_save_status(db, VIDEOS_SYNC_STATUS_KEY, result)
raise
finally:
_videos_lock.release()
class ChannelHasNoSubscriptionId(Exception):
pass
def unsubscribe_channel(db: Session, channel: Channel) -> None:
"""Actually unsubscribes on YouTube (subscriptions.delete) -- deliberately
out of the original MVP scope, added on explicit user request. Requires
the write-capable OAuth scope; see google_oauth.SCOPES."""
if not channel.youtube_subscription_id:
raise ChannelHasNoSubscriptionId(
"This channel was synced before subscription ids were tracked; run a subscriptions sync first"
)
credentials = google_oauth.get_credentials(db)
youtube_client.unsubscribe(credentials, channel.youtube_subscription_id)
channel.subscribed = False
db.commit()