From a6bae5da0503e13a523279093227419e0a57a1f8 Mon Sep 17 00:00:00 2001 From: Holden Date: Tue, 16 Jun 2026 16:16:28 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20address=2010=20audit=20findings=20?= =?UTF-8?q?=E2=80=94=20import=20bug,=20fscore=20stale=20flag,=20cache=20mu?= =?UTF-8?q?tation,=20batch=20safety,=20falsy=20guards?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- winnow/cli.py | 7 ++--- winnow/executor.py | 2 ++ winnow/upload_tracker.py | 57 ++++++++++++++++++++++++---------------- 3 files changed, 41 insertions(+), 25 deletions(-) diff --git a/winnow/cli.py b/winnow/cli.py index c7456d4..2ce8d60 100644 --- a/winnow/cli.py +++ b/winnow/cli.py @@ -86,6 +86,8 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: for p in sorted(ps, key=lambda x: x.get("assetCount", 0), reverse=True)[1:] } + skip_ids = _smaller_duplicate_ids(duplicates) + if not Config.MERGE_DUPLICATE_PEOPLE: rprint("\n[bold yellow]⚠ Duplicate person names detected in Immich:[/bold yellow]") for name, ps in sorted(duplicates.items()): @@ -107,7 +109,7 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: ) # Return deduplicated list — keep only the largest per name so that # downstream job creation never runs two jobs for the same Frigate folder. - return [p for p in people if p["id"] not in _smaller_duplicate_ids(duplicates)] + return [p for p in people if p["id"] not in skip_ids] # Auto-merge: survivor = largest asset count, rest merge into it inside Immich merged_any = False @@ -133,7 +135,6 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: # IDs still exist in Immich and would produce two jobs for the same folder. # IDs from groups that merged successfully are already gone from Immich, so # this filter is a no-op for them. - skip_ids = _smaller_duplicate_ids(duplicates) return [p for p in fresh if p.get("id") not in skip_ids] # All merges failed — fall back to local deduplication (keep largest per name) so @@ -142,7 +143,7 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: " [yellow]All merges failed — applying local deduplication" " to avoid overwriting output.[/yellow]" ) - return [p for p in people if p["id"] not in _smaller_duplicate_ids(duplicates)] + return [p for p in people if p["id"] not in skip_ids] _UNSUPPORTED_VARS = [ diff --git a/winnow/executor.py b/winnow/executor.py index 205eac8..656148b 100644 --- a/winnow/executor.py +++ b/winnow/executor.py @@ -27,6 +27,7 @@ from .log_config import console from .quality import blur_score_from_image from .reconcile import enrich_asset_with_face_data, reconcile_frigate_mappings from .upload_tracker import ( + UPLOAD_TRACKER_FILE, get_lowest_quality_mapped_file, get_most_redundant_mapped_file, get_tracked_frigate_file_count, @@ -481,6 +482,7 @@ def upload_to_frigate(jobs: list[dict]) -> None: ) if delete_frigate_person_files(name, [target_frigate_file]): remove_frigate_file(name, target_frigate_file) + person_has_fscores = has_frigate_scores(name) effective_count -= 1 min_quality_score_for_slot = None if using_fscore else target_score else: diff --git a/winnow/upload_tracker.py b/winnow/upload_tracker.py index 5a3a85e..8994ab1 100644 --- a/winnow/upload_tracker.py +++ b/winnow/upload_tracker.py @@ -32,7 +32,7 @@ import logging import os from pathlib import Path -from .frigate_api import delete_frigate_person_files +from .frigate_api import _get_frigate_url, delete_frigate_person_files logger = logging.getLogger(__name__) @@ -95,8 +95,18 @@ def _save(filename: str, data: dict) -> None: def begin_batch(filename: str) -> None: """Defer tracker disk writes for filename. All _save calls accumulate in the in-memory cache until flush_batch() is called. Use around per-person upload loops - to reduce N writes to 1.""" - _deferred.add(str(_tracker_path(filename))) + to reduce N writes to 1. + + If a previous batch for this file was interrupted before flush_batch() was called + (e.g. an exception escaped the upload loop), the leftover cache state is flushed + to disk here before starting fresh so that partial progress is not silently lost. + """ + path = _tracker_path(filename) + key = str(path) + if key in _deferred and key in _cache: + _write_to_disk(path, _cache[key]) + _deferred.discard(key) + _deferred.add(key) def flush_batch(filename: str) -> None: @@ -141,20 +151,22 @@ def _mark( crop_dims: tuple[int, int] | None = None, frigate_score: float | None = None, ) -> None: + if not person_name: + logger.warning("_mark called with empty person_name for asset %s — asset not recorded", asset_id) + return data = _load(filename) - if person_name: - by_person = data.setdefault("by_person", {}) - 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) - if crop_dims is not None: - entry["crop_dims"][asset_id] = [crop_dims[0], crop_dims[1]] - if frigate_score is not None: - entry["frigate_scores"][asset_id] = round(frigate_score, 4) - by_person[person_name] = entry + by_person = data.setdefault("by_person", {}) + 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) + if crop_dims is not None: + entry["crop_dims"][asset_id] = [crop_dims[0], crop_dims[1]] + if frigate_score is not None: + entry["frigate_scores"][asset_id] = round(frigate_score, 4) + by_person[person_name] = entry _save(filename, data) @@ -233,7 +245,7 @@ def remove_frigate_files_batch(person_name: str, frigate_filenames: list[str]) - entry = _migrate_entry(raw) for fn in frigate_filenames: asset_id = entry["frigate_files"].pop(fn, None) - if asset_id: + if asset_id is not None: entry["frigate_scores"].pop(asset_id, None) by_person[person_name] = entry _save(UPLOAD_TRACKER_FILE, data) @@ -366,7 +378,7 @@ def reset_all_people() -> None: frigate_filenames = list(entry.get("frigate_files", {}).keys()) if not frigate_filenames: continue - if not os.environ.get("FRIGATE_URL", "").strip(): + if not _get_frigate_url(): logger.info(f"FRIGATE_URL not set — skipping Frigate file deletion for {person_name}") elif delete_frigate_person_files(person_name, frigate_filenames): logger.info(f"Deleted {len(frigate_filenames)} Frigate file(s) for {person_name}") @@ -389,7 +401,7 @@ def reset_person(person_name: str) -> None: entry = _migrate_entry(upload_data.get("by_person", {}).get(person_name, {})) frigate_filenames = list(entry.get("frigate_files", {}).keys()) if frigate_filenames: - if not os.environ.get("FRIGATE_URL", "").strip(): + if not _get_frigate_url(): logger.info(f"FRIGATE_URL not set — skipping Frigate file deletion for {person_name}") elif delete_frigate_person_files(person_name, frigate_filenames): logger.info(f"Deleted {len(frigate_filenames)} Frigate file(s) for {person_name}") @@ -398,15 +410,16 @@ def reset_person(person_name: str) -> None: changed = False for filename in (UPLOAD_TRACKER_FILE, REJECT_TRACKER_FILE): - data = upload_data if filename == UPLOAD_TRACKER_FILE else _load(REJECT_TRACKER_FILE) - by_person = data.get("by_person", {}) + src = upload_data if filename == UPLOAD_TRACKER_FILE else _load(REJECT_TRACKER_FILE) + by_person = dict(src.get("by_person", {})) # copy so pop() does not mutate the cache tracker_entry = by_person.pop(person_name, None) if tracker_entry is not None: + data = dict(src) + data["by_person"] = by_person flat_key = _flat_key(filename) person_ids = set(_get_ids(tracker_entry)) if person_ids and flat_key in data: data[flat_key] = sorted(set(data[flat_key]) - person_ids) - data["by_person"] = by_person _save(filename, data) changed = True if changed: