From dc2efb5ac4a921538f02495db66644eccfed378c Mon Sep 17 00:00:00 2001 From: Holden Date: Tue, 16 Jun 2026 15:21:12 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20v0.6.1=20=E2=80=94=20tracker=20integrity?= =?UTF-8?q?,=20quality=20replacement=20correctness,=20code=20review=20fixe?= =?UTF-8?q?s?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Catch OSError alongside PIL.UnidentifiedImageError for corrupt thumbnails - Fix quality replacement mode flip mid-loop (person_has_fscores no longer re-evaluated) - reset_person rebuilds flat list from remaining entries instead of subtracting - _save cache updated only after os.replace succeeds (prevents cache/disk split-brain) - Stale Frigate file cleanup uses remove_frigate_files_batch (N writes → 1) - _migrate_entry deep-copies nested dicts so .pop() cannot mutate the cache - find_by_crop_dimension and _pick_mapped_file consistent on duplicate asset→file mapping - Atomic JSON write (tmp + os.replace) guards against truncated files on crash - get_person_summary uses _migrate_entry instead of three isinstance guards - Quality floor check allows None-scored candidates through (don't block freed slots) - Fix comment-only if body (IndentationError on import) in full-res download path - Merge duplicate if-stale guard into one block - _flat_key uses constant equality instead of substring match - remove_frigate_file returns early when person absent (no ghost entries) - skip_ids extracted to _smaller_duplicate_ids() helper (was duplicated 3×) - blur_score_from_image returns None on error instead of 0.0 --- CHANGELOG.md | 36 +++++++++++++++ pyproject.toml | 2 +- winnow/cli.py | 28 +++++------- winnow/executor.py | 40 ++++++++--------- winnow/quality.py | 7 +-- winnow/upload_tracker.py | 96 ++++++++++++++++++++++++---------------- 6 files changed, 128 insertions(+), 81 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a42a83f..558f920 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,42 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [0.6.1] - 2026-06-16 + +### Fixed + +- **Corrupt or truncated full-res thumbnails now marked rejected**: `OSError` (truncated file) is caught alongside `PIL.UnidentifiedImageError` in the thumbnail path so persistently bad assets are tombstoned instead of retried forever. Full-res download failures (`USE_FULL_RESOLUTION=true`) remain transient — not marked rejected — so a Immich blip doesn't permanently blacklist valid assets. + +- **Quality replacement mode no longer flips mid-loop**: `person_has_fscores` was re-evaluated after each file deletion, which could switch the remaining replacements from Frigate-score mode to blur-score mode if the deleted file was the last scored one. The mode is now fixed for the duration of the upload loop. + +- **`reset_person` no longer removes shared asset IDs**: the flat `uploaded_asset_ids` list is now rebuilt from all remaining `by_person` entries rather than subtracting the reset person's IDs. Previously, resetting Alice could remove an asset ID that also appeared under Bob, making it re-eligible for upload. + +- **`_save` cache updated only after successful write**: the in-memory tracker cache is now updated after `os.replace` succeeds rather than before. A disk-full or permission error no longer leaves the cache permanently ahead of the on-disk file. + +- **Stale Frigate file cleanup batched**: the per-file `remove_frigate_file` loop is replaced with a single `remove_frigate_files_batch` call, reducing N tracker writes to 1 when stale mappings are cleaned up. + +- **`_migrate_entry` no longer mutates the cache through nested dict aliases**: all five nested dicts (`asset_ids`, `scores`, `frigate_scores`, `frigate_files`, `crop_dims`) are now individually copied so `.pop()` calls in write paths cannot reach the in-memory cache. + +- **`find_by_crop_dimension` and `_pick_mapped_file` now agree on duplicate asset→file handling**: both use first-seen-wins when the same `asset_id` maps to multiple Frigate filenames, preventing inconsistent replacement decisions. + +- **Non-atomic JSON write**: tracker files are written to a `.tmp` sibling then renamed with `os.replace` so a crash mid-write never leaves a truncated file. + +- **`get_person_summary` uses `_migrate_entry`**: replaced three ad-hoc `isinstance` guards with a single `_migrate_entry` call, making old-format (list) entries consistent with every other read path. + +- **Quality replacement floor check**: a candidate with a `None` blur score (PIL error during scoring) no longer blocks a freed slot — the `<=` floor comparison is only applied when a score is actually available. + +- **`executor.py` syntax error**: the `if img is None:` block in the full-res download path was comment-only and would have raised `IndentationError` on import. Added `pass`. + +- **Duplicate `if stale:` guard**: two consecutive identical guards around stale-cleanup and its log print were merged into one. + +- **`_flat_key` uses constant equality** instead of substring match, removing a latent routing bug for any filename that happens to contain "uploaded". + +- **`remove_frigate_file` no longer creates ghost entries**: returns early when the person is absent rather than writing an empty stub. + +- **`skip_ids` extracted to helper**: the identical set comprehension in `_handle_duplicate_people` that appeared in three branches is now a single `_smaller_duplicate_ids()` inner function. + +- **`blur_score_from_image` returns `None` on error** instead of `0.0`, so callers can distinguish a failed measurement from a legitimately near-zero Laplacian variance score. + ## [0.6.0] - 2026-06-15 ### Changed diff --git a/pyproject.toml b/pyproject.toml index e8a2694..2889a90 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "winnow" -version = "0.6.0" +version = "0.6.1" description = "Selects diverse, high-quality photos from Immich as training data for Frigate face recognition." license = "AGPL-3.0-or-later" requires-python = ">=3.13" diff --git a/winnow/cli.py b/winnow/cli.py index 945e96f..986960b 100644 --- a/winnow/cli.py +++ b/winnow/cli.py @@ -78,6 +78,14 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: if not duplicates: return people + def _smaller_duplicate_ids(groups: dict) -> set[str]: + """IDs of all but the largest person in each duplicate group.""" + return { + p["id"] + for ps in groups.values() + for p in sorted(ps, key=lambda x: x.get("assetCount", 0), reverse=True)[1:] + } + 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()): @@ -99,12 +107,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. - skip_ids = { - p["id"] - for ps in duplicates.values() - for p in sorted(ps, key=lambda x: x.get("assetCount", 0), reverse=True)[1:] - } - return [p for p in people if p["id"] not in skip_ids] + return [p for p in people if p["id"] not in _smaller_duplicate_ids(duplicates)] # Auto-merge: survivor = largest asset count, rest merge into it inside Immich merged_any = False @@ -130,11 +133,7 @@ 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 = { - p["id"] - for ps in duplicates.values() - for p in sorted(ps, key=lambda x: x.get("assetCount", 0), reverse=True)[1:] - } + 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 @@ -143,12 +142,7 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]: " [yellow]All merges failed — applying local deduplication" " to avoid overwriting output.[/yellow]" ) - skip_ids = { - p["id"] - for ps in duplicates.values() - for p in sorted(ps, key=lambda x: x.get("assetCount", 0), reverse=True)[1:] - } - return [p for p in people if p["id"] not in skip_ids] + return [p for p in people if p["id"] not in _smaller_duplicate_ids(duplicates)] _UNSUPPORTED_VARS = [ diff --git a/winnow/executor.py b/winnow/executor.py index 03c885a..10afd60 100644 --- a/winnow/executor.py +++ b/winnow/executor.py @@ -35,6 +35,7 @@ from .upload_tracker import ( mark_rejected, mark_uploaded, remove_frigate_file, + remove_frigate_files_batch, ) logger = logging.getLogger(__name__) @@ -151,9 +152,9 @@ def execute_jobs(jobs: list[dict]) -> None: if use_full_res: img = fetch_full_image(asset["id"]) if img is None: - # Both original and preview fallback failed — mark rejected - # so this asset isn't retried on every future run. - mark_rejected(asset["id"], person_name=name) + # Full-res download failed — could be a transient network + # error, so don't mark rejected; it will be retried next run. + pass else: resp = requests.get( f"{Config.IMMICH_URL}/api/assets/{asset['id']}/thumbnail?size=preview&format=JPEG", @@ -163,11 +164,11 @@ def execute_jobs(jobs: list[dict]) -> None: if resp.ok: try: img = Image.open(BytesIO(resp.content)) - except PIL.UnidentifiedImageError: - # Pillow cannot identify the format — genuinely corrupt - # Immich thumbnail. Mark rejected so this asset isn't - # retried indefinitely. OSError/truncation errors are - # transient and intentionally not caught here. + except (PIL.UnidentifiedImageError, OSError): + # Pillow cannot identify the format or the content is + # truncated. The download already succeeded (resp.ok), + # so this is a data problem, not a transient network + # error — mark rejected so it isn't retried forever. logger.warning("Invalid image data for asset %s — marking rejected", asset["id"]) mark_rejected(asset["id"], person_name=name) img = None @@ -345,9 +346,8 @@ def upload_to_frigate(jobs: list[dict]) -> None: # (manually deleted, or cleaned up outside winnow). This corrects the # effective_count so those slots are available for new uploads. stale = get_tracked_frigate_filenames(name) - known_frigate_files_at_start - for stale_fn in stale: - remove_frigate_file(name, stale_fn) if stale: + remove_frigate_files_batch(name, list(stale)) progress.console.print( f" [dim]{name}: cleared {len(stale)} stale mapping(s)" " (file(s) no longer in Frigate)[/dim]" @@ -372,10 +372,9 @@ def upload_to_frigate(jobs: list[dict]) -> None: # freed slot isn't filled with something worse than what we removed. if min_quality_score_for_slot is not None: file_score = score_map.get(fname) - if file_score is None or file_score <= min_quality_score_for_slot: - score_str = f"{file_score:.3f}" if file_score is not None else "N/A" + if file_score is not None and file_score <= min_quality_score_for_slot: progress.console.print( - f" [dim]⏭ {fname}: score {score_str} ≤ freed slot floor" + f" [dim]⏭ {fname}: score {file_score:.3f} ≤ freed slot floor" f" {min_quality_score_for_slot:.3f}, skipping[/dim]" ) progress.advance(upload_task) @@ -446,11 +445,13 @@ def upload_to_frigate(jobs: list[dict]) -> None: get_target = get_most_redundant_mapped_file score_label, better_note = "frigate", " (more novel)" no_score_msg = "Frigate recognize unavailable, skipping replacement" + is_better_than = lambda c, t: c < t else: candidate_score = score_map.get(fname) get_target = get_lowest_quality_mapped_file score_label, better_note = "blur", "" no_score_msg = "no quality score, skipping replacement" + is_better_than = lambda c, t: c > t if candidate_score is None: progress.console.print(f" [dim]⏭ {fname}: {no_score_msg}[/dim]") @@ -458,28 +459,25 @@ def upload_to_frigate(jobs: list[dict]) -> None: continue target = get_target(name, exclude=failed_deletes) - not_better = target is None or ( - candidate_score >= target[2] if using_fscore else candidate_score <= target[2] - ) + not_better = target is None or not is_better_than(candidate_score, target[2]) if not_better: target_str = f"{target[2]:.3f}" if target is not None else "N/A" - op = "<" if using_fscore else ">" + cmp_op = "<" if using_fscore else ">" progress.console.print( f" [dim]⏭ {fname}: {score_label} {candidate_score:.3f}" - f" not {op} {target_str}, skipping[/dim]" + f" not {cmp_op} {target_str}, skipping[/dim]" ) progress.advance(upload_task) continue target_frigate_file, _target_asset_id, target_score = target - op = "<" if using_fscore else ">" + cmp_op = "<" if using_fscore else ">" progress.console.print( - f" 🔄 {fname}: {score_label} {candidate_score:.3f} {op} {target_score:.3f}," + f" 🔄 {fname}: {score_label} {candidate_score:.3f} {cmp_op} {target_score:.3f}," f" replacing {target_frigate_file}{better_note}" ) 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/quality.py b/winnow/quality.py index b4e1724..9247f88 100644 --- a/winnow/quality.py +++ b/winnow/quality.py @@ -139,14 +139,15 @@ def assess_quality( return QualityResult(passed=len(reasons) == 0, reasons=reasons, blur_score=blur_score) -def blur_score_from_image(img: Image.Image, max_dim: int = 1440) -> float: +def blur_score_from_image(img: Image.Image, max_dim: int = 1440) -> float | None: """Compute Laplacian-variance blur score, capped at max_dim px to normalise scale. Caps resolution so full-res and thumbnail scores are comparable — Laplacian variance grows with pixel count, making uncapped full-res scores much larger than thumbnail scores for the same perceived sharpness. - Returns 0.0 on any error so callers can treat the result as lowest quality. + Returns None on error so callers can distinguish a failed measurement from a + legitimately low (near-zero) score. """ try: score_img = img.convert("RGB") if img.mode != "RGB" else img @@ -156,5 +157,5 @@ def blur_score_from_image(img: Image.Image, max_dim: int = 1440) -> float: return float(assess_quality(score_img).blur_score) except Exception as exc: logger.debug("blur_score_from_image failed: %s", exc) - return 0.0 + return None diff --git a/winnow/upload_tracker.py b/winnow/upload_tracker.py index 076a157..abf2b55 100644 --- a/winnow/upload_tracker.py +++ b/winnow/upload_tracker.py @@ -71,14 +71,20 @@ def _load(filename: str) -> dict: def _save(filename: str, data: dict) -> None: path = _tracker_path(filename) - _cache[str(path)] = data # keep cache consistent with what we write path.parent.mkdir(parents=True, exist_ok=True) - with open(path, "w") as f: - json.dump(data, f, indent=2) + tmp = path.with_suffix(".tmp") + try: + with open(tmp, "w") as f: + json.dump(data, f, indent=2) + os.replace(tmp, path) + _cache[str(path)] = data # update only after the file is safely on disk + except Exception: + tmp.unlink(missing_ok=True) + raise def _flat_key(filename: str) -> str: - return "uploaded_asset_ids" if "uploaded" in filename else "rejected_asset_ids" + return "uploaded_asset_ids" if filename == UPLOAD_TRACKER_FILE else "rejected_asset_ids" def _load_flat(filename: str) -> set[str]: @@ -96,12 +102,14 @@ 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": {}, "frigate_scores": {}, "frigate_files": {}, "crop_dims": {}} - entry.setdefault("asset_ids", []) - entry.setdefault("scores", {}) - entry.setdefault("frigate_scores", {}) - entry.setdefault("frigate_files", {}) - entry.setdefault("crop_dims", {}) - return entry + # Copy top-level and all nested dicts so callers' mutations never reach the cache. + result = dict(entry) + result["asset_ids"] = list(result.get("asset_ids", [])) + result["scores"] = dict(result.get("scores", {})) + result["frigate_scores"] = dict(result.get("frigate_scores", {})) + result["frigate_files"] = dict(result.get("frigate_files", {})) + result["crop_dims"] = dict(result.get("crop_dims", {})) + return result def _mark( @@ -160,15 +168,10 @@ def mark_rejected(asset_id: str, person_name: str | None = None) -> None: + def record_frigate_file(person_name: str, frigate_filename: str, asset_id: str) -> None: - """Record the mapping from a Frigate training filename to an Immich asset ID.""" - data = _load(UPLOAD_TRACKER_FILE) - by_person = data.setdefault("by_person", {}) - entry = _migrate_entry(by_person.get(person_name, {})) - entry["frigate_files"][frigate_filename] = asset_id - by_person[person_name] = entry - _save(UPLOAD_TRACKER_FILE, data) - logger.debug(f"Mapped Frigate file {frigate_filename} → {asset_id} ({person_name})") + """Record a single Frigate filename → asset_id mapping.""" + record_frigate_files_batch(person_name, {frigate_filename: asset_id}) def record_frigate_files_batch(person_name: str, mappings: dict[str, str]) -> None: @@ -190,15 +193,24 @@ def remove_frigate_file(person_name: str, frigate_filename: str) -> None: Does NOT unmark the source asset_id — the deletion was deliberate and we don't want to re-upload the inferior image on the next run. """ + remove_frigate_files_batch(person_name, [frigate_filename]) + + +def remove_frigate_files_batch(person_name: str, frigate_filenames: list[str]) -> None: + """Remove multiple Frigate filenames in a single load/save.""" data = _load(UPLOAD_TRACKER_FILE) by_person = data.get("by_person", {}) - entry = _migrate_entry(by_person.get(person_name, {})) - asset_id = entry["frigate_files"].pop(frigate_filename, None) - if asset_id: - entry["frigate_scores"].pop(asset_id, None) + raw = by_person.get(person_name) + if raw is None: + return + entry = _migrate_entry(raw) + for fn in frigate_filenames: + asset_id = entry["frigate_files"].pop(fn, None) + if asset_id: + entry["frigate_scores"].pop(asset_id, None) by_person[person_name] = entry _save(UPLOAD_TRACKER_FILE, data) - logger.debug(f"Removed Frigate file mapping {frigate_filename} ({person_name})") + logger.debug(f"Removed {len(frigate_filenames)} Frigate file mapping(s) for {person_name}") def get_tracked_frigate_file_count(person_name: str) -> int: @@ -238,11 +250,12 @@ def _pick_mapped_file( data = _load(UPLOAD_TRACKER_FILE) entry = _migrate_entry(data.get("by_person", {}).get(person_name, {})) scores = entry.get(score_key, {}) - candidates = [ - (ff, asset_id, scores[asset_id]) - for ff, asset_id in entry.get("frigate_files", {}).items() - if (exclude is None or ff not in exclude) and asset_id in scores - ] + seen_assets: set[str] = set() + candidates = [] + for ff, asset_id in entry.get("frigate_files", {}).items(): + if (exclude is None or ff not in exclude) and asset_id in scores and asset_id not in seen_assets: + seen_assets.add(asset_id) + candidates.append((ff, asset_id, scores[asset_id])) if not candidates: return None return max(candidates, key=lambda x: x[2]) if highest else min(candidates, key=lambda x: x[2]) @@ -285,7 +298,9 @@ def find_by_crop_dimension(size: int) -> list[dict]: entry = _migrate_entry(raw_entry) scores = entry.get("scores", {}) frigate_files = entry.get("frigate_files", {}) - asset_to_frigate = {v: k for k, v in frigate_files.items()} + asset_to_frigate: dict[str, str] = {} + for fn, aid in frigate_files.items(): + asset_to_frigate.setdefault(aid, fn) # first-seen wins; plain inversion silently drops duplicates frigate_scores = entry.get("frigate_scores", {}) for asset_id, dims in entry.get("crop_dims", {}).items(): w, h = dims[0], dims[1] @@ -338,9 +353,12 @@ def reset_person(person_name: str) -> None: by_person = data.get("by_person", {}) tracker_entry = by_person.pop(person_name, None) if tracker_entry is not None: - person_ids = set(_get_ids(tracker_entry)) - flat = set(data.get(flat_key, [])) - person_ids - data[flat_key] = sorted(flat) + # Rebuild from remaining entries rather than subtracting, so IDs that + # appear under another person aren't incorrectly removed from the flat list. + remaining_ids: set[str] = set() + for other_entry in by_person.values(): + remaining_ids.update(_get_ids(other_entry)) + data[flat_key] = sorted(remaining_ids) data["by_person"] = by_person _save(filename, data) changed = True @@ -357,14 +375,14 @@ def get_person_summary() -> dict[str, dict]: 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, {}) + u_entry = _migrate_entry(uploaded_data.get(name, {})) + r_entry = _migrate_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 {}, - "frigate_files": u_entry.get("frigate_files", {}) if isinstance(u_entry, dict) else {}, + "uploaded": len(u_entry["asset_ids"]), + "rejected": len(r_entry["asset_ids"]), + "frigate_count": u_entry.get("frigate_count"), + "scores": u_entry["scores"], + "frigate_files": u_entry["frigate_files"], } return result