Files
Magent/backend/app/services/newsletter_catalog.py
T
Assclaw b0f8c89db7
Magent CI/CD / verify (push) Successful in 10m54s
Magent CI/CD / deploy-prod (push) Skipped
Magent CI/CD / deploy-beta (push) Successful in 1m14s
Add Grizzlyflix newsletters with curated editions and weekly delivery
2026-09-09 23:42:52 +12:00

203 lines
11 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Bounded Jellyfin arrival snapshots, recipient access checks and email-safe posters."""
import asyncio
import hashlib
import io
import time
from collections import OrderedDict
from datetime import datetime, timezone
import httpx
from PIL import Image
from .insights_artwork import item_id
from .jellyfin_identity import source_key
MAX_ITEMS = 5000
PAGE_SIZE = 200
MAX_TITLES = 60
_posters = OrderedDict()
_poster_lock = asyncio.Semaphore(4)
class CatalogError(Exception):
pass
def date(value) -> datetime | None:
try:
result = datetime.fromisoformat(str(value).replace('Z', '+00:00'))
return result.replace(tzinfo=timezone.utc) if result.tzinfo is None else result.astimezone(timezone.utc)
except (ValueError, TypeError):
return None
async def get_json(client, runtime, path, params=None):
try:
response = await client.get(runtime.jellyfin_base_url.rstrip('/') + path,
headers={'X-Emby-Token': runtime.jellyfin_api_key}, params=params)
response.raise_for_status()
return response.json()
except (httpx.HTTPError, ValueError) as exc:
raise CatalogError('Jellyfin is temporarily unavailable. Please try again.') from exc
def group_arrivals(items: list[dict], start: datetime, end: datetime) -> list[dict]:
groups = {}
seen = set()
for row in items:
identity = item_id(row.get('Id'))
added = date(row.get('DateCreated'))
if (not identity or identity in seen or not added or not start <= added < end
or row.get('LocationType') == 'Virtual' or row.get('IsPlaceHolder')):
continue
kind = row.get('Type')
if kind not in {'Movie', 'Episode'}:
continue
parent = item_id(row.get('SeriesId')) if kind == 'Episode' else identity
if not parent:
continue
seen.add(identity)
title = str((row.get('SeriesName') if kind == 'Episode' else row.get('Name')) or '').strip()
if not title:
continue
entry = groups.setdefault(parent, {'id': parent, 'type': 'series' if kind == 'Episode' else 'movie',
'title': title[:250], 'year': row.get('ProductionYear') if kind == 'Movie' else None,
'overview': str(row.get('Overview') or '')[:500] if kind == 'Movie' else '',
'added_at': added.isoformat(), 'has_artwork': False, 'items': [], 'selected': False, 'featured': False})
entry['added_at'] = max(entry['added_at'], added.isoformat())
entry['has_artwork'] |= bool(row.get('SeriesPrimaryImageTag') if kind == 'Episode' else (row.get('ImageTags') or {}).get('Primary'))
entry['items'].append({'id': identity, 'season': row.get('ParentIndexNumber') if kind == 'Episode' else None,
'number': row.get('IndexNumber') if kind == 'Episode' else None})
return sorted(groups.values(), key=lambda row: (row['added_at'], row['id']), reverse=True)
async def collect(runtime, start: datetime, end: datetime, limit: int = 12) -> dict:
if not runtime.jellyfin_base_url or not runtime.jellyfin_api_key:
raise CatalogError('Connect Jellyfin before collecting new arrivals.')
rows, seen = [], set()
exhausted = False
async with httpx.AsyncClient(timeout=20) as client:
info = await get_json(client, runtime, '/System/Info')
server_id = item_id(info.get('Id')) if isinstance(info, dict) else None
if not server_id:
raise CatalogError('Jellyfin did not return its server identity.')
for offset in range(0, MAX_ITEMS, PAGE_SIZE):
payload = await get_json(client, runtime, '/Items', {'Recursive': 'true', 'IncludeItemTypes': 'Movie,Episode',
'SortBy': 'DateCreated,SortName', 'SortOrder': 'Descending', 'Fields': 'DateCreated,Overview',
'EnableUserData': 'false', 'IsMissing': 'false', 'IsPlaceHolder': 'false', 'Limit': PAGE_SIZE, 'StartIndex': offset})
if not isinstance(payload, dict) or not isinstance(payload.get('Items'), list):
raise CatalogError('Jellyfin returned an incomplete arrival list.')
page = payload['Items']
total = payload.get('TotalRecordCount')
if not isinstance(total, int) or total < offset + len(page):
raise CatalogError('Jellyfin returned an incomplete arrival count.')
for row in page:
if not isinstance(row, dict) or not item_id(row.get('Id')) or not date(row.get('DateCreated')):
raise CatalogError('Jellyfin returned an arrival without a valid identity or added date.')
identity = item_id(row['Id'])
if identity in seen:
raise CatalogError('The library changed during collection. Refresh arrivals to try again.')
seen.add(identity)
if rows and date(row['DateCreated']) > date(rows[-1]['DateCreated']):
raise CatalogError('The library changed during collection. Refresh arrivals to try again.')
rows.append(row)
if (not page or len(page) < PAGE_SIZE) and offset + len(page) < total:
raise CatalogError('Jellyfin returned an incomplete arrival page.')
if not page or any(date(row['DateCreated']) < start for row in page) or offset + len(page) >= total:
exhausted = True
break
if not exhausted:
raise CatalogError('More than 5,000 recent items were found. Choose a shorter arrival period; no partial edition was created.')
titles = group_arrivals(rows, start, end)
total = len(titles)
titles = titles[:MAX_TITLES]
for index, title in enumerate(titles):
title['selected'] = index < limit
return {'source': source_key(runtime.jellyfin_base_url), 'server_id': server_id,
'period_start': start.isoformat(), 'period_end': end.isoformat(), 'total_titles': total, 'titles': titles}
async def for_recipient(runtime, content: dict, jellyfin_id: str) -> dict:
"""Scope every ID lookup to a view Jellyfin permits this user to browse.
Jellyfin 10.11's AddUserToQuery skips its default library filter when ItemIds
is present. UserId alone is insufficient; ParentId supplies the allowed scope.
"""
if not item_id(jellyfin_id):
raise CatalogError('The recipient does not have a valid Jellyfin identity.')
selected = [entry for entry in content['titles'] if entry['selected']]
ids = sorted({identity for entry in selected for identity in [entry['id'], *(item['id'] for item in entry['items'])]})
allowed = set()
async with httpx.AsyncClient(timeout=20) as client:
info = await get_json(client, runtime, '/System/Info')
if not isinstance(info, dict) or source_key(runtime.jellyfin_base_url) != content['source'] or item_id(info.get('Id')) != content['server_id']:
raise CatalogError('The Jellyfin server changed. Create a new edition for the current library.')
user = await get_json(client, runtime, '/Users/' + jellyfin_id)
if not isinstance(user, dict) or item_id(user.get('Id')) != item_id(jellyfin_id) or not isinstance(user.get('Policy'), dict):
raise CatalogError('Could not verify the recipients Jellyfin account.')
if user['Policy'].get('IsDisabled') or user['Policy'].get('EnableMediaPlayback') is False:
return {**content, 'titles': [], 'recipient_disabled': True}
views = await get_json(client, runtime, '/UserViews', {'UserId': jellyfin_id, 'IncludeHidden': 'true', 'IncludeExternalContent': 'false'})
if not isinstance(views, dict) or not isinstance(views.get('Items'), list) or len(views['Items']) > 32:
raise CatalogError('Could not check the recipients library access.')
for view in views['Items']:
parent = item_id(view.get('Id')) if isinstance(view, dict) else None
if not parent:
raise CatalogError('Jellyfin returned a library without a valid identity.')
for offset in range(0, len(ids), 100):
chunk = ids[offset:offset + 100]
payload = await get_json(client, runtime, '/Items', {'UserId': jellyfin_id, 'ParentId': parent, 'Ids': ','.join(chunk),
'Recursive': 'true', 'Limit': len(chunk), 'EnableUserData': 'false', 'EnableImages': 'false',
'IsMissing': 'false', 'IsPlaceHolder': 'false'})
if not isinstance(payload, dict) or not isinstance(payload.get('Items'), list):
raise CatalogError('Could not check the recipients library access.')
allowed.update(item_id(item.get('Id')) for item in payload['Items'] if isinstance(item, dict))
titles = []
for entry in selected:
accessible = [item for item in entry['items'] if item['id'] in allowed]
if entry['id'] in allowed and accessible:
titles.append({**entry, 'items': accessible})
return {**content, 'titles': titles}
async def poster(runtime, identity: str) -> bytes | None:
if not item_id(identity):
return None
key = (source_key(runtime.jellyfin_base_url), hashlib.sha256(runtime.jellyfin_api_key.encode()).hexdigest(), identity)
async with _poster_lock:
cached = _posters.get(key)
if cached and cached[0] > time.monotonic():
_posters.move_to_end(key)
return cached[1]
result = None
try:
async with httpx.AsyncClient(timeout=10) as client:
async with client.stream('GET', runtime.jellyfin_base_url.rstrip('/') + f'/Items/{identity}/Images/Primary',
headers={'X-Emby-Token': runtime.jellyfin_api_key}, params={'maxWidth': 160, 'maxHeight': 240, 'quality': 82, 'format': 'Jpg'}) as response:
response.raise_for_status()
data = bytearray()
async for chunk in response.aiter_bytes():
data.extend(chunk)
if len(data) > 512 * 1024:
raise ValueError('Poster too large')
with Image.open(io.BytesIO(data)) as image:
if image.width * image.height > 4_000_000:
raise ValueError('Poster dimensions too large')
image.thumbnail((160, 240))
target = io.BytesIO()
image.convert('RGB').save(target, format='JPEG', quality=82)
result = target.getvalue()
except (httpx.HTTPError, ValueError, OSError, Image.DecompressionBombError):
pass
_posters[key] = (time.monotonic() + (1800 if result else 60), result)
while len(_posters) > 128:
_posters.popitem(last=False)
return result
async def posters(runtime, content: dict) -> dict:
titles = [entry for entry in content['titles'] if entry['selected'] and entry['has_artwork']]
results = await asyncio.gather(*(poster(runtime, entry['id']) for entry in titles))
return {entry['id']: data for entry, data in zip(titles, results) if data}