-
Notifications
You must be signed in to change notification settings - Fork 7
Expand file tree
/
Copy pathmain.py
More file actions
8437 lines (7695 loc) · 410 KB
/
Copy pathmain.py
File metadata and controls
8437 lines (7695 loc) · 410 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#main.py
import asyncio
import dataclasses
import hashlib
from html import escape as _html_escape
import urllib.parse
import hmac
import io
import logging
import os
import re
import time
import httpx
import numpy as np
from datetime import datetime, timedelta, timezone
import zoneinfo
from collections import OrderedDict
from concurrent.futures import ThreadPoolExecutor
from contextlib import asynccontextmanager, suppress
from dataclasses import dataclass, field
from functools import lru_cache
from urllib.parse import parse_qsl, urlencode
from fastapi import FastAPI, HTTPException, Request
from fastapi.responses import JSONResponse, Response, HTMLResponse
from fastapi.staticfiles import StaticFiles
from PIL import Image, ImageDraw, ImageFilter, ImageFont, ImageOps
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
force=True,
)
# ElfHosted fork: optional structured JSON logs (LOG_FORMAT=json) for shipping
# to Loki/Elasticsearch. Read straight from the env so logging is configured
# before config.py is imported. Falls back to the text format above if
# python-json-logger isn't installed. Default (text) is upstream behaviour.
if os.environ.get("LOG_FORMAT", "text").strip().lower() == "json":
try:
from pythonjsonlogger.json import JsonFormatter
except Exception:
try:
from pythonjsonlogger.jsonlogger import JsonFormatter # older versions
except Exception:
JsonFormatter = None
if JsonFormatter is not None:
_json_fmt = JsonFormatter(
"%(asctime)s %(levelname)s %(name)s %(message)s",
datefmt="%Y-%m-%dT%H:%M:%S%z",
)
for _h in logging.getLogger().handlers:
_h.setFormatter(_json_fmt)
else:
logging.getLogger(__name__).warning(
"LOG_FORMAT=json but python-json-logger is not installed — "
"falling back to text logs."
)
# Pull uvicorn's loggers into our root handler so all output shares the same format.
for _uv_name in ("uvicorn", "uvicorn.error", "uvicorn.access"):
_uv_logger = logging.getLogger(_uv_name)
_uv_logger.handlers = []
_uv_logger.propagate = True
class _TruncateUrlFilter(logging.Filter):
"""
Redact API keys and truncate long URL paths in log records.
Two responsibilities:
1. For uvicorn.access records, truncate the request path so long URLs
don't fill the log.
2. For ALL records, redact every common API-key query parameter pattern
in both record.msg and record.args. This catches keys that slip
through when an httpx exception is logged (its __str__ includes the
full upstream URL with our outbound api_key=) as well as anything
else that might inadvertently include a key.
"""
_MAX = 80
# Match query params we hold (tmdb_key, mdblist_key, access_key) AND the
# upstream parameter names we forward keys under (api_key, apikey).
_KEY_RE = re.compile(
r'((?:tmdb_key|mdblist_key|access_key|api_key|apikey)=)[^&\s\'\"]*',
re.IGNORECASE,
)
@classmethod
def _redact(cls, value):
if isinstance(value, str):
return cls._KEY_RE.sub(r'\1***', value)
return value
def filter(self, record: logging.LogRecord) -> bool:
# uvicorn.access records: args = (client_addr, method, path, http_version, status_code, ...)
if (
record.name == "uvicorn.access"
and isinstance(record.args, tuple)
and len(record.args) >= 3
):
path = record.args[2]
if isinstance(path, str):
path = self._KEY_RE.sub(r'\1***', path)
if len(path) > self._MAX:
path = path[: self._MAX] + "…"
record.args = (record.args[0], record.args[1], path) + record.args[3:]
# Generic redaction for every other record (application logs).
# We redact in msg and args so the formatted output is safe regardless
# of whether the record uses % substitution or pre-formatted strings.
if isinstance(record.msg, str):
record.msg = self._redact(record.msg)
if isinstance(record.args, tuple):
record.args = tuple(self._redact(a) for a in record.args)
elif isinstance(record.args, dict):
record.args = {k: self._redact(v) for k, v in record.args.items()}
# Tracebacks (logger.exception / exc_info=True) are formatted lazily
# by the handler. Pre-format and redact exc_text here so the
# downstream formatter uses our sanitised copy rather than re-rendering.
if record.exc_info and not record.exc_text:
import traceback
record.exc_text = self._redact(
"".join(traceback.format_exception(*record.exc_info))
)
elif record.exc_text:
record.exc_text = self._redact(record.exc_text)
return True
# Attach to the root handler, not the root logger — propagation calls
# callHandlers() directly on parent loggers, skipping their logger-level filters.
_url_filter = _TruncateUrlFilter()
for _handler in logging.getLogger().handlers:
_handler.addFilter(_url_filter)
# httpx logs every outbound HTTP request at INFO level, including full URLs with
# API keys in query strings. Raise its level to WARNING so those lines are never
# written to the log — our own try/except blocks capture errors explicitly.
logging.getLogger("httpx").setLevel(logging.WARNING)
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Request coalescing
# ---------------------------------------------------------------------------
# Maps final_cache_key -> Future[bytes] for in-flight renders.
# When multiple requests arrive simultaneously for the same uncached poster
# (common during a burst from AIOMetadata loading a library), only the first
# runs the full pipeline; the rest await its Future and get the result for free.
# This dict is per-worker-process — cross-process deduplication would require
# a shared store like Redis, but intra-process coalescing handles the common
# burst pattern well enough at this scale.
# The future carries (jpeg_bytes, provisional): a coalesced request has to know
# whether the render it is riding on was one the pipeline refused to persist, or
# it would hand out a validator for a poster the server never committed to.
_render_inflight: dict[str, "asyncio.Future[tuple[bytes, bool]]"] = {}
# Coalesces concurrent fetch_poster_metadata calls for the same (tmdb_id,
# media_type, language) tuple. Without this, simultaneous /poster + /logo
# requests for the same cold title each fire their own TMDB API call.
_metadata_inflight: dict[str, "asyncio.Future[tuple]"] = {}
async def _coalesced_fetch_poster_metadata(
client: "httpx.AsyncClient",
tmdb_id: str,
tmdb_key: str,
media_type: str,
lang: str,
secondary_lang: str = "",
) -> tuple:
endpoint = "tv" if media_type in ("tv", "series") else "movie"
inflight_key = tmdb_metadata_cache_key(endpoint, tmdb_id, lang, secondary_lang)
existing = _metadata_inflight.get(inflight_key)
if existing is not None:
logger.debug(f"Coalescing metadata fetch for {media_type}/{tmdb_id} ({lang})")
return await existing
fut: "asyncio.Future[tuple]" = asyncio.get_running_loop().create_future()
fut.add_done_callback(
lambda f: f.exception() if not f.cancelled() and f.exception() else None
)
_metadata_inflight[inflight_key] = fut
try:
result = await fetch_poster_metadata(
client, tmdb_id, tmdb_key, media_type, lang, secondary_lang
)
fut.set_result(result)
return result
except Exception as exc:
if not fut.done():
fut.set_exception(exc)
raise
except BaseException:
if not fut.done():
fut.cancel()
raise
finally:
_metadata_inflight.pop(inflight_key, None)
# ---------------------------------------------------------------------------
# Background quality fetching
# ---------------------------------------------------------------------------
# Quality data (AIOStreams / scrapers) is fetched in the background so poster
# responses are never blocked by a slow scraper call. The poster is served
# immediately without quality badges on a cache miss; the next request for the
# same title will find the quality cached and render badges normally.
#
# _quality_bg_inflight: tracks imdb_ids with an active background fetch so
# scroll bursts don't launch duplicate fetches for the same title.
# _quality_bg_semaphore: caps concurrent AIOStreams calls so a large burst
# doesn't hammer the scrapers with hundreds of simultaneous requests.
_quality_bg_inflight: set[str] = set()
_quality_bg_semaphore: "asyncio.Semaphore | None" = None # created inside event loop
_quality_source_backoff_until: dict[str, float] = {}
_quality_source_fail_count: dict[str, int] = {}
# ---------------------------------------------------------------------------
# Rating fetch deduplication
# ---------------------------------------------------------------------------
# Prevents concurrent requests for the same imdb_id (different raw_params /
# final_cache_key) from triggering duplicate MDBlist API calls. The most
# common burst: AIOMetadata requests many posters simultaneously; several
# share an uncached title with different user-config hashes so render
# coalescing alone doesn't protect them.
#
# _rating_fetch_inflight: maps imdb_id -> asyncio.Event that fires once the
# first fetch completes. Subsequent requests wait, then re-read the DB.
# _rating_backoff: maps (imdb_id, API key) -> loop-time after which a new
# attempt is allowed. Scoping by key lets a rotated or replaced key retry
# the same title immediately. Network failures use an escalating ladder
# (30s/2m/8m/1h); rate-limit responses use Retry-After, else the key's
# daily-quota reset (X-RateLimit-Reset), else 1h flat.
_rating_fetch_inflight: dict[str, asyncio.Event] = {}
_rating_backoff: dict[tuple[str, str], float] = {}
_rating_fail_count: dict[tuple[str, str], int] = {}
_mdblist_semaphore: "asyncio.Semaphore | None" = None # caps concurrent MDBlist HTTP calls; created inside event loop
# Caps parallel burned-in-text scans. Each slot owns an independent RapidOCR
# session in a dedicated executor, so cold-cache OCR cannot occupy render workers.
# Created inside the event loop.
_detect_semaphore: "asyncio.Semaphore | None" = None
_detect_executor: "ThreadPoolExecutor | None" = None
# Maps immutable image/detector keys to active OCR tasks. Different poster
# configurations often render the same source image during a burst; they should
# share one scan even when their final composite cache keys differ.
_text_detection_inflight: dict[str, "asyncio.Task[bool | None]"] = {}
_foreground_detection_count = 0
_active_poster_renders = 0
# Admission control for fresh renders (POSTER_RENDER_CONCURRENCY). Cache hits and
# coalesced waiters bypass it; see _get_render_semaphore(). Created inside the
# event loop. _renders_queued counts requests parked on it, for /stats.
_render_semaphore: "asyncio.Semaphore | None" = None
_renders_queued = 0
_background_detection_queue: "asyncio.Queue[_DeferredTextDetection] | None" = None
_background_detection_keys: set[str] = set()
_background_detection_task: "asyncio.Task[None] | None" = None
@dataclass(frozen=True)
class _DeferredTextDetection:
cache_key: str
image_cache_key: str
title: tuple[str, ...]
source: str
tmdb_id: str
media_type: str
image_path: str
vote_count: int | None
source_key: str
def _get_detect_semaphore() -> "asyncio.Semaphore":
"""Lazily create the detection-admission semaphore inside the event loop."""
global _detect_semaphore
if _detect_semaphore is None:
_detect_semaphore = asyncio.Semaphore(_cfg.TEXTLESS_DETECTION_CONCURRENCY)
return _detect_semaphore
def _get_detect_executor() -> ThreadPoolExecutor:
"""Dedicated workers so OCR bursts cannot starve poster compositing."""
global _detect_executor
if _detect_executor is None:
_detect_executor = ThreadPoolExecutor(
max_workers=_cfg.TEXTLESS_DETECTION_CONCURRENCY,
thread_name_prefix="text-detect",
)
return _detect_executor
def _shutdown_detect_executor() -> None:
global _detect_executor
if _detect_executor is not None:
_detect_executor.shutdown(wait=True, cancel_futures=True)
_detect_executor = None
def _get_render_semaphore() -> "asyncio.Semaphore":
"""Lazily create the render-admission semaphore inside the event loop.
Bounds how many uncached /poster renders run at once so a burst from a cold
catalog grid queues here, in order, instead of oversubscribing the shared
HTTP pool and failing with PoolTimeout. Only the render pipeline itself is
gated: a composite cache hit returns before this is touched, and a request
coalesced onto an in-flight render waits on that render's future, never on a
slot of its own. A slot is also not held while waiting on another request's
MDBList fetch (the rating-coalescing event), so a holder can never be
blocked on a request that is itself queued for a slot.
"""
global _render_semaphore
if _render_semaphore is None:
_render_semaphore = asyncio.Semaphore(_cfg.POSTER_RENDER_CONCURRENCY)
return _render_semaphore
def _reserve_foreground_detection() -> None:
global _foreground_detection_count
_foreground_detection_count += 1
def _release_foreground_detection() -> None:
global _foreground_detection_count
_foreground_detection_count = max(0, _foreground_detection_count - 1)
def _start_text_detection(
cache_key: str,
image: Image.Image,
*,
title: tuple[str, ...],
source: str,
tmdb_id: str,
vote_count: int | None,
source_key: str,
media_type: str | None = None,
image_path: str | None = None,
foreground: bool = True,
foreground_reserved: bool = False,
) -> "asyncio.Task[bool | None]":
"""Start or join one OCR scan for an immutable source image."""
cached = get_cached_text_detection(cache_key)
if cached is not None:
if foreground and foreground_reserved:
_release_foreground_detection()
async def _cached_result() -> bool:
return cached
return asyncio.create_task(_cached_result())
existing = _text_detection_inflight.get(cache_key)
if existing is not None:
if foreground and foreground_reserved:
_release_foreground_detection()
logger.info(
f"Coalescing burned-in text scan for {tmdb_id} "
f"(votes={vote_count}, source={source_key})"
)
return existing
if foreground and not foreground_reserved:
_reserve_foreground_detection()
async def _scan() -> bool | None:
from text_detect import poster_has_burned_in_text
try:
async with _get_detect_semaphore():
result = await asyncio.get_running_loop().run_in_executor(
_get_detect_executor(),
lambda: poster_has_burned_in_text(
image,
conf=_cfg.PPOCR_BOX_THRESHOLD,
title=title,
source=source,
debug=True,
),
)
if result is not None:
set_cached_text_detection(cache_key, result)
if result is True and source == "poster" and media_type and image_path:
from textless_report import report_fake_textless_poster
report_fake_textless_poster(
media_type=media_type,
tmdb_id=tmdb_id,
image_path=image_path,
vote_count=vote_count,
)
return result
finally:
if foreground:
_release_foreground_detection()
logger.info(
f"Scanning textless poster {tmdb_id} for burned-in text "
f"(votes={vote_count}, source={source_key}, "
f"priority={'foreground' if foreground else 'background'})"
)
task = asyncio.create_task(_scan())
_text_detection_inflight[cache_key] = task
def _cleanup(done: "asyncio.Task[bool | None]") -> None:
if _text_detection_inflight.get(cache_key) is done:
_text_detection_inflight.pop(cache_key, None)
if not done.cancelled():
done.exception()
task.add_done_callback(_cleanup)
return task
def _queue_background_text_detection(item: _DeferredTextDetection) -> None:
"""Queue one vote-gated scan without retaining its decoded image."""
if get_cached_text_detection(item.cache_key) is not None:
return
if item.cache_key in _background_detection_keys:
return
if _background_detection_queue is None:
logger.warning(
f"Background text-detection queue unavailable for {item.tmdb_id}; "
"scan will retry on the next request"
)
return
_background_detection_keys.add(item.cache_key)
_background_detection_queue.put_nowait(item)
logger.info(
f"Queued vote-gated text scan for {item.tmdb_id} "
f"(votes={item.vote_count}, pending={_background_detection_queue.qsize()})"
)
def _load_detection_image(image_cache_key: str) -> Image.Image | None:
cached_bytes = get_cached_tmdb_poster(image_cache_key)
if not cached_bytes:
return None
return Image.open(io.BytesIO(cached_bytes)).convert("RGBA")
async def _background_text_detection_worker() -> None:
"""Drain vote-gated scans only while no foreground scan is queued or running."""
assert _background_detection_queue is not None
while True:
item = await _background_detection_queue.get()
try:
if get_cached_text_detection(item.cache_key) is not None:
continue
while _foreground_detection_count > 0 or _active_poster_renders > 0:
await asyncio.sleep(0.1)
image = await asyncio.get_running_loop().run_in_executor(
None, _load_detection_image, item.image_cache_key
)
if image is None:
logger.warning(
f"Deferred text scan source unavailable for {item.tmdb_id}; "
"scan will retry on the next request"
)
continue
# A poster render may have arrived while the image was loading.
while _foreground_detection_count > 0 or _active_poster_renders > 0:
await asyncio.sleep(0.1)
await asyncio.shield(_start_text_detection(
item.cache_key,
image,
title=item.title,
source=item.source,
tmdb_id=item.tmdb_id,
vote_count=item.vote_count,
media_type=item.media_type,
image_path=item.image_path,
source_key=item.source_key,
foreground=False,
))
except asyncio.CancelledError:
raise
except Exception as exc:
logger.warning(
f"Deferred text scan failed for {item.tmdb_id}: {exc}"
)
finally:
_background_detection_keys.discard(item.cache_key)
_background_detection_queue.task_done()
# Per-key cooldown timestamps (event-loop time). Keyed by the API key string so
# rotation is independent — a rate-limited key stands down while the other serves.
_mdblist_key_cooldown: dict[str, float] = {}
# Index into _cfg.SERVER_MDBLIST_KEYS for the currently active server-side key.
_mdblist_active_key_idx: int = 0
def _quality_backoff_remaining(now: float | None = None) -> float:
if now is None:
now = asyncio.get_running_loop().time()
return max(0.0, _quality_source_backoff_until.get(active_quality_source(), 0.0) - now)
def _record_quality_result(result) -> None:
# QUALITY_PENDING means the source answered and is healthy — it just has no
# value for this title yet. It is neither a success to reset the failure
# count on nor a failure to count, so the backoff state is left untouched.
if result is QUALITY_PENDING:
return
source = active_quality_source()
if result is not FETCH_FAILED:
_quality_source_backoff_until.pop(source, None)
_quality_source_fail_count.pop(source, None)
return
now = asyncio.get_running_loop().time()
if _quality_source_backoff_until.get(source, 0.0) > now:
return
failures = _quality_source_fail_count.get(source, 0) + 1
_quality_source_fail_count[source] = failures
delay = min(30.0 * (4 ** (failures - 1)), 1800.0)
_quality_source_backoff_until[source] = now + delay
logger.warning(f"Quality source {source} unavailable; backing off for {delay:.0f}s")
def _next_mdblist_server_key(current_key: str, now: float | None = None) -> str | None:
"""Select a healthy configured server key after *current_key*."""
global _mdblist_active_key_idx
keys = _cfg.SERVER_MDBLIST_KEYS
if len(keys) < 2 or current_key not in keys:
return None
if now is None:
now = asyncio.get_running_loop().time()
start = keys.index(current_key)
for offset in range(1, len(keys)):
idx = (start + offset) % len(keys)
candidate = keys[idx]
if now >= _mdblist_key_cooldown.get(candidate, 0.0):
_mdblist_active_key_idx = idx
return candidate
return None
def _mdblist_server_key_number(key: str | None) -> int | None:
if not key:
return None
for idx, candidate in enumerate(_cfg.SERVER_MDBLIST_KEYS):
if candidate == key:
return idx + 1
return None
def _mdblist_server_key_label(key: str | None) -> str:
number = _mdblist_server_key_number(key)
if number is None:
return "request-supplied key"
return f"configured key #{number}"
def _mark_mdblist_rate_limit(
canonical_id: str, key: str, result
) -> tuple[float, str | None]:
"""Cool down a rate-limited key and select a healthy configured fallback.
MDBList's limit is a daily quota, and a quota 429 carries no Retry-After —
only X-RateLimit-Reset. Retrying hourly until then just burns log lines,
so the key sleeps until the reset (capped at a day in case the header is
nonsense) and the fallback key takes over meanwhile.
"""
reset_at = getattr(result, "reset_at", None)
if result.retry_after:
backoff_secs = min(float(result.retry_after), 3600.0)
elif reset_at:
backoff_secs = min(max(float(reset_at) - time.time(), 60.0), 86400.0)
else:
backoff_secs = 3600.0
now = asyncio.get_running_loop().time()
_mdblist_key_cooldown[key] = now + backoff_secs
_rating_backoff[_rating_retry_key(canonical_id, key)] = now + backoff_secs
return backoff_secs, _next_mdblist_server_key(key, now)
def _fleet_cooldown_id(mdblist_key: str) -> str | None:
"""Coordination key for a credential's shared cooldown, or None if it has
no fleet-wide meaning.
Only the OPERATOR's configured keys are shared across replicas, so only
they get one. A request-supplied key belongs to one tenant: its quota
running out says nothing about anyone else's, and cooling the fleet on it
would let any user with an exhausted key stop rating lookups for everyone.
The identity is hashed so the credential never reaches Redis in clear.
"""
if not mdblist_key or mdblist_key not in _cfg.SERVER_MDBLIST_KEYS:
return None
return hashlib.sha256(mdblist_key.encode("utf-8")).hexdigest()[:16]
async def _publish_fleet_mdblist_cooldown(mdblist_key: str, backoff_secs: float) -> None:
"""ElfHosted fork: raise a fleet-wide cooldown for ONE MDBList credential.
_mark_mdblist_rate_limit records the 429 in this process's dicts, which is
all a single-instance deploy needs. Across replicas sharing the same keys,
every replica would otherwise have to discover the same 429 for itself —
one wasted MDBList call each, against a key that has already asked us to
stop. Publishing it means the first replica to be refused backs the rest
off too.
Per credential, not fleet-wide-for-everything: MDBList quota is per key and
v1.2.0 rotates between several. A single shared flag would let one
exhausted key disable the healthy siblings rotation had just selected —
turning a working failover into a fleet-wide outage for up to the 24h a
quota cooldown can last.
A no-op beyond the per-process cooldown on the in-process coordinator, and
never fatal: failing to publish a cooldown must not fail the request that
was merely unlucky enough to hit the 429.
"""
cooldown_id = _fleet_cooldown_id(mdblist_key)
if cooldown_id is None:
return
with suppress(Exception):
await coord.set_backoff(
coord.NS_MDBLIST_GLOBAL_COOLDOWN, cooldown_id,
ttl_seconds=backoff_secs,
)
async def _fleet_mdblist_cooling(mdblist_key: str) -> bool:
"""True when a sibling replica has already been 429'd on this credential."""
cooldown_id = _fleet_cooldown_id(mdblist_key)
if cooldown_id is None:
return False
try:
return await coord.is_backoff_active(
coord.NS_MDBLIST_GLOBAL_COOLDOWN, cooldown_id
)
except Exception:
return False
def _warm_mdblist_key_with_quota(current_key: str, now: float, reserve: int) -> str | None:
"""
Pick a configured key the cache warmer may still spend: not cooling down,
and with unknown or above-reserve daily quota. Tries *current_key* first,
then its siblings in order.
Deliberately does not touch _mdblist_active_key_idx: a key at the reserve
floor is fine for live requests (that is what the reserve is for), so the
warmer moving on to a sibling must not drag live traffic along with it.
"""
keys = _cfg.SERVER_MDBLIST_KEYS
if current_key not in keys:
return current_key if now >= _mdblist_key_cooldown.get(current_key, 0.0) else None
start = keys.index(current_key)
for offset in range(len(keys)):
candidate = keys[(start + offset) % len(keys)]
if now < _mdblist_key_cooldown.get(candidate, 0.0):
continue
remaining = mdblist_quota_remaining(candidate)
if remaining is not None and remaining <= reserve:
continue
return candidate
return None
async def _background_quality_fetch(
quality_id: str,
media_type: str,
season: int,
episode: int,
release_date: str | None,
) -> None:
"""Fetch quality tokens from the configured quality source and cache them. Never raises."""
global _quality_bg_semaphore
if _quality_bg_semaphore is None:
_quality_bg_semaphore = asyncio.Semaphore(_cfg.QUALITY_BG_CONCURRENCY)
try:
async with _quality_bg_semaphore:
if _HTTP_CLIENT is None:
return
remaining = _quality_backoff_remaining()
if remaining > 0:
logger.debug(
f"Quality fetch skipped for {quality_id}; source cooldown has {remaining:.0f}s remaining"
)
return
result = await _with_retry(
fetch_quality,
_HTTP_CLIENT, quality_id, media_type, season, episode, release_date,
)
_record_quality_result(result)
if result is QUALITY_PENDING:
# QualiCache is collecting in the background; the next request
# for this title picks up the value once it lands.
logger.info(f"Background quality fetch pending for {quality_id}")
elif result is not FETCH_FAILED:
logger.info(f"Background quality fetch complete for {quality_id}")
except Exception as exc:
_record_quality_result(FETCH_FAILED)
logger.warning(f"Background quality fetch failed for {quality_id}: {exc}")
finally:
_quality_bg_inflight.discard(quality_id)
# Local imports
from age_badge import draw_quality_age_badge, draw_quality_corner_bookmark, draw_tier_bar, _score_points
from landscape import build_landscape
from awards import _dominant_cluster, _is_skin_tone, dominant_frost_rgb
from awards import FETCH_FAILED, _RateLimited, draw_award_badge, draw_award_sash, parse_mdblist_awards, reconcile_cached_awards
from festivals import match_festival_keyword
from i18n import load_languages, translate_genre, translate_sash
from cache import (
get_cached_quality,
get_cached_rating,
get_cached_final_poster,
get_cached_final_poster_entry,
get_cached_final_poster_url,
get_cached_final_poster_redirect,
is_cached_final_poster_fresh,
set_cached_final_poster,
delete_cached_final_poster,
get_cached_tmdb_poster,
get_cached_tmdb_metadata,
get_cached_text_detection,
set_cached_text_detection,
get_cached_release_status,
get_cached_imdb_to_tmdb,
delete_cached_imdb_to_tmdb,
init_db,
is_digital_release,
set_cached_rating,
delete_cached_tmdb_metadata,
prune_caches,
prune_local_caches,
release_status_ttl_seconds,
get_cache_stats,
get_app_state,
set_app_state,
# ElfHosted fork: upstream reaches for cache.get_db() here and runs the
# composite request_params SELECT inline, which only works against SQLite.
# The query lives behind a backend function so it works on Postgres too.
list_composite_request_params,
close as close_db,
BACKEND_KIND as _STORAGE_KIND,
)
import blobstore
import coordination as coord
import metrics as _metrics
from digital_release import digital_release_poll_loop
import imdb_dataset
from imdb_dataset import imdb_dataset_refresh_loop
import config as _cfg
from discovery import (
ALL_PRIORITY_SLOTS,
DiscoveryMeta,
extract_discovery_meta,
pick_sash,
)
from quality import (
QUALITY_PENDING,
QUALITY_SOURCES,
BadgeItem,
active_quality_source,
fetch_quality,
get_resized_badge,
parse_quality,
quality_source_configured,
render_badges_left,
)
from presets import get_preset, preset_names, preset_catalog
from ratings import (
CustomScorePalette,
MDBLIST_QUOTA,
calculate_weighted_score,
draw_frosted_bar,
draw_score_bar,
fetch_rating,
is_anime_rated,
mdblist_quota_remaining,
parse_custom_score_palette,
score_color_for_mode,
_draw_solid_pip,
_score_color,
_score_color_alt,
_score_color_metal,
)
from tmdb import composite_logo, logo_centre_y, fetch_logo, image_language_order, fetch_poster_metadata, resolve_imdb_to_tmdb, fetch_poster_image, fetch_backdrop_image, fetch_landscape_image, fetch_trending_rank, fetch_trending_candidates, fetch_popular_candidates, fetch_supplemental_candidates, fetch_catalog_candidates, fetch_release_status, fetch_upcoming_movie_release, fetch_recent_movie_digital_release_date, recent_digital_release_from_cache, svg_logo_supported, tmdb_metadata_cache_key, _CROP_VERSION, _fetch_metahub_logo, LOGO_ABS_MAX_H, TEXT_FORWARD_PRIORITIES as _TEXT_FORWARD_LOGO_PRIORITIES
# Logo priorities that consult the secondary preferred language ("custom").
# Elsewhere the secondary language is inert and must be kept out of the image
# fetch / cache key so single-language requests keep their existing cache entry.
_SECONDARY_LANGUAGE_PRIORITIES = frozenset({
"native_custom_text",
"native_custom_original_text",
})
import tvdb
import anime
# ---------------------------------------------------------------------------
# Persistent HTTP client
# ---------------------------------------------------------------------------
# One client for the lifetime of the process. httpx keeps TCP connections
# alive in its connection pool, so repeated requests to the same host
# (TMDB, MDblist, AIOStreams) reuse the existing socket rather than paying
# TLS + TCP handshake overhead on every poster request.
#
# Timeouts are split:
# connect=5s — fail fast when a host is unreachable
# read=12s — allow slow responses from external APIs
# pool=10s — don't block forever waiting for a pool slot, but do wait:
# a request that queues a few seconds behind a burst still
# renders, one that gives up is a 504 the client may cache
#
# The pool is sized from POSTER_RENDER_CONCURRENCY so an operator who raises
# the render cap doesn't silently reintroduce pool exhaustion: each render
# fans out to ~4 upstream calls at its peak, and the background loops
# (quality fetches, TVDB, trending, warm cycles) share the same pool.
_HTTP_CLIENT: httpx.AsyncClient | None = None
# Peak concurrent upstream calls per render (art + logo + rating + trending).
_HTTP_CALLS_PER_RENDER = 4
_HTTP_POOL_MIN_CONNECTIONS = 40
# Connections left for the loops that don't go through /poster: TVDB, trending,
# digital-release sync, the IMDb dataset refresh.
_HTTP_POOL_BACKGROUND_HEADROOM = 8
def _http_pool_size(render_concurrency: int) -> int:
return max(
_HTTP_POOL_MIN_CONNECTIONS,
render_concurrency * _HTTP_CALLS_PER_RENDER
+ _cfg.QUALITY_BG_CONCURRENCY
+ _HTTP_POOL_BACKGROUND_HEADROOM,
)
def _make_http_client() -> httpx.AsyncClient:
_max_connections = _http_pool_size(_cfg.POSTER_RENDER_CONCURRENCY)
return httpx.AsyncClient(
timeout=httpx.Timeout(connect=5.0, read=12.0, write=5.0, pool=10.0),
limits=httpx.Limits(
max_connections=_max_connections,
max_keepalive_connections=max(20, _max_connections // 2),
keepalive_expiry=30,
),
headers={
"Accept-Encoding": "identity",
"User-Agent": (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) "
"Chrome/124.0.0.0 Safari/537.36"
),
},
http2=False, # most poster APIs don't support h2; skip the negotiation
)
# ---------------------------------------------------------------------------
# Input validation
# ---------------------------------------------------------------------------
_TMDB_ID_RE = re.compile(r'^\d{1,10}$')
_IMDB_ID_RE = re.compile(r'^tt\d{1,10}$')
_VALID_TYPES = frozenset({"movie", "tv", "series"})
def _check_tmdb_id(val: str) -> None:
if not _TMDB_ID_RE.match(val):
raise HTTPException(status_code=400, detail="Invalid tmdb_id")
def _check_imdb_id(val: str) -> None:
if not _IMDB_ID_RE.match(val):
raise HTTPException(status_code=400, detail="Invalid imdb_id")
def _normalise_optional_id(raw: str | None, name: str) -> str:
"""Trim an optional id param, reading an unsubstituted placeholder as absent.
A template pasted into a metadata provider arrives with the placeholder
still in it when that provider has no id for the title — AIOMetadata's
optional "{name?}" form is left verbatim by older builds, and some addons
reject the "?" syntax outright so operators write the plain form. Either way
the value is "no id", not a malformed one, and 400ing it would take down
every poster served through that template.
Deliberately narrow: only this parameter's own two literals. Accepting any
brace-wrapped value would silently swallow genuine typos.
"""
value = (raw or "").strip()
if value in ("{" + name + "}", "{" + name + "?}"):
return ""
return value
def _canonical_rating_id(imdb_id: str, anime_key: str, tmdb_id: str) -> str:
"""The immutable cache/coalescing identity for a request.
Chosen once, before any metadata is fetched, and never revised: it keys the
rating cache read, the rating cache write, the coalescing map and the
back-off tables, so a value that changed mid-request would read one row and
write another — turning every subsequent request for that title into a fresh
MDBList call, permanently.
An IMDb id discovered later from TMDB metadata therefore never lands here.
See _quality_identity() for the identity that may use it.
The "tmdb:" form can't collide with a bare TMDB id or a tt-prefixed IMDb id,
and matches the namespacing the anime path already stores in these columns —
so no migration is needed.
"""
return imdb_id or anime_key or f"tmdb:{tmdb_id}"
def _merge_imdb_dataset_rating(
ratings_dict, effective_imdb_id: str | None, rcfg: "RequestConfig"
):
"""Supply the "imdb" entry in *ratings_dict* from the local IMDb dataset.
Three modes, per rcfg.imdb_rating_source:
"mdblist" (default) — no-op. The IMDb weight comes from MDBList like
every other source.
"dataset" — always override. MDBList is not consulted for
this one source at all, so the weight works
with no MDBList key configured.
"fallback" — backfill only. MDBList's answer wins whenever
it has one; the dataset fills the gap when it
does not.
"fallback" is the mode that pays off for operators who do use MDBList,
and it covers two distinct gaps with one rule. A hard gap: MDBList was
rate-limited, timed out, or every configured key is cooling down, so
the failure path arrives here with an empty dict. And a soft gap:
MDBList answered fine but carried no IMDb score for this title, or
carried one that RATING_MIN_VOTES filtered out. Both look identical
from here — "imdb" is absent — so both are covered by testing for its
absence rather than by inspecting how the fetch went.
Note the deliberate asymmetry with fallback_to_imdb, which is a
*scoring* fallback: it fires when no weighted source scored at all and
reaches for whatever "imdb" value is present. This one fires earlier,
at the point the ratings are assembled, and is what puts a value there
for it to find. The two compose: dataset backfill, then weighting,
then fallback_to_imdb if the weights still produced nothing.
In every no-op case the existing MDBList behaviour is left exactly as
it was — this only ever adds or replaces the single "imdb" key, and
only when it has something to put there.
"""
if not isinstance(ratings_dict, dict):
return ratings_dict
mode = rcfg.imdb_rating_source
if mode not in ("dataset", "fallback"):
return ratings_dict
if mode == "fallback" and ratings_dict.get("imdb") is not None:
return ratings_dict
if not effective_imdb_id:
return ratings_dict
value = imdb_dataset.get_rating(effective_imdb_id)
if value is None:
return ratings_dict
return {**ratings_dict, "imdb": value}
def _ratings_base(ratings_dict):
"""Normalise "MDBList was never asked" to an empty dict.
An instance with no MDBList key and no cached rating row resolves its
rating tuple from cached_ratings_dict, which is None rather than {}.
That is precisely the configuration the dataset / direct-TMDB sources
exist to serve, so it has to reach the merge helpers as a dict —
otherwise they are skipped, no MDBList-free source is ever consulted,
and the score is left as None (which the debug JSON then fails to
int()).
Anything that is already a dict, and any non-None sentinel such as the
"N/A" string, is passed through untouched: those are real answers, not
an absence of one.
"""
return {} if ratings_dict is None else ratings_dict
def _merge_direct_tmdb_rating(ratings_dict, tmdb_data: dict, rcfg: "RequestConfig"):
"""Supply the "tmdb" entry in *ratings_dict* from TMDB's own vote_average.
Same three modes as _merge_imdb_dataset_rating, per
rcfg.tmdb_rating_source: "mdblist" (default, no-op), "direct" (always
override), "fallback" (backfill only when MDBList has no "tmdb" value,
whether because the fetch failed or because it simply carried none).
See that function for why absence is the right thing to test.
"fallback" is close to free here in a way it is not for the IMDb
dataset: there is no download, no table and no readiness window, so an
MDBList outage is covered for this source on an instance that has
opted into nothing else.
Unlike the IMDb dataset, this costs nothing extra: vote_average rides
along in the same TMDB details call already made for genre, year and
credits, so this is a pure re-use of data already in hand — no schedule,
no local storage, no separate opt-in infrastructure. It works with no
MDBList key configured at all, same rationale as the IMDb dataset source.