Implement Phases 1-5: skeleton, OAuth, categories, video sync/feed, playback

- FastAPI + PostgreSQL + Alembic + React/Vite skeleton, Docker Compose, healthcheck
- Google OAuth (single allowed account), encrypted refresh token storage
- Subscriptions sync with pagination, uploads playlist batch fetch
- Categories CRUD, many-to-many channel assignment, category filtering
- Video sync (playlistItems + videos.list batching), cached feed with cursor
  pagination, background scheduler (APScheduler)
- Video detail page with YouTube embed player
- SPA fallback routing, optimistic UI updates, client-side query caching

40 backend tests covering OAuth allow-list, sync idempotency, cascade deletes,
cursor pagination, and category filtering.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
vrubelroman 2026-09-16 18:44:30 +00:00
commit 0ed20bb838
90 changed files with 8545 additions and 0 deletions

0
backend/app/__init__.py Normal file
View file

View file

96
backend/app/api/auth.py Normal file
View file

@ -0,0 +1,96 @@
import logging
import secrets
from fastapi import APIRouter, BackgroundTasks, Depends, Request
from fastapi.responses import RedirectResponse
from sqlalchemy.orm import Session
from app.config import settings
from app.core.auth_dependency import require_session
from app.db import SessionLocal, get_db
from app.services import google_oauth, sync
logger = logging.getLogger(__name__)
router = APIRouter()
@router.get("/auth/status")
def auth_status(request: Request, db: Session = Depends(get_db)) -> dict:
connected = google_oauth.is_connected(db)
return {
"authenticated": bool(request.session.get("authenticated")),
"connected": connected,
"email": google_oauth.get_connected_email(db) if connected else None,
}
@router.get("/auth/google/start")
def google_start(request: Request):
auth_url, state = google_oauth.build_authorization_url()
request.session["oauth_state"] = state
return RedirectResponse(auth_url)
def _run_initial_sync() -> None:
db = SessionLocal()
try:
sync.sync_subscriptions(db)
sync.sync_videos(db)
except Exception:
logger.exception("Initial sync after OAuth failed")
finally:
db.close()
@router.get("/auth/google/callback")
def google_callback(
request: Request,
background_tasks: BackgroundTasks,
code: str | None = None,
state: str | None = None,
error: str | None = None,
db: Session = Depends(get_db),
):
expected_state = request.session.pop("oauth_state", None)
if error:
logger.warning("Google OAuth returned error: %s", error)
return RedirectResponse(f"{settings.app_base_url}/?auth_error=google_error")
if not code or not state or not expected_state or not secrets.compare_digest(state, expected_state):
logger.warning("Google OAuth callback with invalid/missing state")
return RedirectResponse(f"{settings.app_base_url}/?auth_error=invalid_state")
try:
credentials = google_oauth.exchange_code(code, state)
except Exception:
logger.exception("Failed to exchange Google OAuth code")
return RedirectResponse(f"{settings.app_base_url}/?auth_error=exchange_failed")
try:
userinfo = google_oauth.fetch_userinfo(credentials.token)
except Exception:
logger.exception("Failed to fetch Google userinfo")
return RedirectResponse(f"{settings.app_base_url}/?auth_error=userinfo_failed")
email = (userinfo.get("email") or "").lower()
allowed_email = settings.allowed_google_email.lower()
if not allowed_email or email != allowed_email:
logger.warning("Rejected Google OAuth login for disallowed account")
google_oauth.revoke_token(credentials.refresh_token or credentials.token)
return RedirectResponse(f"{settings.app_base_url}/?auth_error=account_not_allowed")
google_oauth.store_credentials(db, email, credentials)
request.session["authenticated"] = True
background_tasks.add_task(_run_initial_sync)
return RedirectResponse(f"{settings.app_base_url}/")
@router.post("/auth/logout")
def logout(request: Request, _: None = Depends(require_session)):
request.session.clear()
return {"ok": True}

View file

@ -0,0 +1,134 @@
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy import func
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.core.auth_dependency import require_session
from app.core.slugify import unique_slugify
from app.db import get_db
from app.models.category import Category
from app.models.channel_category import channel_categories
router = APIRouter(dependencies=[Depends(require_session)])
class CategoryCreate(BaseModel):
name: str = Field(min_length=1, max_length=255)
class CategoryUpdate(BaseModel):
name: str = Field(min_length=1, max_length=255)
class CategoryReorder(BaseModel):
category_ids: list[int]
def _existing_slugs(db: Session, exclude_id: int | None = None) -> set[str]:
query = db.query(Category.slug)
if exclude_id is not None:
query = query.filter(Category.id != exclude_id)
return {row[0] for row in query.all()}
def _name_taken(db: Session, name: str, exclude_id: int | None = None) -> bool:
query = db.query(Category.name)
if exclude_id is not None:
query = query.filter(Category.id != exclude_id)
target = name.casefold()
return any(row[0].casefold() == target for row in query.all())
def _serialize(db: Session, category: Category, counts: dict[int, int]) -> dict:
return {
"id": category.id,
"name": category.name,
"slug": category.slug,
"sort_order": category.sort_order,
"channel_count": counts.get(category.id, 0),
}
@router.get("/categories")
def list_categories(db: Session = Depends(get_db)) -> list[dict]:
categories = db.query(Category).order_by(Category.sort_order.asc(), Category.id.asc()).all()
count_rows = (
db.query(channel_categories.c.category_id, func.count(channel_categories.c.channel_id))
.group_by(channel_categories.c.category_id)
.all()
)
counts = dict(count_rows)
return [_serialize(db, c, counts) for c in categories]
@router.post("/categories", status_code=201)
def create_category(payload: CategoryCreate, db: Session = Depends(get_db)) -> dict:
name = payload.name.strip()
if not name:
raise HTTPException(status_code=400, detail="Category name must not be empty")
if _name_taken(db, name):
raise HTTPException(status_code=409, detail="Category with this name already exists")
slug = unique_slugify(name, _existing_slugs(db))
max_sort_order = db.query(func.max(Category.sort_order)).scalar() or 0
category = Category(name=name, slug=slug, sort_order=max_sort_order + 1)
db.add(category)
try:
db.commit()
except IntegrityError:
db.rollback()
raise HTTPException(status_code=409, detail="Category with this name already exists")
return _serialize(db, category, {})
@router.patch("/categories/{category_id}")
def update_category(category_id: int, payload: CategoryUpdate, db: Session = Depends(get_db)) -> dict:
category = db.get(Category, category_id)
if category is None:
raise HTTPException(status_code=404, detail="Category not found")
name = payload.name.strip()
if not name:
raise HTTPException(status_code=400, detail="Category name must not be empty")
if _name_taken(db, name, exclude_id=category_id):
raise HTTPException(status_code=409, detail="Category with this name already exists")
category.name = name
category.slug = unique_slugify(name, _existing_slugs(db, exclude_id=category_id))
try:
db.commit()
except IntegrityError:
db.rollback()
raise HTTPException(status_code=409, detail="Category with this name already exists")
return _serialize(db, category, {})
@router.delete("/categories/{category_id}", status_code=204)
def delete_category(category_id: int, db: Session = Depends(get_db)) -> None:
category = db.get(Category, category_id)
if category is None:
raise HTTPException(status_code=404, detail="Category not found")
db.delete(category)
db.commit()
@router.post("/categories/reorder")
def reorder_categories(payload: CategoryReorder, db: Session = Depends(get_db)) -> list[dict]:
categories = {c.id: c for c in db.query(Category).all()}
if set(payload.category_ids) != set(categories.keys()):
raise HTTPException(status_code=400, detail="category_ids must contain exactly all existing category ids")
for index, category_id in enumerate(payload.category_ids):
categories[category_id].sort_order = index
db.commit()
ordered = db.query(Category).order_by(Category.sort_order.asc(), Category.id.asc()).all()
return [_serialize(db, c, {}) for c in ordered]

106
backend/app/api/channels.py Normal file
View file

