Resolve Race Condition and Optimize Session Switching for Atmos Downloads
This commit is contained in:
+58
-23
@@ -2,10 +2,9 @@ import json
|
|||||||
import os
|
import os
|
||||||
import shutil
|
import shutil
|
||||||
from collections.abc import Callable
|
from collections.abc import Callable
|
||||||
from contextlib import contextmanager
|
|
||||||
from json import JSONDecodeError
|
from json import JSONDecodeError
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from threading import Event
|
from threading import Event, Lock
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import tidalapi
|
import tidalapi
|
||||||
@@ -103,6 +102,19 @@ class Tidal(BaseConfig, metaclass=SingletonMeta):
|
|||||||
self.cls_model = ModelToken
|
self.cls_model = ModelToken
|
||||||
tidal_config: tidalapi.Config = tidalapi.Config(item_limit=10000)
|
tidal_config: tidalapi.Config = tidalapi.Config(item_limit=10000)
|
||||||
self.session = tidalapi.Session(tidal_config)
|
self.session = tidalapi.Session(tidal_config)
|
||||||
|
self.original_client_id = self.session.config.client_id
|
||||||
|
self.original_client_secret = self.session.config.client_secret
|
||||||
|
# Lock to ensure session-switching is thread-safe.
|
||||||
|
# This lock protects against a race condition where one thread
|
||||||
|
# changes the session credentials while another is using them.
|
||||||
|
# It is intentionally held by Download._get_stream_info
|
||||||
|
# for the *entire* duration of the credential switch AND
|
||||||
|
# the get_stream() call.
|
||||||
|
self.stream_lock = Lock()
|
||||||
|
# State-tracking flag to prevent redundant, expensive
|
||||||
|
# session re-authentication when the session is already in the
|
||||||
|
# correct mode (Atmos or Normal).
|
||||||
|
self.is_atmos_session = False
|
||||||
# self.session.config.client_id = "km8T1xS355y7dd3H"
|
# self.session.config.client_id = "km8T1xS355y7dd3H"
|
||||||
# self.session.config.client_secret = "vcmeGW1OuZ0fWYMCSZ6vNvSLJlT3XEpW0ambgYt5ZuI="
|
# self.session.config.client_secret = "vcmeGW1OuZ0fWYMCSZ6vNvSLJlT3XEpW0ambgYt5ZuI="
|
||||||
self.file_path = path_file_token()
|
self.file_path = path_file_token()
|
||||||
@@ -116,7 +128,8 @@ class Tidal(BaseConfig, metaclass=SingletonMeta):
|
|||||||
if settings:
|
if settings:
|
||||||
self.settings = settings
|
self.settings = settings
|
||||||
|
|
||||||
self.session.audio_quality = tidalapi.Quality(self.settings.data.quality_audio)
|
if not self.is_atmos_session:
|
||||||
|
self.session.audio_quality = tidalapi.Quality(self.settings.data.quality_audio)
|
||||||
self.session.video_quality = tidalapi.VideoQuality.high
|
self.session.video_quality = tidalapi.VideoQuality.high
|
||||||
|
|
||||||
return True
|
return True
|
||||||
@@ -162,33 +175,55 @@ class Tidal(BaseConfig, metaclass=SingletonMeta):
|
|||||||
self.set_option("expiry_time", self.session.expiry_time)
|
self.set_option("expiry_time", self.session.expiry_time)
|
||||||
self.save()
|
self.save()
|
||||||
|
|
||||||
@contextmanager
|
def switch_to_atmos_session(self) -> bool:
|
||||||
def atmos_session_context(self):
|
"""
|
||||||
|
Switches the shared session to Dolby Atmos credentials.
|
||||||
|
Only re-authenticates if not already in Atmos mode.
|
||||||
|
"""
|
||||||
|
# If we are already in Atmos mode, do nothing.
|
||||||
|
if self.is_atmos_session:
|
||||||
|
return True
|
||||||
|
|
||||||
if not self.session.check_login():
|
print("Switching session context to Dolby Atmos...")
|
||||||
print("Not logged in.")
|
self.session.config.client_id = ATMOS_CLIENT_ID
|
||||||
|
self.session.config.client_secret = ATMOS_CLIENT_SECRET
|
||||||
|
self.session.audio_quality = ATMOS_REQUEST_QUALITY
|
||||||
|
|
||||||
original_client_id = self.session.config.client_id
|
# Re-login with new credentials
|
||||||
original_client_secret = self.session.config.client_secret
|
if not self.login_token(do_pkce=self.is_pkce):
|
||||||
original_audio_quality = self.session.audio_quality
|
print("Warning: Atmos session authentication failed.")
|
||||||
|
# Try to switch back to normal to be safe
|
||||||
|
self.restore_normal_session(force=True)
|
||||||
|
return False
|
||||||
|
|
||||||
try:
|
self.is_atmos_session = True # Set the flag
|
||||||
self.session.config.client_id = ATMOS_CLIENT_ID
|
print("Session is now in Atmos mode.")
|
||||||
self.session.config.client_secret = ATMOS_CLIENT_SECRET
|
return True
|
||||||
self.session.audio_quality = ATMOS_REQUEST_QUALITY
|
|
||||||
|
|
||||||
if not self.login_token(do_pkce=self.is_pkce):
|
def restore_normal_session(self, force: bool = False) -> bool:
|
||||||
print("Warning: Session restore failed.")
|
"""
|
||||||
|
Restores the shared session to the original user credentials.
|
||||||
|
Only re-authenticates if not already in Normal mode.
|
||||||
|
"""
|
||||||
|
# If we are already in Normal mode (and not forced), do nothing.
|
||||||
|
if not self.is_atmos_session and not force:
|
||||||
|
return True
|
||||||
|
|
||||||
yield
|
print("Restoring session context to Normal...")
|
||||||
|
self.session.config.client_id = self.original_client_id
|
||||||
|
self.session.config.client_secret = self.original_client_secret
|
||||||
|
|
||||||
finally:
|
# Re-apply user's quality setting
|
||||||
self.session.config.client_id = original_client_id
|
self.settings_apply()
|
||||||
self.session.config.client_secret = original_client_secret
|
|
||||||
self.session.audio_quality = original_audio_quality
|
|
||||||
|
|
||||||
if not self.login_token(do_pkce=self.is_pkce):
|
# Re-login with original credentials
|
||||||
print("Warning: Restoring the original session context failed. Please restart the application.")
|
if not self.login_token(do_pkce=self.is_pkce):
|
||||||
|
print("Warning: Restoring the original session context failed. Please restart the application.")
|
||||||
|
return False
|
||||||
|
|
||||||
|
self.is_atmos_session = False # Set the flag
|
||||||
|
print("Session is now in Normal mode.")
|
||||||
|
return True
|
||||||
|
|
||||||
def login(self, fn_print: Callable) -> bool:
|
def login(self, fn_print: Callable) -> bool:
|
||||||
is_token = self.login_token()
|
is_token = self.login_token()
|
||||||
|
|||||||
+76
-22
@@ -794,43 +794,97 @@ class Download:
|
|||||||
stream_manifest: StreamManifest | None = None
|
stream_manifest: StreamManifest | None = None
|
||||||
media_stream: Stream | None = None
|
media_stream: Stream | None = None
|
||||||
do_flac_extract: bool = False
|
do_flac_extract: bool = False
|
||||||
|
file_extension: str = ""
|
||||||
|
|
||||||
if isinstance(media, Track):
|
# CRITICAL: This lock is intentionally broad and serializes all
|
||||||
|
# stream-fetching (Phase 1) to prevent a critical race condition.
|
||||||
|
#
|
||||||
|
# THE PROBLEM:
|
||||||
|
# The single, shared session (self.tidal.session) must change its
|
||||||
|
# credentials to switch between Atmos and Hi-Res/Normal streams.
|
||||||
|
#
|
||||||
|
# THE RACE CONDITION IT FIXES:
|
||||||
|
# If this lock is released *before* get_stream() is called,
|
||||||
|
# another thread could change the session (e.g., back to "Normal")
|
||||||
|
# right after this thread switched it to "Atmos". This would
|
||||||
|
# cause this thread to call get_stream() with the wrong credentials,
|
||||||
|
# resulting in the API returning AAC 320 instead of Atmos.
|
||||||
|
#
|
||||||
|
# THE TRADEOFF:
|
||||||
|
# This creates a "tollbooth" bottleneck, serializing the get_stream()
|
||||||
|
# calls. However, the *actual* segment downloads (Phase 2)
|
||||||
|
# still run in parallel, governed by `downloads_concurrent_max`.
|
||||||
|
#
|
||||||
|
# DO NOT "OPTIMIZE" THIS by making the lock more granular.
|
||||||
|
# Correctness > Performance.
|
||||||
|
|
||||||
|
with self.tidal.stream_lock:
|
||||||
try:
|
try:
|
||||||
if (
|
if isinstance(media, Track):
|
||||||
self.settings.data.download_dolby_atmos
|
stream_manifest, file_extension, do_flac_extract, media_stream = self._get_track_stream_info(media)
|
||||||
and hasattr(media, "audio_modes")
|
if stream_manifest is None:
|
||||||
and AudioMode.dolby_atmos.value in media.audio_modes
|
return None, "", False, None
|
||||||
):
|
|
||||||
with self.tidal.atmos_session_context():
|
|
||||||
atmos_track = self.session.track(media.id)
|
|
||||||
media_stream = atmos_track.get_stream()
|
|
||||||
else:
|
|
||||||
media_stream = media.get_stream()
|
|
||||||
|
|
||||||
stream_manifest = media_stream.get_stream_manifest()
|
elif isinstance(media, Video):
|
||||||
|
# Videos always require the normal session
|
||||||
|
if not self.tidal.restore_normal_session():
|
||||||
|
self.fn_logger.error(f"Failed to restore normal session for video: {media.id}")
|
||||||
|
return None, "", False, None
|
||||||
|
|
||||||
|
file_extension = AudioExtensions.MP4 if self.settings.data.video_convert_mp4 else VideoExtensions.TS
|
||||||
|
|
||||||
|
stream_manifest = None
|
||||||
|
media_stream = None
|
||||||
|
do_flac_extract = False
|
||||||
|
|
||||||
|
else:
|
||||||
|
self.fn_logger.error(f"Unknown media type for stream info: {type(media)}")
|
||||||
|
return None, "", False, None
|
||||||
|
|
||||||
except TooManyRequests:
|
except TooManyRequests:
|
||||||
self.fn_logger.exception(
|
self.fn_logger.exception(
|
||||||
f"Too many requests against TIDAL backend. Skipping '{name_builder_item(media)}'. "
|
f"Too many requests against TIDAL backend. Skipping '{name_builder_item(media)}'. "
|
||||||
f"Consider to activate delay between downloads."
|
f"Consider to activate delay between downloads."
|
||||||
)
|
)
|
||||||
|
|
||||||
return None, "", False, None
|
return None, "", False, None
|
||||||
|
|
||||||
except Exception:
|
except Exception:
|
||||||
self.fn_logger.exception(f"Something went wrong. Skipping '{name_builder_item(media)}'.")
|
self.fn_logger.exception(f"Something went wrong. Skipping '{name_builder_item(media)}'.")
|
||||||
|
|
||||||
return None, "", False, None
|
return None, "", False, None
|
||||||
|
|
||||||
file_extension = stream_manifest.file_extension
|
return stream_manifest, file_extension, do_flac_extract, media_stream
|
||||||
|
|
||||||
if self.settings.data.extract_flac and (
|
def _get_track_stream_info(self, media: Track) -> tuple[StreamManifest | None, str, bool, Stream | None]:
|
||||||
stream_manifest.codecs.upper() == Codec.FLAC and file_extension != AudioExtensions.FLAC
|
"""
|
||||||
):
|
Gets stream info for a Track, handling Atmos/Normal session switching.
|
||||||
file_extension = AudioExtensions.FLAC
|
This is a helper for _get_stream_info to reduce complexity.
|
||||||
do_flac_extract = True
|
"""
|
||||||
elif isinstance(media, Video):
|
want_atmos = (
|
||||||
file_extension = AudioExtensions.MP4 if self.settings.data.video_convert_mp4 else VideoExtensions.TS
|
self.settings.data.download_dolby_atmos
|
||||||
|
and hasattr(media, "audio_modes")
|
||||||
|
and AudioMode.dolby_atmos.value in media.audio_modes
|
||||||
|
)
|
||||||
|
|
||||||
|
if want_atmos:
|
||||||
|
if not self.tidal.switch_to_atmos_session():
|
||||||
|
self.fn_logger.error(f"Failed to switch to Atmos session for track: {media.id}")
|
||||||
|
return None, "", False, None
|
||||||
|
else:
|
||||||
|
if not self.tidal.restore_normal_session():
|
||||||
|
self.fn_logger.error(f"Failed to restore normal session for track: {media.id}")
|
||||||
|
return None, "", False, None
|
||||||
|
|
||||||
|
media_stream = self.session.track(media.id).get_stream() if want_atmos else media.get_stream()
|
||||||
|
|
||||||
|
stream_manifest = media_stream.get_stream_manifest()
|
||||||
|
file_extension = stream_manifest.file_extension
|
||||||
|
do_flac_extract = False
|
||||||
|
|
||||||
|
if self.settings.data.extract_flac and (
|
||||||
|
stream_manifest.codecs.upper() == Codec.FLAC and file_extension != AudioExtensions.FLAC
|
||||||
|
):
|
||||||
|
file_extension = AudioExtensions.FLAC
|
||||||
|
do_flac_extract = True
|
||||||
|
|
||||||
return stream_manifest, file_extension, do_flac_extract, media_stream
|
return stream_manifest, file_extension, do_flac_extract, media_stream
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user