fix: hold one connection to Lidarr open, and retry what deserves retrying #4
@@ -67,6 +67,25 @@ Losing an artist's albums does not cost their tracks, which come from a
|
||||
different endpoint with a different mapper, so matching is unaffected. A cull
|
||||
would not be, and the report says so.
|
||||
|
||||
### Talking to Lidarr
|
||||
|
||||
Indexing is two requests per artist, and more when the album fallback fires. On
|
||||
a large library that is thousands of requests in a few minutes. `urllib` opens a
|
||||
new TCP connection and performs a new DNS lookup for every one of them, which is
|
||||
enough to exhaust a container's resolver and produce `[Errno -3] Try again` on
|
||||
everything at once. The client therefore holds one connection open per host and
|
||||
resolves once.
|
||||
|
||||
Transient failures — a dropped connection, a resolver hiccup, `429`, `502`,
|
||||
`503`, `504` — are retried with a backoff. An HTTP `500` is not: it is an
|
||||
unhandled exception inside Lidarr's own serialisation and will be raised again
|
||||
identically. That distinction also decides whether a failure is worth
|
||||
investigating; a library-wide outage is not probed artist by artist, because
|
||||
doing so multiplies the load that caused it.
|
||||
|
||||
A local address is preferable to a public hostname here. It removes DNS, the
|
||||
reverse proxy and its timeouts from a path that needs none of them.
|
||||
|
||||
Tracks and files have no unfiltered endpoint — Lidarr rejects a call with no
|
||||
filter — so they stay per artist. If one artist cannot be served, that artist is
|
||||
skipped and the run continues, but the count is recorded and the coverage report
|
||||
|
||||
+152
-8
@@ -17,6 +17,8 @@ cursor, so an interrupted run resumes from what it actually has.
|
||||
|
||||
import argparse
|
||||
import fcntl
|
||||
import http.client
|
||||
import io
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
@@ -48,6 +50,11 @@ PAGE_SIZE = 200
|
||||
RETRYABLE_ERRORS = {8, 11, 16, 29}
|
||||
RETRYABLE_STATUS = {429, 500, 502, 503, 504}
|
||||
|
||||
# Lidarr's own list, and 500 is deliberately absent. A 500 from Lidarr is an
|
||||
# unhandled exception inside its serialisation, not a busy server; it will be
|
||||
# raised again identically, and retrying only delays finding that out.
|
||||
LIDARR_RETRYABLE_STATUS = {429, 502, 503, 504}
|
||||
|
||||
# Last.fm asks for no more than five requests a second averaged over five
|
||||
# minutes. A full backfill is thousands of requests, so it is worth staying
|
||||
# well inside that rather than discovering error 29 halfway through.
|
||||
@@ -207,7 +214,18 @@ class LastfmError(Exception):
|
||||
|
||||
|
||||
class LidarrError(Exception):
|
||||
"""A Lidarr request that failed."""
|
||||
"""A Lidarr request that failed.
|
||||
|
||||
`transient` separates "the network or the server had a moment" from "this
|
||||
request will fail identically forever". The distinction matters twice: only
|
||||
the first is worth retrying, and only the second is worth investigating,
|
||||
since probing a library-wide outage artist by artist multiplies the load
|
||||
that caused it.
|
||||
"""
|
||||
|
||||
def __init__(self, message, transient=False):
|
||||
super().__init__(message)
|
||||
self.transient = transient
|
||||
|
||||
|
||||
def normalise(text):
|
||||
@@ -645,35 +663,119 @@ def sync_loved(client, store, user):
|
||||
return len(rows)
|
||||
|
||||
|
||||
class KeepAlive:
|
||||
"""A transport that holds one connection open per host.
|
||||
|
||||
urllib opens a fresh TCP connection -- and performs a fresh DNS lookup --
|
||||
for every request it makes. Indexing a library is two requests per artist,
|
||||
which on a large collection is thousands of lookups inside a few minutes.
|
||||
That is enough to exhaust a container's resolver, and the failure it
|
||||
produces is `[Errno -3] Try again` on everything at once. Resolving once and
|
||||
reusing the socket removes the cause rather than papering over it, and is
|
||||
considerably faster besides.
|
||||
"""
|
||||
|
||||
def __init__(self, timeout=60):
|
||||
self.timeout = timeout
|
||||
self._connections = {}
|
||||
|
||||
def __call__(self, url, timeout=None, headers=None):
|
||||
parsed = urllib.parse.urlparse(url)
|
||||
key = (parsed.scheme, parsed.hostname, parsed.port)
|
||||
target = parsed.path + (f"?{parsed.query}" if parsed.query else "")
|
||||
request_headers = {**(headers or {}), "Accept": "application/json"}
|
||||
|
||||
# Two attempts, because a kept-alive connection the server has since
|
||||
# closed fails on use rather than announcing itself. The second attempt
|
||||
# is on a fresh socket.
|
||||
for attempt in (1, 2):
|
||||
connection = self._connections.get(key)
|
||||
if connection is None:
|
||||
connection = self._connect(parsed, timeout or self.timeout)
|
||||
self._connections[key] = connection
|
||||
try:
|
||||
connection.request("GET", target, headers=request_headers)
|
||||
response = connection.getresponse()
|
||||
body = response.read()
|
||||
except (http.client.HTTPException, OSError) as error:
|
||||
self.close(key)
|
||||
if attempt == 2:
|
||||
raise urllib.error.URLError(error) from error
|
||||
continue
|
||||
|
||||
if response.status >= 300:
|
||||
# Includes redirects: this client does not follow them, and one
|
||||
# here means the URL is pointing somewhere unintended.
|
||||
raise urllib.error.HTTPError(
|
||||
url, response.status, response.reason, response.headers, io.BytesIO(body)
|
||||
)
|
||||
return body.decode("utf-8")
|
||||
raise urllib.error.URLError("unreachable")
|
||||
|
||||
@staticmethod
|
||||
def _connect(parsed, timeout):
|
||||
if parsed.scheme == "https":
|
||||
return http.client.HTTPSConnection(parsed.hostname, parsed.port, timeout=timeout)
|
||||
return http.client.HTTPConnection(parsed.hostname, parsed.port, timeout=timeout)
|
||||
|
||||
def close(self, key=None):
|
||||
for handle in [self._connections.pop(key, None)] if key else self._connections.values():
|
||||
if handle is not None:
|
||||
handle.close()
|
||||
if key is None:
|
||||
self._connections.clear()
|
||||
|
||||
|
||||
class Lidarr:
|
||||
"""Minimal read-only Lidarr client.
|
||||
|
||||
No retries: Lidarr is on the same LAN as this, and a failure there means it
|
||||
is down or the key is wrong, neither of which improves on a second attempt.
|
||||
Retries only what is worth retrying. A dropped connection or a resolver
|
||||
hiccup is transient; an HTTP 500 out of Lidarr is an exception in its own
|
||||
serialisation and will be thrown again identically, so spending three
|
||||
attempts on it only slows down finding out.
|
||||
"""
|
||||
|
||||
def __init__(self, url, api_key, timeout=60, transport=None):
|
||||
def __init__(self, url, api_key, timeout=60, attempts=3, backoff=1.0, transport=None):
|
||||
self.root = url.rstrip("/")
|
||||
self.api_key = api_key
|
||||
self.timeout = timeout
|
||||
self.transport = transport or http_get
|
||||
self.attempts = attempts
|
||||
self.backoff = backoff
|
||||
self.transport = transport or KeepAlive(timeout)
|
||||
|
||||
def get(self, path, params=None):
|
||||
"""Return the decoded response for one API path."""
|
||||
query = urllib.parse.urlencode(params or {})
|
||||
url = f"{self.root}/api/v1/{path}" + (f"?{query}" if query else "")
|
||||
|
||||
for attempt in range(1, self.attempts + 1):
|
||||
try:
|
||||
body = self.transport(url, timeout=self.timeout, headers={"X-Api-Key": self.api_key})
|
||||
body = self.transport(
|
||||
url, timeout=self.timeout, headers={"X-Api-Key": self.api_key}
|
||||
)
|
||||
except urllib.error.HTTPError as error:
|
||||
detail = error_detail(error)
|
||||
raise LidarrError(f"GET {url}: HTTP {error.code}{': ' + detail if detail else ''}")
|
||||
message = f"GET {url}: HTTP {error.code}{': ' + detail if detail else ''}"
|
||||
if error.code not in LIDARR_RETRYABLE_STATUS:
|
||||
raise LidarrError(message)
|
||||
if attempt >= self.attempts:
|
||||
raise LidarrError(message, transient=True)
|
||||
except (urllib.error.URLError, TimeoutError) as error:
|
||||
raise LidarrError(f"GET {url}: {error}") from error
|
||||
message = f"GET {url}: {error}"
|
||||
if attempt >= self.attempts:
|
||||
raise LidarrError(message, transient=True) from error
|
||||
else:
|
||||
try:
|
||||
return json.loads(body)
|
||||
except json.JSONDecodeError as error:
|
||||
raise LidarrError(f"GET {url}: malformed response") from error
|
||||
|
||||
pause = min(self.backoff * 2 ** (attempt - 1), BACKOFF_CEILING_SECONDS)
|
||||
logger.warning("%s; retrying in %.0fs", message, pause)
|
||||
time.sleep(pause)
|
||||
|
||||
raise LidarrError(f"GET {url}: gave up after {self.attempts} attempts", transient=True)
|
||||
|
||||
|
||||
def error_detail(error, limit=300):
|
||||
"""Return whatever the server said about a failure, for the log.
|
||||
@@ -707,6 +809,36 @@ def parse_added(value):
|
||||
return None
|
||||
|
||||
|
||||
def find_bad_albums(client, artist_id, artist_name):
|
||||
"""Name the specific albums Lidarr cannot serialise for one artist.
|
||||
|
||||
Runs only once that artist's album fetch has already failed, so the extra
|
||||
requests are spent on a problem that already exists. Tracks come from an
|
||||
endpoint that still works, and their album ids give a list to probe one at a
|
||||
time; the ones that throw are the culprits. A track title from each is
|
||||
enough to recognise the album in the UI, which the id alone is not.
|
||||
"""
|
||||
try:
|
||||
tracks = client.get("track", {"artistId": artist_id})
|
||||
except LidarrError as error:
|
||||
logger.warning("could not probe %s for the offending album: %s", artist_name, error)
|
||||
return []
|
||||
|
||||
sample = {}
|
||||
for track in tracks:
|
||||
sample.setdefault(track.get("albumId"), track.get("title") or "")
|
||||
|
||||
bad = []
|
||||
for album_id, title in sorted(sample.items(), key=lambda item: item[0] or 0):
|
||||
if not album_id:
|
||||
continue
|
||||
try:
|
||||
client.get("album", {"albumIds": album_id})
|
||||
except LidarrError:
|
||||
bad.append((album_id, title))
|
||||
return bad
|
||||
|
||||
|
||||
def fetch_albums(client, artists):
|
||||
"""Return albums grouped by artist id, plus the artists whose albums failed.
|
||||
|
||||
@@ -738,6 +870,18 @@ def fetch_albums(client, artists):
|
||||
name = artist.get("artistName") or str(artist_id)
|
||||
logger.warning("could not fetch albums for %s: %s", name, error)
|
||||
failed.append(name)
|
||||
# Only worth probing a deterministic failure. When the network or
|
||||
# the resolver is the problem, every artist fails, and probing each
|
||||
# of them album by album multiplies the load that caused it.
|
||||
if not error.transient:
|
||||
for album_id, sample in find_bad_albums(client, artist_id, name):
|
||||
logger.warning(
|
||||
" album id %d is the one Lidarr cannot serialise (it holds the"
|
||||
" track %r). Open it in Lidarr and leave exactly one release"
|
||||
" monitored.",
|
||||
album_id,
|
||||
sample,
|
||||
)
|
||||
return grouped, failed
|
||||
|
||||
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
import http.server
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
|
||||
@@ -201,6 +203,19 @@ class FakeLidarr:
|
||||
url, 500, "Internal Server Error", {}, io.BytesIO(b'{"message": "boom"}')
|
||||
)
|
||||
|
||||
album_ids = query.get("albumIds")
|
||||
if path == "album" and album_ids:
|
||||
album_id = int(album_ids)
|
||||
if ("albumid", album_id) in self.fail:
|
||||
raise urllib.error.HTTPError(
|
||||
url,
|
||||
500,
|
||||
"Internal Server Error",
|
||||
{},
|
||||
io.BytesIO(b'{"message": "Sequence contains more than one element"}'),
|
||||
)
|
||||
return json.dumps([row for row in self.albums if row["id"] == album_id])
|
||||
|
||||
if path == "artist":
|
||||
return json.dumps(self.artists)
|
||||
source = {"album": self.albums, "track": self.tracks, "trackfile": self.files}[path]
|
||||
@@ -228,3 +243,53 @@ def now_playing():
|
||||
"mbid": "",
|
||||
"@attr": {"nowplaying": "true"},
|
||||
}
|
||||
|
||||
|
||||
class _CountingServer(http.server.ThreadingHTTPServer):
|
||||
"""Counts accepted connections, which is what connection reuse is about."""
|
||||
|
||||
daemon_threads = True
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.connections = 0
|
||||
super().__init__(*args, **kwargs)
|
||||
|
||||
def process_request(self, request, client_address):
|
||||
self.connections += 1
|
||||
super().process_request(request, client_address)
|
||||
|
||||
|
||||
class _Handler(http.server.BaseHTTPRequestHandler):
|
||||
# Without HTTP/1.1 the server closes after every response and no client
|
||||
# could reuse anything, which would make the test prove nothing.
|
||||
protocol_version = "HTTP/1.1"
|
||||
|
||||
def do_GET(self):
|
||||
if self.path.startswith("/boom"):
|
||||
body, status = b'{"message": "boom"}', 500
|
||||
else:
|
||||
body = json.dumps(
|
||||
{"path": self.path, "key": self.headers.get("X-Api-Key")}
|
||||
).encode()
|
||||
status = 200
|
||||
self.send_response(status)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(body)))
|
||||
self.end_headers()
|
||||
self.wfile.write(body)
|
||||
|
||||
def log_message(self, *args):
|
||||
pass
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def http_server():
|
||||
"""A real local HTTP server, for the one component that talks sockets."""
|
||||
server = _CountingServer(("127.0.0.1", 0), _Handler)
|
||||
thread = threading.Thread(target=server.serve_forever, daemon=True)
|
||||
thread.start()
|
||||
try:
|
||||
yield server, f"http://127.0.0.1:{server.server_port}"
|
||||
finally:
|
||||
server.shutdown()
|
||||
server.server_close()
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import json
|
||||
import urllib.error
|
||||
|
||||
import pytest
|
||||
@@ -584,3 +585,105 @@ def test_a_lidarr_error_carries_the_url_and_what_the_server_said():
|
||||
assert "http://lidarr:8686/api/v1/album" in message
|
||||
assert "HTTP 500" in message
|
||||
assert "boom" in message
|
||||
|
||||
|
||||
def test_the_offending_album_is_named_not_just_its_artist(tmp_path, caplog):
|
||||
"""An artist's whole discography is too much to click through by hand."""
|
||||
# AC/DC is artist 2; its only album is id 201, holding "Hells Bells".
|
||||
api = FakeLidarr(LIBRARY, fail=[("album", 0), ("album", 2), ("albumid", 201)])
|
||||
store = store_at(tmp_path)
|
||||
|
||||
with caplog.at_level("WARNING"):
|
||||
music_curator.index_library(
|
||||
music_curator.Lidarr("http://lidarr", "key", transport=api), store
|
||||
)
|
||||
|
||||
assert "album id 201" in caplog.text
|
||||
assert "Hells Bells" in caplog.text
|
||||
|
||||
|
||||
def test_a_transient_failure_is_not_probed_album_by_album(tmp_path, caplog):
|
||||
"""When the resolver is the problem every artist fails, and probing each of
|
||||
them multiplies the load that caused it."""
|
||||
|
||||
def unresolvable(url, timeout=None, headers=None):
|
||||
if "artistId" in url:
|
||||
raise urllib.error.URLError("[Errno -3] Try again")
|
||||
return json.dumps(FakeLidarr(LIBRARY).artists) if url.endswith("artist") else "[]"
|
||||
|
||||
store = store_at(tmp_path)
|
||||
client = music_curator.Lidarr("http://lidarr", "key", transport=unresolvable, backoff=0)
|
||||
|
||||
with caplog.at_level("WARNING"):
|
||||
music_curator.index_library(client, store)
|
||||
|
||||
assert "Try again" in caplog.text
|
||||
assert "cannot serialise" not in caplog.text
|
||||
|
||||
|
||||
def test_a_transient_failure_is_retried():
|
||||
attempts = []
|
||||
|
||||
def flaky(url, timeout=None, headers=None):
|
||||
attempts.append(url)
|
||||
if len(attempts) < 3:
|
||||
raise urllib.error.URLError("[Errno -3] Try again")
|
||||
return "[]"
|
||||
|
||||
client = music_curator.Lidarr("http://lidarr", "key", transport=flaky, backoff=0)
|
||||
|
||||
assert client.get("artist") == []
|
||||
assert len(attempts) == 3
|
||||
|
||||
|
||||
def test_a_lidarr_500_is_not_retried():
|
||||
"""It is an exception inside Lidarr's serialisation, not a busy server."""
|
||||
api = FakeLidarr(LIBRARY, fail=[("album", 0)])
|
||||
client = music_curator.Lidarr("http://lidarr", "key", transport=api, backoff=0)
|
||||
|
||||
with pytest.raises(music_curator.LidarrError) as raised:
|
||||
client.get("album")
|
||||
|
||||
assert raised.value.transient is False
|
||||
assert len(api.calls) == 1
|
||||
|
||||
|
||||
def test_keep_alive_uses_one_connection_for_many_requests(http_server):
|
||||
"""The point of the whole class: one DNS lookup and one socket, not N."""
|
||||
server, base = http_server
|
||||
transport = music_curator.KeepAlive()
|
||||
|
||||
try:
|
||||
for index in range(5):
|
||||
body = transport(f"{base}/api/v1/artist?n={index}", headers={"X-Api-Key": "key"})
|
||||
assert json.loads(body)["key"] == "key"
|
||||
finally:
|
||||
transport.close()
|
||||
|
||||
assert server.connections == 1
|
||||
|
||||
|
||||
def test_keep_alive_maps_an_error_status_onto_httperror(http_server):
|
||||
_, base = http_server
|
||||
transport = music_curator.KeepAlive()
|
||||
|
||||
try:
|
||||
with pytest.raises(urllib.error.HTTPError) as raised:
|
||||
transport(f"{base}/boom", headers={"X-Api-Key": "key"})
|
||||
assert raised.value.code == 500
|
||||
assert music_curator.error_detail(raised.value) == "boom"
|
||||
finally:
|
||||
transport.close()
|
||||
|
||||
|
||||
def test_lidarr_talks_to_a_real_server_through_keep_alive(http_server):
|
||||
server, base = http_server
|
||||
client = music_curator.Lidarr(base, "secret")
|
||||
|
||||
try:
|
||||
assert client.get("artist", {"x": 1})["key"] == "secret"
|
||||
assert client.get("album")["path"] == "/api/v1/album"
|
||||
finally:
|
||||
client.transport.close()
|
||||
|
||||
assert server.connections == 1
|
||||
|
||||
Reference in New Issue
Block a user