@ -0,0 +1,106 @@
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.core.auth_dependency import require_session
from app.db import get_db
from app.models.category import Category
from app.models.channel import Channel
from app.models.channel_category import channel_categories
router = APIRouter(dependencies=[Depends(require_session)])
class ChannelCategoriesUpdate(BaseModel):
category_ids: list[int]
def _category_ids_by_channel(db: Session, channel_ids: list[int]) -> dict[int, list[int]]:
if not channel_ids:
return {}
rows = db.execute(
select(channel_categories.c.channel_id, channel_categories.c.category_id).where(
channel_categories.c.channel_id.in_(channel_ids)
)
).all()
result: dict[int, list[int]] = {}
for channel_id, category_id in rows:
result.setdefault(channel_id, []).append(category_id)
return result
def _serialize(channel: Channel, category_ids: list[int]) -> dict:
return {
"id": channel.id,
"youtube_channel_id": channel.youtube_channel_id,
"title": channel.title,
"description": channel.description,
"thumbnail_url": channel.thumbnail_url,
"uploads_playlist_id": channel.uploads_playlist_id,
"subscribed": channel.subscribed,
"last_synced_at": channel.last_synced_at,
"category_ids": category_ids,
}
@router.get("/channels")
def list_channels(
subscribed: bool | None = None,
search: str | None = None,
category_id: int | None = None,
uncategorized: bool = False,
db: Session = Depends(get_db),
) -> list[dict]:
query = db.query(Channel)
if subscribed is not None:
query = query.filter(Channel.subscribed == subscribed)
if search:
query = query.filter(Channel.title.ilike(f"%{search}%"))
if uncategorized:
categorized_ids = select(channel_categories.c.channel_id)
query = query.filter(~Channel.id.in_(categorized_ids))
elif category_id is not None:
channel_ids_in_category = select(channel_categories.c.channel_id).where(
channel_categories.c.category_id == category_id
)
query = query.filter(Channel.id.in_(channel_ids_in_category))
channels = query.order_by(Channel.title.asc()).all()
category_map = _category_ids_by_channel(db, [c.id for c in channels])
return [_serialize(c, category_map.get(c.id, [])) for c in channels]
@router.get("/channels/{channel_id}")
def get_channel(channel_id: int, db: Session = Depends(get_db)) -> dict:
channel = db.get(Channel, channel_id)
if channel is None:
raise HTTPException(status_code=404, detail="Channel not found")
category_map = _category_ids_by_channel(db, [channel_id])
return _serialize(channel, category_map.get(channel_id, []))
@router.put("/channels/{channel_id}/categories")
def set_channel_categories(channel_id: int, payload: ChannelCategoriesUpdate, db: Session = Depends(get_db)) -> dict:
channel = db.get(Channel, channel_id)
if channel is None:
raise HTTPException(status_code=404, detail="Channel not found")
unique_ids = set(payload.category_ids)
if unique_ids:
found = db.query(Category.id).filter(Category.id.in_(unique_ids)).all()
found_ids = {row[0] for row in found}
missing = unique_ids - found_ids
if missing:
raise HTTPException(status_code=400, detail=f"Unknown category ids: {sorted(missing)}")
db.execute(channel_categories.delete().where(channel_categories.c.channel_id == channel_id))
if unique_ids:
db.execute(
channel_categories.insert(),
[{"channel_id": channel_id, "category_id": cid} for cid in unique_ids],
)
db.commit()
return _serialize(channel, sorted(unique_ids))

87
backend/app/api/feed.py Normal file
View file

@ -0,0 +1,87 @@
import base64
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.core.auth_dependency import require_session
from app.db import get_db
from app.models.channel import Channel
from app.models.channel_category import channel_categories
from app.models.video import Video
from app.services.video_presentation import channel_categories_map, serialize_video
router = APIRouter(dependencies=[Depends(require_session)])
DEFAULT_LIMIT = 30
MAX_LIMIT = 100
def _encode_cursor(published_at: datetime, video_id: int) -> str:
raw = f"{published_at.isoformat()}|{video_id}"
return base64.urlsafe_b64encode(raw.encode()).decode()
def _decode_cursor(cursor: str) -> tuple[datetime, int]:
try:
raw = base64.urlsafe_b64decode(cursor.encode()).decode()
published_at_raw, video_id_raw = raw.rsplit("|", 1)
published_at = datetime.fromisoformat(published_at_raw)
if published_at.tzinfo is None:
published_at = published_at.replace(tzinfo=timezone.utc)
return published_at, int(video_id_raw)
except Exception:
raise HTTPException(status_code=400, detail="Invalid cursor")
@router.get("/feed")
def get_feed(
category_id: int | None = None,
uncategorized: bool = False,
channel_id: int | None = None,
limit: int = Query(DEFAULT_LIMIT, ge=1, le=MAX_LIMIT),
cursor: str | None = None,
db: Session = Depends(get_db),
) -> dict:
query = db.query(Video)
if channel_id is not None:
query = query.filter(Video.channel_id == channel_id)
elif uncategorized:
categorized_channel_ids = select(channel_categories.c.channel_id)
query = query.filter(~Video.channel_id.in_(categorized_channel_ids))
elif category_id is not None:
channel_ids_in_category = select(channel_categories.c.channel_id).where(
channel_categories.c.category_id == category_id
)
query = query.filter(Video.channel_id.in_(channel_ids_in_category))
if cursor:
cursor_published_at, cursor_id = _decode_cursor(cursor)
query = query.filter(
(Video.published_at < cursor_published_at)
| ((Video.published_at == cursor_published_at) & (Video.id < cursor_id))
)
query = query.order_by(Video.published_at.desc(), Video.id.desc())
rows = query.limit(limit + 1).all()
next_cursor = None
if len(rows) > limit:
last_kept = rows[limit - 1]
next_cursor = _encode_cursor(last_kept.published_at, last_kept.id)
rows = rows[:limit]
channel_ids = list({v.channel_id for v in rows})
channels = {c.id: c for c in db.query(Channel).filter(Channel.id.in_(channel_ids)).all()}
categories_map = channel_categories_map(db, channel_ids)
items = []
for video in rows:
channel = channels.get(video.channel_id)
if channel is None:
continue
items.append(serialize_video(video, channel, categories_map.get(channel.id, [])))
return {"items": items, "next_cursor": next_cursor}

