diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 97584986..18251df8 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -1,8 +1,9 @@ -## Contributing to Tube Archivist +# Contributing to Tube Archivist Welcome, and thanks for showing interest in improving Tube Archivist! ## Table of Content +- [Next Steps](#next-steps) - [How to open an issue](#how-to-open-an-issue) - [Bug Report](#bug-report) - [Feature Request](#feature-request) @@ -14,8 +15,18 @@ Welcome, and thanks for showing interest in improving Tube Archivist! - [Development Environment](#development-environment) --- +## Next Steps +Going forward, this project will focus on developing a new modern frontend. + +- For the time being, don't open any new PRs that are not towards the new frontend. +- New features requests likely won't get accepted during this process. +- Depending on the severity, bug reports may or may not get fixed during this time. +- When in doubt, reach out. + +Join us on [Discord](https://tubearchivist.com/discord) if you want to help with that process. + ## How to open an issue -Please read this carefully before opening any [issue](https://github.com/tubearchivist/tubearchivist/issues) on GitHub. +Please read this carefully before opening any [issue](https://github.com/tubearchivist/tubearchivist/issues) on GitHub. Make sure you read [Next Steps](#next-steps) above. **Do**: - Do provide details and context, this matters a lot and makes it easier for people to help. @@ -66,6 +77,8 @@ IMPORTANT: When receiving help, contribute back to the community by improving th ## How to make a Pull Request +Make sure you read [Next Steps](#next-steps) above. + Thank you for contributing and helping improve this project. Focus for the foreseeable future is on improving and building on existing functionality, *not* on adding and expanding the application. This is a quick checklist to help streamline the process: diff --git a/Dockerfile b/Dockerfile index 1f7f505a..7e8e0d14 100644 --- a/Dockerfile +++ b/Dockerfile @@ -3,7 +3,7 @@ # First stage to build python wheel -FROM python:3.11.3-slim-bullseye AS builder +FROM python:3.11.8-slim-bookworm AS builder ARG TARGETPLATFORM RUN apt-get update && apt-get install -y --no-install-recommends \ @@ -14,12 +14,12 @@ COPY ./tubearchivist/requirements.txt /requirements.txt RUN pip install --user -r requirements.txt # build ffmpeg -FROM python:3.11.3-slim-bullseye as ffmpeg-builder +FROM python:3.11.8-slim-bookworm as ffmpeg-builder COPY docker_assets/ffmpeg_download.py ffmpeg_download.py RUN python ffmpeg_download.py $TARGETPLATFORM # build final image -FROM python:3.11.3-slim-bullseye as tubearchivist +FROM python:3.11.8-slim-bookworm as tubearchivist ARG TARGETPLATFORM ARG INSTALL_DEBUG diff --git a/docker_assets/run.sh b/docker_assets/run.sh index 4b60e585..8cda03c0 100644 --- a/docker_assets/run.sh +++ b/docker_assets/run.sh @@ -17,7 +17,7 @@ python manage.py ta_startup # start all tasks nginx & -celery -A home.tasks worker --loglevel=INFO --max-tasks-per-child 10 & +celery -A home.celery worker --loglevel=INFO --max-tasks-per-child 10 & celery -A home beat --loglevel=INFO \ - -s "${BEAT_SCHEDULE_PATH:-${cachedir}/celerybeat-schedule}" & + --scheduler django_celery_beat.schedulers:DatabaseScheduler & uwsgi --ini uwsgi.ini diff --git a/tubearchivist/api/urls.py b/tubearchivist/api/urls.py index b03547cd..4ad57a73 100644 --- a/tubearchivist/api/urls.py +++ b/tubearchivist/api/urls.py @@ -121,6 +121,16 @@ urlpatterns = [ views.TaskIDView.as_view(), name="api-task-id", ), + path( + "schedule/", + views.ScheduleView.as_view(), + name="api-schedule", + ), + path( + "schedule/notification/", + views.ScheduleNotification.as_view(), + name="api-schedule-notification", + ), path( "config/user/", views.UserConfigView.as_view(), diff --git a/tubearchivist/api/views.py b/tubearchivist/api/views.py index 77c345ab..1daca9dc 100644 --- a/tubearchivist/api/views.py +++ b/tubearchivist/api/views.py @@ -10,6 +10,7 @@ from api.src.aggs import ( WatchProgress, ) from api.src.search_processor import SearchProcess +from home.models import CustomPeriodicTask from home.src.download.queue import PendingInteract from home.src.download.subscriptions import ( ChannelSubscription, @@ -27,13 +28,14 @@ from home.src.index.playlist import YoutubePlaylist from home.src.index.reindex import ReindexProgress from home.src.index.video import SponsorBlock, YoutubeVideo from home.src.ta.config import AppConfig, ReleaseVersion +from home.src.ta.notify import Notifications, get_all_notifications from home.src.ta.settings import EnvironmentSettings from home.src.ta.ta_redis import RedisArchivist +from home.src.ta.task_config import TASK_CONFIG from home.src.ta.task_manager import TaskCommand, TaskManager from home.src.ta.urlparser import Parser from home.src.ta.users import UserConfig from home.tasks import ( - BaseTask, check_reindex, download_pending, extrac_dl, @@ -592,7 +594,7 @@ class DownloadApiView(ApiBaseView): """ search_base = "ta_download/_doc/" - valid_status = ["pending", "ignore", "priority"] + valid_status = ["pending", "ignore", "ignore-force", "priority"] permission_classes = [AdminOnly] def get(self, request, video_id): @@ -609,6 +611,11 @@ class DownloadApiView(ApiBaseView): print(message) return Response({"message": message}, status=400) + if item_status == "ignore-force": + extrac_dl.delay(video_id, status="ignore") + message = f"{video_id}: set status to ignore" + return Response(request.data) + _, status_code = PendingInteract(video_id).get_item() if status_code == 404: message = f"{video_id}: item not found {status_code}" @@ -911,7 +918,7 @@ class TaskNameListView(ApiBaseView): def get(self, request, task_name): """handle get request""" # pylint: disable=unused-argument - if task_name not in BaseTask.TASK_CONFIG: + if task_name not in TASK_CONFIG: message = {"message": "invalid task name"} return Response(message, status=404) @@ -926,12 +933,12 @@ class TaskNameListView(ApiBaseView): 400 if task can't be started here without argument """ # pylint: disable=unused-argument - task_config = BaseTask.TASK_CONFIG.get(task_name) + task_config = TASK_CONFIG.get(task_name) if not task_config: message = {"message": "invalid task name"} return Response(message, status=404) - if not task_config.get("api-start"): + if not task_config.get("api_start"): message = {"message": "can not start task through this endpoint"} return Response(message, status=400) @@ -970,16 +977,16 @@ class TaskIDView(ApiBaseView): message = {"message": "task id not found"} return Response(message, status=404) - task_conf = BaseTask.TASK_CONFIG.get(task_result.get("name")) + task_conf = TASK_CONFIG.get(task_result.get("name")) if command == "stop": - if not task_conf.get("api-stop"): + if not task_conf.get("api_stop"): message = {"message": "task can not be stopped"} return Response(message, status=400) message_key = self._build_message_key(task_conf, task_id) TaskCommand().stop(task_id, message_key) if command == "kill": - if not task_conf.get("api-stop"): + if not task_conf.get("api_stop"): message = {"message": "task can not be killed"} return Response(message, status=400) @@ -992,6 +999,56 @@ class TaskIDView(ApiBaseView): return f"message:{task_conf.get('group')}:{task_id.split('-')[0]}" +class ScheduleView(ApiBaseView): + """resolves to /api/schedule/ + DEL: delete schedule for task + """ + + permission_classes = [AdminOnly] + + def delete(self, request): + """delete schedule by task_name query""" + task_name = request.data.get("task_name") + try: + task = CustomPeriodicTask.objects.get(name=task_name) + except CustomPeriodicTask.DoesNotExist: + message = {"message": "task_name not found"} + return Response(message, status=404) + + _ = task.delete() + + return Response({"success": True}) + + +class ScheduleNotification(ApiBaseView): + """resolves to /api/schedule/notification/ + GET: get all schedule notifications + DEL: delete notification + """ + + def get(self, request): + """handle get request""" + + return Response(get_all_notifications()) + + def delete(self, request): + """handle delete""" + + task_name = request.data.get("task_name") + url = request.data.get("url") + + if not TASK_CONFIG.get(task_name): + message = {"message": "task_name not found"} + return Response(message, status=404) + + if url: + response, status_code = Notifications(task_name).remove_url(url) + else: + response, status_code = Notifications(task_name).remove_task() + + return Response({"response": response, "status_code": status_code}) + + class RefreshView(ApiBaseView): """resolves to /api/refresh/ GET: get refresh progress diff --git a/tubearchivist/config/management/commands/ta_startup.py b/tubearchivist/config/management/commands/ta_startup.py index 16d43d2f..80f8edf1 100644 --- a/tubearchivist/config/management/commands/ta_startup.py +++ b/tubearchivist/config/management/commands/ta_startup.py @@ -5,18 +5,24 @@ Functionality: """ import os +from random import randint from time import sleep +from django.conf import settings from django.core.management.base import BaseCommand, CommandError +from django_celery_beat.models import CrontabSchedule +from home.models import CustomPeriodicTask from home.src.es.connect import ElasticWrap from home.src.es.index_setup import ElasitIndexWrap from home.src.es.snapshot import ElasticSnapshot from home.src.ta.config import AppConfig, ReleaseVersion +from home.src.ta.config_schedule import ScheduleBuilder from home.src.ta.helper import clear_dl_cache +from home.src.ta.notify import Notifications from home.src.ta.settings import EnvironmentSettings from home.src.ta.ta_redis import RedisArchivist +from home.src.ta.task_config import TASK_CONFIG from home.src.ta.task_manager import TaskManager -from home.src.ta.users import UserConfig TOPIC = """ @@ -40,12 +46,12 @@ class Command(BaseCommand): self._clear_redis_keys() self._clear_tasks() self._clear_dl_cache() - self._mig_clear_failed_versioncheck() self._version_check() self._mig_index_setup() self._mig_snapshot_check() - self._mig_move_users_to_es() + self._mig_schedule_store() self._mig_custom_playlist() + self._create_default_schedules() def _sync_redis_state(self): """make sure redis gets new config.json values""" @@ -151,102 +157,134 @@ class Command(BaseCommand): self.stdout.write("[MIGRATION] setup snapshots") ElasticSnapshot().setup() - def _mig_clear_failed_versioncheck(self): - """hotfix for v0.4.5, clearing faulty versioncheck""" - ReleaseVersion().clear_fail() - - def _mig_move_users_to_es(self): # noqa: C901 - """migration: update from 0.4.1 to 0.4.2 move user config to ES""" - self.stdout.write("[MIGRATION] move user configuration to ES") - redis = RedisArchivist() - - # 1: Find all users in Redis - users = {i.split(":")[0] for i in redis.list_keys("[0-9]*:")} - if not users: - self.stdout.write(" no users needed migrating to ES") + def _mig_schedule_store(self): + """ + update from 0.4.7 to 0.4.8 + migrate schedule task store to CustomCronSchedule + """ + self.stdout.write("[MIGRATION] migrate schedule store") + config = AppConfig().config + current_schedules = config.get("scheduler") + if not current_schedules: + self.stdout.write( + self.style.SUCCESS(" no schedules to migrate") + ) return - # 2: Write all Redis user settings to ES - # 3: Remove user settings from Redis - try: - for user in users: - new_conf = UserConfig(user) + self._mig_update_subscribed(current_schedules) + self._mig_download_pending(current_schedules) + self._mig_check_reindex(current_schedules) + self._mig_thumbnail_check(current_schedules) + self._mig_run_backup(current_schedules) + self._mig_version_check() - stylesheet_key = f"{user}:color" - stylesheet = redis.get_message(stylesheet_key).get("status") - if stylesheet: - new_conf.set_value("stylesheet", stylesheet) - redis.del_message(stylesheet_key) + del config["scheduler"] + RedisArchivist().set_message("config", config, save=True) - sort_by_key = f"{user}:sort_by" - sort_by = redis.get_message(sort_by_key).get("status") - if sort_by: - new_conf.set_value("sort_by", sort_by) - redis.del_message(sort_by_key) + def _mig_update_subscribed(self, current_schedules): + """create update_subscribed schedule""" + task_name = "update_subscribed" + update_subscribed_schedule = current_schedules.get(task_name) + if update_subscribed_schedule: + self._create_task(task_name, update_subscribed_schedule) - page_size_key = f"{user}:page_size" - page_size = redis.get_message(page_size_key).get("status") - if page_size: - new_conf.set_value("page_size", page_size) - redis.del_message(page_size_key) + self._create_notifications(task_name, current_schedules) - sort_order_key = f"{user}:sort_order" - sort_order = redis.get_message(sort_order_key).get("status") - if sort_order: - new_conf.set_value("sort_order", sort_order) - redis.del_message(sort_order_key) + def _mig_download_pending(self, current_schedules): + """create download_pending schedule""" + task_name = "download_pending" + download_pending_schedule = current_schedules.get(task_name) + if download_pending_schedule: + self._create_task(task_name, download_pending_schedule) - grid_items_key = f"{user}:grid_items" - grid_items = redis.get_message(grid_items_key).get("status") - if grid_items: - new_conf.set_value("grid_items", grid_items) - redis.del_message(grid_items_key) + self._create_notifications(task_name, current_schedules) - hide_watch_key = f"{user}:hide_watched" - hide_watch = redis.get_message(hide_watch_key).get("status") - if hide_watch: - new_conf.set_value("hide_watched", hide_watch) - redis.del_message(hide_watch_key) + def _mig_check_reindex(self, current_schedules): + """create check_reindex schedule""" + task_name = "check_reindex" + check_reindex_schedule = current_schedules.get(task_name) + if check_reindex_schedule: + task_config = {} + days = current_schedules.get("check_reindex_days") + if days: + task_config.update({"days": days}) - ignore_only_key = f"{user}:show_ignored_only" - ignore_only = redis.get_message(ignore_only_key).get("status") - if ignore_only: - new_conf.set_value("show_ignored_only", ignore_only) - redis.del_message(ignore_only_key) - - subed_only_key = f"{user}:show_subed_only" - subed_only = redis.get_message(subed_only_key).get("status") - if subed_only: - new_conf.set_value("show_subed_only", subed_only) - redis.del_message(subed_only_key) - - for view in ["channel", "playlist", "home", "downloads"]: - view_key = f"{user}:view:{view}" - view_style = redis.get_message(view_key).get("status") - if view_style: - new_conf.set_value(f"view_style_{view}", view_style) - redis.del_message(view_key) - - self.stdout.write( - self.style.SUCCESS( - f" ✓ Settings for user '{user}' migrated to ES" - ) - ) - except Exception as err: - message = " 🗙 user migration to ES failed" - self.stdout.write(self.style.ERROR(message)) - self.stdout.write(self.style.ERROR(err)) - sleep(60) - raise CommandError(message) from err - else: - self.stdout.write( - self.style.SUCCESS( - " ✓ Settings for all users migrated to ES" - ) + self._create_task( + task_name, + check_reindex_schedule, + task_config=task_config, ) + self._create_notifications(task_name, current_schedules) + + def _mig_thumbnail_check(self, current_schedules): + """create thumbnail_check schedule""" + thumbnail_check_schedule = current_schedules.get("thumbnail_check") + if thumbnail_check_schedule: + self._create_task("thumbnail_check", thumbnail_check_schedule) + + def _mig_run_backup(self, current_schedules): + """create run_backup schedule""" + run_backup_schedule = current_schedules.get("run_backup") + if run_backup_schedule: + task_config = False + rotate = current_schedules.get("run_backup_rotate") + if rotate: + task_config = {"rotate": rotate} + + self._create_task( + "run_backup", run_backup_schedule, task_config=task_config + ) + + def _mig_version_check(self): + """create version_check schedule""" + version_check_schedule = { + "minute": randint(0, 59), + "hour": randint(0, 23), + "day_of_week": "*", + } + self._create_task("version_check", version_check_schedule) + + def _create_task(self, task_name, schedule, task_config=False): + """create task""" + description = TASK_CONFIG[task_name].get("title") + schedule, _ = CrontabSchedule.objects.get_or_create(**schedule) + schedule.timezone = settings.TIME_ZONE + schedule.save() + + task, _ = CustomPeriodicTask.objects.get_or_create( + crontab=schedule, + name=task_name, + description=description, + task=task_name, + ) + if task_config: + task.task_config = task_config + task.save() + + self.stdout.write( + self.style.SUCCESS(f" ✓ new task created: '{task}'") + ) + + def _create_notifications(self, task_name, current_schedules): + """migrate notifications of task""" + notifications = current_schedules.get(f"{task_name}_notify") + if not notifications: + return + + urls = [i.strip() for i in notifications.split()] + if not urls: + return + + self.stdout.write( + self.style.SUCCESS(f" ✓ migrate notifications: '{urls}'") + ) + handler = Notifications(task_name) + for url in urls: + handler.add_url(url) + def _mig_custom_playlist(self): - """migration for custom playlist""" + """add playlist_type for migration from v0.4.6 to v0.4.7""" self.stdout.write("[MIGRATION] custom playlist") data = { "query": { @@ -277,3 +315,54 @@ class Command(BaseCommand): self.stdout.write(response) sleep(60) raise CommandError(message) + + def _create_default_schedules(self) -> None: + """ + create default schedules for new installations + needs to be called after _mig_schedule_store + """ + self.stdout.write("[7] create initial schedules") + init_has_run = CustomPeriodicTask.objects.filter( + name="version_check" + ).exists() + + if init_has_run: + self.stdout.write( + self.style.SUCCESS( + " schedule init already done, skipping..." + ) + ) + return + + builder = ScheduleBuilder() + check_reindex = builder.get_set_task( + "check_reindex", schedule=builder.SCHEDULES["check_reindex"] + ) + check_reindex.task_config.update({"days": 90}) + check_reindex.save() + self.stdout.write( + self.style.SUCCESS( + f" ✓ created new default schedule: {check_reindex}" + ) + ) + + thumbnail_check = builder.get_set_task( + "thumbnail_check", schedule=builder.SCHEDULES["thumbnail_check"] + ) + self.stdout.write( + self.style.SUCCESS( + f" ✓ created new default schedule: {thumbnail_check}" + ) + ) + daily_random = f"{randint(0, 59)} {randint(0, 23)} *" + version_check = builder.get_set_task( + "version_check", schedule=daily_random + ) + self.stdout.write( + self.style.SUCCESS( + f" ✓ created new default schedule: {version_check}" + ) + ) + self.stdout.write( + self.style.SUCCESS(" ✓ all default schedules created") + ) diff --git a/tubearchivist/config/settings.py b/tubearchivist/config/settings.py index 3be6177b..ac41eaa0 100644 --- a/tubearchivist/config/settings.py +++ b/tubearchivist/config/settings.py @@ -38,6 +38,7 @@ ALLOWED_HOSTS, CSRF_TRUSTED_ORIGINS = ta_host_parser(environ["TA_HOST"]) # Application definition INSTALLED_APPS = [ + "django_celery_beat", "home.apps.HomeConfig", "django.contrib.admin", "django.contrib.auth", diff --git a/tubearchivist/home/__init__.py b/tubearchivist/home/__init__.py index 736385d6..2e00ac76 100644 --- a/tubearchivist/home/__init__.py +++ b/tubearchivist/home/__init__.py @@ -1,5 +1,7 @@ -""" handle celery startup """ +"""start celery app""" -from .tasks import app as celery_app +from __future__ import absolute_import, unicode_literals + +from home.celery import app as celery_app __all__ = ("celery_app",) diff --git a/tubearchivist/home/admin.py b/tubearchivist/home/admin.py index 2bcb7011..3c6e83c5 100644 --- a/tubearchivist/home/admin.py +++ b/tubearchivist/home/admin.py @@ -2,6 +2,7 @@ from django.contrib import admin from django.contrib.auth.admin import UserAdmin as BaseUserAdmin +from django_celery_beat import models as BeatModels from .models import Account @@ -34,3 +35,12 @@ class HomeAdmin(BaseUserAdmin): admin.site.register(Account, HomeAdmin) +admin.site.unregister( + [ + BeatModels.ClockedSchedule, + BeatModels.CrontabSchedule, + BeatModels.IntervalSchedule, + BeatModels.PeriodicTask, + BeatModels.SolarSchedule, + ] +) diff --git a/tubearchivist/home/celery.py b/tubearchivist/home/celery.py new file mode 100644 index 00000000..9f65795b --- /dev/null +++ b/tubearchivist/home/celery.py @@ -0,0 +1,24 @@ +"""initiate celery""" + +import os + +from celery import Celery +from home.src.ta.config import AppConfig +from home.src.ta.settings import EnvironmentSettings + +CONFIG = AppConfig().config +REDIS_HOST = EnvironmentSettings.REDIS_HOST +REDIS_PORT = EnvironmentSettings.REDIS_PORT + +os.environ.setdefault("DJANGO_SETTINGS_MODULE", "config.settings") +app = Celery( + "tasks", + broker=f"redis://{REDIS_HOST}:{REDIS_PORT}", + backend=f"redis://{REDIS_HOST}:{REDIS_PORT}", + result_extended=True, +) +app.config_from_object( + "django.conf:settings", namespace=EnvironmentSettings.REDIS_NAME_SPACE +) +app.autodiscover_tasks() +app.conf.timezone = EnvironmentSettings.TZ diff --git a/tubearchivist/home/config.json b/tubearchivist/home/config.json index 1b655f47..184721b4 100644 --- a/tubearchivist/home/config.json +++ b/tubearchivist/home/config.json @@ -26,18 +26,5 @@ }, "application": { "enable_snapshot": true - }, - "scheduler": { - "update_subscribed": false, - "update_subscribed_notify": false, - "download_pending": false, - "download_pending_notify": false, - "check_reindex": {"minute": "0", "hour": "12", "day_of_week": "*"}, - "check_reindex_notify": false, - "check_reindex_days": 90, - "thumbnail_check": {"minute": "0", "hour": "17", "day_of_week": "*"}, - "run_backup": false, - "run_backup_rotate": 5, - "version_check": "rand-d" } } diff --git a/tubearchivist/home/migrations/0002_customperiodictask.py b/tubearchivist/home/migrations/0002_customperiodictask.py new file mode 100644 index 00000000..cc584980 --- /dev/null +++ b/tubearchivist/home/migrations/0002_customperiodictask.py @@ -0,0 +1,23 @@ +# Generated by Django 4.2.7 on 2023-12-05 13:47 + +from django.db import migrations, models +import django.db.models.deletion + + +class Migration(migrations.Migration): + + dependencies = [ + ('django_celery_beat', '0018_improve_crontab_helptext'), + ('home', '0001_initial'), + ] + + operations = [ + migrations.CreateModel( + name='CustomPeriodicTask', + fields=[ + ('periodictask_ptr', models.OneToOneField(auto_created=True, on_delete=django.db.models.deletion.CASCADE, parent_link=True, primary_key=True, serialize=False, to='django_celery_beat.periodictask')), + ('task_config', models.JSONField(default=dict)), + ], + bases=('django_celery_beat.periodictask',), + ), + ] diff --git a/tubearchivist/home/models.py b/tubearchivist/home/models.py index 2dacd0c7..3f0c376b 100644 --- a/tubearchivist/home/models.py +++ b/tubearchivist/home/models.py @@ -6,6 +6,7 @@ from django.contrib.auth.models import ( PermissionsMixin, ) from django.db import models +from django_celery_beat.models import PeriodicTask class AccountManager(BaseUserManager): @@ -52,3 +53,9 @@ class Account(AbstractBaseUser, PermissionsMixin): USERNAME_FIELD = "name" REQUIRED_FIELDS = ["password"] + + +class CustomPeriodicTask(PeriodicTask): + """add custom metadata to to task""" + + task_config = models.JSONField(default=dict) diff --git a/tubearchivist/home/src/download/queue.py b/tubearchivist/home/src/download/queue.py index 3c733ca6..661f39b9 100644 --- a/tubearchivist/home/src/download/queue.py +++ b/tubearchivist/home/src/download/queue.py @@ -7,10 +7,7 @@ Functionality: import json from datetime import datetime -from home.src.download.subscriptions import ( - ChannelSubscription, - PlaylistSubscription, -) +from home.src.download.subscriptions import ChannelSubscription from home.src.download.thumbnails import ThumbManager from home.src.download.yt_dlp_base import YtWrap from home.src.es.connect import ElasticWrap, IndexPaginate @@ -196,7 +193,6 @@ class PendingList(PendingIndex): self._parse_channel(entry["url"], vid_type) elif entry["type"] == "playlist": self._parse_playlist(entry["url"]) - PlaylistSubscription().process_url_str([entry], subscribed=False) else: raise ValueError(f"invalid url_type: {entry}") @@ -227,15 +223,18 @@ class PendingList(PendingIndex): def _parse_playlist(self, url): """add all videos of playlist to list""" playlist = YoutubePlaylist(url) - playlist.build_json() - if not playlist.json_data: + is_active = playlist.update_playlist() + if not is_active: message = f"{playlist.youtube_id}: failed to extract metadata" print(message) raise ValueError(message) - video_results = playlist.json_data.get("playlist_entries") - youtube_ids = [i["youtube_id"] for i in video_results] - for video_id in youtube_ids: + entries = playlist.json_data["playlist_entries"] + to_add = [i["youtube_id"] for i in entries if not i["downloaded"]] + if not to_add: + return + + for video_id in to_add: # match vid_type later self._add_video(video_id, VideoTypeEnum.UNKNOWN) @@ -245,6 +244,7 @@ class PendingList(PendingIndex): bulk_list = [] total = len(self.missing_videos) + videos_added = [] for idx, (youtube_id, vid_type) in enumerate(self.missing_videos): if self.task and self.task.is_stopped(): break @@ -268,6 +268,7 @@ class PendingList(PendingIndex): url = video_details["vid_thumb_url"] ThumbManager(youtube_id).download_video_thumb(url) + videos_added.append(youtube_id) if len(bulk_list) >= 20: self._ingest_bulk(bulk_list) @@ -275,6 +276,8 @@ class PendingList(PendingIndex): self._ingest_bulk(bulk_list) + return videos_added + def _ingest_bulk(self, bulk_list): """add items to queue in bulk""" if not bulk_list: diff --git a/tubearchivist/home/src/download/subscriptions.py b/tubearchivist/home/src/download/subscriptions.py index c6d85a9c..a276dcee 100644 --- a/tubearchivist/home/src/download/subscriptions.py +++ b/tubearchivist/home/src/download/subscriptions.py @@ -4,14 +4,15 @@ Functionality: - handle playlist subscriptions """ -from home.src.download import queue # partial import from home.src.download.thumbnails import ThumbManager from home.src.download.yt_dlp_base import YtWrap from home.src.es.connect import IndexPaginate from home.src.index.channel import YoutubeChannel from home.src.index.playlist import YoutubePlaylist +from home.src.index.video import YoutubeVideo from home.src.index.video_constants import VideoTypeEnum from home.src.ta.config import AppConfig +from home.src.ta.helper import is_missing from home.src.ta.urlparser import Parser @@ -105,10 +106,6 @@ class ChannelSubscription: if not all_channels: return False - pending = queue.PendingList() - pending.get_download() - pending.get_indexed() - missing_videos = [] total = len(all_channels) @@ -118,22 +115,22 @@ class ChannelSubscription: last_videos = self.get_last_youtube_videos(channel_id) if last_videos: + ids_to_add = is_missing([i[0] for i in last_videos]) for video_id, _, vid_type in last_videos: - if video_id not in pending.to_skip: + if video_id in ids_to_add: missing_videos.append((video_id, vid_type)) if not self.task: continue - if self.task: - if self.task.is_stopped(): - self.task.send_progress(["Received Stop signal."]) - break + if self.task.is_stopped(): + self.task.send_progress(["Received Stop signal."]) + break - self.task.send_progress( - message_lines=[f"Scanning Channel {idx + 1}/{total}"], - progress=(idx + 1) / total, - ) + self.task.send_progress( + message_lines=[f"Scanning Channel {idx + 1}/{total}"], + progress=(idx + 1) / total, + ) return missing_videos @@ -174,10 +171,6 @@ class PlaylistSubscription: def process_url_str(self, new_playlists, subscribed=True): """process playlist subscribe form url_str""" - data = {"query": {"match_all": {}}, "_source": ["youtube_id"]} - all_indexed = IndexPaginate("ta_video", data).get_results() - all_youtube_ids = [i["youtube_id"] for i in all_indexed] - for idx, playlist in enumerate(new_playlists): playlist_id = playlist["url"] if not playlist["type"] == "playlist": @@ -185,7 +178,6 @@ class PlaylistSubscription: continue playlist_h = YoutubePlaylist(playlist_id) - playlist_h.all_youtube_ids = all_youtube_ids playlist_h.build_json() if not playlist_h.json_data: message = f"{playlist_h.youtube_id}: failed to extract data" @@ -223,27 +215,15 @@ class PlaylistSubscription: playlist.json_data["playlist_subscribed"] = subscribe_status playlist.upload_to_es() - @staticmethod - def get_to_ignore(): - """get all youtube_ids already downloaded or ignored""" - pending = queue.PendingList() - pending.get_download() - pending.get_indexed() - - return pending.to_skip - def find_missing(self): """find videos in subscribed playlists not downloaded yet""" all_playlists = [i["playlist_id"] for i in self.get_playlists()] if not all_playlists: return False - to_ignore = self.get_to_ignore() - missing_videos = [] total = len(all_playlists) for idx, playlist_id in enumerate(all_playlists): - size_limit = self.config["subscriptions"]["channel_size"] playlist = YoutubePlaylist(playlist_id) is_active = playlist.update_playlist() if not is_active: @@ -251,27 +231,29 @@ class PlaylistSubscription: continue playlist_entries = playlist.json_data["playlist_entries"] + size_limit = self.config["subscriptions"]["channel_size"] if size_limit: del playlist_entries[size_limit:] - all_missing = [i for i in playlist_entries if not i["downloaded"]] - - for video in all_missing: - youtube_id = video["youtube_id"] - if youtube_id not in to_ignore: - missing_videos.append(youtube_id) + to_check = [ + i["youtube_id"] + for i in playlist_entries + if i["downloaded"] is False + ] + needs_downloading = is_missing(to_check) + missing_videos.extend(needs_downloading) if not self.task: continue - if self.task: - self.task.send_progress( - message_lines=[f"Scanning Playlists {idx + 1}/{total}"], - progress=(idx + 1) / total, - ) - if self.task.is_stopped(): - self.task.send_progress(["Received Stop signal."]) - break + if self.task.is_stopped(): + self.task.send_progress(["Received Stop signal."]) + break + + self.task.send_progress( + message_lines=[f"Scanning Playlists {idx + 1}/{total}"], + progress=(idx + 1) / total, + ) return missing_videos @@ -359,8 +341,10 @@ class SubscriptionHandler: if item["type"] == "video": # extract channel id from video - vid = queue.PendingList().get_youtube_details(item["url"]) - channel_id = vid["channel_id"] + video = YoutubeVideo(item["url"]) + video.get_from_youtube() + video.process_youtube_meta() + channel_id = video.channel_id elif item["type"] == "channel": channel_id = item["url"] else: diff --git a/tubearchivist/home/src/download/thumbnails.py b/tubearchivist/home/src/download/thumbnails.py index a2698fd7..cf4c485f 100644 --- a/tubearchivist/home/src/download/thumbnails.py +++ b/tubearchivist/home/src/download/thumbnails.py @@ -11,6 +11,7 @@ from time import sleep import requests from home.src.es.connect import ElasticWrap, IndexPaginate +from home.src.ta.helper import is_missing from home.src.ta.settings import EnvironmentSettings from mutagen.mp4 import MP4, MP4Cover from PIL import Image, ImageFile, ImageFilter, UnidentifiedImageError @@ -326,7 +327,7 @@ class ThumbValidator: }, ] - def __init__(self, task): + def __init__(self, task=False): self.task = task def validate(self): @@ -346,6 +347,89 @@ class ThumbValidator: ) _ = paginate.get_results() + def clean_up(self): + """clean up all thumbs""" + self._clean_up_vids() + self._clean_up_channels() + self._clean_up_playlists() + + def _clean_up_vids(self): + """clean unneeded vid thumbs""" + video_dir = os.path.join(EnvironmentSettings.CACHE_DIR, "videos") + video_folders = os.listdir(video_dir) + for video_folder in video_folders: + folder_path = os.path.join(video_dir, video_folder) + thumbs_is = {i.split(".")[0] for i in os.listdir(folder_path)} + thumbs_should = self._get_vid_thumbs_should(video_folder) + to_delete = thumbs_is - thumbs_should + for thumb in to_delete: + delete_path = os.path.join(folder_path, f"{thumb}.jpg") + os.remove(delete_path) + + if to_delete: + message = ( + f"[thumbs][video][{video_folder}] " + + f"delete {len(to_delete)} unused thumbnails" + ) + print(message) + if self.task: + self.task.send_progress([message]) + + @staticmethod + def _get_vid_thumbs_should(video_folder: str) -> set[str]: + """get indexed""" + should_list = [ + {"prefix": {"youtube_id": {"value": video_folder.lower()}}}, + {"prefix": {"youtube_id": {"value": video_folder.upper()}}}, + ] + data = { + "query": {"bool": {"should": should_list}}, + "_source": ["youtube_id"], + } + result = IndexPaginate("ta_video,ta_download", data).get_results() + thumbs_should = {i["youtube_id"] for i in result} + + return thumbs_should + + def _clean_up_channels(self): + """clean unneeded channel thumbs""" + channel_dir = os.path.join(EnvironmentSettings.CACHE_DIR, "channels") + channel_art = os.listdir(channel_dir) + thumbs_is = {"_".join(i.split("_")[:-1]) for i in channel_art} + to_delete = is_missing(list(thumbs_is), "ta_channel", "channel_id") + for channel_thumb in channel_art: + if channel_thumb[:24] in to_delete: + delete_path = os.path.join(channel_dir, channel_thumb) + os.remove(delete_path) + + if to_delete: + message = ( + "[thumbs][channel] " + + f"delete {len(to_delete)} unused channel art" + ) + print(message) + if self.task: + self.task.send_progress([message]) + + def _clean_up_playlists(self): + """clean up unneeded playlist thumbs""" + playlist_dir = os.path.join(EnvironmentSettings.CACHE_DIR, "playlists") + playlist_art = os.listdir(playlist_dir) + thumbs_is = {i.split(".")[0] for i in playlist_art} + to_delete = is_missing(list(thumbs_is), "ta_playlist", "playlist_id") + for playlist_id in to_delete: + delete_path = os.path.join(playlist_dir, f"{playlist_id}.jpg") + os.remove(delete_path) + + if to_delete: + message = ( + "[thumbs][playlist] " + + f"delete {len(to_delete)} unused playlist art" + ) + print(message) + if self.task: + self.task.send_progress([message]) + @staticmethod def _get_total(index_name): """get total documents in index""" diff --git a/tubearchivist/home/src/download/yt_dlp_handler.py b/tubearchivist/home/src/download/yt_dlp_handler.py index cb76ac34..c1a150c2 100644 --- a/tubearchivist/home/src/download/yt_dlp_handler.py +++ b/tubearchivist/home/src/download/yt_dlp_handler.py @@ -20,182 +20,69 @@ from home.src.index.playlist import YoutubePlaylist from home.src.index.video import YoutubeVideo, index_new_video from home.src.index.video_constants import VideoTypeEnum from home.src.ta.config import AppConfig -from home.src.ta.helper import ignore_filelist +from home.src.ta.helper import get_channel_overwrites, ignore_filelist from home.src.ta.settings import EnvironmentSettings +from home.src.ta.ta_redis import RedisQueue -class DownloadPostProcess: - """handle task to run after download queue finishes""" +class DownloaderBase: + """base class for shared config""" - def __init__(self, download): - self.download = download - self.now = int(datetime.now().timestamp()) - self.pending = False + CACHE_DIR = EnvironmentSettings.CACHE_DIR + MEDIA_DIR = EnvironmentSettings.MEDIA_DIR + CHANNEL_QUEUE = "download:channel" + PLAYLIST_QUEUE = "download:playlist:full" + PLAYLIST_QUICK = "download:playlist:quick" + VIDEO_QUEUE = "download:video" - def run(self): - """run all functions""" - self.pending = PendingList() - self.pending.get_download() - self.pending.get_channels() - self.pending.get_indexed() - self.auto_delete_all() - self.auto_delete_overwrites() - self.validate_playlists() - self.get_comments() - - def auto_delete_all(self): - """handle auto delete""" - autodelete_days = self.download.config["downloads"]["autodelete_days"] - if not autodelete_days: - return - - print(f"auto delete older than {autodelete_days} days") - now_lte = str(self.now - autodelete_days * 24 * 60 * 60) - data = { - "query": {"range": {"player.watched_date": {"lte": now_lte}}}, - "sort": [{"player.watched_date": {"order": "asc"}}], - } - self._auto_delete_watched(data) - - def auto_delete_overwrites(self): - """handle per channel auto delete from overwrites""" - for channel_id, value in self.pending.channel_overwrites.items(): - if "autodelete_days" in value: - autodelete_days = value.get("autodelete_days") - print(f"{channel_id}: delete older than {autodelete_days}d") - now_lte = str(self.now - autodelete_days * 24 * 60 * 60) - must_list = [ - {"range": {"player.watched_date": {"lte": now_lte}}}, - {"term": {"channel.channel_id": {"value": channel_id}}}, - ] - data = { - "query": {"bool": {"must": must_list}}, - "sort": [{"player.watched_date": {"order": "desc"}}], - } - self._auto_delete_watched(data) - - @staticmethod - def _auto_delete_watched(data): - """delete watched videos after x days""" - to_delete = IndexPaginate("ta_video", data).get_results() - if not to_delete: - return - - for video in to_delete: - youtube_id = video["youtube_id"] - print(f"{youtube_id}: auto delete video") - YoutubeVideo(youtube_id).delete_media_file() - - print("add deleted to ignore list") - vids = [{"type": "video", "url": i["youtube_id"]} for i in to_delete] - pending = PendingList(youtube_ids=vids) - pending.parse_url_list() - pending.add_to_pending(status="ignore") - - def validate_playlists(self): - """look for playlist needing to update""" - for id_c, channel_id in enumerate(self.download.channels): - channel = YoutubeChannel(channel_id, task=self.download.task) - overwrites = self.pending.channel_overwrites.get(channel_id, False) - if overwrites and overwrites.get("index_playlists"): - # validate from remote - channel.index_channel_playlists() - continue - - # validate from local - playlists = channel.get_indexed_playlists(active_only=True) - all_channel_playlist = [i["playlist_id"] for i in playlists] - self._validate_channel_playlist(all_channel_playlist, id_c) - - def _validate_channel_playlist(self, all_channel_playlist, id_c): - """scan channel for playlist needing update""" - all_youtube_ids = [i["youtube_id"] for i in self.pending.all_videos] - for id_p, playlist_id in enumerate(all_channel_playlist): - playlist = YoutubePlaylist(playlist_id) - playlist.all_youtube_ids = all_youtube_ids - playlist.build_json(scrape=True) - if not playlist.json_data: - playlist.deactivate() - continue - - playlist.add_vids_to_playlist() - playlist.upload_to_es() - self._notify_playlist_progress(all_channel_playlist, id_c, id_p) - - def _notify_playlist_progress(self, all_channel_playlist, id_c, id_p): - """notify to UI""" - if not self.download.task: - return - - total_channel = len(self.download.channels) - total_playlist = len(all_channel_playlist) - - message = [ - f"Post Processing Channels: {id_c}/{total_channel}", - f"Validate Playlists {id_p + 1}/{total_playlist}", - ] - progress = (id_c + 1) / total_channel - self.download.task.send_progress(message, progress=progress) - - def get_comments(self): - """get comments from youtube""" - CommentList(self.download.videos, task=self.download.task).index() - - -class VideoDownloader: - """ - handle the video download functionality - if not initiated with list, take from queue - """ - - def __init__(self, youtube_id_list=False, task=False): - self.obs = False - self.video_overwrites = False - self.youtube_id_list = youtube_id_list + def __init__(self, task): self.task = task self.config = AppConfig().config - self.cache_dir = EnvironmentSettings.CACHE_DIR - self.media_dir = EnvironmentSettings.MEDIA_DIR - self._build_obs() - self.channels = set() - self.videos = set() + self.channel_overwrites = get_channel_overwrites() + self.now = int(datetime.now().timestamp()) - def run_queue(self, auto_only=False): + +class VideoDownloader(DownloaderBase): + """handle the video download functionality""" + + def __init__(self, task=False): + super().__init__(task) + self.obs = False + self._build_obs() + + def run_queue(self, auto_only=False) -> int: """setup download queue in redis loop until no more items""" - self._get_overwrites() + downloaded = 0 while True: video_data = self._get_next(auto_only) if self.task.is_stopped() or not video_data: self._reset_auto() break - youtube_id = video_data.get("youtube_id") + youtube_id = video_data["youtube_id"] + channel_id = video_data["channel_id"] print(f"{youtube_id}: Downloading video") self._notify(video_data, "Validate download format") - success = self._dl_single_vid(youtube_id) + success = self._dl_single_vid(youtube_id, channel_id) if not success: continue self._notify(video_data, "Add video metadata to index", progress=1) - - vid_dict = index_new_video( - youtube_id, - video_overwrites=self.video_overwrites, - video_type=VideoTypeEnum(video_data["vid_type"]), - ) - self.channels.add(vid_dict["channel"]["channel_id"]) - self.videos.add(vid_dict["youtube_id"]) + video_type = VideoTypeEnum(video_data["vid_type"]) + vid_dict = index_new_video(youtube_id, video_type=video_type) + RedisQueue(self.CHANNEL_QUEUE).add(channel_id) + RedisQueue(self.VIDEO_QUEUE).add(youtube_id) self._notify(video_data, "Move downloaded file to archive") self.move_to_archive(vid_dict) self._delete_from_pending(youtube_id) + downloaded += 1 # post processing - self._add_subscribed_channels() - DownloadPostProcess(self).run() + DownloadPostProcess(self.task).run() - return self.videos + return downloaded def _notify(self, video_data, message, progress=False): """send progress notification to task""" @@ -230,13 +117,6 @@ class VideoDownloader: return response["hits"]["hits"][0]["_source"] - def _get_overwrites(self): - """get channel overwrites""" - pending = PendingList() - pending.get_download() - pending.get_channels() - self.video_overwrites = pending.video_overwrites - def _progress_hook(self, response): """process the progress_hooks from yt_dlp""" progress = False @@ -267,7 +147,7 @@ class VideoDownloader: """initial obs""" self.obs = { "merge_output_format": "mp4", - "outtmpl": (self.cache_dir + "/download/%(id)s.mp4"), + "outtmpl": (self.CACHE_DIR + "/download/%(id)s.mp4"), "progress_hooks": [self._progress_hook], "noprogress": True, "continuedl": True, @@ -327,22 +207,17 @@ class VideoDownloader: self.obs["postprocessors"] = postprocessors - def get_format_overwrites(self, youtube_id): - """get overwrites from single video""" - overwrites = self.video_overwrites.get(youtube_id, False) - if overwrites: - return overwrites.get("download_format", False) + def _set_overwrites(self, obs: dict, channel_id: str) -> None: + """add overwrites to obs""" + overwrites = self.channel_overwrites.get(channel_id) + if overwrites and overwrites.get("download_format"): + obs["format"] = overwrites.get("download_format") - return False - - def _dl_single_vid(self, youtube_id): + def _dl_single_vid(self, youtube_id: str, channel_id: str) -> bool: """download single video""" obs = self.obs.copy() - format_overwrite = self.get_format_overwrites(youtube_id) - if format_overwrite: - obs["format"] = format_overwrite - - dl_cache = self.cache_dir + "/download/" + self._set_overwrites(obs, channel_id) + dl_cache = os.path.join(self.CACHE_DIR, "download") # check if already in cache to continue from there all_cached = ignore_filelist(os.listdir(dl_cache)) @@ -376,7 +251,7 @@ class VideoDownloader: host_gid = EnvironmentSettings.HOST_GID # make folder folder = os.path.join( - self.media_dir, vid_dict["channel"]["channel_id"] + self.MEDIA_DIR, vid_dict["channel"]["channel_id"] ) if not os.path.exists(folder): os.makedirs(folder) @@ -384,8 +259,8 @@ class VideoDownloader: os.chown(folder, host_uid, host_gid) # move media file media_file = vid_dict["youtube_id"] + ".mp4" - old_path = os.path.join(self.cache_dir, "download", media_file) - new_path = os.path.join(self.media_dir, vid_dict["media_url"]) + old_path = os.path.join(self.CACHE_DIR, "download", media_file) + new_path = os.path.join(self.MEDIA_DIR, vid_dict["media_url"]) # move media file and fix permission shutil.move(old_path, new_path, copy_function=shutil.copyfile) if host_uid and host_gid: @@ -397,18 +272,6 @@ class VideoDownloader: path = f"ta_download/_doc/{youtube_id}?refresh=true" _, _ = ElasticWrap(path).delete() - def _add_subscribed_channels(self): - """add all channels subscribed to refresh""" - all_subscribed = PlaylistSubscription().get_playlists() - if not all_subscribed: - return - - channel_ids = [i["playlist_channel_id"] for i in all_subscribed] - for channel_id in channel_ids: - self.channels.add(channel_id) - - return - def _reset_auto(self): """reset autostart to defaults after queue stop""" path = "ta_download/_update_by_query" @@ -423,3 +286,169 @@ class VideoDownloader: updated = response.get("updated") if updated: print(f"[download] reset auto start on {updated} videos.") + + +class DownloadPostProcess(DownloaderBase): + """handle task to run after download queue finishes""" + + def run(self): + """run all functions""" + self.auto_delete_all() + self.auto_delete_overwrites() + self.refresh_playlist() + self.match_videos() + self.get_comments() + + def auto_delete_all(self): + """handle auto delete""" + autodelete_days = self.config["downloads"]["autodelete_days"] + if not autodelete_days: + return + + print(f"auto delete older than {autodelete_days} days") + now_lte = str(self.now - autodelete_days * 24 * 60 * 60) + data = { + "query": {"range": {"player.watched_date": {"lte": now_lte}}}, + "sort": [{"player.watched_date": {"order": "asc"}}], + } + self._auto_delete_watched(data) + + def auto_delete_overwrites(self): + """handle per channel auto delete from overwrites""" + for channel_id, value in self.channel_overwrites.items(): + if "autodelete_days" in value: + autodelete_days = value.get("autodelete_days") + print(f"{channel_id}: delete older than {autodelete_days}d") + now_lte = str(self.now - autodelete_days * 24 * 60 * 60) + must_list = [ + {"range": {"player.watched_date": {"lte": now_lte}}}, + {"term": {"channel.channel_id": {"value": channel_id}}}, + ] + data = { + "query": {"bool": {"must": must_list}}, + "sort": [{"player.watched_date": {"order": "desc"}}], + } + self._auto_delete_watched(data) + + @staticmethod + def _auto_delete_watched(data): + """delete watched videos after x days""" + to_delete = IndexPaginate("ta_video", data).get_results() + if not to_delete: + return + + for video in to_delete: + youtube_id = video["youtube_id"] + print(f"{youtube_id}: auto delete video") + YoutubeVideo(youtube_id).delete_media_file() + + print("add deleted to ignore list") + vids = [{"type": "video", "url": i["youtube_id"]} for i in to_delete] + pending = PendingList(youtube_ids=vids) + pending.parse_url_list() + _ = pending.add_to_pending(status="ignore") + + def refresh_playlist(self) -> None: + """match videos with playlists""" + self.add_playlists_to_refresh() + + queue = RedisQueue(self.PLAYLIST_QUEUE) + while True: + total = queue.max_score() + playlist_id, idx = queue.get_next() + if not playlist_id or not idx or not total: + break + + playlist = YoutubePlaylist(playlist_id) + playlist.update_playlist(skip_on_empty=True) + + if not self.task: + continue + + channel_name = playlist.json_data["playlist_channel"] + playlist_title = playlist.json_data["playlist_name"] + message = [ + f"Post Processing Playlists for: {channel_name}", + f"{playlist_title} [{idx}/{total}]", + ] + progress = idx / total + self.task.send_progress(message, progress=progress) + + def add_playlists_to_refresh(self) -> None: + """add playlists to refresh""" + if self.task: + message = ["Post Processing Playlists", "Scanning for Playlists"] + self.task.send_progress(message) + + self._add_playlist_sub() + self._add_channel_playlists() + self._add_video_playlists() + + def _add_playlist_sub(self): + """add subscribed playlists to refresh""" + subs = PlaylistSubscription().get_playlists() + to_add = [i["playlist_id"] for i in subs] + RedisQueue(self.PLAYLIST_QUEUE).add_list(to_add) + + def _add_channel_playlists(self): + """add playlists from channels to refresh""" + queue = RedisQueue(self.CHANNEL_QUEUE) + while True: + channel_id, _ = queue.get_next() + if not channel_id: + break + + channel = YoutubeChannel(channel_id) + channel.get_from_es() + overwrites = channel.get_overwrites() + if "index_playlists" in overwrites: + channel.get_all_playlists() + to_add = [i[0] for i in channel.all_playlists] + RedisQueue(self.PLAYLIST_QUEUE).add_list(to_add) + + def _add_video_playlists(self): + """add other playlists for quick sync""" + all_playlists = RedisQueue(self.PLAYLIST_QUEUE).get_all() + must_not = [{"terms": {"playlist_id": all_playlists}}] + video_ids = RedisQueue(self.VIDEO_QUEUE).get_all() + must = [{"terms": {"playlist_entries.youtube_id": video_ids}}] + data = { + "query": {"bool": {"must_not": must_not, "must": must}}, + "_source": ["playlist_id"], + } + playlists = IndexPaginate("ta_playlist", data).get_results() + to_add = [i["playlist_id"] for i in playlists] + RedisQueue(self.PLAYLIST_QUICK).add_list(to_add) + + def match_videos(self) -> None: + """scan rest of indexed playlists to match videos""" + queue = RedisQueue(self.PLAYLIST_QUICK) + while True: + total = queue.max_score() + playlist_id, idx = queue.get_next() + if not playlist_id or not idx or not total: + break + + playlist = YoutubePlaylist(playlist_id) + playlist.get_from_es() + playlist.add_vids_to_playlist() + playlist.remove_vids_from_playlist() + + if not self.task: + continue + + message = [ + "Post Processing Playlists.", + f"Validate Playlists: - {idx}/{total}", + ] + progress = idx / total + self.task.send_progress(message, progress=progress) + + def get_comments(self): + """get comments from youtube""" + video_queue = RedisQueue(self.VIDEO_QUEUE) + comment_list = CommentList(task=self.task) + comment_list.add(video_ids=video_queue.get_all()) + + video_queue.clear() + comment_list.index() diff --git a/tubearchivist/home/src/es/backup.py b/tubearchivist/home/src/es/backup.py index 1c3778bb..3f134472 100644 --- a/tubearchivist/home/src/es/backup.py +++ b/tubearchivist/home/src/es/backup.py @@ -10,6 +10,7 @@ import os import zipfile from datetime import datetime +from home.models import CustomPeriodicTask from home.src.es.connect import ElasticWrap, IndexPaginate from home.src.ta.config import AppConfig from home.src.ta.helper import get_mapping, ignore_filelist @@ -197,7 +198,12 @@ class ElasticBackup: def rotate_backup(self): """delete old backups if needed""" - rotate = self.config["scheduler"]["run_backup_rotate"] + try: + task = CustomPeriodicTask.objects.get(name="run_backup") + except CustomPeriodicTask.DoesNotExist: + return + + rotate = task.task_config.get("rotate") if not rotate: return diff --git a/tubearchivist/home/src/frontend/forms.py b/tubearchivist/home/src/frontend/forms.py index f8fcb798..7f42797c 100644 --- a/tubearchivist/home/src/frontend/forms.py +++ b/tubearchivist/home/src/frontend/forms.py @@ -159,50 +159,6 @@ class ApplicationSettingsForm(forms.Form): ) -class SchedulerSettingsForm(forms.Form): - """handle scheduler settings""" - - HELP_TEXT = "Add Apprise notification URLs, one per line" - - update_subscribed = forms.CharField(required=False) - update_subscribed_notify = forms.CharField( - label=False, - widget=forms.Textarea( - attrs={ - "rows": 2, - "placeholder": HELP_TEXT, - } - ), - required=False, - ) - download_pending = forms.CharField(required=False) - download_pending_notify = forms.CharField( - label=False, - widget=forms.Textarea( - attrs={ - "rows": 2, - "placeholder": HELP_TEXT, - } - ), - required=False, - ) - check_reindex = forms.CharField(required=False) - check_reindex_notify = forms.CharField( - label=False, - widget=forms.Textarea( - attrs={ - "rows": 2, - "placeholder": HELP_TEXT, - } - ), - required=False, - ) - check_reindex_days = forms.IntegerField(required=False) - thumbnail_check = forms.CharField(required=False) - run_backup = forms.CharField(required=False) - run_backup_rotate = forms.IntegerField(required=False) - - class MultiSearchForm(forms.Form): """multi search form for /search/""" diff --git a/tubearchivist/home/src/frontend/forms_schedule.py b/tubearchivist/home/src/frontend/forms_schedule.py new file mode 100644 index 00000000..f4424fff --- /dev/null +++ b/tubearchivist/home/src/frontend/forms_schedule.py @@ -0,0 +1,101 @@ +""" +Functionality: +- handle schedule forms +- implement form validation +""" + +from celery.schedules import crontab +from django import forms +from home.src.ta.task_config import TASK_CONFIG + + +class CrontabValidator: + """validate crontab""" + + @staticmethod + def validate_fields(cron_fields): + """expect 3 cron fields""" + if not len(cron_fields) == 3: + raise forms.ValidationError("expected three cron schedule fields") + + @staticmethod + def validate_minute(minute_field): + """expect minute int""" + try: + minute_value = int(minute_field) + if not 0 <= minute_value <= 59: + raise forms.ValidationError( + "Invalid value for minutes. Must be between 0 and 59." + ) + except ValueError as err: + raise forms.ValidationError( + "Invalid value for minutes. Must be an integer." + ) from err + + @staticmethod + def validate_cron_tab(minute, hour, day_of_week): + """check if crontab can be created""" + try: + crontab(minute=minute, hour=hour, day_of_week=day_of_week) + except ValueError as err: + raise forms.ValidationError(f"invalid crontab: {err}") from err + + def validate(self, cron_expression): + """create crontab schedule""" + if cron_expression == "auto": + return + + cron_fields = cron_expression.split() + self.validate_fields(cron_fields) + + minute, hour, day_of_week = cron_fields + self.validate_minute(minute) + self.validate_cron_tab(minute, hour, day_of_week) + + +def validate_cron(cron_expression): + """callable for field""" + CrontabValidator().validate(cron_expression) + + +class SchedulerSettingsForm(forms.Form): + """handle scheduler settings""" + + update_subscribed = forms.CharField( + required=False, validators=[validate_cron] + ) + download_pending = forms.CharField( + required=False, validators=[validate_cron] + ) + check_reindex = forms.CharField(required=False, validators=[validate_cron]) + check_reindex_days = forms.IntegerField(required=False) + thumbnail_check = forms.CharField( + required=False, validators=[validate_cron] + ) + run_backup = forms.CharField(required=False, validators=[validate_cron]) + run_backup_rotate = forms.IntegerField(required=False) + + +class NotificationSettingsForm(forms.Form): + """add notification URL""" + + SUPPORTED_TASKS = [ + "update_subscribed", + "extract_download", + "download_pending", + "check_reindex", + ] + TASK_LIST = [(i, TASK_CONFIG[i]["title"]) for i in SUPPORTED_TASKS] + + TASK_CHOICES = [("", "-- select task --")] + TASK_CHOICES.extend(TASK_LIST) + + PLACEHOLDER = "Apprise notification URL" + + task = forms.ChoiceField( + widget=forms.Select, choices=TASK_CHOICES, required=False + ) + notification_url = forms.CharField( + required=False, + widget=forms.TextInput(attrs={"placeholder": PLACEHOLDER}), + ) diff --git a/tubearchivist/home/src/index/channel.py b/tubearchivist/home/src/index/channel.py index 08166f70..dd8f064e 100644 --- a/tubearchivist/home/src/index/channel.py +++ b/tubearchivist/home/src/index/channel.py @@ -8,7 +8,6 @@ import json import os from datetime import datetime -from home.src.download import queue # partial import from home.src.download.thumbnails import ThumbManager from home.src.download.yt_dlp_base import YtWrap from home.src.es.connect import ElasticWrap, IndexPaginate @@ -23,17 +22,16 @@ class YoutubeChannel(YouTubeItem): es_path = False index_name = "ta_channel" yt_base = "https://www.youtube.com/channel/" - yt_obs = {"playlist_items": "0,0"} + yt_obs = { + "playlist_items": "1,0", + "skip_download": True, + } def __init__(self, youtube_id, task=False): super().__init__(youtube_id) self.all_playlists = False self.task = task - def build_yt_url(self): - """overwrite base to use channel about page""" - return f"{self.yt_base}{self.youtube_id}/about" - def build_json(self, upload=False, fallback=False): """get from es or from youtube""" self.get_from_es() @@ -53,14 +51,13 @@ class YoutubeChannel(YouTubeItem): def process_youtube_meta(self): """extract relevant fields""" self.youtube_meta["thumbnails"].reverse() - channel_subs = self.youtube_meta.get("channel_follower_count") or 0 self.json_data = { "channel_active": True, "channel_description": self.youtube_meta.get("description", False), "channel_id": self.youtube_id, "channel_last_refresh": int(datetime.now().timestamp()), "channel_name": self.youtube_meta["uploader"], - "channel_subs": channel_subs, + "channel_subs": self._extract_follower_count(), "channel_subscribed": False, "channel_tags": self._parse_tags(self.youtube_meta.get("tags")), "channel_banner_url": self._get_banner_art(), @@ -69,6 +66,22 @@ class YoutubeChannel(YouTubeItem): "channel_views": self.youtube_meta.get("view_count") or 0, } + def _extract_follower_count(self) -> int: + """workaround for upstream #9893, extract subs from first video""" + subs = self.youtube_meta.get("channel_follower_count") + if subs is not None: + return subs + + entries = self.youtube_meta.get("entries", []) + if entries: + first_entry = entries[0] + if isinstance(first_entry, dict): + subs_entry = first_entry.get("channel_follower_count") + if subs_entry is not None: + return subs_entry + + return 0 + def _parse_tags(self, tags): """parse channel tags""" if not tags: @@ -212,9 +225,7 @@ class YoutubeChannel(YouTubeItem): """delete all indexed playlist from es""" all_playlists = self.get_indexed_playlists() for playlist in all_playlists: - playlist_id = playlist["playlist_id"] - playlist = YoutubePlaylist(playlist_id) - YoutubePlaylist(playlist_id).delete_metadata() + YoutubePlaylist(playlist["playlist_id"]).delete_metadata() def delete_channel(self): """delete channel and all videos""" @@ -253,13 +264,12 @@ class YoutubeChannel(YouTubeItem): print(f"{self.youtube_id}: no playlists found.") return - all_youtube_ids = self.get_all_video_ids() total = len(self.all_playlists) for idx, playlist in enumerate(self.all_playlists): if self.task: self._notify_single_playlist(idx, total) - self._index_single_playlist(playlist, all_youtube_ids) + self._index_single_playlist(playlist) print("add playlist: " + playlist[1]) def _notify_single_playlist(self, idx, total): @@ -272,32 +282,10 @@ class YoutubeChannel(YouTubeItem): self.task.send_progress(message, progress=(idx + 1) / total) @staticmethod - def _index_single_playlist(playlist, all_youtube_ids): + def _index_single_playlist(playlist): """add single playlist if needed""" playlist = YoutubePlaylist(playlist[0]) - playlist.all_youtube_ids = all_youtube_ids - playlist.build_json() - if not playlist.json_data: - return - - entries = playlist.json_data["playlist_entries"] - downloaded = [i for i in entries if i["downloaded"]] - if not downloaded: - return - - playlist.upload_to_es() - playlist.add_vids_to_playlist() - playlist.get_playlist_art() - - @staticmethod - def get_all_video_ids(): - """match all playlists with videos""" - handler = queue.PendingList() - handler.get_download() - handler.get_indexed() - all_youtube_ids = [i["youtube_id"] for i in handler.all_videos] - - return all_youtube_ids + playlist.update_playlist(skip_on_empty=True) def get_channel_videos(self): """get all videos from channel""" @@ -334,9 +322,9 @@ class YoutubeChannel(YouTubeItem): all_playlists = IndexPaginate("ta_playlist", data).get_results() return all_playlists - def get_overwrites(self): + def get_overwrites(self) -> dict: """get all per channel overwrites""" - return self.json_data.get("channel_overwrites", False) + return self.json_data.get("channel_overwrites", {}) def set_overwrites(self, overwrites): """set per channel overwrites""" diff --git a/tubearchivist/home/src/index/comments.py b/tubearchivist/home/src/index/comments.py index 5dced73d..794cbc35 100644 --- a/tubearchivist/home/src/index/comments.py +++ b/tubearchivist/home/src/index/comments.py @@ -10,6 +10,7 @@ from datetime import datetime from home.src.download.yt_dlp_base import YtWrap from home.src.es.connect import ElasticWrap from home.src.ta.config import AppConfig +from home.src.ta.ta_redis import RedisQueue class Comments: @@ -189,20 +190,30 @@ class Comments: class CommentList: """interact with comments in group""" - def __init__(self, video_ids, task=False): - self.video_ids = video_ids + COMMENT_QUEUE = "index:comment" + + def __init__(self, task=False): self.task = task self.config = AppConfig().config - def index(self): - """index comments for list, init with task object to notify""" + def add(self, video_ids: list[str]) -> None: + """add list of videos to get comments, if enabled in config""" if not self.config["downloads"].get("comment_max"): return - total_videos = len(self.video_ids) - for idx, youtube_id in enumerate(self.video_ids): + RedisQueue(self.COMMENT_QUEUE).add_list(video_ids) + + def index(self): + """run comment index""" + queue = RedisQueue(self.COMMENT_QUEUE) + while True: + total = queue.max_score() + youtube_id, idx = queue.get_next() + if not youtube_id or not idx or not total: + break + if self.task: - self.notify(idx, total_videos) + self.notify(idx, total) comment = Comments(youtube_id, config=self.config) comment.build_json() @@ -211,6 +222,6 @@ class CommentList: def notify(self, idx, total_videos): """send notification on task""" - message = [f"Add comments for new videos {idx + 1}/{total_videos}"] - progress = (idx + 1) / total_videos + message = [f"Add comments for new videos {idx}/{total_videos}"] + progress = idx / total_videos self.task.send_progress(message, progress=progress) diff --git a/tubearchivist/home/src/index/filesystem.py b/tubearchivist/home/src/index/filesystem.py index c650cedb..484271a8 100644 --- a/tubearchivist/home/src/index/filesystem.py +++ b/tubearchivist/home/src/index/filesystem.py @@ -89,7 +89,9 @@ class Scanner: ) index_new_video(youtube_id) - CommentList(self.to_index, task=self.task).index() + comment_list = CommentList(task=self.task) + comment_list.add(video_ids=list(self.to_index)) + comment_list.index() def url_fix(self) -> None: """ diff --git a/tubearchivist/home/src/index/generic.py b/tubearchivist/home/src/index/generic.py index e18623ee..8a502bb5 100644 --- a/tubearchivist/home/src/index/generic.py +++ b/tubearchivist/home/src/index/generic.py @@ -17,7 +17,7 @@ class YouTubeItem: es_path = False index_name = "" yt_base = "" - yt_obs = { + yt_obs: dict[str, bool | str] = { "skip_download": True, "noplaylist": True, } diff --git a/tubearchivist/home/src/index/manual.py b/tubearchivist/home/src/index/manual.py index a25b0611..aa65ebdb 100644 --- a/tubearchivist/home/src/index/manual.py +++ b/tubearchivist/home/src/index/manual.py @@ -147,7 +147,9 @@ class ImportFolderScanner: ManualImport(current_video, self.CONFIG).run() video_ids = [i["video_id"] for i in self.to_import] - CommentList(video_ids, task=self.task).index() + comment_list = CommentList(task=self.task) + comment_list.add(video_ids=video_ids) + comment_list.index() def _notify(self, idx, current_video): """send notification back to task""" diff --git a/tubearchivist/home/src/index/playlist.py b/tubearchivist/home/src/index/playlist.py index 196a8844..2c577f57 100644 --- a/tubearchivist/home/src/index/playlist.py +++ b/tubearchivist/home/src/index/playlist.py @@ -8,7 +8,7 @@ import json from datetime import datetime from home.src.download.thumbnails import ThumbManager -from home.src.es.connect import ElasticWrap +from home.src.es.connect import ElasticWrap, IndexPaginate from home.src.index.generic import YouTubeItem from home.src.index.video import YoutubeVideo @@ -28,7 +28,6 @@ class YoutubePlaylist(YouTubeItem): super().__init__(youtube_id) self.all_members = False self.nav = False - self.all_youtube_ids = [] def build_json(self, scrape=False): """collection to create json_data""" @@ -45,7 +44,8 @@ class YoutubePlaylist(YouTubeItem): return self.process_youtube_meta() - self.get_entries() + ids_found = self.get_local_vids() + self.get_entries(ids_found) self.json_data["playlist_entries"] = self.all_members self.json_data["playlist_subscribed"] = subscribed @@ -69,25 +69,31 @@ class YoutubePlaylist(YouTubeItem): "playlist_type": "regular", } - def get_entries(self, playlistend=False): - """get all videos in playlist""" - if playlistend: - # implement playlist end - print(playlistend) + def get_local_vids(self) -> list[str]: + """get local video ids from youtube entries""" + entries = self.youtube_meta["entries"] + data = { + "query": {"terms": {"youtube_id": [i["id"] for i in entries]}}, + "_source": ["youtube_id"], + } + indexed_vids = IndexPaginate("ta_video", data).get_results() + ids_found = [i["youtube_id"] for i in indexed_vids] + + return ids_found + + def get_entries(self, ids_found) -> None: + """get all videos in playlist, match downloaded with ids_found""" all_members = [] for idx, entry in enumerate(self.youtube_meta["entries"]): - if self.all_youtube_ids: - downloaded = entry["id"] in self.all_youtube_ids - else: - downloaded = False if not entry["channel"]: continue + to_append = { "youtube_id": entry["id"], "title": entry["title"], "uploader": entry["channel"], "idx": idx, - "downloaded": downloaded, + "downloaded": entry["id"] in ids_found, } all_members.append(to_append) @@ -128,17 +134,50 @@ class YoutubePlaylist(YouTubeItem): ElasticWrap("_bulk").post(query_str, ndjson=True) - def update_playlist(self): + def remove_vids_from_playlist(self): + """remove playlist ids from videos if needed""" + needed = [i["youtube_id"] for i in self.json_data["playlist_entries"]] + data = { + "query": {"match": {"playlist": self.youtube_id}}, + "_source": ["youtube_id"], + } + result = IndexPaginate("ta_video", data).get_results() + to_remove = [ + i["youtube_id"] for i in result if i["youtube_id"] not in needed + ] + s = "ctx._source.playlist.removeAll(Collections.singleton(params.rm))" + for video_id in to_remove: + query = { + "script": { + "source": s, + "lang": "painless", + "params": {"rm": self.youtube_id}, + }, + "query": {"match": {"youtube_id": video_id}}, + } + path = "ta_video/_update_by_query" + _, status_code = ElasticWrap(path).post(query) + if status_code == 200: + print(f"{self.youtube_id}: removed {video_id} from playlist") + + def update_playlist(self, skip_on_empty=False): """update metadata for playlist with data from YouTube""" - self.get_from_es() - subscribed = self.json_data["playlist_subscribed"] - self.get_from_youtube() + self.build_json(scrape=True) if not self.json_data: # return false to deactivate return False - self.json_data["playlist_subscribed"] = subscribed + if skip_on_empty: + has_item_downloaded = any( + i["downloaded"] for i in self.json_data["playlist_entries"] + ) + if not has_item_downloaded: + return True + self.upload_to_es() + self.add_vids_to_playlist() + self.remove_vids_from_playlist() + self.get_playlist_art() return True def build_nav(self, youtube_id): diff --git a/tubearchivist/home/src/index/reindex.py b/tubearchivist/home/src/index/reindex.py index 10a5596d..1db975e5 100644 --- a/tubearchivist/home/src/index/reindex.py +++ b/tubearchivist/home/src/index/reindex.py @@ -8,8 +8,9 @@ import json import os from datetime import datetime from time import sleep +from typing import Callable, TypedDict -from home.src.download.queue import PendingList +from home.models import CustomPeriodicTask from home.src.download.subscriptions import ChannelSubscription from home.src.download.thumbnails import ThumbManager from home.src.download.yt_dlp_base import CookieHandler @@ -23,10 +24,19 @@ from home.src.ta.settings import EnvironmentSettings from home.src.ta.ta_redis import RedisQueue +class ReindexConfigType(TypedDict): + """represents config type""" + + index_name: str + queue_name: str + active_key: str + refresh_key: str + + class ReindexBase: """base config class for reindex task""" - REINDEX_CONFIG = { + REINDEX_CONFIG: dict[str, ReindexConfigType] = { "video": { "index_name": "ta_video", "queue_name": "reindex:ta_video", @@ -53,25 +63,36 @@ class ReindexBase: def __init__(self): self.config = AppConfig().config self.now = int(datetime.now().timestamp()) - self.total = None - def populate(self, all_ids, reindex_config): + def populate(self, all_ids, reindex_config: ReindexConfigType): """add all to reindex ids to redis queue""" if not all_ids: return RedisQueue(queue_name=reindex_config["queue_name"]).add_list(all_ids) - self.total = None class ReindexPopulate(ReindexBase): """add outdated and recent documents to reindex queue""" + INTERVAL_DEFAIULT: int = 90 + def __init__(self): super().__init__() - self.interval = self.config["scheduler"]["check_reindex_days"] + self.interval = self.INTERVAL_DEFAIULT - def add_recent(self): + def get_interval(self) -> None: + """get reindex days interval from task""" + try: + task = CustomPeriodicTask.objects.get(name="check_reindex") + except CustomPeriodicTask.DoesNotExist: + return + + task_config = task.task_config + if task_config.get("days"): + self.interval = task_config.get("days") + + def add_recent(self) -> None: """add recent videos to refresh""" gte = datetime.fromtimestamp(self.now - self.DAYS3).date().isoformat() must_list = [ @@ -89,10 +110,10 @@ class ReindexPopulate(ReindexBase): return all_ids = [i["_source"]["youtube_id"] for i in hits] - reindex_config = self.REINDEX_CONFIG.get("video") + reindex_config: ReindexConfigType = self.REINDEX_CONFIG["video"] self.populate(all_ids, reindex_config) - def add_outdated(self): + def add_outdated(self) -> None: """add outdated documents""" for reindex_config in self.REINDEX_CONFIG.values(): total_hits = self._get_total_hits(reindex_config) @@ -101,7 +122,7 @@ class ReindexPopulate(ReindexBase): self.populate(all_ids, reindex_config) @staticmethod - def _get_total_hits(reindex_config): + def _get_total_hits(reindex_config: ReindexConfigType) -> int: """get total hits from index""" index_name = reindex_config["index_name"] active_key = reindex_config["active_key"] @@ -113,7 +134,7 @@ class ReindexPopulate(ReindexBase): return len(total) - def _get_daily_should(self, total_hits): + def _get_daily_should(self, total_hits: int) -> int: """calc how many should reindex daily""" daily_should = int((total_hits // self.interval + 1) * self.MULTIPLY) if daily_should >= 10000: @@ -121,7 +142,9 @@ class ReindexPopulate(ReindexBase): return daily_should - def _get_outdated_ids(self, reindex_config, daily_should): + def _get_outdated_ids( + self, reindex_config: ReindexConfigType, daily_should: int + ) -> list[str]: """get outdated from index_name""" index_name = reindex_config["index_name"] refresh_key = reindex_config["refresh_key"] @@ -158,7 +181,7 @@ class ReindexManual(ReindexBase): self.extract_videos = extract_videos self.data = False - def extract_data(self, data): + def extract_data(self, data) -> None: """process data""" self.data = data for key, values in self.data.items(): @@ -169,7 +192,9 @@ class ReindexManual(ReindexBase): self.process_index(reindex_config, values) - def process_index(self, index_config, values): + def process_index( + self, index_config: ReindexConfigType, values: list[str] + ) -> None: """process values per index""" index_name = index_config["index_name"] if index_name == "ta_video": @@ -179,32 +204,35 @@ class ReindexManual(ReindexBase): elif index_name == "ta_playlist": self._add_playlists(values) - def _add_videos(self, values): + def _add_videos(self, values: list[str]) -> None: """add list of videos to reindex queue""" if not values: return - RedisQueue("reindex:ta_video").add_list(values) + queue_name = self.REINDEX_CONFIG["video"]["queue_name"] + RedisQueue(queue_name).add_list(values) - def _add_channels(self, values): + def _add_channels(self, values: list[str]) -> None: """add list of channels to reindex queue""" - RedisQueue("reindex:ta_channel").add_list(values) + queue_name = self.REINDEX_CONFIG["channel"]["queue_name"] + RedisQueue(queue_name).add_list(values) if self.extract_videos: for channel_id in values: all_videos = self._get_channel_videos(channel_id) self._add_videos(all_videos) - def _add_playlists(self, values): + def _add_playlists(self, values: list[str]) -> None: """add list of playlists to reindex queue""" - RedisQueue("reindex:ta_playlist").add_list(values) + queue_name = self.REINDEX_CONFIG["playlist"]["queue_name"] + RedisQueue(queue_name).add_list(values) if self.extract_videos: for playlist_id in values: all_videos = self._get_playlist_videos(playlist_id) self._add_videos(all_videos) - def _get_channel_videos(self, channel_id): + def _get_channel_videos(self, channel_id: str) -> list[str]: """get all videos from channel""" data = { "query": {"term": {"channel.channel_id": {"value": channel_id}}}, @@ -213,7 +241,7 @@ class ReindexManual(ReindexBase): all_results = IndexPaginate("ta_video", data).get_results() return [i["youtube_id"] for i in all_results] - def _get_playlist_videos(self, playlist_id): + def _get_playlist_videos(self, playlist_id: str) -> list[str]: """get all videos from playlist""" data = { "query": {"term": {"playlist.keyword": {"value": playlist_id}}}, @@ -229,14 +257,13 @@ class Reindex(ReindexBase): def __init__(self, task=False): super().__init__() self.task = task - self.all_indexed_ids = False self.processed = { "videos": 0, "channels": 0, "playlists": 0, } - def reindex_all(self): + def reindex_all(self) -> None: """reindex all in queue""" if not self.cookie_is_valid(): print("[reindex] cookie invalid, exiting...") @@ -246,26 +273,26 @@ class Reindex(ReindexBase): if not RedisQueue(index_config["queue_name"]).length(): continue - self.total = RedisQueue(index_config["queue_name"]).length() - while True: - has_next = self.reindex_index(name, index_config) - if not has_next: - break + self.reindex_type(name, index_config) - def reindex_index(self, name, index_config): + def reindex_type(self, name: str, index_config: ReindexConfigType) -> None: """reindex all of a single index""" - reindex = self.get_reindex_map(index_config["index_name"]) - youtube_id = RedisQueue(index_config["queue_name"]).get_next() - if youtube_id: + reindex = self._get_reindex_map(index_config["index_name"]) + queue = RedisQueue(index_config["queue_name"]) + while True: + total = queue.max_score() + youtube_id, idx = queue.get_next() + if not youtube_id or not idx or not total: + break + if self.task: - self._notify(name, index_config) + self._notify(name, total, idx) + reindex(youtube_id) sleep_interval = self.config["downloads"].get("sleep_interval", 0) sleep(sleep_interval) - return bool(youtube_id) - - def get_reindex_map(self, index_name): + def _get_reindex_map(self, index_name: str) -> Callable: """return def to run for index""" def_map = { "ta_video": self._reindex_single_video, @@ -273,20 +300,15 @@ class Reindex(ReindexBase): "ta_playlist": self._reindex_single_playlist, } - return def_map.get(index_name) + return def_map[index_name] - def _notify(self, name, index_config): + def _notify(self, name: str, total: int, idx: int) -> None: """send notification back to task""" - if self.total is None: - self.total = RedisQueue(index_config["queue_name"]).length() - - remaining = RedisQueue(index_config["queue_name"]).length() - idx = self.total - remaining - message = [f"Reindexing {name.title()}s {idx}/{self.total}"] - progress = idx / self.total + message = [f"Reindexing {name.title()}s {idx}/{total}"] + progress = idx / total self.task.send_progress(message, progress=progress) - def _reindex_single_video(self, youtube_id): + def _reindex_single_video(self, youtube_id: str) -> None: """refresh data for single video""" video = YoutubeVideo(youtube_id) @@ -325,9 +347,7 @@ class Reindex(ReindexBase): Comments(youtube_id, config=self.config).reindex_comments() self.processed["videos"] += 1 - return - - def _reindex_single_channel(self, channel_id): + def _reindex_single_channel(self, channel_id: str) -> None: """refresh channel data and sync to videos""" # read current state channel = YoutubeChannel(channel_id) @@ -358,9 +378,8 @@ class Reindex(ReindexBase): ChannelFullScan(channel_id).scan() self.processed["channels"] += 1 - def _reindex_single_playlist(self, playlist_id): + def _reindex_single_playlist(self, playlist_id: str) -> None: """refresh playlist data""" - self._get_all_videos() playlist = YoutubePlaylist(playlist_id) playlist.get_from_es() if ( @@ -369,29 +388,14 @@ class Reindex(ReindexBase): ): return - subscribed = playlist.json_data["playlist_subscribed"] - playlist.all_youtube_ids = self.all_indexed_ids - playlist.build_json(scrape=True) - if not playlist.json_data: + is_active = playlist.update_playlist() + if not is_active: playlist.deactivate() return - playlist.json_data["playlist_subscribed"] = subscribed - playlist.upload_to_es() self.processed["playlists"] += 1 - return - def _get_all_videos(self): - """add all videos for playlist index validation""" - if self.all_indexed_ids: - return - - handler = PendingList() - handler.get_download() - handler.get_indexed() - self.all_indexed_ids = [i["youtube_id"] for i in handler.all_videos] - - def cookie_is_valid(self): + def cookie_is_valid(self) -> bool: """return true if cookie is enabled and valid""" if not self.config["downloads"]["cookie_import"]: # is not activated, continue reindex @@ -400,7 +404,7 @@ class Reindex(ReindexBase): valid = CookieHandler(self.config).validate() return valid - def build_message(self): + def build_message(self) -> str: """build progress message""" message = "" for key, value in self.processed.items(): @@ -430,7 +434,7 @@ class ReindexProgress(ReindexBase): self.request_type = request_type self.request_id = request_id - def get_progress(self): + def get_progress(self) -> dict: """get progress from task""" queue_name, request_type = self._get_queue_name() total = self._get_total_in_queue(queue_name) diff --git a/tubearchivist/home/src/index/video.py b/tubearchivist/home/src/index/video.py index 41f01147..22efc088 100644 --- a/tubearchivist/home/src/index/video.py +++ b/tubearchivist/home/src/index/video.py @@ -125,15 +125,9 @@ class YoutubeVideo(YouTubeItem, YoutubeSubtitle): index_name = "ta_video" yt_base = "https://www.youtube.com/watch?v=" - def __init__( - self, - youtube_id, - video_overwrites=False, - video_type=VideoTypeEnum.VIDEOS, - ): + def __init__(self, youtube_id, video_type=VideoTypeEnum.VIDEOS): super().__init__(youtube_id) self.channel_id = False - self.video_overwrites = video_overwrites self.video_type = video_type self.offline_import = False @@ -147,7 +141,7 @@ class YoutubeVideo(YouTubeItem, YoutubeSubtitle): self.youtube_meta = youtube_meta_overwrite self.offline_import = True - self._process_youtube_meta() + self.process_youtube_meta() self._add_channel() self._add_stats() self.add_file_path() @@ -165,17 +159,16 @@ class YoutubeVideo(YouTubeItem, YoutubeSubtitle): """check if need to run sponsor block""" integrate = self.config["downloads"]["integrate_sponsorblock"] - if self.video_overwrites: - single_overwrite = self.video_overwrites.get(self.youtube_id) - if not single_overwrite: + if overwrite := self.json_data["channel"].get("channel_overwrites"): + if not overwrite: return integrate - if "integrate_sponsorblock" in single_overwrite: - return single_overwrite.get("integrate_sponsorblock") + if "integrate_sponsorblock" in overwrite: + return overwrite.get("integrate_sponsorblock") return integrate - def _process_youtube_meta(self): + def process_youtube_meta(self): """extract relevant fields from youtube""" self._validate_id() # extract @@ -399,13 +392,9 @@ class YoutubeVideo(YouTubeItem, YoutubeSubtitle): _, _ = ElasticWrap(path).post(data=data) -def index_new_video( - youtube_id, video_overwrites=False, video_type=VideoTypeEnum.VIDEOS -): +def index_new_video(youtube_id, video_type=VideoTypeEnum.VIDEOS): """combined classes to create new video in index""" - video = YoutubeVideo( - youtube_id, video_overwrites=video_overwrites, video_type=video_type - ) + video = YoutubeVideo(youtube_id, video_type=video_type) video.build_json() if not video.json_data: raise ValueError("failed to get metadata for " + youtube_id) diff --git a/tubearchivist/home/src/ta/config.py b/tubearchivist/home/src/ta/config.py index 39819322..af342b93 100644 --- a/tubearchivist/home/src/ta/config.py +++ b/tubearchivist/home/src/ta/config.py @@ -5,12 +5,10 @@ Functionality: """ import json -import re from random import randint from time import sleep import requests -from celery.schedules import crontab from django.conf import settings from home.src.ta.ta_redis import RedisArchivist @@ -74,15 +72,6 @@ class AppConfig: RedisArchivist().set_message("config", self.config, save=True) return updated - @staticmethod - def _build_rand_daily(): - """build random daily schedule per installation""" - return { - "minute": randint(0, 59), - "hour": randint(0, 23), - "day_of_week": "*", - } - def load_new_defaults(self): """check config.json for missing defaults""" default_config = self.get_config_file() @@ -91,7 +80,6 @@ class AppConfig: # check for customizations if not redis_config: config = self.get_config() - config["scheduler"]["version_check"] = self._build_rand_daily() RedisArchivist().set_message("config", config) return False @@ -106,13 +94,7 @@ class AppConfig: # missing nested values for sub_key, sub_value in value.items(): - if ( - sub_key not in redis_config[key].keys() - or sub_value == "rand-d" - ): - if sub_value == "rand-d": - sub_value = self._build_rand_daily() - + if sub_key not in redis_config[key].keys(): redis_config[key].update({sub_key: sub_value}) needs_update = True @@ -122,147 +104,6 @@ class AppConfig: return needs_update -class ScheduleBuilder: - """build schedule dicts for beat""" - - SCHEDULES = { - "update_subscribed": "0 8 *", - "download_pending": "0 16 *", - "check_reindex": "0 12 *", - "thumbnail_check": "0 17 *", - "run_backup": "0 18 0", - "version_check": "0 11 *", - } - CONFIG = ["check_reindex_days", "run_backup_rotate"] - NOTIFY = [ - "update_subscribed_notify", - "download_pending_notify", - "check_reindex_notify", - ] - MSG = "message:setting" - - def __init__(self): - self.config = AppConfig().config - - def update_schedule_conf(self, form_post): - """process form post""" - print("processing form, restart container for changes to take effect") - redis_config = self.config - for key, value in form_post.items(): - if key in self.SCHEDULES and value: - try: - to_write = self.value_builder(key, value) - except ValueError: - print(f"failed: {key} {value}") - mess_dict = { - "group": "setting:schedule", - "level": "error", - "title": "Scheduler update failed.", - "messages": ["Invalid schedule input"], - "id": "0000", - } - RedisArchivist().set_message( - self.MSG, mess_dict, expire=True - ) - return - - redis_config["scheduler"][key] = to_write - if key in self.CONFIG and value: - redis_config["scheduler"][key] = int(value) - if key in self.NOTIFY and value: - if value == "0": - to_write = False - else: - to_write = value - redis_config["scheduler"][key] = to_write - - RedisArchivist().set_message("config", redis_config, save=True) - mess_dict = { - "group": "setting:schedule", - "level": "info", - "title": "Scheduler changed.", - "messages": ["Restart container for changes to take effect"], - "id": "0000", - } - RedisArchivist().set_message(self.MSG, mess_dict, expire=True) - - def value_builder(self, key, value): - """validate single cron form entry and return cron dict""" - print(f"change schedule for {key} to {value}") - if value == "0": - # deactivate this schedule - return False - if re.search(r"[\d]{1,2}\/[\d]{1,2}", value): - # number/number cron format will fail in celery - print("number/number schedule formatting not supported") - raise ValueError - - keys = ["minute", "hour", "day_of_week"] - if value == "auto": - # set to sensible default - values = self.SCHEDULES[key].split() - else: - values = value.split() - - if len(keys) != len(values): - print(f"failed to parse {value} for {key}") - raise ValueError("invalid input") - - to_write = dict(zip(keys, values)) - self._validate_cron(to_write) - - return to_write - - @staticmethod - def _validate_cron(to_write): - """validate all fields, raise value error for impossible schedule""" - all_hours = list(re.split(r"\D+", to_write["hour"])) - for hour in all_hours: - if hour.isdigit() and int(hour) > 23: - print("hour can not be greater than 23") - raise ValueError("invalid input") - - all_days = list(re.split(r"\D+", to_write["day_of_week"])) - for day in all_days: - if day.isdigit() and int(day) > 6: - print("day can not be greater than 6") - raise ValueError("invalid input") - - if not to_write["minute"].isdigit(): - print("too frequent: only number in minutes are supported") - raise ValueError("invalid input") - - if int(to_write["minute"]) > 59: - print("minutes can not be greater than 59") - raise ValueError("invalid input") - - def build_schedule(self): - """build schedule dict as expected by app.conf.beat_schedule""" - AppConfig().load_new_defaults() - self.config = AppConfig().config - schedule_dict = {} - - for schedule_item in self.SCHEDULES: - item_conf = self.config["scheduler"][schedule_item] - if not item_conf: - continue - - schedule_dict.update( - { - f"schedule_{schedule_item}": { - "task": schedule_item, - "schedule": crontab( - minute=item_conf["minute"], - hour=item_conf["hour"], - day_of_week=item_conf["day_of_week"], - ), - } - } - ) - - return schedule_dict - - class ReleaseVersion: """compare local version with remote version""" @@ -340,12 +181,3 @@ class ReleaseVersion: return {} return message - - def clear_fail(self) -> None: - """clear key, catch previous error in v0.4.5""" - message = self.get_update() - if not message: - return - - if isinstance(message.get("version"), list): - RedisArchivist().del_message(self.NEW_KEY) diff --git a/tubearchivist/home/src/ta/config_schedule.py b/tubearchivist/home/src/ta/config_schedule.py new file mode 100644 index 00000000..24b81ca9 --- /dev/null +++ b/tubearchivist/home/src/ta/config_schedule.py @@ -0,0 +1,89 @@ +""" +Functionality: +- Handle scheduler config update +""" + +from django_celery_beat.models import CrontabSchedule +from home.models import CustomPeriodicTask +from home.src.ta.config import AppConfig +from home.src.ta.settings import EnvironmentSettings +from home.src.ta.task_config import TASK_CONFIG + + +class ScheduleBuilder: + """build schedule dicts for beat""" + + SCHEDULES = { + "update_subscribed": "0 8 *", + "download_pending": "0 16 *", + "check_reindex": "0 12 *", + "thumbnail_check": "0 17 *", + "run_backup": "0 18 0", + "version_check": "0 11 *", + } + CONFIG = { + "check_reindex_days": "check_reindex", + "run_backup_rotate": "run_backup", + "update_subscribed_notify": "update_subscribed", + "download_pending_notify": "download_pending", + "check_reindex_notify": "check_reindex", + } + MSG = "message:setting" + + def __init__(self): + self.config = AppConfig().config + + def update_schedule_conf(self, form_post): + """process form post, schedules need to be validated before""" + for key, value in form_post.items(): + if not value: + continue + + if key in self.SCHEDULES: + if value == "auto": + value = self.SCHEDULES.get(key) + + _ = self.get_set_task(key, value) + continue + + if key in self.CONFIG: + self.set_config(key, value) + + def get_set_task(self, task_name, schedule=False): + """get task""" + try: + task = CustomPeriodicTask.objects.get(name=task_name) + except CustomPeriodicTask.DoesNotExist: + description = TASK_CONFIG[task_name].get("title") + task = CustomPeriodicTask( + name=task_name, + task=task_name, + description=description, + ) + + if schedule: + task_crontab = self.get_set_cron_tab(schedule) + task.crontab = task_crontab + task.save() + + return task + + @staticmethod + def get_set_cron_tab(schedule): + """needs to be validated before""" + kwargs = dict(zip(["minute", "hour", "day_of_week"], schedule.split())) + kwargs.update({"timezone": EnvironmentSettings.TZ}) + crontab, _ = CrontabSchedule.objects.get_or_create(**kwargs) + + return crontab + + def set_config(self, key, value): + """set task_config""" + task_name = self.CONFIG.get(key) + if not task_name: + raise ValueError("invalid config key") + + task = CustomPeriodicTask.objects.get(name=task_name) + config_key = key.split(f"{task_name}_")[-1] + task.task_config.update({config_key: value}) + task.save() diff --git a/tubearchivist/home/src/ta/helper.py b/tubearchivist/home/src/ta/helper.py index e3089d85..05e1ef0e 100644 --- a/tubearchivist/home/src/ta/helper.py +++ b/tubearchivist/home/src/ta/helper.py @@ -9,9 +9,11 @@ import random import string import subprocess from datetime import datetime +from typing import Any from urllib.parse import urlparse import requests +from home.src.es.connect import IndexPaginate from home.src.ta.settings import EnvironmentSettings @@ -222,3 +224,37 @@ def check_stylesheet(stylesheet: str): return stylesheet return "dark.css" + + +def is_missing( + to_check: str | list[str], + index_name: str = "ta_video,ta_download", + on_key: str = "youtube_id", +) -> list[str]: + """id or list of ids that are missing from index_name""" + if isinstance(to_check, str): + to_check = [to_check] + + data = { + "query": {"terms": {on_key: to_check}}, + "_source": [on_key], + } + result = IndexPaginate(index_name, data=data).get_results() + existing_ids = [i[on_key] for i in result] + dl = [i for i in to_check if i not in existing_ids] + + return dl + + +def get_channel_overwrites() -> dict[str, dict[str, Any]]: + """get overwrites indexed my channel_id""" + data = { + "query": { + "bool": {"must": [{"exists": {"field": "channel_overwrites"}}]} + }, + "_source": ["channel_id", "channel_overwrites"], + } + result = IndexPaginate("ta_channel", data).get_results() + overwrites = {i["channel_id"]: i["channel_overwrites"] for i in result} + + return overwrites diff --git a/tubearchivist/home/src/ta/notify.py b/tubearchivist/home/src/ta/notify.py index 5b1c9c7a..63140775 100644 --- a/tubearchivist/home/src/ta/notify.py +++ b/tubearchivist/home/src/ta/notify.py @@ -1,55 +1,141 @@ """send notifications using apprise""" import apprise -from home.src.ta.config import AppConfig +from home.src.es.connect import ElasticWrap +from home.src.ta.task_config import TASK_CONFIG from home.src.ta.task_manager import TaskManager class Notifications: - """notification handler""" + """store notifications in ES""" - def __init__(self, name: str, task_id: str, task_title: str): - self.name: str = name - self.task_id: str = task_id - self.task_title: str = task_title + GET_PATH = "ta_config/_doc/notify" + UPDATE_PATH = "ta_config/_update/notify/" - def send(self) -> None: + def __init__(self, task_name: str): + self.task_name = task_name + + def send(self, task_id: str, task_title: str) -> None: """send notifications""" apobj = apprise.Apprise() - hooks: str | None = self.get_url() - if not hooks: + urls: list[str] = self.get_urls() + if not urls: return - hook_list: list[str] = self.parse_hooks(hooks=hooks) - title, body = self.build_message() + title, body = self._build_message(task_id, task_title) if not body: return - for hook in hook_list: - apobj.add(hook) + for url in urls: + apobj.add(url) apobj.notify(body=body, title=title) - def get_url(self) -> str | None: - """get apprise urls for task""" - config = AppConfig().config - hooks: str = config["scheduler"].get(f"{self.name}_notify") - - return hooks - - def parse_hooks(self, hooks: str) -> list[str]: - """create list of hooks""" - - hook_list: list[str] = [i.strip() for i in hooks.split()] - - return hook_list - - def build_message(self) -> tuple[str, str | None]: + def _build_message( + self, task_id: str, task_title: str + ) -> tuple[str, str | None]: """build message to send notification""" - task = TaskManager().get_task(self.task_id) + task = TaskManager().get_task(task_id) status = task.get("status") - title: str = f"[TA] {self.task_title} process ended with {status}" + title: str = f"[TA] {task_title} process ended with {status}" body: str | None = task.get("result") return title, body + + def get_urls(self) -> list[str]: + """get stored urls for task""" + response, code = ElasticWrap(self.GET_PATH).get(print_error=False) + if not code == 200: + return [] + + urls = response["_source"].get(self.task_name, []) + + return urls + + def add_url(self, url: str) -> None: + """add url to task notification""" + source = ( + "if (!ctx._source.containsKey(params.task_name)) " + + "{ctx._source[params.task_name] = [params.url]} " + + "else if (!ctx._source[params.task_name].contains(params.url)) " + + "{ctx._source[params.task_name].add(params.url)} " + + "else {ctx.op = 'none'}" + ) + + data = { + "script": { + "source": source, + "lang": "painless", + "params": {"url": url, "task_name": self.task_name}, + }, + "upsert": {self.task_name: [url]}, + } + + _, _ = ElasticWrap(self.UPDATE_PATH).post(data) + + def remove_url(self, url: str) -> tuple[dict, int]: + """remove url from task""" + source = ( + "if (ctx._source.containsKey(params.task_name) " + + "&& ctx._source[params.task_name].contains(params.url)) " + + "{ctx._source[params.task_name]." + + "remove(ctx._source[params.task_name].indexOf(params.url))}" + ) + + data = { + "script": { + "source": source, + "lang": "painless", + "params": {"url": url, "task_name": self.task_name}, + } + } + + response, status_code = ElasticWrap(self.UPDATE_PATH).post(data) + if not self.get_urls(): + _, _ = self.remove_task() + + return response, status_code + + def remove_task(self) -> tuple[dict, int]: + """remove all notifications from task""" + source = ( + "if (ctx._source.containsKey(params.task_name)) " + + "{ctx._source.remove(params.task_name)}" + ) + data = { + "script": { + "source": source, + "lang": "painless", + "params": {"task_name": self.task_name}, + } + } + + response, status_code = ElasticWrap(self.UPDATE_PATH).post(data) + + return response, status_code + + +def get_all_notifications() -> dict[str, list[str]]: + """get all notifications stored""" + path = "ta_config/_doc/notify" + response, status_code = ElasticWrap(path).get(print_error=False) + if not status_code == 200: + return {} + + notifications: dict = {} + source = response.get("_source") + if not source: + return notifications + + for task_id, urls in source.items(): + notifications.update( + { + task_id: { + "urls": urls, + "title": TASK_CONFIG[task_id]["title"], + } + } + ) + + return notifications diff --git a/tubearchivist/home/src/ta/ta_redis.py b/tubearchivist/home/src/ta/ta_redis.py index dd9ac511..1feb0232 100644 --- a/tubearchivist/home/src/ta/ta_redis.py +++ b/tubearchivist/home/src/ta/ta_redis.py @@ -104,6 +104,17 @@ class RedisQueue(RedisBase): dynamically interact with queues in redis using sorted set - low score number is first in queue - add new items with high score number + + queue names in use: + download:channel channels during download + download:playlist:full playlists during dl for full refresh + download:playlist:quick playlists during dl for quick refresh + download:video videos during downloads + index:comment videos needing comment indexing + reindex:ta_video reindex videos + reindex:ta_channel reindex channels + reindex:ta_playlist reindex playlists + """ def __init__(self, queue_name: str): @@ -127,18 +138,48 @@ class RedisQueue(RedisBase): return False + def add(self, to_add: str) -> None: + """add single item to queue""" + if not to_add: + return + + next_score = self._get_next_score() + self.conn.zadd(self.key, {to_add: next_score}) + def add_list(self, to_add: list) -> None: """add list to queue""" - mapping = {i: "+inf" for i in to_add} + if not to_add: + return + + next_score = self._get_next_score() + mapping = {i[1]: next_score + i[0] for i in enumerate(to_add)} self.conn.zadd(self.key, mapping) - def get_next(self) -> str | bool: + def max_score(self) -> int | None: + """get max score""" + last = self.conn.zrange(self.key, -1, -1, withscores=True) + if not last: + return None + + return int(last[0][1]) + + def _get_next_score(self) -> float: + """get next score in queue to append""" + last = self.conn.zrange(self.key, -1, -1, withscores=True) + if not last: + return 1.0 + + return last[0][1] + 1 + + def get_next(self) -> tuple[str | None, int | None]: """return next element in the queue, if available""" result = self.conn.zpopmin(self.key) if not result: - return False + return None, None - return result[0][0] + item, idx = result[0][0], int(result[0][1]) + + return item, idx def clear(self) -> None: """delete list from redis""" diff --git a/tubearchivist/home/src/ta/task_config.py b/tubearchivist/home/src/ta/task_config.py new file mode 100644 index 00000000..5c3dc433 --- /dev/null +++ b/tubearchivist/home/src/ta/task_config.py @@ -0,0 +1,125 @@ +""" +Functionality: +- Static Task config values +- Type definitions +- separate to avoid circular imports +""" + +from typing import TypedDict + + +class TaskItemConfig(TypedDict): + """describes a task item config""" + + title: str + group: str + api_start: bool + api_stop: bool + + +UPDATE_SUBSCRIBED: TaskItemConfig = { + "title": "Rescan your Subscriptions", + "group": "download:scan", + "api_start": True, + "api_stop": True, +} + +DOWNLOAD_PENDING: TaskItemConfig = { + "title": "Downloading", + "group": "download:run", + "api_start": True, + "api_stop": True, +} + +EXTRACT_DOWNLOAD: TaskItemConfig = { + "title": "Add to download queue", + "group": "download:add", + "api_start": False, + "api_stop": True, +} + +CHECK_REINDEX: TaskItemConfig = { + "title": "Reindex Documents", + "group": "reindex:run", + "api_start": False, + "api_stop": False, +} + +MANUAL_IMPORT: TaskItemConfig = { + "title": "Manual video import", + "group": "setting:import", + "api_start": True, + "api_stop": False, +} + +RUN_BACKUP: TaskItemConfig = { + "title": "Index Backup", + "group": "setting:backup", + "api_start": True, + "api_stop": False, +} + +RESTORE_BACKUP: TaskItemConfig = { + "title": "Restore Backup", + "group": "setting:restore", + "api_start": False, + "api_stop": False, +} + +RESCAN_FILESYSTEM: TaskItemConfig = { + "title": "Rescan your Filesystem", + "group": "setting:filesystemscan", + "api_start": True, + "api_stop": False, +} + +THUMBNAIL_CHECK: TaskItemConfig = { + "title": "Check your Thumbnails", + "group": "setting:thumbnailcheck", + "api_start": True, + "api_stop": False, +} + +RESYNC_THUMBS: TaskItemConfig = { + "title": "Sync Thumbnails to Media Files", + "group": "setting:thumbnailsync", + "api_start": True, + "api_stop": False, +} + +INDEX_PLAYLISTS: TaskItemConfig = { + "title": "Index Channel Playlist", + "group": "channel:indexplaylist", + "api_start": False, + "api_stop": False, +} + +SUBSCRIBE_TO: TaskItemConfig = { + "title": "Add Subscription", + "group": "subscription:add", + "api_start": False, + "api_stop": False, +} + +VERSION_CHECK: TaskItemConfig = { + "title": "Look for new Version", + "group": "", + "api_start": False, + "api_stop": False, +} + +TASK_CONFIG: dict[str, TaskItemConfig] = { + "update_subscribed": UPDATE_SUBSCRIBED, + "download_pending": DOWNLOAD_PENDING, + "extract_download": EXTRACT_DOWNLOAD, + "check_reindex": CHECK_REINDEX, + "manual_import": MANUAL_IMPORT, + "run_backup": RUN_BACKUP, + "restore_backup": RESTORE_BACKUP, + "rescan_filesystem": RESCAN_FILESYSTEM, + "thumbnail_check": THUMBNAIL_CHECK, + "resync_thumbs": RESYNC_THUMBS, + "index_playlists": INDEX_PLAYLISTS, + "subscribe_to": SUBSCRIBE_TO, + "version_check": VERSION_CHECK, +} diff --git a/tubearchivist/home/src/ta/task_manager.py b/tubearchivist/home/src/ta/task_manager.py index 3771fdd6..0813ccb5 100644 --- a/tubearchivist/home/src/ta/task_manager.py +++ b/tubearchivist/home/src/ta/task_manager.py @@ -4,8 +4,9 @@ functionality: - handle threads and locks """ -from home import tasks as ta_tasks +from home.celery import app as celery_app from home.src.ta.ta_redis import RedisArchivist, TaskRedis +from home.src.ta.task_config import TASK_CONFIG class TaskManager: @@ -86,7 +87,7 @@ class TaskCommand: def start(self, task_name): """start task by task_name, only pass task that don't take args""" - task = ta_tasks.app.tasks.get(task_name).delay() + task = celery_app.tasks.get(task_name).delay() message = { "task_id": task.id, "status": task.status, @@ -104,7 +105,7 @@ class TaskCommand: handler = TaskRedis() task = handler.get_single(task_id) - if not task["name"] in ta_tasks.BaseTask.TASK_CONFIG: + if not task["name"] in TASK_CONFIG: raise ValueError handler.set_command(task_id, "STOP") @@ -113,4 +114,4 @@ class TaskCommand: def kill(self, task_id): """send kill signal to task_id""" print(f"[task][{task_id}]: received KILL signal.") - ta_tasks.app.control.revoke(task_id, terminate=True) + celery_app.control.revoke(task_id, terminate=True) diff --git a/tubearchivist/home/tasks.py b/tubearchivist/home/tasks.py index e0ab8e7c..89f15a6b 100644 --- a/tubearchivist/home/tasks.py +++ b/tubearchivist/home/tasks.py @@ -1,14 +1,12 @@ """ Functionality: -- initiate celery app - collect tasks -- user config changes won't get applied here - because tasks are initiated at application start +- handle task callbacks +- handle task notifications +- handle task locking """ -import os - -from celery import Celery, Task, shared_task +from celery import Task, shared_task from home.src.download.queue import PendingList from home.src.download.subscriptions import ( SubscriptionHandler, @@ -22,97 +20,19 @@ from home.src.index.channel import YoutubeChannel from home.src.index.filesystem import Scanner from home.src.index.manual import ImportFolderScanner from home.src.index.reindex import Reindex, ReindexManual, ReindexPopulate -from home.src.ta.config import AppConfig, ReleaseVersion, ScheduleBuilder +from home.src.ta.config import ReleaseVersion from home.src.ta.notify import Notifications -from home.src.ta.settings import EnvironmentSettings from home.src.ta.ta_redis import RedisArchivist +from home.src.ta.task_config import TASK_CONFIG from home.src.ta.task_manager import TaskManager from home.src.ta.urlparser import Parser -CONFIG = AppConfig().config -REDIS_HOST = EnvironmentSettings.REDIS_HOST -REDIS_PORT = EnvironmentSettings.REDIS_PORT - -os.environ.setdefault("DJANGO_SETTINGS_MODULE", "config.settings") -app = Celery( - "tasks", - broker=f"redis://{REDIS_HOST}:{REDIS_PORT}", - backend=f"redis://{REDIS_HOST}:{REDIS_PORT}", - result_extended=True, -) -app.config_from_object( - "django.conf:settings", namespace=EnvironmentSettings.REDIS_NAME_SPACE -) -app.autodiscover_tasks() -app.conf.timezone = EnvironmentSettings.TZ - class BaseTask(Task): """base class to inherit each class from""" # pylint: disable=abstract-method - TASK_CONFIG = { - "update_subscribed": { - "title": "Rescan your Subscriptions", - "group": "download:scan", - "api-start": True, - "api-stop": True, - }, - "download_pending": { - "title": "Downloading", - "group": "download:run", - "api-start": True, - "api-stop": True, - }, - "extract_download": { - "title": "Add to download queue", - "group": "download:add", - "api-stop": True, - }, - "check_reindex": { - "title": "Reindex Documents", - "group": "reindex:run", - }, - "manual_import": { - "title": "Manual video import", - "group": "setting:import", - "api-start": True, - }, - "run_backup": { - "title": "Index Backup", - "group": "setting:backup", - "api-start": True, - }, - "restore_backup": { - "title": "Restore Backup", - "group": "setting:restore", - }, - "rescan_filesystem": { - "title": "Rescan your Filesystem", - "group": "setting:filesystemscan", - "api-start": True, - }, - "thumbnail_check": { - "title": "Check your Thumbnails", - "group": "setting:thumbnailcheck", - "api-start": True, - }, - "resync_thumbs": { - "title": "Sync Thumbnails to Media Files", - "group": "setting:thumbnailsync", - "api-start": True, - }, - "index_playlists": { - "title": "Index Channel Playlist", - "group": "channel:indexplaylist", - }, - "subscribe_to": { - "title": "Add Subscription", - "group": "subscription:add", - }, - } - def on_failure(self, exc, task_id, args, kwargs, einfo): """callback for task failure""" print(f"{task_id} Failed callback") @@ -137,8 +57,8 @@ class BaseTask(Task): def after_return(self, status, retval, task_id, args, kwargs, einfo): """callback after task returns""" print(f"{task_id} return callback") - task_title = self.TASK_CONFIG.get(self.name).get("title") - Notifications(self.name, task_id, task_title).send() + task_title = TASK_CONFIG.get(self.name).get("title") + Notifications(self.name).send(task_id, task_title) def send_progress(self, message_lines, progress=False, title=False): """send progress message""" @@ -157,7 +77,7 @@ class BaseTask(Task): def _build_message(self, level="info"): """build message dict""" task_id = self.request.id - message = self.TASK_CONFIG.get(self.name).copy() + message = TASK_CONFIG.get(self.name).copy() message.update({"level": level, "id": task_id}) task_result = TaskManager().get_task(task_id) if task_result: @@ -208,13 +128,13 @@ def download_pending(self, auto_only=False): videos_downloaded = downloader.run_queue(auto_only=auto_only) if videos_downloaded: - return f"downloaded {len(videos_downloaded)} videos." + return f"downloaded {videos_downloaded} video(s)." return None @shared_task(name="extract_download", bind=True, base=BaseTask) -def extrac_dl(self, youtube_ids, auto_start=False): +def extrac_dl(self, youtube_ids, auto_start=False, status="pending"): """parse list passed and add to pending""" TaskManager().init(self) if isinstance(youtube_ids, str): @@ -224,11 +144,18 @@ def extrac_dl(self, youtube_ids, auto_start=False): pending_handler = PendingList(youtube_ids=to_add, task=self) pending_handler.parse_url_list() - pending_handler.add_to_pending(auto_start=auto_start) + videos_added = pending_handler.add_to_pending( + status=status, auto_start=auto_start + ) if auto_start: download_pending.delay(auto_only=True) + if videos_added: + return f"added {len(videos_added)} Videos to Queue" + + return None + @shared_task(bind=True, name="check_reindex", base=BaseTask) def check_reindex(self, data=False, extract_videos=False): @@ -251,6 +178,7 @@ def check_reindex(self, data=False, extract_videos=False): populate = ReindexPopulate() print(f"[task][{self.name}] reindex outdated documents") self.send_progress("Add recent documents to the reindex Queue.") + populate.get_interval() populate.add_recent() self.send_progress("Add outdated documents to the reindex Queue.") populate.add_outdated() @@ -331,7 +259,9 @@ def thumbnail_check(self): return manager.init(self) - ThumbValidator(task=self).validate() + thumnail = ThumbValidator(task=self) + thumnail.validate() + thumnail.clean_up() @shared_task(bind=True, name="resync_thumbs", base=BaseTask) @@ -367,7 +297,3 @@ def index_channel_playlists(self, channel_id): def version_check(): """check for new updates""" ReleaseVersion().check() - - -# start schedule here -app.conf.beat_schedule = ScheduleBuilder().build_schedule() diff --git a/tubearchivist/home/templates/home/channel_id_playlist.html b/tubearchivist/home/templates/home/channel_id_playlist.html index 539d4388..52b99800 100644 --- a/tubearchivist/home/templates/home/channel_id_playlist.html +++ b/tubearchivist/home/templates/home/channel_id_playlist.html @@ -53,13 +53,14 @@
Last refreshed: {{ playlist.playlist_last_refresh }}
- {% if playlist.playlist_subscribed and request.user|has_group:"admin" or request.user.is_staff %} - - {% else %} - + {% if request.user|has_group:"admin" or request.user.is_staff %} + {% if playlist.playlist_subscribed %} + + {% else %} + + {% endif %} {% endif %}Note:
Become a sponsor and join members.tubearchivist.com to get access to real time notifications for new videos uploaded by your favorite channels.
Current rescan schedule: - {% if config.scheduler.update_subscribed %} - {% for key, value in config.scheduler.update_subscribed.items %} - {{ value }} - {% endfor %} + {% if update_subscribed %} + {{ update_subscribed.crontab.minute }} {{ update_subscribed.crontab.hour }} {{ update_subscribed.crontab.day_of_week }} + {% else %} False {% endif %}
-Become a sponsor and join members.tubearchivist.com to get access to real time notifications for new videos uploaded by your favorite channels.
Periodically rescan your subscriptions:
+ {% for error in scheduler_form.update_subscribed.errors %} +{{ error }}
+ {% endfor %} {{ scheduler_form.update_subscribed }}Send notification on task completed:
- {% if config.scheduler.update_subscribed_notify %} -stored notification links
-{{ config.scheduler.update_subscribed_notify|linebreaks }}
-Current notification urls: {{ config.scheduler.update_subscribed_notify }}
- {% endif %} - {{ scheduler_form.update_subscribed_notify }} -Current Download schedule: - {% if config.scheduler.download_pending %} - {% for key, value in config.scheduler.download_pending.items %} - {{ value }} - {% endfor %} + {% if download_pending %} + {{ download_pending.crontab.minute }} {{ download_pending.crontab.hour }} {{ download_pending.crontab.day_of_week }} + {% else %} False {% endif %}
Automatic video download schedule:
+ {% for error in scheduler_form.download_pending.errors %} +{{ error }}
+ {% endfor %} {{ scheduler_form.download_pending }}Send notification on task completed:
- {% if config.scheduler.download_pending_notify %} -stored notification links
-{{ config.scheduler.download_pending_notify|linebreaks }}
-Current notification urls: {{ config.scheduler.download_pending_notify }}
- {% endif %} - {{ scheduler_form.download_pending_notify }} -Current Metadata refresh schedule: - {% if config.scheduler.check_reindex %} - {% for key, value in config.scheduler.check_reindex.items %} - {{ value }} - {% endfor %} + {% if check_reindex %} + {{ check_reindex.crontab.minute }} {{ check_reindex.crontab.hour }} {{ check_reindex.crontab.day_of_week }} + {% else %} False {% endif %} @@ -94,36 +71,29 @@ {{ scheduler_form.check_reindex }}
Current refresh for metadata older than x days: {{ config.scheduler.check_reindex_days }}
+Current refresh for metadata older than x days: {{ check_reindex.task_config.days }}
Refresh older than x days, recommended 90:
+ {% for error in scheduler_form.check_reindex.errors %} +{{ error }}
+ {% endfor %} {{ scheduler_form.check_reindex_days }}Send notification on task completed:
- {% if config.scheduler.check_reindex_notify %} -stored notification links
-{{ config.scheduler.check_reindex_notify|linebreaks }}
-Current notification urls: {{ config.scheduler.check_reindex_notify }}
- {% endif %} - {{ scheduler_form.check_reindex_notify }} -Current thumbnail check schedule: - {% if config.scheduler.thumbnail_check %} - {% for key, value in config.scheduler.thumbnail_check.items %} - {{ value }} - {% endfor %} + {% if thumbnail_check %} + {{ thumbnail_check.crontab.minute }} {{ thumbnail_check.crontab.hour }} {{ thumbnail_check.crontab.day_of_week }} + {% else %} False {% endif %}
Periodically check and cleanup thumbnails:
+ {% for error in scheduler_form.thumbnail_check.errors %} +{{ error }}
+ {% endfor %} {{ scheduler_form.thumbnail_check }}Zip file backups are very slow for large archives and consistency is not guaranteed, use snapshots instead. Make sure no other tasks are running when creating a Zip file backup.
Current index backup schedule: - {% if config.scheduler.run_backup %} - {% for key, value in config.scheduler.run_backup.items %} - {{ value }} - {% endfor %} + {% if run_backup %} + {{ run_backup.crontab.minute }} {{ run_backup.crontab.hour }} {{ run_backup.crontab.day_of_week }} + {% else %} False {% endif %}
Automatically backup metadata to a zip file:
+ {% for error in scheduler_form.run_backup.errors %} +{{ error }}
+ {% endfor %} {{ scheduler_form.run_backup }}Current backup files to keep: {{ config.scheduler.run_backup_rotate }}
+Current backup files to keep: {{ run_backup.task_config.rotate }}
Max auto backups to keep:
{{ scheduler_form.run_backup_rotate }}stored notification links
++ + {{ url }} +
+ {% endfor %} + {% endfor %} +No notifications stored
+ {% endif %} +Send notification on completed tasks with the help of the Apprise library.
+ {{ notification_form.task }} + {{ notification_form.notification_url }} +