* fix: audit hardening — input validation, error handling, and robustness - immich_api: guard person["id"] with .get() + early return on missing field - immich_api: include page number in pagination exception log - immich_api: validate faces response is a list before indexing - executor: wrap Image.open() in try/except for non-image HTTP responses - executor: strip leading 'v' from Frigate version before parsing (v0.16.0 was misread) - config: wrap FRIGATE_SCORE_CEILING float() parse in try/except with warning - config: warn when both DATA_DIR and legacy CWD config files exist simultaneously - scheduler: wrap PID file write in try/except so /tmp failures don't crash startup - scheduler: clamp sleep to 60s max to bound recovery time after NTP clock jumps - frigate_api: log unexpected non-list type in get_frigate_person_files at DEBUG * fix: LIMIT env var crash and symlink guard on person output dir - jobs: wrap int(LIMIT) parse in try/except — bad value (e.g. "30.5", "all") now logs a warning and falls back to the default instead of crashing - executor: check for symlink before shutil.rmtree on person_dir — prevents following a symlink out of OUTPUT_DIR on a shared volume * chore: bump version to 0.5.3
258 lines
8.3 KiB
Python
258 lines
8.3 KiB
Python
"""Immich API client for fetching people, assets, and face data."""
|
|
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from datetime import datetime, timedelta, timezone
|
|
from io import BytesIO
|
|
|
|
import requests
|
|
from PIL import Image, ImageOps
|
|
|
|
from .config import Config, get_headers
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
MAX_PAGES = 1000 # Safety limit for pagination
|
|
_MAX_ASSETS_PER_PERSON = 5000 # Stop fetching after this many — diversity pool is capped at 3000 anyway
|
|
|
|
|
|
@dataclass
|
|
class FaceData:
|
|
"""Pre-computed face data from Immich."""
|
|
|
|
bbox: tuple[float, float, float, float] # (x1, y1, x2, y2)
|
|
confidence: float | None
|
|
image_width: int
|
|
image_height: int
|
|
|
|
|
|
def get_immich_version() -> tuple[int, int, int] | None:
|
|
"""Fetch Immich server version from GET /api/server/version.
|
|
|
|
Returns (major, minor, patch) or None if unreachable or unparseable.
|
|
"""
|
|
try:
|
|
resp = requests.get(
|
|
f"{Config.IMMICH_URL}/api/server/version",
|
|
headers=get_headers(),
|
|
timeout=5,
|
|
)
|
|
if resp.ok:
|
|
data = resp.json()
|
|
return (int(data["major"]), int(data["minor"]), int(data["patch"]))
|
|
return None
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def get_people() -> list[dict]:
|
|
"""Fetch all people from Immich."""
|
|
try:
|
|
resp = requests.get(
|
|
f"{Config.IMMICH_URL}/api/people",
|
|
headers=get_headers(),
|
|
timeout=10,
|
|
)
|
|
if resp.status_code == 401:
|
|
logger.error("Immich API key is invalid or expired (401 Unauthorized). Update API_KEY.")
|
|
return []
|
|
resp.raise_for_status()
|
|
return resp.json().get("people", [])
|
|
except (requests.RequestException, ValueError) as e:
|
|
logger.error("Failed to fetch people from Immich: %s", e)
|
|
return []
|
|
|
|
|
|
def merge_people(survivor_id: str, merge_ids: list[str]) -> bool:
|
|
"""Merge duplicate people into survivor via Immich's merge endpoint.
|
|
|
|
The survivor (identified by survivor_id) absorbs all faces and assets
|
|
from the people in merge_ids, which are then removed from Immich.
|
|
"""
|
|
try:
|
|
resp = requests.put(
|
|
f"{Config.IMMICH_URL}/api/people/{survivor_id}/merge",
|
|
headers={**get_headers(), "Content-Type": "application/json"},
|
|
json={"ids": merge_ids},
|
|
timeout=30,
|
|
)
|
|
resp.raise_for_status()
|
|
return True
|
|
except requests.RequestException as e:
|
|
logger.error("Failed to merge people into %s: %s", survivor_id, e)
|
|
return False
|
|
|
|
|
|
def fetch_all_assets(person: dict) -> list[dict]:
|
|
"""Fetch all assets for a person with pagination."""
|
|
name = person.get("name", "Unknown")
|
|
person_id = person.get("id")
|
|
if not person_id:
|
|
logger.error("Person dict missing 'id' field for %s — skipping asset fetch", name)
|
|
return []
|
|
url = f"{Config.IMMICH_URL}/api/search/metadata"
|
|
page_size = 1000
|
|
|
|
logger.debug("Fetching assets for %s...", name)
|
|
|
|
assets = []
|
|
for page in range(1, MAX_PAGES + 1):
|
|
try:
|
|
resp = requests.post(
|
|
url,
|
|
json={"personIds": [person_id], "size": page_size, "page": page},
|
|
headers=get_headers(),
|
|
timeout=30,
|
|
)
|
|
|
|
if not resp.ok:
|
|
logger.error("Error fetching assets for %s (page %s): %s", name, page, resp.status_code)
|
|
break
|
|
|
|
page_assets = resp.json().get("assets", [])
|
|
if isinstance(page_assets, dict):
|
|
page_assets = page_assets.get("items", [])
|
|
|
|
if not page_assets:
|
|
break
|
|
|
|
assets.extend(a for a in page_assets if isinstance(a, dict))
|
|
logger.debug("Fetched page %s, total: %s", page, len(assets))
|
|
|
|
if len(page_assets) < page_size or len(assets) >= _MAX_ASSETS_PER_PERSON:
|
|
break
|
|
|
|
except (requests.RequestException, ValueError) as e:
|
|
logger.error("Exception fetching assets for %s (page %s): %s", name, page, e)
|
|
break
|
|
|
|
return assets
|
|
|
|
|
|
def fetch_face_data(asset_id: str, person_id: str | None = None) -> FaceData | None:
|
|
"""Fetch pre-computed face data (bbox, confidence) from Immich.
|
|
|
|
Queries GET /api/faces?id={asset_id} to retrieve face detection results
|
|
that Immich already computed using InsightFace Buffalo_L.
|
|
|
|
Args:
|
|
asset_id: The asset to get face data for
|
|
person_id: Optional person ID to match the specific face
|
|
|
|
Returns:
|
|
FaceData with bbox and confidence, or None if unavailable
|
|
"""
|
|
try:
|
|
resp = requests.get(
|
|
f"{Config.IMMICH_URL}/api/faces",
|
|
params={"id": asset_id},
|
|
headers=get_headers(),
|
|
timeout=10,
|
|
)
|
|
|
|
if not resp.ok:
|
|
logger.debug("Face data endpoint returned %s for %s", resp.status_code, asset_id)
|
|
return None
|
|
|
|
faces = resp.json()
|
|
if not isinstance(faces, list) or not faces:
|
|
return None
|
|
|
|
# Match the target person if specified
|
|
face = None
|
|
if person_id:
|
|
face = next(
|
|
(f for f in faces if isinstance(f, dict) and (f.get("person") or {}).get("id") == person_id),
|
|
None,
|
|
)
|
|
if face is None:
|
|
face = faces[0] if isinstance(faces[0], dict) else None
|
|
if face is None:
|
|
return None
|
|
|
|
bbox = (
|
|
face.get("boundingBoxX1", 0),
|
|
face.get("boundingBoxY1", 0),
|
|
face.get("boundingBoxX2", 0),
|
|
face.get("boundingBoxY2", 0),
|
|
)
|
|
|
|
score = face.get("score")
|
|
return FaceData(
|
|
bbox=bbox,
|
|
confidence=score if score is not None else face.get("confidence"),
|
|
image_width=face.get("imageWidth", 0),
|
|
image_height=face.get("imageHeight", 0),
|
|
)
|
|
|
|
except requests.RequestException as e:
|
|
logger.debug("Failed to fetch face data for %s: %s", asset_id, e)
|
|
return None
|
|
except (AttributeError, KeyError, TypeError, ValueError) as e:
|
|
logger.debug("Failed to parse face data for %s: %s", asset_id, e)
|
|
return None
|
|
|
|
|
|
def fetch_full_image(asset_id: str, timeout: int = 60) -> Image.Image | None:
|
|
"""Fetch full-resolution image from Immich, falling back to preview thumbnail.
|
|
|
|
The /original endpoint may return HEIC, RAW, or video files that PIL
|
|
cannot open directly. In that case, we fall back to the JPEG thumbnail.
|
|
"""
|
|
# Try original first
|
|
try:
|
|
resp = requests.get(
|
|
f"{Config.IMMICH_URL}/api/assets/{asset_id}/original",
|
|
headers=get_headers(),
|
|
timeout=timeout,
|
|
)
|
|
if resp.ok:
|
|
try:
|
|
return ImageOps.exif_transpose(Image.open(BytesIO(resp.content)))
|
|
except Exception:
|
|
logger.debug("PIL can't open original for %s, falling back to preview", asset_id)
|
|
except requests.RequestException:
|
|
logger.debug("Original request failed for %s, falling back to preview", asset_id)
|
|
|
|
# Fall back to preview thumbnail (always JPEG)
|
|
try:
|
|
resp = requests.get(
|
|
f"{Config.IMMICH_URL}/api/assets/{asset_id}/thumbnail?size=preview&format=JPEG",
|
|
headers=get_headers(),
|
|
timeout=30,
|
|
)
|
|
if resp.ok:
|
|
return ImageOps.exif_transpose(Image.open(BytesIO(resp.content)))
|
|
except Exception as e:
|
|
logger.error("Failed to fetch image %s: %s", asset_id, e)
|
|
|
|
return None
|
|
|
|
|
|
def filter_recent_assets(assets: list[dict], years: int | None = None) -> list[dict]:
|
|
"""Filter assets to keep only those from the last N years."""
|
|
years = years or Config.YEARS_FILTER
|
|
cutoff = datetime.now(timezone.utc) - timedelta(days=365 * years)
|
|
|
|
logger.debug("Filtering assets older than %s years (%s)", years, cutoff)
|
|
|
|
recent, skipped = [], 0
|
|
for asset in assets:
|
|
created_at_str = asset.get("fileCreatedAt")
|
|
if not created_at_str:
|
|
continue
|
|
|
|
try:
|
|
# Handle ISO8601 with 'Z' suffix
|
|
created_at = datetime.fromisoformat(created_at_str.replace("Z", "+00:00"))
|
|
if created_at > cutoff:
|
|
recent.append(asset)
|
|
else:
|
|
skipped += 1
|
|
except ValueError:
|
|
continue
|
|
|
|
logger.debug("Retained %s assets (filtered %s old assets).", len(recent), skipped)
|
|
return recent
|
|
|