42
backend/app/api/health.py Normal file
View file

@ -0,0 +1,42 @@
import logging
import httpx
from fastapi import APIRouter, Depends
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.config import settings
from app.db import get_db
logger = logging.getLogger(__name__)
router = APIRouter()
def _check_database(db: Session) -> str:
try:
db.execute(text("SELECT 1"))
return "ok"
except Exception:
logger.exception("Database healthcheck failed")
return "error"
def _check_metube() -> str:
try:
response = httpx.get(settings.metube_api_base_url, timeout=5)
if response.status_code < 500:
return "ok"
return "error"
except Exception:
logger.warning("MeTube healthcheck failed", exc_info=True)
return "error"
@router.get("/health")
def health(db: Session = Depends(get_db)) -> dict:
return {
"status": "ok",
"database": _check_database(db),
"metube": _check_metube(),
}

50
backend/app/api/sync.py Normal file
View file

@ -0,0 +1,50 @@
import logging
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from app.core.auth_dependency import require_session
from app.db import get_db
from app.services import sync
from app.services.google_oauth import OAuthNotConnected
from app.services.youtube_client import YouTubeAPIError, YouTubeQuotaExceeded
logger = logging.getLogger(__name__)
router = APIRouter(dependencies=[Depends(require_session)])
@router.post("/sync/subscriptions")
def sync_subscriptions(db: Session = Depends(get_db)) -> dict:
try:
return sync.sync_subscriptions(db)
except sync.SyncInProgress:
raise HTTPException(status_code=409, detail="Subscriptions sync already in progress")
except OAuthNotConnected:
raise HTTPException(status_code=400, detail="Google account is not connected")
except YouTubeQuotaExceeded:
raise HTTPException(status_code=503, detail="YouTube API quota exhausted")
except YouTubeAPIError as exc:
raise HTTPException(status_code=502, detail=f"YouTube API error: {exc}")
@router.post("/sync/videos")
def sync_videos(db: Session = Depends(get_db)) -> dict:
try:
return sync.sync_videos(db)
except sync.SyncInProgress:
raise HTTPException(status_code=409, detail="Videos sync already in progress")
except OAuthNotConnected:
raise HTTPException(status_code=400, detail="Google account is not connected")
except YouTubeQuotaExceeded:
raise HTTPException(status_code=503, detail="YouTube API quota exhausted")
except YouTubeAPIError as exc:
raise HTTPException(status_code=502, detail=f"YouTube API error: {exc}")
@router.get("/sync/status")
def sync_status(db: Session = Depends(get_db)) -> dict:
return {
"subscriptions": sync.get_subscriptions_sync_status(db),
"videos": sync.get_videos_sync_status(db),
}

24
backend/app/api/videos.py Normal file
View file

@ -0,0 +1,24 @@
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from app.core.auth_dependency import require_session
from app.db import get_db
from app.models.channel import Channel
from app.models.video import Video
from app.services.video_presentation import channel_categories_map, serialize_video
router = APIRouter(dependencies=[Depends(require_session)])
@router.get("/videos/{youtube_video_id}")
def get_video(youtube_video_id: str, db: Session = Depends(get_db)) -> dict:
video = db.query(Video).filter(Video.youtube_video_id == youtube_video_id).one_or_none()
if video is None:
raise HTTPException(status_code=404, detail="Video not found")
channel = db.get(Channel, video.channel_id)
if channel is None:
raise HTTPException(status_code=404, detail="Channel not found")
categories = channel_categories_map(db, [channel.id]).get(channel.id, [])
return serialize_video(video, channel, categories)

32
backend/app/config.py Normal file
View file

@ -0,0 +1,32 @@
from pydantic_settings import BaseSettings, SettingsConfigDict
class Settings(BaseSettings):
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
app_env: str = "production"
app_base_url: str = "http://localhost:8080"
app_port: int = 8080
app_secret_key: str
token_encryption_key: str
database_url: str
google_client_id: str = ""
google_client_secret: str = ""
google_redirect_uri: str = ""
allowed_google_email: str = ""
metube_api_base_url: str = "http://127.0.0.1:8081"
metube_public_base_url: str = "http://127.0.0.1:8081"
metube_container_download_dir: str = "/downloads"
metube_request_timeout_seconds: int = 30
subscriptions_sync_interval_hours: int = 6
videos_sync_interval_minutes: int = 60
videos_per_channel_sync: int = 10
log_level: str = "INFO"
settings = Settings()

View file

View file

@ -0,0 +1,6 @@
from fastapi import HTTPException, Request
def require_session(request: Request) -> None:
if not request.session.get("authenticated"):
raise HTTPException(status_code=401, detail="Not authenticated")

View file

