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 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, uploads_playlist_id in uploads.items(): channel = existing.get(channel_id) if channel is not None: channel.uploads_playlist_id = uploads_playlist_id 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} 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 ) candidate_video_ids.update(video_ids) except Exception: logger.exception("Failed to fetch playlist items for channel %s", channel.youtube_channel_id) 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: 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"]) video = existing.get(item["youtube_video_id"]) if video is None: video = 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, ) 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 db.commit() result = { "status": "completed", "started_at": started_at, "finished_at": _now_iso(), "error": None, "videos_added": added, "videos_updated": updated, "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) 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()