From b2337669ca30a8d8f644b43178bffe4331447a6f Mon Sep 17 00:00:00 2001 From: Simon Date: Sun, 16 Feb 2025 11:03:50 +0700 Subject: [PATCH] fix channel tags empty list serializers and migration --- backend/channel/src/index.py | 2 +- .../config/management/commands/ta_startup.py | 276 +++++++++++------- 2 files changed, 172 insertions(+), 106 deletions(-) diff --git a/backend/channel/src/index.py b/backend/channel/src/index.py index 35422611..b2aa688c 100644 --- a/backend/channel/src/index.py +++ b/backend/channel/src/index.py @@ -115,7 +115,7 @@ class YoutubeChannel(YouTubeItem): "channel_tvart_url": False, "channel_id": self.youtube_id, "channel_subscribed": False, - "channel_tags": False, + "channel_tags": [], "channel_description": False, "channel_thumb_url": False, "channel_views": 0, diff --git a/backend/config/management/commands/ta_startup.py b/backend/config/management/commands/ta_startup.py index e2d0991f..ec6d2ffc 100644 --- a/backend/config/management/commands/ta_startup.py +++ b/backend/config/management/commands/ta_startup.py @@ -42,7 +42,6 @@ class Command(BaseCommand): def handle(self, *args, **options): """run all commands""" self.stdout.write(TOPIC) - self._mig_app_settings() self._make_folders() self._clear_redis_keys() self._clear_tasks() @@ -50,9 +49,113 @@ class Command(BaseCommand): self._version_check() self._index_setup() self._snapshot_check() + self._mig_app_settings() self._create_default_schedules() self._update_schedule_tz() self._init_app_config() + self._mig_channel_tags() + self._mig_video_channel_tags() + + def _make_folders(self): + """make expected cache folders""" + self.stdout.write("[1] create expected cache folders") + folders = [ + "backup", + "channels", + "download", + "import", + "playlists", + "videos", + ] + cache_dir = EnvironmentSettings.CACHE_DIR + for folder in folders: + folder_path = os.path.join(cache_dir, folder) + os.makedirs(folder_path, exist_ok=True) + + self.stdout.write(self.style.SUCCESS(" ✓ expected folders created")) + + def _clear_redis_keys(self): + """make sure there are no leftover locks or keys set in redis""" + self.stdout.write("[2] clear leftover keys in redis") + all_keys = [ + "dl_queue_id", + "dl_queue", + "downloading", + "manual_import", + "reindex", + "rescan", + "run_backup", + "startup_check", + "reindex:ta_video", + "reindex:ta_channel", + "reindex:ta_playlist", + ] + + redis_con = RedisArchivist() + has_changed = False + for key in all_keys: + if redis_con.del_message(key): + self.stdout.write( + self.style.SUCCESS(f" ✓ cleared key {key}") + ) + has_changed = True + + if not has_changed: + self.stdout.write(self.style.SUCCESS(" no keys found")) + + def _clear_tasks(self): + """clear tasks and messages""" + self.stdout.write("[3] clear task leftovers") + TaskManager().fail_pending() + redis_con = RedisArchivist() + to_delete = redis_con.list_keys("message:") + if to_delete: + for key in to_delete: + redis_con.del_message(key) + + self.stdout.write( + self.style.SUCCESS(f" ✓ cleared {len(to_delete)} messages") + ) + + def _clear_dl_cache(self): + """clear leftover files from dl cache""" + self.stdout.write("[4] clear leftover files from dl cache") + leftover_files = clear_dl_cache(EnvironmentSettings.CACHE_DIR) + if leftover_files: + self.stdout.write( + self.style.SUCCESS(f" ✓ cleared {leftover_files} files") + ) + else: + self.stdout.write(self.style.SUCCESS(" no files found")) + + def _version_check(self): + """remove new release key if updated now""" + self.stdout.write("[5] check for first run after update") + new_version = ReleaseVersion().is_updated() + if new_version: + self.stdout.write( + self.style.SUCCESS(f" ✓ update to {new_version} completed") + ) + else: + self.stdout.write(self.style.SUCCESS(" no new update found")) + + version_task = CustomPeriodicTask.objects.filter(name="version_check") + if not version_task.exists(): + return + + if not version_task.first().last_run_at: + self.style.SUCCESS(" ✓ send initial version check task") + version_check.delay() + + def _index_setup(self): + """migration: validate index mappings""" + self.stdout.write("[6] validate index mappings") + ElasitIndexWrap().setup() + + def _snapshot_check(self): + """migration setup snapshots""" + self.stdout.write("[7] setup snapshots") + ElasticSnapshot().setup() def _mig_app_settings(self) -> None: """update from v0.4.13 to v0.5.0, migrate application settings""" @@ -87,110 +190,9 @@ class Command(BaseCommand): sleep(60) raise CommandError(message) - def _make_folders(self): - """make expected cache folders""" - self.stdout.write("[2] create expected cache folders") - folders = [ - "backup", - "channels", - "download", - "import", - "playlists", - "videos", - ] - cache_dir = EnvironmentSettings.CACHE_DIR - for folder in folders: - folder_path = os.path.join(cache_dir, folder) - os.makedirs(folder_path, exist_ok=True) - - self.stdout.write(self.style.SUCCESS(" ✓ expected folders created")) - - def _clear_redis_keys(self): - """make sure there are no leftover locks or keys set in redis""" - self.stdout.write("[3] clear leftover keys in redis") - all_keys = [ - "dl_queue_id", - "dl_queue", - "downloading", - "manual_import", - "reindex", - "rescan", - "run_backup", - "startup_check", - "reindex:ta_video", - "reindex:ta_channel", - "reindex:ta_playlist", - ] - - redis_con = RedisArchivist() - has_changed = False - for key in all_keys: - if redis_con.del_message(key): - self.stdout.write( - self.style.SUCCESS(f" ✓ cleared key {key}") - ) - has_changed = True - - if not has_changed: - self.stdout.write(self.style.SUCCESS(" no keys found")) - - def _clear_tasks(self): - """clear tasks and messages""" - self.stdout.write("[4] clear task leftovers") - TaskManager().fail_pending() - redis_con = RedisArchivist() - to_delete = redis_con.list_keys("message:") - if to_delete: - for key in to_delete: - redis_con.del_message(key) - - self.stdout.write( - self.style.SUCCESS(f" ✓ cleared {len(to_delete)} messages") - ) - - def _clear_dl_cache(self): - """clear leftover files from dl cache""" - self.stdout.write("[5] clear leftover files from dl cache") - leftover_files = clear_dl_cache(EnvironmentSettings.CACHE_DIR) - if leftover_files: - self.stdout.write( - self.style.SUCCESS(f" ✓ cleared {leftover_files} files") - ) - else: - self.stdout.write(self.style.SUCCESS(" no files found")) - - def _version_check(self): - """remove new release key if updated now""" - self.stdout.write("[6] check for first run after update") - new_version = ReleaseVersion().is_updated() - if new_version: - self.stdout.write( - self.style.SUCCESS(f" ✓ update to {new_version} completed") - ) - else: - self.stdout.write(self.style.SUCCESS(" no new update found")) - - version_task = CustomPeriodicTask.objects.filter(name="version_check") - if not version_task.exists(): - return - - if not version_task.first().last_run_at: - self.style.SUCCESS(" ✓ send initial version check task") - version_check.delay() - - def _index_setup(self): - """migration: validate index mappings""" - self.stdout.write("[7] validate index mappings") - ElasitIndexWrap().setup() - - def _snapshot_check(self): - """migration setup snapshots""" - self.stdout.write("[8] setup snapshots") - ElasticSnapshot().setup() - def _create_default_schedules(self) -> None: """create default schedules for new installations""" - self.stdout.write("[9] create initial schedules") + self.stdout.write("[8] create initial schedules") init_has_run = CustomPeriodicTask.objects.filter( name="version_check" ).exists() @@ -241,7 +243,7 @@ class Command(BaseCommand): def _update_schedule_tz(self) -> None: """update timezone for Schedule instances""" - self.stdout.write("[10] validate schedules TZ") + self.stdout.write("[9] validate schedules TZ") tz = EnvironmentSettings.TZ to_update = CrontabSchedule.objects.exclude(timezone=tz) @@ -259,7 +261,7 @@ class Command(BaseCommand): def _init_app_config(self) -> None: """init default app config to ES""" - self.stdout.write("[11] Check AppConfig") + self.stdout.write("[10] Check AppConfig") try: _ = AppConfig().config self.stdout.write( @@ -280,3 +282,67 @@ class Command(BaseCommand): self.stdout.write( self.style.SUCCESS(f" Status code: {status_code}") ) + + def _mig_channel_tags(self) -> None: + """update from v0.4.13 to v0.5.0, migrate incorrect data types""" + self.stdout.write("[MIGRATION] fix incorrect channel tags types") + path = "ta_channel/_update_by_query" + data = { + "query": {"match": {"channel_tags": False}}, + "script": { + "source": "ctx._source.channel_tags = []", + "lang": "painless", + }, + } + response, status_code = ElasticWrap(path).post(data) + if status_code in [200, 201]: + updated = response.get("updated") + if updated: + self.stdout.write( + self.style.SUCCESS(f" ✓ fixed {updated} channel tags") + ) + else: + self.stdout.write( + self.style.SUCCESS(" no channel tags needed fixing") + ) + return + + message = " 🗙 failed to fix channel tags" + self.stdout.write(self.style.ERROR(message)) + self.stdout.write(response) + sleep(60) + raise CommandError(message) + + def _mig_video_channel_tags(self) -> None: + """update from v0.4.13 to v0.5.0, migrate incorrect data types""" + self.stdout.write("[MIGRATION] fix incorrect video channel tags types") + path = "ta_video/_update_by_query" + data = { + "query": {"match": {"channel.channel_tags": False}}, + "script": { + "source": "ctx._source.channel.channel_tags = []", + "lang": "painless", + }, + } + response, status_code = ElasticWrap(path).post(data) + if status_code in [200, 201]: + updated = response.get("updated") + if updated: + self.stdout.write( + self.style.SUCCESS( + f" ✓ fixed {updated} video channel tags" + ) + ) + else: + self.stdout.write( + self.style.SUCCESS( + " no video channel tags needed fixing" + ) + ) + return + + message = " 🗙 failed to fix video channel tags" + self.stdout.write(self.style.ERROR(message)) + self.stdout.write(response) + sleep(60) + raise CommandError(message)