@ -0,0 +1,18 @@
from functools import lru_cache
from cryptography.fernet import Fernet
from app.config import settings
@lru_cache
def _fernet() -> Fernet:
return Fernet(settings.token_encryption_key.encode())
def encrypt_token(plaintext: str) -> str:
return _fernet().encrypt(plaintext.encode()).decode()
def decrypt_token(ciphertext: str) -> str:
return _fernet().decrypt(ciphertext.encode()).decode()

View file

@ -0,0 +1,22 @@
import re
_ISO8601_DURATION_RE = re.compile(
r"P(?:(?P<days>\d+)D)?"
r"(?:T(?:(?P<hours>\d+)H)?(?:(?P<minutes>\d+)M)?(?:(?P<seconds>\d+)S)?)?"
)
def parse_iso8601_duration(value: str | None) -> int | None:
if not value:
return None
match = _ISO8601_DURATION_RE.fullmatch(value)
if not match:
return None
parts = match.groupdict()
if not any(parts.values()):
return None
days = int(parts["days"] or 0)
hours = int(parts["hours"] or 0)
minutes = int(parts["minutes"] or 0)
seconds = int(parts["seconds"] or 0)
return days * 86400 + hours * 3600 + minutes * 60 + seconds

View file

@ -0,0 +1,29 @@
import re
_CYRILLIC_MAP = {
"а": "a", "б": "b", "в": "v", "г": "g", "д": "d", "е": "e", "ё": "e",
"ж": "zh", "з": "z", "и": "i", "й": "y", "к": "k", "л": "l", "м": "m",
"н": "n", "о": "o", "п": "p", "р": "r", "с": "s", "т": "t", "у": "u",
"ф": "f", "х": "h", "ц": "ts", "ч": "ch", "ш": "sh", "щ": "sch", "ъ": "",
"ы": "y", "ь": "", "э": "e", "ю": "yu", "я": "ya",
}
def _transliterate(text: str) -> str:
return "".join(_CYRILLIC_MAP.get(ch, ch) for ch in text)
def slugify(text: str) -> str:
text = _transliterate(text.lower())
text = re.sub(r"[^a-z0-9]+", "-", text).strip("-")
return text or "category"
def unique_slugify(text: str, existing_slugs: set[str]) -> str:
base = slugify(text)
if base not in existing_slugs:
return base
suffix = 2
while f"{base}-{suffix}" in existing_slugs:
suffix += 1
return f"{base}-{suffix}"

19
backend/app/db.py Normal file
View file

@ -0,0 +1,19 @@
from sqlalchemy import create_engine
from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker
from app.config import settings
engine = create_engine(settings.database_url, pool_pre_ping=True)
SessionLocal = sessionmaker(bind=engine, autoflush=False, autocommit=False)
class Base(DeclarativeBase):
pass
def get_db() -> Session:
db = SessionLocal()
try:
yield db
finally:
db.close()

69
backend/app/main.py Normal file
View file

@ -0,0 +1,69 @@
import logging
from contextlib import asynccontextmanager
from pathlib import Path
from fastapi import FastAPI
from fastapi.staticfiles import StaticFiles
from starlette.exceptions import HTTPException as StarletteHTTPException
from starlette.middleware.sessions import SessionMiddleware
from starlette.types import Scope
class SPAStaticFiles(StaticFiles):
"""Falls back to index.html for unknown paths so client-side routes
(e.g. /categories) work on direct navigation/refresh, not just /."""
async def get_response(self, path: str, scope: Scope):
try:
return await super().get_response(path, scope)
except StarletteHTTPException as exc:
if exc.status_code == 404:
return await super().get_response("index.html", scope)
raise
from app.api.auth import router as auth_router
from app.api.categories import router as categories_router
from app.api.channels import router as channels_router
from app.api.feed import router as feed_router
from app.api.health import router as health_router
from app.api.sync import router as sync_router
from app.api.videos import router as videos_router
from app.config import settings
from app.services.scheduler import create_scheduler
logging.basicConfig(level=settings.log_level)
logger = logging.getLogger(__name__)
@asynccontextmanager
async def lifespan(_: FastAPI):
logger.info("Application startup")
scheduler = create_scheduler()
scheduler.start()
yield
scheduler.shutdown(wait=False)
logger.info("Application shutdown")
app = FastAPI(title="MyYouTube", lifespan=lifespan)
app.add_middleware(
SessionMiddleware,
secret_key=settings.app_secret_key,
session_cookie="myyoutube_session",
same_site="lax",
https_only=settings.app_base_url.startswith("https"),
max_age=60 * 60 * 24 * 30,
)
app.include_router(health_router, prefix="/api")
app.include_router(auth_router, prefix="/api")
app.include_router(sync_router, prefix="/api")
app.include_router(channels_router, prefix="/api")
app.include_router(categories_router, prefix="/api")
app.include_router(feed_router, prefix="/api")
app.include_router(videos_router, prefix="/api")
STATIC_DIR = Path(__file__).resolve().parent.parent / "static"
if STATIC_DIR.is_dir():
app.mount("/", SPAStaticFiles(directory=STATIC_DIR, html=True), name="static")

