myYouTube/backend/app/services/sync.py

272 lines
9.4 KiB
Python
Raw 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 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 = f"https://www.youtube.com/watch?v={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()