myYouTube/backend/app/services/sync.py

287 lines
10 KiB
Python
Raw Permalink Normal View History

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()