fix: address 10 full-codebase audit findings + lint
Correctness:
- jobs: cap auto-diversity limit for brand-new people (was never capped,
could exceed MAX_AUTO_IMAGES on first run)
- image_processing: separate None/0 guard for imageWidth/imageHeight so
missing field is explicit rather than silently aliased to img_w
- upload_tracker (_mark, update_frigate_count): copy-before-mutate so
exceptions between cache access and _save don't corrupt in-process state
- jobs: reject LIMIT=0 on no-embedding path (was silently empty run)
- jobs: add STRATEGY=skip to strategy_map so env var is honoured
- embeddings: convert to RGB before cvtColor so RGBA/grayscale thumbnails
don't raise cv2.error and silently drop from diversity selection
- config: use falsy guard for OUTPUT_DIR so blank env var falls through
to config file value
- reconcile: _ts() returns float("inf") on parse failure so unrecognised
filenames sort last instead of collapsing to 0.0 and corrupting FIFO mapping
- diversity: remove dead selected_set (never read; -np.inf sentinel already
prevents re-selection)
Lint (ruff):
- executor: sort upload_tracker import block (I001)
- executor: replace lambda is_better_than with operator.lt/gt (E731 x2)
- executor, upload_tracker: wrap long logger.warning calls (E501 x4)
This commit is contained in:
+1
-1
@@ -173,7 +173,7 @@ class _Config:
|
|||||||
data = json.loads(config_file.read_text())
|
data = json.loads(config_file.read_text())
|
||||||
if not self.IMMICH_URL:
|
if not self.IMMICH_URL:
|
||||||
self.IMMICH_URL = data.get("IMMICH_URL")
|
self.IMMICH_URL = data.get("IMMICH_URL")
|
||||||
if os.getenv("OUTPUT_DIR") is None:
|
if not os.getenv("OUTPUT_DIR"):
|
||||||
self.OUTPUT_DIR = data.get("OUTPUT_DIR", self.OUTPUT_DIR)
|
self.OUTPUT_DIR = data.get("OUTPUT_DIR", self.OUTPUT_DIR)
|
||||||
except (json.JSONDecodeError, OSError) as e:
|
except (json.JSONDecodeError, OSError) as e:
|
||||||
logging.warning("Failed to load config file: %s", e)
|
logging.warning("Failed to load config file: %s", e)
|
||||||
|
|||||||
@@ -508,7 +508,6 @@ def _cluster_aware_selection(
|
|||||||
|
|
||||||
medoid_indices, cluster_labels = _kmedoids(dist_matrix, k)
|
medoid_indices, cluster_labels = _kmedoids(dist_matrix, k)
|
||||||
selected = list(medoid_indices)
|
selected = list(medoid_indices)
|
||||||
selected_set = set(selected)
|
|
||||||
|
|
||||||
logger.debug("Selected %s cluster medoids as initial picks.", len(selected))
|
logger.debug("Selected %s cluster medoids as initial picks.", len(selected))
|
||||||
|
|
||||||
@@ -541,7 +540,6 @@ def _cluster_aware_selection(
|
|||||||
break
|
break
|
||||||
|
|
||||||
selected.append(best_idx)
|
selected.append(best_idx)
|
||||||
selected_set.add(best_idx)
|
|
||||||
|
|
||||||
# Update min distances
|
# Update min distances
|
||||||
dists_to_new = dist_matrix[best_idx]
|
dists_to_new = dist_matrix[best_idx]
|
||||||
|
|||||||
@@ -181,8 +181,9 @@ def get_face_embedding(img_pil: Image.Image) -> np.ndarray | None:
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# InsightFace expects BGR cv2 image
|
# InsightFace expects BGR cv2 image; normalise mode first so RGBA/grayscale don't
|
||||||
img_bgr = cv2.cvtColor(np.asarray(img_pil), cv2.COLOR_RGB2BGR)
|
# raise a channel-count error inside cvtColor.
|
||||||
|
img_bgr = cv2.cvtColor(np.asarray(img_pil.convert("RGB")), cv2.COLOR_RGB2BGR)
|
||||||
|
|
||||||
# Suppress scikit-image FutureWarning from InsightFace's face_align.py
|
# Suppress scikit-image FutureWarning from InsightFace's face_align.py
|
||||||
with warnings.catch_warnings():
|
with warnings.catch_warnings():
|
||||||
|
|||||||
+20
-8
@@ -1,6 +1,7 @@
|
|||||||
"""Execution phase: image processing and Frigate upload."""
|
"""Execution phase: image processing and Frigate upload."""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
import operator
|
||||||
import os
|
import os
|
||||||
import shutil
|
import shutil
|
||||||
from io import BytesIO
|
from io import BytesIO
|
||||||
@@ -27,8 +28,10 @@ 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,
|
REJECT_TRACKER_FILE,
|
||||||
|
UPLOAD_TRACKER_FILE,
|
||||||
|
begin_batch,
|
||||||
|
flush_batch,
|
||||||
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,
|
||||||
@@ -36,8 +39,6 @@ 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,
|
remove_frigate_files_batch,
|
||||||
)
|
)
|
||||||
@@ -452,13 +453,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
|
is_better_than = operator.lt
|
||||||
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
|
is_better_than = operator.gt
|
||||||
|
|
||||||
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]")
|
||||||
@@ -489,7 +490,10 @@ def upload_to_frigate(jobs: list[dict]) -> None:
|
|||||||
effective_count -= 1
|
effective_count -= 1
|
||||||
min_quality_score_for_slot = None if using_fscore else target_score
|
min_quality_score_for_slot = None if using_fscore else target_score
|
||||||
else:
|
else:
|
||||||
logger.warning("Failed to delete %s for %s, skipping replacement", target_frigate_file, name)
|
logger.warning(
|
||||||
|
"Failed to delete %s for %s, skipping replacement",
|
||||||
|
target_frigate_file, name,
|
||||||
|
)
|
||||||
failed_deletes.add(target_frigate_file)
|
failed_deletes.add(target_frigate_file)
|
||||||
progress.advance(upload_task)
|
progress.advance(upload_task)
|
||||||
continue
|
continue
|
||||||
@@ -603,11 +607,19 @@ def upload_to_frigate(jobs: list[dict]) -> None:
|
|||||||
try:
|
try:
|
||||||
flush_batch(UPLOAD_TRACKER_FILE)
|
flush_batch(UPLOAD_TRACKER_FILE)
|
||||||
except Exception as _flush_exc:
|
except Exception as _flush_exc:
|
||||||
logger.warning("flush_batch failed during cleanup — batch will be recovered on next begin_batch: %s", _flush_exc)
|
logger.warning(
|
||||||
|
"flush_batch failed during cleanup"
|
||||||
|
" — batch will be recovered on next begin_batch: %s",
|
||||||
|
_flush_exc,
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
flush_batch(REJECT_TRACKER_FILE)
|
flush_batch(REJECT_TRACKER_FILE)
|
||||||
except Exception as _flush_exc:
|
except Exception as _flush_exc:
|
||||||
logger.warning("flush_batch failed during cleanup — batch will be recovered on next begin_batch: %s", _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:
|
||||||
|
|||||||
@@ -92,11 +92,14 @@ def process_face_mode(
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
img_w, img_h = img.size
|
img_w, img_h = img.size
|
||||||
meta_w = face_info.get("imageWidth") or img_w
|
meta_w = face_info.get("imageWidth") or 0
|
||||||
meta_h = face_info.get("imageHeight") or img_h
|
meta_h = face_info.get("imageHeight") or 0
|
||||||
|
|
||||||
# Scale bounding box to actual image dimensions
|
# Scale bounding box from detection-image space to actual image dimensions.
|
||||||
scale_x, scale_y = img_w / meta_w, img_h / meta_h
|
# Fall back to 1.0 if Immich omits the field — bbox is assumed to already
|
||||||
|
# be in image space (correct for thumbnails, wrong for full-res).
|
||||||
|
scale_x = img_w / meta_w if meta_w else 1.0
|
||||||
|
scale_y = img_h / meta_h if meta_h else 1.0
|
||||||
x1 = face_info["boundingBoxX1"] * scale_x
|
x1 = face_info["boundingBoxX1"] * scale_x
|
||||||
y1 = face_info["boundingBoxY1"] * scale_y
|
y1 = face_info["boundingBoxY1"] * scale_y
|
||||||
x2 = face_info["boundingBoxX2"] * scale_x
|
x2 = face_info["boundingBoxX2"] * scale_x
|
||||||
|
|||||||
+7
-3
@@ -66,7 +66,11 @@ def _get_strategy_choice(has_embedding: bool) -> tuple[int | str, str]:
|
|||||||
def _resolve_strategy(strategy: str, has_embedding: bool) -> tuple[int | str, str]:
|
def _resolve_strategy(strategy: str, has_embedding: bool) -> tuple[int | str, str]:
|
||||||
"""Resolve env var strategy to (limit, selection_mode) without prompts."""
|
"""Resolve env var strategy to (limit, selection_mode) without prompts."""
|
||||||
if not has_embedding:
|
if not has_embedding:
|
||||||
return _getenv_int("LIMIT", 30), "time"
|
limit = _getenv_int("LIMIT", 30)
|
||||||
|
if limit <= 0:
|
||||||
|
logger.warning("LIMIT=%s is invalid — ignoring and using default 30", limit)
|
||||||
|
limit = 30
|
||||||
|
return limit, "time"
|
||||||
|
|
||||||
custom_limit = _getenv_optional_int("LIMIT")
|
custom_limit = _getenv_optional_int("LIMIT")
|
||||||
if custom_limit is not None:
|
if custom_limit is not None:
|
||||||
@@ -77,6 +81,7 @@ def _resolve_strategy(strategy: str, has_embedding: bool) -> tuple[int | str, st
|
|||||||
strategy_map = {
|
strategy_map = {
|
||||||
"adaptive": ("auto", "smart"),
|
"adaptive": ("auto", "smart"),
|
||||||
"auto": ("auto", "smart"), # legacy alias for adaptive
|
"auto": ("auto", "smart"), # legacy alias for adaptive
|
||||||
|
"skip": (0, "skip"),
|
||||||
"standard": (30, "smart"),
|
"standard": (30, "smart"),
|
||||||
"broad": (100, "smart"),
|
"broad": (100, "smart"),
|
||||||
}
|
}
|
||||||
@@ -297,8 +302,7 @@ def auto_configure(people: list[dict]) -> list[dict]:
|
|||||||
if limit == "auto":
|
if limit == "auto":
|
||||||
# Switch from open-ended auto to a fixed budget at remaining capacity
|
# Switch from open-ended auto to a fixed budget at remaining capacity
|
||||||
# so the diversity selector itself stops at the right count instead of
|
# so the diversity selector itself stops at the right count instead of
|
||||||
# selecting MAX_AUTO_IMAGES and then discarding the excess by position.
|
# selecting more than MAX_AUTO_IMAGES and overflowing the cap.
|
||||||
if already_uploaded > 0:
|
|
||||||
limit = capacity
|
limit = capacity
|
||||||
else:
|
else:
|
||||||
limit = min(limit, capacity)
|
limit = min(limit, capacity)
|
||||||
|
|||||||
+1
-1
@@ -61,7 +61,7 @@ def reconcile_frigate_mappings(
|
|||||||
try:
|
try:
|
||||||
return float(fname.rsplit("_", 1)[-1].rsplit(".", 1)[0])
|
return float(fname.rsplit("_", 1)[-1].rsplit(".", 1)[0])
|
||||||
except (ValueError, IndexError):
|
except (ValueError, IndexError):
|
||||||
return 0.0
|
return float("inf")
|
||||||
|
|
||||||
logger.debug(
|
logger.debug(
|
||||||
"%s: mapping %s file(s) by filename timestamp — assumes Frigate processes"
|
"%s: mapping %s file(s) by filename timestamp — assumes Frigate processes"
|
||||||
|
|||||||
@@ -109,7 +109,11 @@ def begin_batch(filename: str) -> None:
|
|||||||
try:
|
try:
|
||||||
_write_to_disk(path, _cache[key])
|
_write_to_disk(path, _cache[key])
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("begin_batch: could not flush leftover deferred state for %s — partial progress may be lost", path)
|
logger.warning(
|
||||||
|
"begin_batch: could not flush leftover deferred state for %s"
|
||||||
|
" — partial progress may be lost",
|
||||||
|
path,
|
||||||
|
)
|
||||||
_deferred.discard(key)
|
_deferred.discard(key)
|
||||||
_dirty.discard(key)
|
_dirty.discard(key)
|
||||||
_deferred.add(key)
|
_deferred.add(key)
|
||||||
@@ -162,7 +166,7 @@ def _mark(
|
|||||||
logger.warning("_mark called with empty person_name for asset %s — asset not recorded", asset_id)
|
logger.warning("_mark called with empty person_name for asset %s — asset not recorded", asset_id)
|
||||||
return
|
return
|
||||||
data = _load(filename)
|
data = _load(filename)
|
||||||
by_person = data.setdefault("by_person", {})
|
by_person = dict(data.get("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"])
|
||||||
ids.add(asset_id)
|
ids.add(asset_id)
|
||||||
@@ -174,7 +178,9 @@ def _mark(
|
|||||||
if frigate_score is not None:
|
if frigate_score is not None:
|
||||||
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)
|
new_data = dict(data)
|
||||||
|
new_data["by_person"] = by_person
|
||||||
|
_save(filename, new_data)
|
||||||
logger.debug("Marked %s in %s (%s)", asset_id, filename, person_name)
|
logger.debug("Marked %s in %s (%s)", asset_id, filename, person_name)
|
||||||
|
|
||||||
|
|
||||||
@@ -371,11 +377,13 @@ def find_by_crop_dimension(size: int) -> list[dict]:
|
|||||||
def update_frigate_count(person_name: str, count: int) -> None:
|
def update_frigate_count(person_name: str, count: int) -> None:
|
||||||
"""Record Frigate's authoritative training image count for a person."""
|
"""Record Frigate's authoritative training image count for a person."""
|
||||||
data = _load(UPLOAD_TRACKER_FILE)
|
data = _load(UPLOAD_TRACKER_FILE)
|
||||||
by_person = data.setdefault("by_person", {})
|
by_person = dict(data.get("by_person", {}))
|
||||||
entry = _migrate_entry(by_person.get(person_name, {}))
|
entry = _migrate_entry(by_person.get(person_name, {}))
|
||||||
entry["frigate_count"] = count
|
entry["frigate_count"] = count
|
||||||
by_person[person_name] = entry
|
by_person[person_name] = entry
|
||||||
_save(UPLOAD_TRACKER_FILE, data)
|
new_data = dict(data)
|
||||||
|
new_data["by_person"] = by_person
|
||||||
|
_save(UPLOAD_TRACKER_FILE, new_data)
|
||||||
|
|
||||||
|
|
||||||
def reset_all_people() -> None:
|
def reset_all_people() -> None:
|
||||||
|
|||||||
Reference in New Issue
Block a user