Compare commits

...
12 Commits
Author SHA1 Message Date
flan 3c0ef47fdc release: v0.6.3 2026-06-16 18:13:42 +00:00
flan cf7660595d chore: update lockfile 2026-06-16 18:13:42 +00:00
flan 14f759e960 fix: address 3 quality review findings — record_frigate_files_batch cache mutation, tracker_ok flag, LIMIT guard
- record_frigate_files_batch: copy-before-mutate so a write failure
  doesn't leave cache ahead of disk (same fix as remove_frigate_files_batch)
- executor: replace tracker_ok boolean with try/else
- jobs: collapse duplicate custom_limit is not None checks into one guard

Bump version to 0.6.3.
2026-06-16 18:13:36 +00:00
flan e8cb390fe4 fix: address 2 quality review findings — begin_batch dirty guard, LIMIT<=0 warning 2026-06-16 18:04:38 +00:00
flan 54b52b0a73 fix: address 3 quality review findings — batch reject tracker, skip flush when clean, hoist frigate url check 2026-06-16 17:55:34 +00:00
flan b622e58f1b fix: address 3 quality review findings — flush_batch finally guard, _laplacian_var helper, has_frigate_scores no-copy 2026-06-16 17:09:00 +00:00
flan f3622b8d41 fix: address 4 quality review findings — flush_batch order, batch finally guard, cache copy, LIMIT=0 fallthrough 2026-06-16 16:59:43 +00:00
flan 817fa17e41 fix: address 3 quality review findings — tracker_ok gate, LIMIT=0 warning, cache write log level 2026-06-16 16:41:32 +00:00
flan 8846a4f1df fix: address 3 post-fix audit findings — begin_batch flush guard, misleading debug log, shared asset_id score deletion 2026-06-16 16:19:54 +00:00
flan a6bae5da05 fix: address 10 audit findings — import bug, fscore stale flag, cache mutation, batch safety, falsy guards 2026-06-16 16:16:28 +00:00
flan 4cdd4657d6 fix: v0.6.2 — structural tracker refactor, batch writes, multi-instance prep
- Drop flat list as primary storage; derive uploaded/rejected IDs from by_person
  (single source of truth). Legacy flat lists in existing files still read for
  backward compat. Removes dual-representation sync hazard.
- Add begin_batch/flush_batch: per-person upload loop now does 1 os.replace
  instead of N (one per mark_uploaded call). Benefit on slow storage.
- reset_all_people(): RESET_PERSON=* is now O(1) disk writes instead of O(P^2).
- blur_score_from_image inlines cv2.Laplacian directly, removing assess_quality
  call overhead and decoupling from the full quality pipeline.
