From e10df8dcbd18ef43ce089238dfa92703922cfd1f Mon Sep 17 00:00:00 2001 From: vrubelroman Date: Thu, 17 Sep 2026 21:07:33 +0000 Subject: [PATCH] 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 --- .env.example | 13 +- .opencode/agent/orchestrator.md | 5 +- AGENT_TEAM.md | 3 +- backend/app/api/categories.py | 21 +- backend/app/api/channels.py | 39 +- backend/app/api/feed.py | 13 +- backend/app/config.py | 15 +- backend/app/core/auth_dependency.py | 5 + backend/app/models/channel.py | 3 +- backend/app/services/scheduler.py | 20 +- backend/app/services/sync.py | 63 +-- backend/app/services/sync_trigger.py | 91 ++++ backend/app/services/youtube_client.py | 143 +++++- frontend/src/App.css | 13 +- frontend/src/api/client.ts | 7 + frontend/src/components/AppShell.tsx | 5 +- frontend/src/components/ChannelCard.tsx | 1 + frontend/src/pages/ChannelVideos.tsx | 6 +- frontend/src/pages/Feed.tsx | 8 +- frontend/src/utils/format.ts | 7 + .../versions/0008_channel_subscriber_count.py | 26 ++ tests/test_categories.py | 77 ++++ tests/test_channels.py | 54 +++ tests/test_feed.py | 127 ++++++ tests/test_sync.py | 406 +++++++++++++++++- tests/test_sync_trigger.py | 159 +++++++ tests/test_youtube_client.py | 197 ++++++++- 27 files changed, 1420 insertions(+), 107 deletions(-) create mode 100644 backend/app/services/sync_trigger.py create mode 100644 migrations/versions/0008_channel_subscriber_count.py create mode 100644 tests/test_sync_trigger.py diff --git a/.env.example b/.env.example index 7714907..7d8fd4f 100644 --- a/.env.example +++ b/.env.example @@ -26,7 +26,16 @@ METUBE_CONTAINER_DOWNLOAD_DIR=/downloads METUBE_REQUEST_TIMEOUT_SECONDS=30 SUBSCRIPTIONS_SYNC_INTERVAL_HOURS=6 -VIDEOS_SYNC_INTERVAL_MINUTES=60 -VIDEOS_PER_CHANNEL_SYNC=10 + +# Необязательные настройки +# Авто-синк видео запускается по активности (запросы к API) после простоя +# в VIDEOS_SYNC_IDLE_HOURS с последнего завершённого синка. +VIDEOS_SYNC_IDLE_HOURS=2 +# Максимум элементов плейлиста, рассматриваемых на канал за один синк. +VIDEOS_BACKFILL_CAP=200 +# Остановить пагинацию плейлиста канала после стольких уже известных +# видео подряд (дальше — уже синхронизированная история). +VIDEOS_KNOWN_STOP_THRESHOLD=50 +NEW_VIDEOS_WINDOW_DAYS=2 LOG_LEVEL=INFO diff --git a/.opencode/agent/orchestrator.md b/.opencode/agent/orchestrator.md index 4400499..2ab2287 100644 --- a/.opencode/agent/orchestrator.md +++ b/.opencode/agent/orchestrator.md @@ -15,11 +15,12 @@ You are the only agent who normally speaks with the user. Coordinate the opencod 4. Submit the task to Coder with the Task tool (subagent_type `coder`). Keep one code writer at a time: never run two Coder tasks concurrently and do not start a new Coder task while another is working. 5. Read Coder's report (ends with `CODER_DONE`) and inspect the diff yourself. Then run Reviewer and Tester on the result. Reviewer must not edit; Tester must not fix. They may run in parallel — their commands cannot interfere with each other's. 6. Collect Critical, Major, and Minor findings with evidence. Send actionable findings back to Coder as a follow-up task. Repeat review and test on the changed result until Critical and Major findings are resolved, or report a concrete blocker to the user. -7. Give the user a concise final report: changes, affected files, review findings resolved or remaining, tests run and results, known limits, and worktree/branch. Do not claim visual or integration checks that were not performed. +6.5. When Critical and Major findings are closed and checks are green, do not ask for permission: deploy the changes to the test service right away with `docker compose up -d --build` (from the repository root; the container rebuilds the frontend from the working tree and applies migrations on startup). Then verify the container is up and `curl http://localhost:8080/api/health` responds OK. +7. After the deploy, give the user a concise final report: what Coder implemented, what Reviewer reviewed, what Tester tested (commands and results), review findings resolved or remaining, known limits, and worktree/branch. Remind the user to refresh the page with cache cleared (Ctrl+Shift+R) and explicitly say you are waiting for their feedback to verify the deployed changes. Do not claim visual or integration checks that were not performed. ## Working rules -- Do not merge, deploy, publish, or commit unless the user requested it or existing authorization covers it. +- Do not merge, publish, or commit unless the user requested it or existing authorization covers it. Deploying to the test service is authorized by this workflow (step 6.5); deploying elsewhere still requires the user's go-ahead. - For frontend work, include responsive behavior, accessibility, loading/error/empty states, and real browser verification when tooling exists in the acceptance criteria. - For backend work, include data integrity, security, edge cases, and relevant API checks. diff --git a/AGENT_TEAM.md b/AGENT_TEAM.md index bfb8954..130b201 100644 --- a/AGENT_TEAM.md +++ b/AGENT_TEAM.md @@ -9,7 +9,8 @@ - декомпозирует задачу и отправляет её Coder'у через Task tool; - после Coder'а запускает Reviewer (только чтение) и Tester (проверка без правок) параллельно; - возвращает Coder'у actionable findings и повторяет цикл, пока Critical/Major не закрыты; -- отдаёт финальный отчёт с изменениями, проверками и оставшимися рисками. +- когда Critical/Major закрыты и проверки зелёные — без вопроса деплоит на тестовый сервис (`docker compose up -d --build`, затем `curl http://localhost:8080/api/health`); +- отдаёт финальный отчёт уже после деплоя: что написано, просмотрено и протестировано, оставшиеся риски и просьба проверить в браузере с очисткой кэша (Ctrl+Shift+R). Вручную писать субагентам не нужно — их вызывает только Orchestrator. diff --git a/backend/app/api/categories.py b/backend/app/api/categories.py index ec02a08..2f55135 100644 --- a/backend/app/api/categories.py +++ b/backend/app/api/categories.py @@ -1,14 +1,19 @@ +from datetime import datetime, timedelta, timezone + 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.config import settings 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 import Channel from app.models.channel_category import channel_categories +from app.models.video import Video router = APIRouter(dependencies=[Depends(require_session)]) @@ -40,13 +45,15 @@ def _name_taken(db: Session, name: str, exclude_id: int | None = None) -> bool: return any(row[0].casefold() == target for row in query.all()) -def _serialize(db: Session, category: Category, counts: dict[int, int]) -> dict: +def _serialize(db: Session, category: Category, counts: dict[int, int], new_videos: dict[int, int] | None = None) -> dict: + new_videos = new_videos or {} return { "id": category.id, "name": category.name, "slug": category.slug, "sort_order": category.sort_order, "channel_count": counts.get(category.id, 0), + "new_videos_count": new_videos.get(category.id, 0), } @@ -59,7 +66,17 @@ def list_categories(db: Session = Depends(get_db)) -> list[dict]: .all() ) counts = dict(count_rows) - return [_serialize(db, c, counts) for c in categories] + since = datetime.now(timezone.utc) - timedelta(days=settings.new_videos_window_days) + new_videos_rows = ( + db.query(channel_categories.c.category_id, func.count(Video.id)) + .join(Channel, Channel.id == channel_categories.c.channel_id) + .join(Video, Video.channel_id == Channel.id) + .filter(Channel.subscribed.is_(True), Video.published_at >= since) + .group_by(channel_categories.c.category_id) + .all() + ) + new_videos = dict(new_videos_rows) + return [_serialize(db, c, counts, new_videos) for c in categories] @router.post("/categories", status_code=201) diff --git a/backend/app/api/channels.py b/backend/app/api/channels.py index ea6cf6e..f3a87dd 100644 --- a/backend/app/api/channels.py +++ b/backend/app/api/channels.py @@ -1,15 +1,18 @@ import logging +from datetime import datetime, timedelta, timezone from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel -from sqlalchemy import select +from sqlalchemy import func, select from sqlalchemy.orm import Session +from app.config import settings 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 +from app.models.video import Video from app.services import sync from app.services.google_oauth import OAuthNotConnected from app.services.youtube_client import YouTubeAPIError, YouTubeInsufficientScope, YouTubeQuotaExceeded @@ -37,7 +40,22 @@ def _category_ids_by_channel(db: Session, channel_ids: list[int]) -> dict[int, l return result -def _serialize(channel: Channel, category_ids: list[int]) -> dict: +def _new_videos_counts(db: Session, channel_ids: list[int]) -> dict[int, int]: + """Videos published within the new-videos window, per channel, in one + aggregate query.""" + if not channel_ids: + return {} + since = datetime.now(timezone.utc) - timedelta(days=settings.new_videos_window_days) + rows = ( + db.query(Video.channel_id, func.count(Video.id)) + .filter(Video.channel_id.in_(channel_ids), Video.published_at >= since) + .group_by(Video.channel_id) + .all() + ) + return {channel_id: count for channel_id, count in rows} + + +def _serialize(channel: Channel, category_ids: list[int], new_videos_count: int = 0) -> dict: return { "id": channel.id, "youtube_channel_id": channel.youtube_channel_id, @@ -45,9 +63,11 @@ def _serialize(channel: Channel, category_ids: list[int]) -> dict: "description": channel.description, "thumbnail_url": channel.thumbnail_url, "uploads_playlist_id": channel.uploads_playlist_id, + "subscriber_count": channel.subscriber_count, "subscribed": channel.subscribed, "last_synced_at": channel.last_synced_at, "category_ids": category_ids, + "new_videos_count": new_videos_count, } @@ -75,8 +95,10 @@ def list_channels( 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] + channel_ids = [c.id for c in channels] + category_map = _category_ids_by_channel(db, channel_ids) + new_videos_map = _new_videos_counts(db, channel_ids) + return [_serialize(c, category_map.get(c.id, []), new_videos_map.get(c.id, 0)) for c in channels] @router.get("/channels/{channel_id}") @@ -85,7 +107,8 @@ def get_channel(channel_id: int, db: Session = Depends(get_db)) -> dict: 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, [])) + new_videos_map = _new_videos_counts(db, [channel_id]) + return _serialize(channel, category_map.get(channel_id, []), new_videos_map.get(channel_id, 0)) @router.put("/channels/{channel_id}/categories") @@ -110,7 +133,8 @@ def set_channel_categories(channel_id: int, payload: ChannelCategoriesUpdate, db ) db.commit() - return _serialize(channel, sorted(unique_ids)) + new_videos_map = _new_videos_counts(db, [channel_id]) + return _serialize(channel, sorted(unique_ids), new_videos_map.get(channel_id, 0)) @router.post("/channels/{channel_id}/unsubscribe") @@ -139,4 +163,5 @@ def unsubscribe_channel(channel_id: int, db: Session = Depends(get_db)) -> dict: raise HTTPException(status_code=502, detail="YouTube is unavailable") category_map = _category_ids_by_channel(db, [channel_id]) - return _serialize(channel, category_map.get(channel_id, [])) + new_videos_map = _new_videos_counts(db, [channel_id]) + return _serialize(channel, category_map.get(channel_id, []), new_videos_map.get(channel_id, 0)) diff --git a/backend/app/api/feed.py b/backend/app/api/feed.py index a50e81e..e94f97c 100644 --- a/backend/app/api/feed.py +++ b/backend/app/api/feed.py @@ -1,10 +1,11 @@ import base64 -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy import func, or_, select from sqlalchemy.orm import Session +from app.config import settings from app.core.auth_dependency import require_session from app.db import get_db from app.models.channel import Channel @@ -80,6 +81,7 @@ def get_feed( channel_id: int | None = None, downloaded: bool = False, search: str | None = Query(None, max_length=200), + new_only: bool = False, limit: int = Query(DEFAULT_LIMIT, ge=1, le=MAX_LIMIT), cursor: str | None = None, db: Session = Depends(get_db), @@ -97,6 +99,15 @@ def get_feed( ) query = query.filter(Video.channel_id.in_(channel_ids_in_category)) + if new_only: + since = datetime.now(timezone.utc) - timedelta(days=settings.new_videos_window_days) + query = query.filter(Video.published_at >= since) + # Align with the sidebar badge count (categories.py), which only + # counts subscribed channels: after unsubscribing, a channel's + # fresh videos must disappear from "new only" too. + subscribed_channel_ids = select(Channel.id).where(Channel.subscribed.is_(True)) + query = query.filter(Video.channel_id.in_(subscribed_channel_ids)) + if downloaded: query = query.filter(Video.id.in_(_completed_video_ids())) diff --git a/backend/app/config.py b/backend/app/config.py index d6e5d37..1bab4e9 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -32,8 +32,19 @@ class Settings(BaseSettings): metube_request_timeout_seconds: int = 30 subscriptions_sync_interval_hours: int = 6 - videos_sync_interval_minutes: int = 60 - videos_per_channel_sync: int = 10 + # Videos sync is no longer scheduled: it is triggered by user activity + # (any authenticated request) after this much idle time since the last + # completed sync. + videos_sync_idle_hours: int = 2 + # Hard history-depth limit per channel per videos sync: at most the + # newest N playlist items are considered; anything older is never + # backfilled by design (full channel history is out of scope). + videos_backfill_cap: int = 200 + # Stop playlistItems pagination for a channel once this many already + # known video ids are seen in a row (the rest of the playlist is history + # we have already synced). + videos_known_stop_threshold: int = 50 + new_videos_window_days: int = 2 log_level: str = "INFO" diff --git a/backend/app/core/auth_dependency.py b/backend/app/core/auth_dependency.py index 70b0483..1b7716e 100644 --- a/backend/app/core/auth_dependency.py +++ b/backend/app/core/auth_dependency.py @@ -1,6 +1,11 @@ from fastapi import HTTPException, Request +from app.services.sync_trigger import maybe_trigger_videos_sync + def require_session(request: Request) -> None: if not request.session.get("authenticated"): raise HTTPException(status_code=401, detail="Not authenticated") + # Every authenticated request may trigger a background videos sync when it + # is due (cheap check: lock + one app_settings row read). + maybe_trigger_videos_sync() diff --git a/backend/app/models/channel.py b/backend/app/models/channel.py index 7384bdf..7cbb5ef 100644 --- a/backend/app/models/channel.py +++ b/backend/app/models/channel.py @@ -1,6 +1,6 @@ from datetime import datetime -from sqlalchemy import Boolean, DateTime, String, Text, func +from sqlalchemy import BigInteger, Boolean, DateTime, String, Text, func from sqlalchemy.orm import Mapped, mapped_column from app.db import Base @@ -16,6 +16,7 @@ class Channel(Base): 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) + subscriber_count: Mapped[int | None] = mapped_column(BigInteger, 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) diff --git a/backend/app/services/scheduler.py b/backend/app/services/scheduler.py index 301203a..92d5654 100644 --- a/backend/app/services/scheduler.py +++ b/backend/app/services/scheduler.py @@ -21,19 +21,9 @@ def _run_subscriptions_sync() -> None: 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: + # Videos sync is not scheduled anymore: it runs on user activity via + # services.sync_trigger.maybe_trigger_videos_sync (see require_session). scheduler = BackgroundScheduler(timezone="UTC") scheduler.add_job( _run_subscriptions_sync, @@ -41,10 +31,4 @@ def create_scheduler() -> BackgroundScheduler: 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 diff --git a/backend/app/services/sync.py b/backend/app/services/sync.py index c15778f..95ec767 100644 --- a/backend/app/services/sync.py +++ b/backend/app/services/sync.py @@ -52,6 +52,10 @@ 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") @@ -108,10 +112,13 @@ def sync_subscriptions(db: Session) -> dict: 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(): + for channel_id, data in uploads.items(): channel = existing.get(channel_id) if channel is not None: - channel.uploads_playlist_id = uploads_playlist_id + 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") @@ -166,21 +173,38 @@ def sync_videos(db: Session) -> dict: ) 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: - video_ids = youtube_client.fetch_playlist_video_ids( - credentials, channel.uploads_playlist_id, settings.videos_per_channel_sync + # 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(video_ids) + 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)) - existing = {v.youtube_video_id: v for v in db.query(Video).all()} added = 0 - updated = 0 skipped = 0 for item in details: @@ -197,9 +221,8 @@ def sync_videos(db: Session) -> dict: 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( + db.add( + Video( youtube_video_id=item["youtube_video_id"], channel_id=channel.id, title=item["title"], @@ -209,18 +232,8 @@ def sync_videos(db: Session) -> dict: 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 + ) + added += 1 db.commit() @@ -230,12 +243,14 @@ def sync_videos(db: Session) -> dict: "finished_at": _now_iso(), "error": None, "videos_added": added, - "videos_updated": updated, + # 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 updated=%d skipped=%d", added, updated, skipped) + logger.info("Videos sync completed: added=%d skipped=%d", added, skipped) return result except Exception as exc: diff --git a/backend/app/services/sync_trigger.py b/backend/app/services/sync_trigger.py new file mode 100644 index 0000000..d941fc6 --- /dev/null +++ b/backend/app/services/sync_trigger.py @@ -0,0 +1,91 @@ +"""Activity-triggered videos sync. + +The hourly scheduled videos sync is gone; instead, every authenticated request +(require_session) calls maybe_trigger_videos_sync(). The check must stay cheap: +if the sync is not running and enough idle time has passed since the last +completed sync, a background thread is started which calls sync.sync_videos +with its own DB session. +""" + +import json +import logging +import threading +from datetime import datetime, timedelta, timezone + +from app.config import settings +from app.db import SessionLocal +from app.services import sync +from app.services.state import get_setting + +logger = logging.getLogger(__name__) + + +def is_videos_sync_due(finished_at_iso: str | None, now: datetime, idle_hours: int) -> bool: + """Pure predicate: should an automatic videos sync start now? + + True when the last completed sync's finished_at is missing/unparseable or + older than `idle_hours`.""" + if not finished_at_iso: + return True + if not isinstance(finished_at_iso, str): + return True + try: + finished_at = datetime.fromisoformat(finished_at_iso.replace("Z", "+00:00")) + except (ValueError, TypeError): + return True + if finished_at.tzinfo is None: + finished_at = finished_at.replace(tzinfo=timezone.utc) + if now.tzinfo is None: + now = now.replace(tzinfo=timezone.utc) + return now - finished_at >= timedelta(hours=idle_hours) + + +def _extract_finished_at(raw: str | None) -> str | None: + if not raw: + return None + try: + payload = json.loads(raw) + except (ValueError, TypeError): + return None + if not isinstance(payload, dict): + return None + return payload.get("finished_at") + + +def _run_videos_sync() -> None: + db = SessionLocal() + try: + sync.sync_videos(db) + except sync.SyncInProgress: + # Lost the race against a manual sync or another trigger: fine, the + # other run will record the status. + logger.info("Triggered videos sync skipped: already running") + except Exception: + logger.exception("Triggered videos sync failed") + finally: + db.close() + + +def maybe_trigger_videos_sync() -> None: + """Cheap per-request check; never raises. Starts a background videos sync + when due, or does nothing when a sync is already running.""" + try: + if sync.is_videos_sync_running(): + return + + db = SessionLocal() + try: + raw = get_setting(db, sync.VIDEOS_SYNC_STATUS_KEY) + finished_at = _extract_finished_at(raw) + if not is_videos_sync_due(finished_at, datetime.now(timezone.utc), settings.videos_sync_idle_hours): + return + finally: + db.close() + + logger.info( + "Videos sync due (last finished: %s), starting in background thread", finished_at + ) + threading.Thread(target=_run_videos_sync, daemon=True).start() + except Exception: + # Never break the request because of the trigger itself. + logger.exception("Failed to evaluate videos sync trigger") diff --git a/backend/app/services/youtube_client.py b/backend/app/services/youtube_client.py index fc8aae0..0a64f07 100644 --- a/backend/app/services/youtube_client.py +++ b/backend/app/services/youtube_client.py @@ -9,6 +9,12 @@ logger = logging.getLogger(__name__) BATCH_SIZE = 50 +# Hard safety cap for playlistItems pagination: at most this many pages +# (BATCH_SIZE items each, i.e. 500 ids) per playlist, so the pageToken loop +# can never run away. Raise both the cap and the caller's max_results +# together if more is ever needed. +MAX_PLAYLIST_PAGES = 10 + class YouTubeQuotaExceeded(Exception): pass @@ -102,26 +108,108 @@ def fetch_subscriptions(credentials: Credentials) -> list[dict]: 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"{settings.youtube_api_base_url}/playlistItems", params=params, headers=_headers(credentials)) - if response.status_code == 404: - return [] - _raise_for_status(response) - data = response.json() + """Low-level primitive: fetch up to `max_results` playlist item ids, + newest first, with no knowledge of what is already synced. The videos sync + uses fetch_playlist_video_ids_incremental instead; this stays as the plain + paginated helper (kept for tests and any future non-incremental callers).""" + video_ids: list[str] = [] + page_token: str | None = None + pages_fetched = 0 + + with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: + while pages_fetched < MAX_PLAYLIST_PAGES and len(video_ids) < max_results: + params = { + "part": "contentDetails", + "playlistId": playlist_id, + "maxResults": min(max_results - len(video_ids), BATCH_SIZE), + } + if page_token: + params["pageToken"] = page_token + + response = client.get(f"{settings.youtube_api_base_url}/playlistItems", params=params, headers=_headers(credentials)) + if response.status_code == 404: + return [] + _raise_for_status(response) + data = response.json() + + for item in data.get("items", []): + video_id = item.get("contentDetails", {}).get("videoId") + if video_id: + video_ids.append(video_id) + + pages_fetched += 1 + page_token = data.get("nextPageToken") + if not page_token: + break - 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_playlist_video_ids_incremental( + credentials: Credentials, + playlist_id: str, + known_ids: set[str], + max_results: int, + stop_threshold: int, +) -> list[str]: + """Fetch up to `max_results` previously-unknown video ids from a playlist, + stopping pagination early once `stop_threshold` consecutive ids that are + already in `known_ids` are encountered. Uploads playlists are ordered + newest-first, so a long run of known ids means we reached history that was + already synced and there is nothing new further down. + + `max_results` is a hard history-depth limit: at most the newest + `max_results` videos of the channel are ever considered, and everything + older than that is intentionally not backfilled (we do not mirror full + channel history). The early stop on known ids only saves pages on repeated + syncs within that window. + + Returns only the unknown ids, in playlist order.""" + new_ids: list[str] = [] + seen_new: set[str] = set() + consecutive_known = 0 + page_token: str | None = None + pages_fetched = 0 + + with httpx.Client(timeout=settings.metube_request_timeout_seconds) as client: + while pages_fetched < MAX_PLAYLIST_PAGES and len(new_ids) < max_results: + params = { + "part": "contentDetails", + "playlistId": playlist_id, + "maxResults": BATCH_SIZE, + } + if page_token: + params["pageToken"] = page_token + + response = client.get(f"{settings.youtube_api_base_url}/playlistItems", params=params, headers=_headers(credentials)) + if response.status_code == 404: + return new_ids + _raise_for_status(response) + data = response.json() + + for item in data.get("items", []): + video_id = item.get("contentDetails", {}).get("videoId") + if not video_id: + continue + if video_id in known_ids or video_id in seen_new: + consecutive_known += 1 + if consecutive_known >= stop_threshold: + return new_ids + else: + seen_new.add(video_id) + new_ids.append(video_id) + consecutive_known = 0 + if len(new_ids) >= max_results: + return new_ids + + pages_fetched += 1 + page_token = data.get("nextPageToken") + if not page_token: + break + + return new_ids + + def fetch_videos_details(credentials: Credentials, video_ids: list[str]) -> list[dict]: results: list[dict] = [] @@ -171,14 +259,24 @@ def unsubscribe(credentials: Credentials, youtube_subscription_id: str) -> None: _raise_for_status(response) -def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> dict[str, str]: - result: dict[str, str] = {} +def _parse_subscriber_count(statistics: dict) -> int | None: + raw = statistics.get("subscriberCount") + if raw is None: + return None + try: + return int(raw) + except (TypeError, ValueError): + return None + + +def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> dict[str, dict]: + result: dict[str, dict] = {} 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", + "part": "snippet,contentDetails,statistics", "id": ",".join(batch), "maxResults": BATCH_SIZE, } @@ -191,7 +289,10 @@ def fetch_uploads_playlists(credentials: Credentials, channel_ids: list[str]) -> uploads = ( item.get("contentDetails", {}).get("relatedPlaylists", {}).get("uploads") ) - if channel_id and uploads: - result[channel_id] = uploads + if channel_id: + result[channel_id] = { + "uploads_playlist_id": uploads or None, + "subscriber_count": _parse_subscriber_count(item.get("statistics", {})), + } return result diff --git a/frontend/src/App.css b/frontend/src/App.css index bd94692..b2cbc02 100644 --- a/frontend/src/App.css +++ b/frontend/src/App.css @@ -33,7 +33,12 @@ .category-dot { display: inline-block; width: 8px; height: 8px; border-radius: 50%; flex: none; background: var(--accent); box-shadow: 0 0 0 3px var(--accent-soft); } .category-link { gap: 15px; padding-left: 18px; } .category-link.active .category-dot { background: var(--accent); } -.sidebar-count { margin-left: auto; font-size: 11px; color: var(--subtle); } +.sidebar-category-row { display: flex; align-items: center; gap: 4px; } +.sidebar-category-row .sidebar-link { flex: 1; min-width: 0; } +.sidebar-badge { margin-left: auto; flex: none; padding: 1px 8px; border-radius: 999px; background: rgba(255,107,104,.13); color: var(--error); font-size: 11px; font-weight: 650; line-height: 1.5; white-space: nowrap; } +.sidebar-badge-link { cursor: pointer; text-decoration: none; transition: background .16s; } +.sidebar-badge-link:hover { background: rgba(255,107,104,.26); } +.sidebar-badge-link:focus-visible { outline: 2px solid var(--accent); outline-offset: 2px; } .sidebar-add { color: var(--accent); font-size: 13px; margin-top: 4px; } .sidebar-hint { margin: 0; padding: 5px 13px 9px; color: var(--subtle); font-size: 12px; } .sidebar-bottom { margin-top: auto; padding-bottom: 0; } @@ -46,6 +51,8 @@ h1, h2, h3, p { margin-top: 0; } h1 { color: var(--text); font-size: clamp(25px, 2.25vw, 34px); line-height: 1.18; letter-spacing: -.035em; font-weight: 700; margin-bottom: 8px; } h2 { color: var(--text); font-size: 18px; line-height: 1.3; font-weight: 650; } .page-subtitle { margin: 0; color: var(--muted); font-size: 14px; line-height: 1.5; } +.page-subtitle a { color: var(--accent); text-decoration: none; font-weight: 650; } +.page-subtitle a:hover { text-decoration: underline; } button, .button-primary, .button-secondary, .button-quiet, .button-danger, .button-link { transition: background .16s, color .16s, border-color .16s; } .button-primary, .button-secondary, .button-danger, .button-link { min-height: 42px; display: inline-flex; align-items: center; justify-content: center; gap: 8px; border-radius: 8px; padding: 0 15px; text-decoration: none; font-size: 13px; font-weight: 650; white-space: nowrap; } .button-primary { border: 1px solid var(--accent); background: var(--accent); color: #081521; } @@ -125,6 +132,7 @@ button, .button-primary, .button-secondary, .button-quiet, .button-danger, .butt .channel-title { font-size: 15px; font-weight: 650; text-decoration: none; } .channel-title:hover { color: var(--accent); } .badge { font-size: 11px; color: var(--subtle); } +.badge-new { padding: 2px 8px; border-radius: 999px; background: rgba(255,107,104,.13); color: var(--error); font-weight: 650; white-space: nowrap; } .unsubscribe-button { margin-left: auto; min-height: 28px; border: 0; border-radius: 6px; padding: 0 8px; color: var(--subtle); background: transparent; font-size: 11px; } .unsubscribe-button:hover { color: var(--error); background: rgba(255,107,104,.1); } .channel-categories { display: flex; align-items: center; flex-wrap: wrap; gap: 7px; margin-top: 10px; } @@ -271,7 +279,8 @@ button, .button-primary, .button-secondary, .button-quiet, .button-danger, .butt } @media (min-width: 621px) and (max-width: 699px) { .video-grid { grid-template-columns: 1fr; } } @media (pointer: coarse) { - .video-actions .button-link, .video-actions .button-secondary, .download-badge, .category-nav button, .local-filter-nav a, .mobile-category-nav a, .chip, .chip-edit, .category-checkbox, .unsubscribe-button { min-height: 44px; } + .video-actions .button-link, .video-actions .button-secondary, .download-badge, .category-nav button, .local-filter-nav a, .mobile-category-nav a, .chip, .chip-edit, .category-checkbox, .unsubscribe-button, .sidebar-badge { min-height: 44px; } + .sidebar-badge { display: inline-flex; align-items: center; } .reorder-buttons .icon-button { width: 40px; height: 40px; } .popover-create input, .popover-create button { height: 44px; } .popover-create button { width: 44px; } diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index cf106c9..8480ae7 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -58,9 +58,11 @@ export interface ChannelDto { description: string | null thumbnail_url: string | null uploads_playlist_id: string | null + subscriber_count: number | null subscribed: boolean last_synced_at: string | null category_ids: number[] + new_videos_count: number } export function getChannel(channelId: number) { @@ -96,6 +98,7 @@ export interface CategoryDto { slug: string sort_order: number channel_count: number + new_videos_count: number } export function listCategories() { @@ -211,6 +214,8 @@ export function getFeed( downloaded?: boolean search?: string cursor?: string + limit?: number + newOnly?: boolean } = {}, ) { const qs = new URLSearchParams() @@ -219,7 +224,9 @@ export function getFeed( if (params.channelId != null) qs.set('channel_id', String(params.channelId)) if (params.downloaded) qs.set('downloaded', 'true') if (params.search) qs.set('search', params.search) + if (params.newOnly) qs.set('new_only', 'true') if (params.cursor) qs.set('cursor', params.cursor) + if (params.limit != null) qs.set('limit', String(params.limit)) const suffix = qs.toString() ? `?${qs.toString()}` : '' return request<{ items: FeedVideoDto[]; next_cursor: string | null }>(`/api/feed${suffix}`) } diff --git a/frontend/src/components/AppShell.tsx b/frontend/src/components/AppShell.tsx index f7bba3e..b063d5c 100644 --- a/frontend/src/components/AppShell.tsx +++ b/frontend/src/components/AppShell.tsx @@ -49,7 +49,10 @@ function Sidebar({ close }: { close: () => void }) {
Категории
- {categories.map((category) => `sidebar-link category-link ${isActive ? 'active' : ''}`}>{category.name}{category.channel_count})} + {categories.map((category) =>
+ `sidebar-link category-link ${isActive ? 'active' : ''}`}>{category.name} + {category.new_videos_count > 0 && {category.new_videos_count}} +
)} {categories.length === 0 && !categoriesQuery.isLoading &&

Пока нет категорий

} Новая категория
diff --git a/frontend/src/components/ChannelCard.tsx b/frontend/src/components/ChannelCard.tsx index cd5ab4b..559749c 100644 --- a/frontend/src/components/ChannelCard.tsx +++ b/frontend/src/components/ChannelCard.tsx @@ -82,6 +82,7 @@ function ChannelCard({ channel, categories }: Props) { {channel.thumbnail_url ? : {channel.title.charAt(0).toUpperCase()}}
{channel.title} + {channel.new_videos_count > 0 && {channel.new_videos_count} новых} {!channel.subscribed && отписан} {channel.subscribed && ( }