View file

@ -0,0 +1,8 @@
from app.models.app_settings import AppSetting
from app.models.category import Category
from app.models.channel import Channel
from app.models.channel_category import channel_categories
from app.models.oauth_credentials import OAuthCredentials
from app.models.video import Video
__all__ = ["AppSetting", "Category", "Channel", "channel_categories", "OAuthCredentials", "Video"]

View file

@ -0,0 +1,17 @@
from datetime import datetime
from sqlalchemy import DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
class AppSetting(Base):
__tablename__ = "app_settings"
id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
key: Mapped[str] = mapped_column(String(255), unique=True, nullable=False)
value: Mapped[str | None] = mapped_column(String, nullable=True)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)

View file

@ -0,0 +1,19 @@
from datetime import datetime
from sqlalchemy import DateTime, Integer, String, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
class Category(Base):
__tablename__ = "categories"
id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
name: Mapped[str] = mapped_column(String(255), unique=True, nullable=False)
slug: Mapped[str] = mapped_column(String(255), unique=True, nullable=False)
sort_order: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), nullable=False)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)

View file

@ -0,0 +1,23 @@
from datetime import datetime
from sqlalchemy import Boolean, DateTime, String, Text, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
class Channel(Base):
__tablename__ = "channels"
id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
youtube_channel_id: Mapped[str] = mapped_column(String(64), unique=True, nullable=False, index=True)
title: Mapped[str] = mapped_column(String(255), nullable=False)
description: Mapped[str | None] = mapped_column(Text, nullable=True)
thumbnail_url: Mapped[str | None] = mapped_column(String, nullable=True)
uploads_playlist_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
subscribed: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
last_synced_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), nullable=False)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)

View file

@ -0,0 +1,10 @@
from sqlalchemy import Column, ForeignKey, Integer, Table
from app.db import Base
channel_categories = Table(
"channel_categories",
Base.metadata,
Column("channel_id", Integer, ForeignKey("channels.id"), primary_key=True),
Column("category_id", Integer, ForeignKey("categories.id", ondelete="CASCADE"), primary_key=True),
)

View file

@ -0,0 +1,21 @@
from datetime import datetime
from sqlalchemy import DateTime, Integer, String, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
SINGLETON_ID = 1
class OAuthCredentials(Base):
__tablename__ = "oauth_credentials"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
google_email: Mapped[str] = mapped_column(String(255), nullable=False)
encrypted_refresh_token: Mapped[str] = mapped_column(String, nullable=False)
access_token_expires_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), nullable=False)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)

View file

@ -0,0 +1,24 @@
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, Integer, String, Text, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db import Base
class Video(Base):
__tablename__ = "videos"
id: Mapped[int] = mapped_column(primary_key=True, autoincrement=True)
youtube_video_id: Mapped[str] = mapped_column(String(32), unique=True, nullable=False, index=True)
channel_id: Mapped[int] = mapped_column(ForeignKey("channels.id"), nullable=False, index=True)
title: Mapped[str] = mapped_column(String(500), nullable=False)
description: Mapped[str | None] = mapped_column(Text, nullable=True)
thumbnail_url: Mapped[str | None] = mapped_column(String, nullable=True)
published_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, index=True)
duration_seconds: Mapped[int | None] = mapped_column(Integer, nullable=True)
youtube_url: Mapped[str] = mapped_column(String, nullable=False)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), nullable=False)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
)

View file

View file

@ -0,0 +1,142 @@
import logging
from datetime import datetime, timezone
import httpx
from google.auth.transport.requests import Request as GoogleAuthRequest
from google.oauth2.credentials import Credentials
from google_auth_oauthlib.flow import Flow
from sqlalchemy.orm import Session
from app.config import settings
from app.core.crypto import decrypt_token, encrypt_token
from app.models.oauth_credentials import SINGLETON_ID, OAuthCredentials
logger = logging.getLogger(__name__)
SCOPES = [
"https://www.googleapis.com/auth/youtube.readonly",
"openid",
"https://www.googleapis.com/auth/userinfo.email",
"https://www.googleapis.com/auth/userinfo.profile",
]
AUTH_URI = "https://accounts.google.com/o/oauth2/auth"
TOKEN_URI = "https://oauth2.googleapis.com/token"
USERINFO_URI = "https://www.googleapis.com/oauth2/v3/userinfo"
REVOKE_URI = "https://oauth2.googleapis.com/revoke"
class OAuthNotConnected(Exception):
pass
def _client_config() -> dict:
return {
"web": {
"client_id": settings.google_client_id,
"client_secret": settings.google_client_secret,
"auth_uri": AUTH_URI,
"token_uri": TOKEN_URI,
"redirect_uris": [settings.google_redirect_uri],
}
}
def _build_flow(state: str | None = None) -> Flow:
return Flow.from_client_config(
_client_config(),
scopes=SCOPES,
state=state,
redirect_uri=settings.google_redirect_uri,
)
def build_authorization_url() -> tuple[str, str]:
flow = _build_flow()
auth_url, state = flow.authorization_url(
access_type="offline",
prompt="consent",
include_granted_scopes="true",
)
return auth_url, state
def exchange_code(code: str, state: str) -> Credentials:
flow = _build_flow(state=state)
flow.fetch_token(code=code)
return flow.credentials
def fetch_userinfo(access_token: str) -> dict:
response = httpx.get(
USERINFO_URI,
headers={"Authorization": f"Bearer {access_token}"},
timeout=settings.metube_request_timeout_seconds,
)
response.raise_for_status()
return response.json()
def revoke_token(token: str) -> None:
try:
httpx.post(REVOKE_URI, params={"token": token}, timeout=10)
except Exception:
logger.warning("Failed to revoke Google token", exc_info=True)
def store_credentials(db: Session, google_email: str, credentials: Credentials) -> None:
encrypted = encrypt_token(credentials.refresh_token)
expires_at = credentials.expiry
if expires_at is not None and expires_at.tzinfo is None:
expires_at = expires_at.replace(tzinfo=timezone.utc)
row = db.get(OAuthCredentials, SINGLETON_ID)
if row is None:
row = OAuthCredentials(id=SINGLETON_ID, google_email=google_email, encrypted_refresh_token=encrypted)
db.add(row)
else:
row.google_email = google_email
row.encrypted_refresh_token = encrypted
row.access_token_expires_at = expires_at
db.commit()
def clear_credentials(db: Session) -> None:
row = db.get(OAuthCredentials, SINGLETON_ID)
if row is not None:
db.delete(row)
db.commit()
def is_connected(db: Session) -> bool:
return db.get(OAuthCredentials, SINGLETON_ID) is not None
def get_connected_email(db: Session) -> str | None:
row = db.get(OAuthCredentials, SINGLETON_ID)
return row.google_email if row else None
def get_credentials(db: Session) -> Credentials:
row = db.get(OAuthCredentials, SINGLETON_ID)
if row is None:
raise OAuthNotConnected("Google account is not connected")
refresh_token = decrypt_token(row.encrypted_refresh_token)
credentials = Credentials(
token=None,
refresh_token=refresh_token,
token_uri=TOKEN_URI,
client_id=settings.google_client_id,
client_secret=settings.google_client_secret,
scopes=SCOPES,
)
credentials.refresh(GoogleAuthRequest())
expires_at = credentials.expiry
if expires_at is not None and expires_at.tzinfo is None:
expires_at = expires_at.replace(tzinfo=timezone.utc)
row.access_token_expires_at = expires_at
db.commit()
return credentials

View file

