203 lines
11 KiB
Python
203 lines
11 KiB
Python
"""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 recipient’s 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 recipient’s 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 recipient’s 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}
|