diff --git a/CHANGELOG.md b/CHANGELOG.md index d0c0092..2655076 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,22 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [0.2.2] - 2026-06-12 + +### Added + +- **Confidence scores in upload tracker**: Immich face confidence scores are now stored per asset in `frigate_uploaded_ids.json` under `by_person[name].scores`. Lays the groundwork for future replacement logic (remove low-confidence uploads when better images are found). +- **Frigate-authoritative capacity tracking**: At startup, `GET /api/faces` is queried on the Frigate host to retrieve the actual number of trained images per person from the `train` directory (pending/unclassified queue is excluded). This count is stored as `frigate_count` in the tracker JSON so it survives Frigate downtime. +- **Lifetime cap uses Frigate count**: `MAX_AUTO_IMAGES` is now enforced against Frigate's live training image count rather than the local uploaded-asset tally. Fallback priority: live Frigate API → last cached `frigate_count` in JSON → local uploaded count. +- **Startup summary shows Frigate count**: Tracker summary at startup now includes the last known Frigate training count per person (e.g. `78 uploaded, 2 rejected, 42 in Frigate`). +- **`winnow/frigate_api.py`**: new module encapsulating Frigate API helpers; currently exposes `get_frigate_face_counts()`. + +### Changed + +- `upload_tracker.py`: `by_person` entries migrated from flat list to `{asset_ids, scores, frigate_count}` dict. Old list format is read and migrated transparently on first write. +- `mark_uploaded()` now accepts an optional `score` keyword argument. +- `get_person_summary()` now returns `frigate_count` and `scores` fields alongside `uploaded` and `rejected`. + ## [0.2.1] - 2026-06-12 ### Fixed diff --git a/pyproject.toml b/pyproject.toml index 9b873fb..65ce249 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "winnow" -version = "0.2.1" +version = "0.2.2" description = "Immich to Frigate training sets" license = "MIT" requires-python = ">=3.12" diff --git a/winnow/cli.py b/winnow/cli.py index 56a2eb2..d052595 100644 --- a/winnow/cli.py +++ b/winnow/cli.py @@ -48,7 +48,15 @@ def main() -> None: if summary: rprint("\n[dim]Tracker summary:[/dim]") for person_name, counts in summary.items(): - rprint(f" [dim]{person_name}: {counts['uploaded']} uploaded, {counts['rejected']} rejected[/dim]") + frigate_part = ( + f", {counts['frigate_count']} in Frigate" + if counts.get("frigate_count") is not None + else "" + ) + rprint( + f" [dim]{person_name}: {counts['uploaded']} uploaded," + f" {counts['rejected']} rejected{frigate_part}[/dim]" + ) people = get_people() if not people: diff --git a/winnow/executor.py b/winnow/executor.py index c59ce95..5f9e85e 100644 --- a/winnow/executor.py +++ b/winnow/executor.py @@ -58,6 +58,7 @@ def _enrich_asset_with_face_data(asset: dict, person: dict) -> dict: # Inject into asset so process_face_mode can find it via asset["people"] asset["people"] = [{"id": person_id, "faces": [face_info]}] + asset["face_confidence"] = face_data.confidence return asset @@ -96,8 +97,9 @@ def execute_jobs(jobs: list[dict]) -> None: shutil.rmtree(person_dir) os.makedirs(person_dir, exist_ok=True) - # Track filename → asset_id mapping for upload dedup + # Track filename → asset_id and filename → confidence score asset_map: dict[str, str] = {} + score_map: dict[str, float | None] = {} count = 0 for asset in assets: @@ -132,11 +134,13 @@ def execute_jobs(jobs: list[dict]) -> None: # Record which asset produced which output file filename = f"{count}.jpg" asset_map[filename] = asset["id"] + score_map[filename] = asset.get("face_confidence") # Also record object-mode variant filenames if mode == "object": for f in sorted(os.listdir(person_dir)): if f.startswith(f"{count}_") and f not in asset_map: asset_map[f] = asset["id"] + score_map[f] = asset.get("face_confidence") count += 1 else: @@ -149,8 +153,9 @@ def execute_jobs(jobs: list[dict]) -> None: progress.advance(job_task) progress.advance(overall_task) - # Store asset_map on the job so upload_to_frigate can use it + # Store maps on the job so upload_to_frigate can use them job["asset_map"] = asset_map + job["score_map"] = score_map progress.remove_task(job_task) @@ -230,6 +235,7 @@ def upload_to_frigate(jobs: list[dict]) -> None: continue asset_map = filename_to_asset_id.get(name, {}) + score_map = job.get("score_map", {}) person_files = sorted(asset_map.keys()) if not person_files: @@ -257,7 +263,7 @@ def upload_to_frigate(jobs: list[dict]) -> None: # Mark this asset as uploaded so it's skipped on future runs asset_id = asset_map.get(fname) if asset_id: - mark_uploaded(asset_id, person_name=name) + mark_uploaded(asset_id, person_name=name, score=score_map.get(fname)) break else: diff --git a/winnow/frigate_api.py b/winnow/frigate_api.py new file mode 100644 index 0000000..4409f05 --- /dev/null +++ b/winnow/frigate_api.py @@ -0,0 +1,28 @@ +"""Frigate API helpers for querying face training state.""" + +import logging +import os + +import requests + +logger = logging.getLogger(__name__) + + +def get_frigate_face_counts() -> dict[str, int] | None: + """Return {person_name: training_image_count} from Frigate's train directory. + + Returns None if FRIGATE_URL is not set or the API is unreachable, so callers + can distinguish "API unavailable" from "person has 0 images." + """ + frigate_url = os.environ.get("FRIGATE_URL", "").rstrip("/") + if not frigate_url: + return None + try: + resp = requests.get(f"{frigate_url}/api/faces", timeout=10) + resp.raise_for_status() + data = resp.json() + train = data.get("train", {}) + return {name: len(files) for name, files in train.items() if isinstance(files, list)} + except Exception as e: + logger.warning(f"Could not query Frigate face counts: {e}") + return None diff --git a/winnow/jobs.py b/winnow/jobs.py index 15cc1eb..55af9df 100644 --- a/winnow/jobs.py +++ b/winnow/jobs.py @@ -11,9 +11,10 @@ from rich.table import Table from .config import Config from .diversity import select_diverse_assets from .embeddings import is_embedding_available, load_embedding_model +from .frigate_api import get_frigate_face_counts from .immich_api import fetch_all_assets, filter_recent_assets from .logging import console -from .upload_tracker import filter_already_uploaded +from .upload_tracker import filter_already_uploaded, get_person_summary, update_frigate_count logger = logging.getLogger(__name__) @@ -236,6 +237,12 @@ def auto_configure(people: list[dict]) -> list[dict]: f" ≥{min_face_count} assets (MIN_FACE_COUNT={min_face_count})" ) + frigate_counts = get_frigate_face_counts() + # Persist each count to tracker so the last known value survives Frigate downtime + if frigate_counts is not None: + for pname, count in frigate_counts.items(): + update_frigate_count(pname, count) + upload_summary = get_person_summary() jobs = [] for person in valid_people: name = person["name"] @@ -263,9 +270,34 @@ def auto_configure(people: list[dict]) -> list[dict]: rprint(f" [dim]Skipping {name} (0 new images after dedup).[/dim]") continue + # Enforce MAX_AUTO_IMAGES as a lifetime cap per person. + # Priority: live Frigate count → last cached Frigate count → local uploaded count. + person_summary = upload_summary.get(name, {}) + if frigate_counts is not None: + already_uploaded = frigate_counts.get(name, 0) + else: + already_uploaded = ( + person_summary.get("frigate_count") + or person_summary.get("uploaded", 0) + ) + capacity = Config.MAX_AUTO_IMAGES - already_uploaded + if capacity <= 0: + rprint( + f" [dim]Skipping {name} (at lifetime cap:" + f" {already_uploaded}/{Config.MAX_AUTO_IMAGES} trained).[/dim]" + ) + continue + has_embedding = is_embedding_available(entity_type) limit, selection_mode = _resolve_strategy(strategy, has_embedding) + # Cap selection to remaining capacity + if limit == "auto": + if already_uploaded > 0: + limit = capacity # partially filled — select exactly what remains + else: + limit = min(limit, capacity) + if selection_mode == "skip": continue diff --git a/winnow/upload_tracker.py b/winnow/upload_tracker.py index 957ee28..0317639 100644 --- a/winnow/upload_tracker.py +++ b/winnow/upload_tracker.py @@ -8,6 +8,13 @@ Both are excluded from future candidate pools. To reset: - All: delete both files - One person: call reset_person("Name") or set RESET_PERSON=Name - Rejects only: delete frigate_rejected_ids.json, or set RETRY_REJECTED=true + +by_person schema (frigate_uploaded_ids.json): + { + "asset_ids": ["immich-id-1", ...], # all assets we attempted to upload + "scores": {"immich-id-1": 0.953}, # Immich face confidence at upload time + "frigate_count": 42 # last known Frigate training image count + } """ import json @@ -55,7 +62,23 @@ def _load_flat(filename: str) -> set[str]: return set(_load(filename).get(_flat_key(filename), [])) -def _mark(filename: str, asset_id: str, person_name: str | None) -> None: +def _get_ids(entry: list | dict) -> list[str]: + """Extract asset_ids from either the old list format or the new dict format.""" + if isinstance(entry, list): + return entry + return entry.get("asset_ids", []) + + +def _migrate_entry(entry: list | dict) -> dict: + """Ensure by_person entry is in the current dict format.""" + if isinstance(entry, list): + return {"asset_ids": sorted(entry), "scores": {}} + entry.setdefault("asset_ids", []) + entry.setdefault("scores", {}) + return entry + + +def _mark(filename: str, asset_id: str, person_name: str | None, score: float | None = None) -> None: data = _load(filename) flat_key = _flat_key(filename) flat = set(data.get(flat_key, [])) @@ -63,9 +86,13 @@ def _mark(filename: str, asset_id: str, person_name: str | None) -> None: data[flat_key] = sorted(flat) if person_name: by_person = data.setdefault("by_person", {}) - person_ids = set(by_person.get(person_name, [])) - person_ids.add(asset_id) - by_person[person_name] = sorted(person_ids) + entry = _migrate_entry(by_person.get(person_name, {})) + ids = set(entry["asset_ids"]) + ids.add(asset_id) + entry["asset_ids"] = sorted(ids) + if score is not None: + entry["scores"][asset_id] = round(score, 4) + by_person[person_name] = entry _save(filename, data) @@ -79,8 +106,8 @@ def load_rejected_ids() -> set[str]: return _load_flat(REJECT_TRACKER_FILE) -def mark_uploaded(asset_id: str, person_name: str | None = None) -> None: - _mark(UPLOAD_TRACKER_FILE, asset_id, person_name) +def mark_uploaded(asset_id: str, person_name: str | None = None, score: float | None = None) -> None: + _mark(UPLOAD_TRACKER_FILE, asset_id, person_name, score=score) logger.debug(f"Marked {asset_id} as uploaded ({person_name})") @@ -89,14 +116,25 @@ def mark_rejected(asset_id: str, person_name: str | None = None) -> None: logger.debug(f"Marked {asset_id} as rejected ({person_name})") +def update_frigate_count(person_name: str, count: int) -> None: + """Record Frigate's authoritative training image count for a person.""" + data = _load(UPLOAD_TRACKER_FILE) + by_person = data.setdefault("by_person", {}) + entry = _migrate_entry(by_person.get(person_name, {})) + entry["frigate_count"] = count + by_person[person_name] = entry + _save(UPLOAD_TRACKER_FILE, data) + + def reset_person(person_name: str) -> None: """Remove all uploaded and rejected records for a given person.""" for filename in (UPLOAD_TRACKER_FILE, REJECT_TRACKER_FILE): data = _load(filename) flat_key = _flat_key(filename) by_person = data.get("by_person", {}) - person_ids = set(by_person.pop(person_name, [])) - if person_ids: + entry = by_person.pop(person_name, None) + if entry is not None: + person_ids = set(_get_ids(entry)) flat = set(data.get(flat_key, [])) - person_ids data[flat_key] = sorted(flat) data["by_person"] = by_person @@ -104,18 +142,22 @@ def reset_person(person_name: str) -> None: logger.info(f"Reset tracking data for {person_name}") -def get_person_summary() -> dict[str, dict[str, int]]: - """Return {person_name: {uploaded: N, rejected: N}} for display.""" - uploaded_by = _load(UPLOAD_TRACKER_FILE).get("by_person", {}) - rejected_by = _load(REJECT_TRACKER_FILE).get("by_person", {}) - names = set(uploaded_by) | set(rejected_by) - return { - name: { - "uploaded": len(uploaded_by.get(name, [])), - "rejected": len(rejected_by.get(name, [])), +def get_person_summary() -> dict[str, dict]: + """Return {person_name: {uploaded, rejected, frigate_count, scores}} for display/capacity.""" + uploaded_data = _load(UPLOAD_TRACKER_FILE).get("by_person", {}) + rejected_data = _load(REJECT_TRACKER_FILE).get("by_person", {}) + names = set(uploaded_data) | set(rejected_data) + result = {} + for name in sorted(names): + u_entry = uploaded_data.get(name, {}) + r_entry = rejected_data.get(name, {}) + result[name] = { + "uploaded": len(_get_ids(u_entry)), + "rejected": len(_get_ids(r_entry)), + "frigate_count": u_entry.get("frigate_count") if isinstance(u_entry, dict) else None, + "scores": u_entry.get("scores", {}) if isinstance(u_entry, dict) else {}, } - for name in sorted(names) - } + return result def filter_already_uploaded(