@ -0,0 +1,50 @@
import logging
from apscheduler.schedulers.background import BackgroundScheduler
from app.config import settings
from app.db import SessionLocal
from app.services import sync
logger = logging.getLogger(__name__)
def _run_subscriptions_sync() -> None:
db = SessionLocal()
try:
sync.sync_subscriptions(db)
except sync.SyncInProgress:
logger.info("Scheduled subscriptions sync skipped: already running")
except Exception:
logger.exception("Scheduled subscriptions sync failed")
finally:
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:
scheduler = BackgroundScheduler(timezone="UTC")
scheduler.add_job(
_run_subscriptions_sync,
"interval",
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

View file

@ -0,0 +1,18 @@
from sqlalchemy.orm import Session
from app.models.app_settings import AppSetting
def get_setting(db: Session, key: str) -> str | None:
row = db.query(AppSetting).filter(AppSetting.key == key).one_or_none()
return row.value if row else None
def set_setting(db: Session, key: str, value: str) -> None:
row = db.query(AppSetting).filter(AppSetting.key == key).one_or_none()
if row is None:
row = AppSetting(key=key, value=value)
db.add(row)
else:
row.value = value
db.commit()

View file

@ -0,0 +1,251 @@
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"],
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.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()

View file

@ -0,0 +1,46 @@
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.models.category import Category
from app.models.channel import Channel
from app.models.channel_category import channel_categories
from app.models.video import Video
def channel_categories_map(db: Session, channel_ids: list[int]) -> dict[int, list[dict]]:
if not channel_ids:
return {}
rows = db.execute(
select(channel_categories.c.channel_id, Category.id, Category.name)
.join(Category, Category.id == channel_categories.c.category_id)
.where(channel_categories.c.channel_id.in_(channel_ids))
).all()
result: dict[int, list[dict]] = {}
for channel_id, category_id, category_name in rows:
result.setdefault(channel_id, []).append({"id": category_id, "name": category_name})
return result
def serialize_video(video: Video, channel: Channel, categories: list[dict]) -> dict:
return {
"youtube_video_id": video.youtube_video_id,
"title": video.title,
"description": video.description,
"channel": {
"id": channel.id,
"youtube_channel_id": channel.youtube_channel_id,
"title": channel.title,
"thumbnail_url": channel.thumbnail_url,
},
"thumbnail_url": video.thumbnail_url,
"published_at": video.published_at,
"duration_seconds": video.duration_seconds,
"youtube_url": video.youtube_url,
"categories": categories,
"local": {
"available": False,
"status": "not_downloaded",
"progress_percent": None,
"media_url": None,
},
}

View file

@ -0,0 +1,168 @@
import logging
import httpx
from google.oauth2.credentials import Credentials
from app.config import settings
logger = logging.getLogger(__name__)
API_BASE = "https://www.googleapis.com/youtube/v3"
BATCH_SIZE = 50
class YouTubeQuotaExceeded(Exception):
pass
class YouTubeAPIError(Exception):
pass
def _headers(credentials: Credentials) -> dict:
return {"Authorization": f"Bearer {credentials.token}"}
def _raise_for_status(response: httpx.Response) -> None:
if response.status_code == 200:
return
try:
payload = response.json()
reason = payload.get("error", {}).get("errors", [{}])[0].get("reason", "")
message = payload.get("error", {}).get("message", response.text)
except Exception:
reason = ""
message = response.text
if response.status_code == 403 and reason in ("quotaExceeded", "dailyLimitExceeded", "rateLimitExceeded"):
raise YouTubeQuotaExceeded(message)
logger.error("YouTube API error %s: %s", response.status_code, message)
raise YouTubeAPIError(f"{response.status_code}: {message}")
def fetch_subscriptions(credentials: Credentials) -> list[dict]:
subscriptions: list[dict] = []
page_token: str | None = None
with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client:
while True:
params = {
"part": "snippet,contentDetails",
"mine": "true",
"maxResults": 50,
}
if page_token:
params["pageToken"] = page_token
response = client.get(f"{API_BASE}/subscriptions", params=params, headers=_headers(credentials))
_raise_for_status(response)
data = response.json()
for item in data.get("items", []):
snippet = item.get("snippet", {})
resource_id = snippet.get("resourceId", {})
channel_id = resource_id.get("channelId")
if not channel_id:
continue
thumbnails = snippet.get("thumbnails", {})
thumbnail = (
thumbnails.get("high") or thumbnails.get("medium") or thumbnails.get("default") or {}
).get("url")
subscriptions.append(
{
"youtube_channel_id": channel_id,
"title": snippet.get("title", ""),
"description": snippet.get("description", ""),
"thumbnail_url": thumbnail,
}
)
page_token = data.get("nextPageToken")
if not page_token:
break
return subscriptions
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"{API_BASE}/playlistItems", params=params, headers=_headers(credentials))
if response.status_code == 404:
return []
_raise_for_status(response)
data = response.json()
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_videos_details(credentials: Credentials, video_ids: list[str]) -> list[dict]:
results: list[dict] = []
with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client:
for i in range(0, len(video_ids), BATCH_SIZE):
batch = video_ids[i : i + BATCH_SIZE]
params = {
"part": "snippet,contentDetails,status",
"id": ",".join(batch),
"maxResults": BATCH_SIZE,
}
response = client.get(f"{API_BASE}/videos", params=params, headers=_headers(credentials))
_raise_for_status(response)
data = response.json()
for item in data.get("items", []):
snippet = item.get("snippet", {})
thumbnails = snippet.get("thumbnails", {})
thumbnail = (
thumbnails.get("high") or thumbnails.get("medium") or thumbnails.get("default") or {}
).get("url")
results.append(
{
"youtube_video_id": item.get("id"),
"youtube_channel_id": snippet.get("channelId"),
"title": snippet.get("title", ""),
"description": snippet.get("description", ""),
"thumbnail_url": thumbnail,
"published_at": snippet.get("publishedAt"),
"duration_iso8601": item.get("contentDetails", {}).get("duration"),
}
)
return results
def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> dict[str, str]:
result: dict[str, str] = {}
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",
"id": ",".join(batch),
"maxResults": BATCH_SIZE,
}
response = client.get(f"{API_BASE}/channels", params=params, headers=_headers(credentials))
_raise_for_status(response)
data = response.json()
for item in data.get("items", []):
channel_id = item.get("id")
uploads = (
item.get("contentDetails", {}).get("relatedPlaylists", {}).get("uploads")
)
if channel_id and uploads:
result[channel_id] = uploads
return result