2026-06-16 15:44:41 +00:00
flan dc2efb5ac4 fix: v0.6.1 — tracker integrity, quality replacement correctness, code review fixes
- 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
2026-06-16 15:21:12 +00:00
9 changed files with 496 additions and 329 deletions
+58
View File
@@ -7,6 +7,64 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased] ## [Unreleased]
## [0.6.3] - 2026-06-16
### Fixed
- **`record_frigate_files_batch` no longer mutates the tracker cache before write**: the function shared the same cache-corruption-on-write-failure bug that was fixed in `remove_frigate_files_batch` in v0.6.1 — `data.setdefault("by_person", {})` mutated the cached dict in-place, so a disk-full or permission error left the in-memory cache ahead of the on-disk file. Now uses the same copy-before-mutate pattern (shallow copies of the top-level dict and `by_person` sub-dict) so a failed write leaves cache and disk in sync.
- **`tracker_ok` boolean flag replaced with try/else**: the intermediate boolean was a misleading placeholder — the `True` initial value suggested success before the operation ran. The control flow is now expressed directly with a try/except/else block.
- **`LIMIT` env var guard simplified**: the two adjacent `if custom_limit is not None` checks in `_resolve_strategy` are collapsed into a single `if custom_limit is not None:` with nested branches, removing redundant evaluation.
## [0.6.2] - 2026-06-16
### Changed
- **Flat `uploaded_asset_ids` / `rejected_asset_ids` lists dropped as primary storage**: asset IDs are now derived on read from `by_person` entries, which are the single source of truth. The legacy flat lists in existing tracker files are still read (union) so no assets become re-eligible after upgrading. New writes no longer maintain the flat lists. This removes the dual-representation sync hazard and paves the way for multi-instance support (per-instance `by_person` keying in a future release).
- **Tracker writes batched per person**: `mark_uploaded` calls inside the per-person upload loop are now accumulated in memory (`begin_batch`) and flushed in a single `os.replace` write at the end of each person's loop (`flush_batch`), reducing N tracker writes per person to 1. Benefits users on slow storage (NAS, SD card, spinning disks).
- **`RESET_PERSON=*` is now O(1) disk writes**: replaced the per-person `reset_person` loop with `reset_all_people()`, which makes one Frigate API call per person for file deletion and then clears both tracker files in two writes. Previously it was O(P²) iterations and 2P writes.
- **`blur_score_from_image` inlines Laplacian computation**: replaced the `assess_quality()` call (which ran grayscale, exposure, and confidence checks whose results were discarded) with a direct `cv2.Laplacian` computation. The function is now self-contained and does not silently inherit future costs added to the full quality pipeline.
## [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 ## [0.6.0] - 2026-06-15
### Changed ### Changed
+1 -1
View File
@@ -1,6 +1,6 @@
[project] [project]
name = "winnow" name = "winnow"
version = "0.6.0" version = "0.6.3"
description = "Selects diverse, high-quality photos from Immich as training data for Frigate face recognition." description = "Selects diverse, high-quality photos from Immich as training data for Frigate face recognition."
license = "AGPL-3.0-or-later" license = "AGPL-3.0-or-later"
requires-python = ">=3.13" requires-python = ">=3.13"
Generated
+1 -1
View File
@@ -862,7 +862,7 @@ wheels = [
[[package]] [[package]]
name = "winnow" name = "winnow"
version = "0.6.0" version = "0.6.2"
source = { editable = "." } source = { editable = "." }
dependencies = [ dependencies = [
{ name = "croniter" }, { name = "croniter" },
+1 -1
View File
@@ -90,7 +90,7 @@ class EmbeddingCache:
np.save(tmp, embedding) np.save(tmp, embedding)
os.replace(tmp, final) os.replace(tmp, final)
except Exception as e: except Exception as e:
logger.debug("Cache write failed for %s: %s", asset_id, e) logger.warning("Cache write failed for %s: %s", asset_id, e)
try: try:
os.remove(tmp) os.remove(tmp)
except OSError: except OSError:
+12 -18
View File
@@ -12,7 +12,7 @@ from .executor import execute_jobs, upload_to_frigate
from .immich_api import get_immich_version, get_people, merge_people from .immich_api import get_immich_version, get_people, merge_people
from .jobs import _show_preview, auto_configure, interactive_configure from .jobs import _show_preview, auto_configure, interactive_configure
from .log_config import console, setup_logging from .log_config import console, setup_logging
from .upload_tracker import find_by_crop_dimension, get_person_summary, reset_person from .upload_tracker import find_by_crop_dimension, get_person_summary, reset_all_people, reset_person
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -78,6 +78,16 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]:
if not duplicates: if not duplicates:
return people 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:]
}
skip_ids = _smaller_duplicate_ids(duplicates)
if not Config.MERGE_DUPLICATE_PEOPLE: if not Config.MERGE_DUPLICATE_PEOPLE:
rprint("\n[bold yellow]⚠ Duplicate person names detected in Immich:[/bold yellow]") rprint("\n[bold yellow]⚠ Duplicate person names detected in Immich:[/bold yellow]")
for name, ps in sorted(duplicates.items()): for name, ps in sorted(duplicates.items()):
@@ -99,11 +109,6 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]:
) )
# Return deduplicated list — keep only the largest per name so that # Return deduplicated list — keep only the largest per name so that
# downstream job creation never runs two jobs for the same Frigate folder. # 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 skip_ids]
# Auto-merge: survivor = largest asset count, rest merge into it inside Immich # Auto-merge: survivor = largest asset count, rest merge into it inside Immich
@@ -130,11 +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 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 # IDs from groups that merged successfully are already gone from Immich, so
# this filter is a no-op for them. # 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:]
}
return [p for p in fresh if p.get("id") not in skip_ids] 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 # All merges failed — fall back to local deduplication (keep largest per name) so
@@ -143,11 +143,6 @@ def _handle_duplicate_people(people: list[dict]) -> list[dict]:
" [yellow]All merges failed — applying local deduplication" " [yellow]All merges failed — applying local deduplication"
" to avoid overwriting output.[/yellow]" " 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 skip_ids]
@@ -213,8 +208,7 @@ def main() -> None:
"and will be reset along with everyone else.[/yellow]" "and will be reset along with everyone else.[/yellow]"
) )
if names: if names:
for name in names: reset_all_people()
reset_person(name)
rprint(f"[bold yellow]Reset tracking data for all {len(names)} people.[/bold yellow]") rprint(f"[bold yellow]Reset tracking data for all {len(names)} people.[/bold yellow]")
else: else:
rprint("[dim]No tracking data to reset.[/dim]") rprint("[dim]No tracking data to reset.[/dim]")
+37 -20
View File
@@ -27,6 +27,8 @@ from .log_config import console
from .quality import blur_score_from_image from .quality import blur_score_from_image
from .reconcile import enrich_asset_with_face_data, reconcile_frigate_mappings from .reconcile import enrich_asset_with_face_data, reconcile_frigate_mappings
from .upload_tracker import ( from .upload_tracker import (
UPLOAD_TRACKER_FILE,
REJECT_TRACKER_FILE,
get_lowest_quality_mapped_file, get_lowest_quality_mapped_file,
get_most_redundant_mapped_file, get_most_redundant_mapped_file,
get_tracked_frigate_file_count, get_tracked_frigate_file_count,
@@ -34,7 +36,10 @@ from .upload_tracker import (
has_frigate_scores, has_frigate_scores,
mark_rejected, mark_rejected,
mark_uploaded, mark_uploaded,
begin_batch,
flush_batch,
remove_frigate_file, remove_frigate_file,
remove_frigate_files_batch,
) )
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -151,9 +156,9 @@ def execute_jobs(jobs: list[dict]) -> None:
if use_full_res: if use_full_res:
img = fetch_full_image(asset["id"]) img = fetch_full_image(asset["id"])
if img is None: if img is None:
# Both original and preview fallback failed — mark rejected # Full-res download failed — could be a transient network
# so this asset isn't retried on every future run. # error, so don't mark rejected; it will be retried next run.
mark_rejected(asset["id"], person_name=name) pass
else: else:
resp = requests.get( resp = requests.get(
f"{Config.IMMICH_URL}/api/assets/{asset['id']}/thumbnail?size=preview&format=JPEG", f"{Config.IMMICH_URL}/api/assets/{asset['id']}/thumbnail?size=preview&format=JPEG",
@@ -163,11 +168,11 @@ def execute_jobs(jobs: list[dict]) -> None:
if resp.ok: if resp.ok:
try: try:
img = Image.open(BytesIO(resp.content)) img = Image.open(BytesIO(resp.content))
except PIL.UnidentifiedImageError: except (PIL.UnidentifiedImageError, OSError):
# Pillow cannot identify the format — genuinely corrupt # Pillow cannot identify the format or the content is
# Immich thumbnail. Mark rejected so this asset isn't # truncated. The download already succeeded (resp.ok),
# retried indefinitely. OSError/truncation errors are # so this is a data problem, not a transient network
# transient and intentionally not caught here. # error — mark rejected so it isn't retried forever.
logger.warning("Invalid image data for asset %s — marking rejected", asset["id"]) logger.warning("Invalid image data for asset %s — marking rejected", asset["id"])
mark_rejected(asset["id"], person_name=name) mark_rejected(asset["id"], person_name=name)
img = None img = None
@@ -345,9 +350,8 @@ def upload_to_frigate(jobs: list[dict]) -> None:
# (manually deleted, or cleaned up outside winnow). This corrects the # (manually deleted, or cleaned up outside winnow). This corrects the
# effective_count so those slots are available for new uploads. # effective_count so those slots are available for new uploads.
stale = get_tracked_frigate_filenames(name) - known_frigate_files_at_start stale = get_tracked_frigate_filenames(name) - known_frigate_files_at_start
for stale_fn in stale:
remove_frigate_file(name, stale_fn)
if stale: if stale:
remove_frigate_files_batch(name, list(stale))
progress.console.print( progress.console.print(
f" [dim]{name}: cleared {len(stale)} stale mapping(s)" f" [dim]{name}: cleared {len(stale)} stale mapping(s)"
" (file(s) no longer in Frigate)[/dim]" " (file(s) no longer in Frigate)[/dim]"
@@ -364,6 +368,9 @@ def upload_to_frigate(jobs: list[dict]) -> None:
min_quality_score_for_slot: float | None = None min_quality_score_for_slot: float | None = None
person_has_fscores: bool = has_frigate_scores(name) person_has_fscores: bool = has_frigate_scores(name)
begin_batch(UPLOAD_TRACKER_FILE)
begin_batch(REJECT_TRACKER_FILE)
try:
for fname in person_files: for fname in person_files:
fpath = os.path.join(person_dir, fname) fpath = os.path.join(person_dir, fname)
@@ -372,10 +379,9 @@ def upload_to_frigate(jobs: list[dict]) -> None:
# freed slot isn't filled with something worse than what we removed. # freed slot isn't filled with something worse than what we removed.
if min_quality_score_for_slot is not None: if min_quality_score_for_slot is not None:
file_score = score_map.get(fname) file_score = score_map.get(fname)
if file_score is None or file_score <= min_quality_score_for_slot: if file_score is not None and file_score <= min_quality_score_for_slot:
score_str = f"{file_score:.3f}" if file_score is not None else "N/A"
progress.console.print( 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]" f" {min_quality_score_for_slot:.3f}, skipping[/dim]"
) )
progress.advance(upload_task) progress.advance(upload_task)
@@ -446,11 +452,13 @@ def upload_to_frigate(jobs: list[dict]) -> None:
get_target = get_most_redundant_mapped_file get_target = get_most_redundant_mapped_file
score_label, better_note = "frigate", " (more novel)" score_label, better_note = "frigate", " (more novel)"
no_score_msg = "Frigate recognize unavailable, skipping replacement" no_score_msg = "Frigate recognize unavailable, skipping replacement"
is_better_than = lambda c, t: c < t
else: else:
candidate_score = score_map.get(fname) candidate_score = score_map.get(fname)
get_target = get_lowest_quality_mapped_file get_target = get_lowest_quality_mapped_file
score_label, better_note = "blur", "" score_label, better_note = "blur", ""
no_score_msg = "no quality score, skipping replacement" no_score_msg = "no quality score, skipping replacement"
is_better_than = lambda c, t: c > t
if candidate_score is None: if candidate_score is None:
progress.console.print(f" [dim]⏭ {fname}: {no_score_msg}[/dim]") progress.console.print(f" [dim]⏭ {fname}: {no_score_msg}[/dim]")
@@ -458,23 +466,21 @@ def upload_to_frigate(jobs: list[dict]) -> None:
continue continue
target = get_target(name, exclude=failed_deletes) target = get_target(name, exclude=failed_deletes)
not_better = target is None or ( not_better = target is None or not is_better_than(candidate_score, target[2])
candidate_score >= target[2] if using_fscore else candidate_score <= target[2]
)
if not_better: if not_better:
target_str = f"{target[2]:.3f}" if target is not None else "N/A" 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( progress.console.print(
f" [dim]⏭ {fname}: {score_label} {candidate_score:.3f}" 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) progress.advance(upload_task)
continue continue
target_frigate_file, _target_asset_id, target_score = target target_frigate_file, _target_asset_id, target_score = target
op = "<" if using_fscore else ">" cmp_op = "<" if using_fscore else ">"
progress.console.print( 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}" f" replacing {target_frigate_file}{better_note}"
) )
if delete_frigate_person_files(name, [target_frigate_file]): if delete_frigate_person_files(name, [target_frigate_file]):
@@ -520,6 +526,7 @@ def upload_to_frigate(jobs: list[dict]) -> None:
" but asset may be re-selected next run: %s", " but asset may be re-selected next run: %s",
fname, tracker_exc, fname, tracker_exc,
) )
else:
if pre_fscore is not None: if pre_fscore is not None:
person_has_fscores = True person_has_fscores = True
actually_uploaded.append((fname, asset_id)) actually_uploaded.append((fname, asset_id))
@@ -592,6 +599,16 @@ def upload_to_frigate(jobs: list[dict]) -> None:
" was not filled this run — will be available next run" " was not filled this run — will be available next run"
) )
finally:
try:
flush_batch(UPLOAD_TRACKER_FILE)
except Exception as _flush_exc:
logger.warning("flush_batch failed during cleanup — batch will be recovered on next begin_batch: %s", _flush_exc)
try:
flush_batch(REJECT_TRACKER_FILE)
except Exception as _flush_exc:
logger.warning("flush_batch failed during cleanup — batch will be recovered on next begin_batch: %s", _flush_exc)
# Batch-map Frigate filenames to asset IDs now that all uploads are done. # Batch-map Frigate filenames to asset IDs now that all uploads are done.
if actually_uploaded and not _skip_reconcile: if actually_uploaded and not _skip_reconcile:
reconcile_frigate_mappings(name, known_frigate_files_at_start, actually_uploaded) reconcile_frigate_mappings(name, known_frigate_files_at_start, actually_uploaded)
+2
View File
@@ -70,7 +70,9 @@ def _resolve_strategy(strategy: str, has_embedding: bool) -> tuple[int | str, st
custom_limit = _getenv_optional_int("LIMIT") custom_limit = _getenv_optional_int("LIMIT")
if custom_limit is not None: if custom_limit is not None:
if custom_limit > 0:
return custom_limit, "smart" return custom_limit, "smart"
logger.warning("LIMIT=%s is invalid — ignoring and using auto strategy", custom_limit)
strategy_map = { strategy_map = {
"adaptive": ("auto", "smart"), "adaptive": ("auto", "smart"),
+12 -8
View File
@@ -14,6 +14,11 @@ from PIL import Image
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _laplacian_var(img_np: np.ndarray) -> float:
gray = cv2.cvtColor(img_np, cv2.COLOR_RGB2GRAY) if img_np.ndim == 3 else img_np
return float(cv2.Laplacian(gray, cv2.CV_64F).var())
@dataclass @dataclass
class QualityResult: class QualityResult:
"""Result of quality assessment on a face/image crop.""" """Result of quality assessment on a face/image crop."""
@@ -32,8 +37,7 @@ def check_blur(img_np: np.ndarray, threshold: float = 100.0) -> tuple[bool, str]
Lower variance = blurrier image. ArcFace needs clear facial features. Lower variance = blurrier image. ArcFace needs clear facial features.
""" """
gray = cv2.cvtColor(img_np, cv2.COLOR_RGB2GRAY) if img_np.ndim == 3 else img_np variance = _laplacian_var(img_np)
variance = cv2.Laplacian(gray, cv2.CV_64F).var()
if variance < threshold: if variance < threshold:
return False, f"Blurry (laplacian={variance:.1f}, threshold={threshold})" return False, f"Blurry (laplacian={variance:.1f}, threshold={threshold})"
return True, "" return True, ""
@@ -115,8 +119,7 @@ def assess_quality(
reasons = [] reasons = []
# Compute laplacian variance once (used by check_blur and stored as blur_score) # Compute laplacian variance once (used by check_blur and stored as blur_score)
gray = cv2.cvtColor(img_np, cv2.COLOR_RGB2GRAY) if img_np.ndim == 3 else img_np blur_score = _laplacian_var(img_np)
blur_score = float(cv2.Laplacian(gray, cv2.CV_64F).var())
checks = [ checks = [
( (
@@ -139,22 +142,23 @@ def assess_quality(
return QualityResult(passed=len(reasons) == 0, reasons=reasons, blur_score=blur_score) 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. """Compute Laplacian-variance blur score, capped at max_dim px to normalise scale.
Caps resolution so full-res and thumbnail scores are comparable — Laplacian Caps resolution so full-res and thumbnail scores are comparable — Laplacian
variance grows with pixel count, making uncapped full-res scores much larger variance grows with pixel count, making uncapped full-res scores much larger
than thumbnail scores for the same perceived sharpness. 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: try:
score_img = img.convert("RGB") if img.mode != "RGB" else img score_img = img.convert("RGB") if img.mode != "RGB" else img
if score_img.width > max_dim or score_img.height > max_dim: if score_img.width > max_dim or score_img.height > max_dim:
score_img = score_img.copy() score_img = score_img.copy()
score_img.thumbnail((max_dim, max_dim), Image.LANCZOS) score_img.thumbnail((max_dim, max_dim), Image.LANCZOS)
return float(assess_quality(score_img).blur_score) return _laplacian_var(np.array(score_img))
except Exception as exc: except Exception as exc:
logger.debug("blur_score_from_image failed: %s", exc) logger.debug("blur_score_from_image failed: %s", exc)
return 0.0 return None
+157 -65
View File
@@ -32,7 +32,7 @@ import logging
import os import os
from pathlib import Path 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__) logger = logging.getLogger(__name__)
@@ -43,6 +43,8 @@ REJECT_TRACKER_FILE = "frigate_rejected_ids.json"
# Reduces per-call JSON reads from O(calls) to O(1) after the first load. # Reduces per-call JSON reads from O(calls) to O(1) after the first load.
# Keyed by full path so tests with isolated tmp dirs never share entries. # Keyed by full path so tests with isolated tmp dirs never share entries.
_cache: dict[str, dict] = {} _cache: dict[str, dict] = {}
_deferred: set[str] = set() # paths whose disk writes are batched until flush_batch()
_dirty: set[str] = set() # deferred paths that received at least one _save during the batch
def _tracker_path(filename: str) -> Path: def _tracker_path(filename: str) -> Path:
@@ -69,20 +71,62 @@ def _load(filename: str) -> dict:
return data return data
def _write_to_disk(path: Path, data: dict) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(".tmp")
try:
with open(tmp, "w") as f:
json.dump(data, f, indent=2)
os.replace(tmp, path)
except Exception:
tmp.unlink(missing_ok=True)
raise
def _save(filename: str, data: dict) -> None: def _save(filename: str, data: dict) -> None:
path = _tracker_path(filename) path = _tracker_path(filename)
_cache[str(path)] = data # keep cache consistent with what we write key = str(path)
path.parent.mkdir(parents=True, exist_ok=True) if key in _deferred:
with open(path, "w") as f: _cache[key] = data # accumulate in cache; disk write deferred until flush_batch()
json.dump(data, f, indent=2) _dirty.add(key)
return
_write_to_disk(path, data)
_cache[key] = data # update cache only after successful write
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.
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 _dirty:
try:
_write_to_disk(path, _cache[key])
except Exception:
logger.warning("begin_batch: could not flush leftover deferred state for %s — partial progress may be lost", path)
_deferred.discard(key)
_dirty.discard(key)
_deferred.add(key)
def flush_batch(filename: str) -> None:
"""Write the accumulated cache state for filename to disk."""
path = _tracker_path(filename)
key = str(path)
if key in _dirty and key in _cache:
_write_to_disk(path, _cache[key])
_deferred.discard(key)
_dirty.discard(key)
def _flat_key(filename: str) -> str: 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]:
return set(_load(filename).get(_flat_key(filename), []))
def _get_ids(entry: list | dict) -> list[str]: def _get_ids(entry: list | dict) -> list[str]:
@@ -96,12 +140,14 @@ def _migrate_entry(entry: list | dict) -> dict:
"""Ensure by_person entry is in the current dict format.""" """Ensure by_person entry is in the current dict format."""
if isinstance(entry, list): if isinstance(entry, list):
return {"asset_ids": sorted(entry), "scores": {}, "frigate_scores": {}, "frigate_files": {}, "crop_dims": {}} return {"asset_ids": sorted(entry), "scores": {}, "frigate_scores": {}, "frigate_files": {}, "crop_dims": {}}
entry.setdefault("asset_ids", []) # Copy top-level and all nested dicts so callers' mutations never reach the cache.
entry.setdefault("scores", {}) result = dict(entry)
entry.setdefault("frigate_scores", {}) result["asset_ids"] = list(result.get("asset_ids", []))
entry.setdefault("frigate_files", {}) result["scores"] = dict(result.get("scores", {}))
entry.setdefault("crop_dims", {}) result["frigate_scores"] = dict(result.get("frigate_scores", {}))
return entry result["frigate_files"] = dict(result.get("frigate_files", {}))
result["crop_dims"] = dict(result.get("crop_dims", {}))
return result
def _mark( def _mark(
@@ -112,12 +158,10 @@ def _mark(
crop_dims: tuple[int, int] | None = None, crop_dims: tuple[int, int] | None = None,
frigate_score: float | None = None, frigate_score: float | None = 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) data = _load(filename)
flat_key = _flat_key(filename)
flat = set(data.get(flat_key, []))
flat.add(asset_id)
data[flat_key] = sorted(flat)
if person_name:
by_person = data.setdefault("by_person", {}) by_person = data.setdefault("by_person", {})
entry = _migrate_entry(by_person.get(person_name, {})) entry = _migrate_entry(by_person.get(person_name, {}))
ids = set(entry["asset_ids"]) ids = set(entry["asset_ids"])
@@ -131,16 +175,27 @@ def _mark(
entry["frigate_scores"][asset_id] = round(frigate_score, 4) entry["frigate_scores"][asset_id] = round(frigate_score, 4)
by_person[person_name] = entry by_person[person_name] = entry
_save(filename, data) _save(filename, data)
logger.debug("Marked %s in %s (%s)", asset_id, filename, person_name)
# ── Public API ──────────────────────────────────────────────────────────────── # ── Public API ────────────────────────────────────────────────────────────────
def load_uploaded_ids() -> set[str]: def load_uploaded_ids() -> set[str]:
return _load_flat(UPLOAD_TRACKER_FILE) """Return all asset IDs recorded as uploaded. Derives from by_person (primary)
plus any legacy flat list still present in old tracker files."""
data = _load(UPLOAD_TRACKER_FILE)
ids = {aid for e in data.get("by_person", {}).values() for aid in _get_ids(e)}
ids.update(data.get("uploaded_asset_ids", [])) # backward compat with pre-0.6.1 files
return ids
def load_rejected_ids() -> set[str]: def load_rejected_ids() -> set[str]:
return _load_flat(REJECT_TRACKER_FILE) """Return all asset IDs recorded as rejected. Derives from by_person (primary)
plus any legacy flat list still present in old tracker files."""
data = _load(REJECT_TRACKER_FILE)
ids = {aid for e in data.get("by_person", {}).values() for aid in _get_ids(e)}
ids.update(data.get("rejected_asset_ids", [])) # backward compat with pre-0.6.1 files
return ids
def mark_uploaded( def mark_uploaded(
@@ -151,35 +206,30 @@ def mark_uploaded(
frigate_score: float | None = None, frigate_score: float | None = None,
) -> None: ) -> None:
_mark(UPLOAD_TRACKER_FILE, asset_id, person_name, score=score, crop_dims=crop_dims, frigate_score=frigate_score) _mark(UPLOAD_TRACKER_FILE, asset_id, person_name, score=score, crop_dims=crop_dims, frigate_score=frigate_score)
logger.debug(f"Marked {asset_id} as uploaded ({person_name})")
def mark_rejected(asset_id: str, person_name: str | None = None) -> None: def mark_rejected(asset_id: str, person_name: str | None = None) -> None:
_mark(REJECT_TRACKER_FILE, asset_id, person_name) _mark(REJECT_TRACKER_FILE, asset_id, person_name)
logger.debug(f"Marked {asset_id} as rejected ({person_name})")
def record_frigate_file(person_name: str, frigate_filename: str, asset_id: str) -> 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.""" """Record a single Frigate filename → asset_id mapping."""
data = _load(UPLOAD_TRACKER_FILE) record_frigate_files_batch(person_name, {frigate_filename: asset_id})
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})")
def record_frigate_files_batch(person_name: str, mappings: dict[str, str]) -> None: def record_frigate_files_batch(person_name: str, mappings: dict[str, str]) -> None:
"""Record multiple Frigate filename → asset_id mappings in a single load/save.""" """Record multiple Frigate filename → asset_id mappings in a single load/save."""
if not mappings: if not mappings:
return return
data = _load(UPLOAD_TRACKER_FILE) src = _load(UPLOAD_TRACKER_FILE)
by_person = data.setdefault("by_person", {}) by_person = dict(src.get("by_person", {}))
entry = _migrate_entry(by_person.get(person_name, {})) entry = _migrate_entry(by_person.get(person_name, {}))
entry["frigate_files"].update(mappings) entry["frigate_files"].update(mappings)
by_person[person_name] = entry by_person[person_name] = entry
data = dict(src)
data["by_person"] = by_person
_save(UPLOAD_TRACKER_FILE, data) _save(UPLOAD_TRACKER_FILE, data)
logger.debug(f"Batch-mapped {len(mappings)} Frigate file(s) for {person_name}") logger.debug(f"Batch-mapped {len(mappings)} Frigate file(s) for {person_name}")
@@ -190,15 +240,26 @@ def remove_frigate_file(person_name: str, frigate_filename: str) -> None:
Does NOT unmark the source asset_id — the deletion was deliberate and 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. we don't want to re-upload the inferior image on the next run.
""" """
data = _load(UPLOAD_TRACKER_FILE) remove_frigate_files_batch(person_name, [frigate_filename])
by_person = data.get("by_person", {})
entry = _migrate_entry(by_person.get(person_name, {}))
asset_id = entry["frigate_files"].pop(frigate_filename, None) def remove_frigate_files_batch(person_name: str, frigate_filenames: list[str]) -> None:
if asset_id: """Remove multiple Frigate filenames in a single load/save."""
src = _load(UPLOAD_TRACKER_FILE)
raw = src.get("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 is not None and asset_id not in entry["frigate_files"].values():
entry["frigate_scores"].pop(asset_id, None) entry["frigate_scores"].pop(asset_id, None)
by_person = dict(src.get("by_person", {})) # copy so assignment does not mutate the cache
by_person[person_name] = entry by_person[person_name] = entry
data = dict(src)
data["by_person"] = by_person
_save(UPLOAD_TRACKER_FILE, data) _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: def get_tracked_frigate_file_count(person_name: str) -> int:
@@ -226,9 +287,11 @@ def get_tracked_frigate_filenames(person_name: str) -> set[str]:
def has_frigate_scores(person_name: str) -> bool: def has_frigate_scores(person_name: str) -> bool:
"""Return True if any mapped file for this person has a stored Frigate recognition score.""" """Return True if any mapped file for this person has a stored Frigate recognition score."""
data = _load(UPLOAD_TRACKER_FILE) data = _load(UPLOAD_TRACKER_FILE)
entry = _migrate_entry(data.get("by_person", {}).get(person_name, {})) raw = data.get("by_person", {}).get(person_name)
frigate_files = entry.get("frigate_files", {}) if not raw or isinstance(raw, list):
frigate_scores = entry.get("frigate_scores", {}) return False
frigate_files = raw.get("frigate_files", {})
frigate_scores = raw.get("frigate_scores", {})
return any(asset_id in frigate_scores for asset_id in frigate_files.values()) return any(asset_id in frigate_scores for asset_id in frigate_files.values())
@@ -238,11 +301,12 @@ def _pick_mapped_file(
data = _load(UPLOAD_TRACKER_FILE) data = _load(UPLOAD_TRACKER_FILE)
entry = _migrate_entry(data.get("by_person", {}).get(person_name, {})) entry = _migrate_entry(data.get("by_person", {}).get(person_name, {}))
scores = entry.get(score_key, {}) scores = entry.get(score_key, {})
candidates = [ seen_assets: set[str] = set()
(ff, asset_id, scores[asset_id]) candidates = []
for ff, asset_id in entry.get("frigate_files", {}).items() for ff, asset_id in entry.get("frigate_files", {}).items():
if (exclude is None or ff not in exclude) and asset_id in scores 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: if not candidates:
return None return None
return max(candidates, key=lambda x: x[2]) if highest else min(candidates, key=lambda x: x[2]) return max(candidates, key=lambda x: x[2]) if highest else min(candidates, key=lambda x: x[2])
@@ -285,7 +349,9 @@ def find_by_crop_dimension(size: int) -> list[dict]:
entry = _migrate_entry(raw_entry) entry = _migrate_entry(raw_entry)
scores = entry.get("scores", {}) scores = entry.get("scores", {})
frigate_files = entry.get("frigate_files", {}) 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", {}) frigate_scores = entry.get("frigate_scores", {})
for asset_id, dims in entry.get("crop_dims", {}).items(): for asset_id, dims in entry.get("crop_dims", {}).items():
w, h = dims[0], dims[1] w, h = dims[0], dims[1]
@@ -312,6 +378,31 @@ def update_frigate_count(person_name: str, count: int) -> None:
_save(UPLOAD_TRACKER_FILE, data) _save(UPLOAD_TRACKER_FILE, data)
def reset_all_people() -> None:
"""Reset all tracking data in two writes (O(P) Frigate API calls, O(1) disk writes).
Preferred over calling reset_person() in a loop when RESET_PERSON=* — that
approach is O(P²) because each call rebuilds the flat list from all remaining entries.
"""
upload_data = _load(UPLOAD_TRACKER_FILE)
frigate_url = _get_frigate_url()
if not frigate_url:
logger.info("FRIGATE_URL not set — skipping Frigate file deletion")
for person_name, raw_entry in upload_data.get("by_person", {}).items():
entry = _migrate_entry(raw_entry)
frigate_filenames = list(entry.get("frigate_files", {}).keys())
if not frigate_filenames:
continue
if frigate_url:
if delete_frigate_person_files(person_name, frigate_filenames):
logger.info(f"Deleted {len(frigate_filenames)} Frigate file(s) for {person_name}")
else:
logger.warning(f"Could not delete Frigate files for {person_name} — tracker reset proceeding anyway")
_save(UPLOAD_TRACKER_FILE, {})
_save(REJECT_TRACKER_FILE, {})
logger.info("Reset all tracking data")
def reset_person(person_name: str) -> None: def reset_person(person_name: str) -> None:
"""Remove all uploaded and rejected records for a given person. """Remove all uploaded and rejected records for a given person.
@@ -324,7 +415,7 @@ def reset_person(person_name: str) -> None:
entry = _migrate_entry(upload_data.get("by_person", {}).get(person_name, {})) entry = _migrate_entry(upload_data.get("by_person", {}).get(person_name, {}))
frigate_filenames = list(entry.get("frigate_files", {}).keys()) frigate_filenames = list(entry.get("frigate_files", {}).keys())
if frigate_filenames: 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}") logger.info(f"FRIGATE_URL not set — skipping Frigate file deletion for {person_name}")
elif delete_frigate_person_files(person_name, frigate_filenames): elif delete_frigate_person_files(person_name, frigate_filenames):
logger.info(f"Deleted {len(frigate_filenames)} Frigate file(s) for {person_name}") logger.info(f"Deleted {len(frigate_filenames)} Frigate file(s) for {person_name}")
@@ -332,16 +423,17 @@ def reset_person(person_name: str) -> None:
logger.warning(f"Could not delete Frigate files for {person_name} — tracker reset proceeding anyway") logger.warning(f"Could not delete Frigate files for {person_name} — tracker reset proceeding anyway")
changed = False changed = False
tracker_files = ((UPLOAD_TRACKER_FILE, upload_data), (REJECT_TRACKER_FILE, _load(REJECT_TRACKER_FILE))) for filename in (UPLOAD_TRACKER_FILE, REJECT_TRACKER_FILE):
for filename, data in tracker_files: src = upload_data if filename == UPLOAD_TRACKER_FILE else _load(REJECT_TRACKER_FILE)
flat_key = _flat_key(filename) by_person = dict(src.get("by_person", {})) # copy so pop() does not mutate the cache
by_person = data.get("by_person", {})
tracker_entry = by_person.pop(person_name, None) tracker_entry = by_person.pop(person_name, None)
if tracker_entry is not None: if tracker_entry is not None:
person_ids = set(_get_ids(tracker_entry)) data = dict(src)
flat = set(data.get(flat_key, [])) - person_ids
data[flat_key] = sorted(flat)
data["by_person"] = by_person 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)
_save(filename, data) _save(filename, data)
changed = True changed = True
if changed: if changed:
@@ -357,14 +449,14 @@ def get_person_summary() -> dict[str, dict]:
names = set(uploaded_data) | set(rejected_data) names = set(uploaded_data) | set(rejected_data)
result = {} result = {}
for name in sorted(names): for name in sorted(names):
u_entry = uploaded_data.get(name, {}) u_entry = _migrate_entry(uploaded_data.get(name, {}))
r_entry = rejected_data.get(name, {}) r_entry = _migrate_entry(rejected_data.get(name, {}))
result[name] = { result[name] = {
"uploaded": len(_get_ids(u_entry)), "uploaded": len(u_entry["asset_ids"]),
"rejected": len(_get_ids(r_entry)), "rejected": len(r_entry["asset_ids"]),
"frigate_count": u_entry.get("frigate_count") if isinstance(u_entry, dict) else None, "frigate_count": u_entry.get("frigate_count"),
"scores": u_entry.get("scores", {}) if isinstance(u_entry, dict) else {}, "scores": u_entry["scores"],
"frigate_files": u_entry.get("frigate_files", {}) if isinstance(u_entry, dict) else {}, "frigate_files": u_entry["frigate_files"],
} }
return result return result