From 6cf215cd409de2bed641005d2ea1ff2a087f208e Mon Sep 17 00:00:00 2001 From: yorah Date: Sun, 2 Aug 2026 16:52:32 +0200 Subject: [PATCH] Fixed subsource being throttled for one hour when its API reports a reset delay of seconds --- .../subliminal_patch/providers/subsource.py | 119 +++++++++++++----- tests/subliminal_patch/test_subsource.py | 76 +++++++++++ 2 files changed, 165 insertions(+), 30 deletions(-) diff --git a/custom_libs/subliminal_patch/providers/subsource.py b/custom_libs/subliminal_patch/providers/subsource.py index d4acec712..55ae4995d 100644 --- a/custom_libs/subliminal_patch/providers/subsource.py +++ b/custom_libs/subliminal_patch/providers/subsource.py @@ -1,13 +1,14 @@ # -*- coding: utf-8 -*- from __future__ import annotations import logging +import math import os import time import io import datetime from typing import Set -from typing import Optional, TYPE_CHECKING +from typing import Callable, Optional, TYPE_CHECKING from babelfish import language_converters from zipfile import ZipFile, is_zipfile @@ -33,6 +34,9 @@ TITLES_EXPIRATION_TIME = datetime.timedelta(hours=6).total_seconds() QUERIES_EXPIRATION_TIME = datetime.timedelta(hours=1).total_seconds() ARCHIVES_EXPIRATION_TIME = datetime.timedelta(minutes=15).total_seconds() +# the longest the per minute window can ask us to wait, so anything above is an hour or day quota +MAX_INLINE_WAIT = int(datetime.timedelta(minutes=1).total_seconds()) + retry_amount = 3 retry_timeout = 5 @@ -169,14 +173,10 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix if season: parameters['season'] = season - results = self.retry( - lambda: self.session.get(self._server_url() + 'movies/search', params=parameters, timeout=30), - amount=retry_amount, - retry_timeout=retry_timeout + results = self.checked( + lambda: self.session.get(self._server_url() + 'movies/search', params=parameters, timeout=30) ) - self._status_raiser(results) - # deserialize results results_dict = results.json()['data'] @@ -191,13 +191,10 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix if season: parameters['season'] = season - results = self.retry( - lambda: self.session.get(self._server_url() + 'movies/search', params=parameters, timeout=30), - amount=retry_amount, - retry_timeout=retry_timeout + results = self.checked( + lambda: self.session.get(self._server_url() + 'movies/search', params=parameters, timeout=30) ) - - self._status_raiser(results) + results_dict = results.json()['data'] def get_alternative_titles(video): @@ -299,24 +296,18 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix # query the server if isinstance(self.video, Episode): parameters += (('seasonNumber', self.video.season), ('episodeNumber', self.video.episode)) - res = self.retry( + res = self.checked( lambda: self.session.get(self._server_url() + 'subtitles', params=parameters, - timeout=30), - amount=retry_amount, - retry_timeout=retry_timeout + timeout=30) ) else: - res = self.retry( + res = self.checked( lambda: self.session.get(self._server_url() + 'subtitles', params=parameters, - timeout=30), - amount=retry_amount, - retry_timeout=retry_timeout + timeout=30) ) - self._status_raiser(res) - subtitles = [] result = res.json() @@ -452,6 +443,77 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix return contributor['displayname'] return '' + @staticmethod + def _retry_after(response: Response) -> Optional[int]: + """ + Extracts the delay before the rate limit that rejected this response resets. The + server states it on every response through the X-RateLimit-* headers, and repeats + it in the body of a 429 as retryAfter. The absolute reset timestamp is preferred + as it stays accurate however long the response waited before being read. + + Several windows are enforced (per minute, hour and day) and whichever one is + currently exhausted gets reported, so no assumption is made about its length. + + :param response: A `Response` object from an HTTP request. + :type response: Response + :return: The number of seconds to wait, or None if the server didn't say. + :rtype: Optional[int] + """ + delay = None + reset_at = response.headers.get('X-RateLimit-Reset') + + if reset_at: + try: + reset_at = datetime.datetime.fromisoformat(reset_at.replace('Z', '+00:00')) + except ValueError: + logger.debug(f'Unparsable X-RateLimit-Reset value: {reset_at}') + else: + delay = (reset_at - datetime.datetime.now(datetime.timezone.utc)).total_seconds() + + if delay is None: + try: + delay = response.json().get('retryAfter') + except (ValueError, AttributeError): + return None + + if not isinstance(delay, (int, float)): + return None + + # a reset already in the past still means we are being refused, so wait a little + return max(int(math.ceil(delay)), 1) + + def checked(self, fn: Callable, is_retry: bool = False) -> Response: + """ + Executes a given callable, retrying it on transient network failures, and turns the + API-related errors it may report into exceptions. + + A rate limit that resets within the minute is waited out and the call retried once, + rather than raised: throttling the whole provider for an hour over a delay of + seconds would drop it from every remaining search. Anything longer is left to raise + so that the throttle map benches the provider instead. + + :param fn: The callable to execute, expected to return a Response object. + :type fn: Callable + :param is_retry: Indicates whether the current execution already waited out a rate + limit. Defaults to False. + :type is_retry: bool, optional + :return: The HTTP response object returned by the callable upon success. + :rtype: Response + """ + response = self.retry(fn, amount=retry_amount, retry_timeout=retry_timeout) + + if response.status_code == 429 and not is_retry: + retry_after = self._retry_after(response) + + if retry_after is not None and retry_after <= MAX_INLINE_WAIT: + logger.debug(f'Rate limit exceeded, waiting {retry_after} seconds and trying again') + time.sleep(retry_after) + return self.checked(fn, is_retry=True) + + self._status_raiser(response) + + return response + @staticmethod def _status_raiser(response: Response): """ @@ -555,8 +617,6 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix r = self._get_subtitles_archive(download_link) - self._status_raiser(r) - if not r: logger.error(f'Could not download subtitle from {download_link}') subtitle.content = None @@ -576,15 +636,14 @@ class SubsourceProvider(ProviderRetryMixin, Provider, ProviderSubtitleArchiveMix """ Fetches a subtitle archive from the given download link. The method uses caching to store the result for a defined expiration period and retries the network - request upon failure due to transient issues. + request upon failure due to transient issues. Errors are raised before the result + is cached, so a failed download is not replayed for the lifetime of the entry. :param download_link: The URL for the subtitles archive to download. :type download_link: str :return: The HTTP response object containing the subtitle archive. :rtype: Response """ - return self.retry( - lambda: self.session.get(download_link, params={'api_key': self.api_key}, timeout=30), - amount=retry_amount, - retry_timeout=retry_timeout + return self.checked( + lambda: self.session.get(download_link, params={'api_key': self.api_key}, timeout=30) ) diff --git a/tests/subliminal_patch/test_subsource.py b/tests/subliminal_patch/test_subsource.py index e6845df55..3cd2af232 100644 --- a/tests/subliminal_patch/test_subsource.py +++ b/tests/subliminal_patch/test_subsource.py @@ -1,7 +1,13 @@ # coding=utf-8 +import datetime + +import pytest from babelfish import Language +from requests import Response from subliminal_patch.converters.subsource import SubsourceConverter +from subliminal_patch.exceptions import TooManyRequests +from subliminal_patch.providers.subsource import MAX_INLINE_WAIT, SubsourceProvider def test_convert_brazilian_portuguese(): @@ -20,3 +26,73 @@ def test_convert_plain_portuguese_still_works(): def test_convert_brazilian_portuguese_round_trip(): converter = SubsourceConverter() assert converter.convert(*converter.reverse("Brazillian Portuguese")) == "Brazillian Portuguese" + + +def _rate_limited(body=b"", reset_in=None): + """A 429 as the API sends it: X-RateLimit-* headers, and retryAfter in the body.""" + response = Response() + response.status_code = 429 + response._content = body + if reset_in is not None: + reset_at = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(seconds=reset_in) + response.headers["X-RateLimit-Reset"] = reset_at.isoformat().replace("+00:00", "Z") + return response + + +def test_retry_after_prefers_the_reset_header(): + # the header is authoritative: it stays right however long the response waited + response = _rate_limited(body=b'{"retryAfter": 3}', reset_in=26) + assert SubsourceProvider._retry_after(response) == 26 + + +@pytest.mark.parametrize("body", [b"too many requests", b"", b"[1, 2]", b'{"retryAfter": "26"}']) +def test_retry_after_is_none_when_the_server_does_not_say(body): + # a 429 from the CDN in front of the API carries neither the headers nor the body + assert SubsourceProvider._retry_after(_rate_limited(body=body)) is None + + +def test_retry_after_is_never_zero(): + # the server still refusing past its own reset must not be retried with no delay at all + assert SubsourceProvider._retry_after(_rate_limited(reset_in=-5)) == 1 + + +@pytest.fixture +def _provider(monkeypatch): + """`checked` sleeps out short rate limits, which no test should actually wait for.""" + monkeypatch.setattr("subliminal_patch.providers.subsource.time.sleep", lambda seconds: None) + return SubsourceProvider(api_key="key") + + +def _responses(*responses): + """A callable standing in for a request, answering with each response in turn.""" + remaining = list(responses) + return lambda: remaining.pop(0) + + +def test_checked_waits_out_a_rate_limit_that_resets_within_the_minute(_provider): + ok = Response() + ok.status_code = 200 + fn = _responses(_rate_limited(reset_in=MAX_INLINE_WAIT), ok) + + assert _provider.checked(fn) is ok + + +def test_checked_throttles_on_a_rate_limit_that_outlasts_the_minute(_provider): + # an hour or day quota is gone, so waiting it out would stall the search for that long + fn = _responses(_rate_limited(reset_in=MAX_INLINE_WAIT + 1)) + + with pytest.raises(TooManyRequests): + _provider.checked(fn) + + +def test_checked_throttles_when_the_server_reports_no_delay(_provider): + with pytest.raises(TooManyRequests): + _provider.checked(_responses(_rate_limited())) + + +def test_checked_waits_only_once(_provider): + # a server refusing again right after its own reset is not going to yield to more waiting + fn = _responses(_rate_limited(reset_in=5), _rate_limited(reset_in=5)) + + with pytest.raises(TooManyRequests): + _provider.checked(fn)