diff --git a/CHANGELOG.md b/CHANGELOG.md index c7f8a5a..e48a2d9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- **`get_all_frigate_person_files` crashed on a non-dict `/api/faces` response**: `data.items()` assumed the JSON body was always an object; a malformed response from Frigate crashed the call instead of degrading gracefully like the rest of the API layer. `_get_faces_data` now validates the response shape and returns `None` on a non-dict body, protecting both `get_all_frigate_person_files` and `get_frigate_person_files`. + +- **Tracker JSON files had no cross-process locking**: two winnow invocations against the same `DATA_DIR` (e.g. a scheduled run overlapping a manual `docker exec`, which the docs explicitly instruct) could race a read-modify-write on `frigate_uploaded_ids.json`/`frigate_rejected_ids.json` and silently lose the loser's marks. Tracker load-mutate-save cycles now hold an exclusive `flock` on a `DATA_DIR` lock file for the duration of the operation. + +- **`begin_batch`/`flush_batch` could lose an entire person's upload marks on a crash**: tracker writes were deferred for the whole per-person upload loop, so a SIGKILL/OOM/host crash mid-batch lost every successful Frigate upload from the tracker even though the files were already live in Frigate, causing duplicate re-uploads on the next run. Batches now flush to disk every 10 marks, bounding the loss to a small, fixed window. + ## [0.6.6] - 2026-06-18 ### Changed diff --git a/tests/test_frigate_api.py b/tests/test_frigate_api.py new file mode 100644 index 0000000..4577da1 --- /dev/null +++ b/tests/test_frigate_api.py @@ -0,0 +1,56 @@ +"""Tests for Frigate API helpers (no network required).""" + +import pytest + + +class _FakeResponse: + def __init__(self, json_body, status_code=200): + self._json_body = json_body + self.status_code = status_code + self.ok = 200 <= status_code < 300 + + def raise_for_status(self): + if not self.ok: + raise Exception(f"HTTP {self.status_code}") + + def json(self): + return self._json_body + + +@pytest.fixture(autouse=True) +def frigate_url(monkeypatch): + monkeypatch.setenv("FRIGATE_URL", "http://frigate.test") + yield + + +def test_get_all_frigate_person_files_happy_path(monkeypatch): + from winnow import frigate_api + body = {"Alice": ["a-1.webp", "a-2.webp"], "train": ["pending.webp"]} + monkeypatch.setattr(frigate_api.requests, "get", lambda *a, **k: _FakeResponse(body)) + result = frigate_api.get_all_frigate_person_files() + assert result == {"Alice": ["a-1.webp", "a-2.webp"]} + + +def test_get_all_frigate_person_files_rejects_non_dict_body(monkeypatch): + """A non-dict /api/faces response must not crash data.items() — it should degrade to None.""" + from winnow import frigate_api + monkeypatch.setattr(frigate_api.requests, "get", lambda *a, **k: _FakeResponse(["not", "a", "dict"])) + assert frigate_api.get_all_frigate_person_files() is None + + +def test_get_frigate_person_files_rejects_non_dict_body(monkeypatch): + from winnow import frigate_api + monkeypatch.setattr(frigate_api.requests, "get", lambda *a, **k: _FakeResponse(["not", "a", "dict"])) + assert frigate_api.get_frigate_person_files("Alice") is None + + +def test_get_frigate_face_counts_rejects_non_dict_body(monkeypatch): + from winnow import frigate_api + monkeypatch.setattr(frigate_api.requests, "get", lambda *a, **k: _FakeResponse("oops")) + assert frigate_api.get_frigate_face_counts() is None + + +def test_get_all_frigate_person_files_rejects_scalar_body(monkeypatch): + from winnow import frigate_api + monkeypatch.setattr(frigate_api.requests, "get", lambda *a, **k: _FakeResponse(42)) + assert frigate_api.get_all_frigate_person_files() is None diff --git a/tests/test_upload_tracker.py b/tests/test_upload_tracker.py index 51c9159..acbe926 100644 --- a/tests/test_upload_tracker.py +++ b/tests/test_upload_tracker.py @@ -1,5 +1,8 @@ """Tests for upload tracker — mark, filter, reset, and summary logic.""" +import fcntl +import json +import os import pytest @@ -260,3 +263,121 @@ def test_get_most_redundant_exclude_all_returns_none(): mark_uploaded("asset-a", person_name="Alice", score=0.5, frigate_score=0.70) record_frigate_file("Alice", "Alice-a.webp", "asset-a") assert get_most_redundant_mapped_file("Alice", exclude={"Alice-a.webp"}) is None + + +# ── cross-process locking ────────────────────────────────────────────────────── + +def test_tracker_lock_blocks_concurrent_exclusive_access(): + """While upload_tracker holds the tracker lock, a second (non-blocking) attempt + to exclusively lock the same file must fail — proving the lock is real and + guards a second winnow process from racing a read-modify-write.""" + from winnow import upload_tracker as ut + + ut._acquire_lock() + try: + lock_path = ut._lock_path() + assert lock_path.exists() + fd = os.open(lock_path, os.O_RDWR) + try: + with pytest.raises(OSError): + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + finally: + os.close(fd) + finally: + ut._release_lock() + + # Released — a second exclusive, non-blocking lock now succeeds immediately. + fd = os.open(ut._lock_path(), os.O_RDWR) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(fd, fcntl.LOCK_UN) + finally: + os.close(fd) + + +def test_tracker_lock_reentrant_within_process(): + """Nested acquisitions (e.g. begin_batch() for both tracker files, or a tracker + call made while a batch is open) must not deadlock the process that holds them.""" + from winnow import upload_tracker as ut + + assert ut._lock_depth == 0 + ut._acquire_lock() + ut._acquire_lock() + assert ut._lock_depth == 2 + ut._release_lock() + assert ut._lock_depth == 1 + ut._release_lock() + assert ut._lock_depth == 0 + + +def test_begin_flush_batch_releases_lock_for_next_caller(): + """begin_batch()/flush_batch() (including a nested mark_uploaded call inside the + batch) must fully release the lock so it doesn't stay held for the rest of the run.""" + from winnow import upload_tracker as ut + from winnow.upload_tracker import ( + REJECT_TRACKER_FILE, + UPLOAD_TRACKER_FILE, + begin_batch, + flush_batch, + mark_uploaded, + ) + + begin_batch(UPLOAD_TRACKER_FILE) + begin_batch(REJECT_TRACKER_FILE) + mark_uploaded("a1", person_name="Alice") + flush_batch(UPLOAD_TRACKER_FILE) + flush_batch(REJECT_TRACKER_FILE) + + assert ut._lock_depth == 0 + fd = os.open(ut._lock_path(), os.O_RDWR) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(fd, fcntl.LOCK_UN) + finally: + os.close(fd) + + +# ── batched writes: incremental flush bounds crash loss ─────────────────────── + +def test_batch_defers_writes_below_flush_threshold(isolated_cache): + """Below the flush threshold, writes stay in memory — batching still avoids + a disk write per mark.""" + from winnow.upload_tracker import UPLOAD_TRACKER_FILE, begin_batch, flush_batch, mark_uploaded + + tracker_path = isolated_cache / UPLOAD_TRACKER_FILE + begin_batch(UPLOAD_TRACKER_FILE) + try: + mark_uploaded("asset-0", person_name="Alice") + assert not tracker_path.exists() + finally: + flush_batch(UPLOAD_TRACKER_FILE) + assert tracker_path.exists() + + +def test_batch_flushes_incrementally_bounding_crash_loss(isolated_cache): + """A crash mid-batch (SIGKILL/OOM/host crash) must lose at most a bounded number + of marks, not the whole per-person batch — verified by reading the file straight + off disk before flush_batch() is ever called.""" + from winnow.upload_tracker import _BATCH_FLUSH_EVERY, UPLOAD_TRACKER_FILE, begin_batch, flush_batch, mark_uploaded + + tracker_path = isolated_cache / UPLOAD_TRACKER_FILE + begin_batch(UPLOAD_TRACKER_FILE) + try: + for i in range(_BATCH_FLUSH_EVERY): + mark_uploaded(f"asset-{i}", person_name="Alice") + # Threshold reached: disk already reflects all marks so far, even though + # flush_batch() has not run yet — simulates surviving a crash here. + on_disk = json.loads(tracker_path.read_text()) + ids = on_disk["by_person"]["Alice"]["asset_ids"] + assert len(ids) == _BATCH_FLUSH_EVERY + + # One more mark past the threshold stays deferred again until the next + # periodic flush or flush_batch(). + mark_uploaded("asset-extra", person_name="Alice") + on_disk = json.loads(tracker_path.read_text()) + assert len(on_disk["by_person"]["Alice"]["asset_ids"]) == _BATCH_FLUSH_EVERY + finally: + flush_batch(UPLOAD_TRACKER_FILE) + + on_disk = json.loads(tracker_path.read_text()) + assert len(on_disk["by_person"]["Alice"]["asset_ids"]) == _BATCH_FLUSH_EVERY + 1 diff --git a/winnow/frigate_api.py b/winnow/frigate_api.py index eba6939..b13d9b5 100644 --- a/winnow/frigate_api.py +++ b/winnow/frigate_api.py @@ -32,14 +32,18 @@ def get_frigate_version() -> str | None: def _get_faces_data() -> dict | None: - """Fetch raw GET /api/faces response. Returns None if unavailable.""" + """Fetch raw GET /api/faces response. Returns None if unavailable or malformed.""" frigate_url = _get_frigate_url() if not frigate_url: return None try: resp = requests.get(f"{frigate_url}/api/faces", timeout=10) resp.raise_for_status() - return resp.json() + data = resp.json() + if not isinstance(data, dict): + logger.warning("Unexpected response shape from Frigate /api/faces: %s", type(data).__name__) + return None + return data except Exception as e: logger.warning("Could not query Frigate faces API: %s", e) return None diff --git a/winnow/upload_tracker.py b/winnow/upload_tracker.py index 04ead4c..b176714 100644 --- a/winnow/upload_tracker.py +++ b/winnow/upload_tracker.py @@ -27,9 +27,11 @@ frigate_files only contains files winnow uploaded — files added manually throu Frigate's UI are never mapped here and are never touched by quality replacement. """ +import fcntl import json import logging import os +from contextlib import contextmanager from pathlib import Path from .frigate_api import _get_frigate_url, delete_frigate_person_files @@ -38,6 +40,14 @@ logger = logging.getLogger(__name__) UPLOAD_TRACKER_FILE = "frigate_uploaded_ids.json" REJECT_TRACKER_FILE = "frigate_rejected_ids.json" +LOCK_FILE = ".tracker.lock" + +# Number of deferred _save calls a batch accumulates before it is flushed to disk +# early. Bounds how many marks a crash mid-batch (SIGKILL/OOM/host crash) can lose — +# without this, begin_batch()/flush_batch() defer every write for an entire +# per-person upload loop, so a crash could lose every mark from that person's batch +# even though the files are already live in Frigate. +_BATCH_FLUSH_EVERY = 10 # Write-through in-memory cache keyed by the resolved file path. # Reduces per-call JSON reads from O(calls) to O(1) after the first load. @@ -45,6 +55,16 @@ REJECT_TRACKER_FILE = "frigate_rejected_ids.json" _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 +_batch_writes: dict[str, int] = {} # deferred _save calls since the last disk write, per path + +# Cross-process lock guarding the load-mutate-save cycle below. Two winnow +# invocations against the same DATA_DIR (e.g. a scheduled run overlapping a +# manual `docker exec`, which the docs explicitly instruct) can otherwise race +# a read-modify-write and silently lose whichever one saves first. Reentrant +# within a single process (depth-counted) so begin_batch()/flush_batch() pairs +# and nested tracker calls made while a batch is open don't self-deadlock. +_lock_fd: int | None = None +_lock_depth = 0 def _tracker_path(filename: str) -> Path: @@ -55,6 +75,49 @@ def _tracker_path(filename: str) -> Path: return Path(filename) +def _lock_path() -> Path: + return _tracker_path(LOCK_FILE) + + +def _acquire_lock() -> None: + """Acquire the cross-process tracker lock. Reentrant within this process.""" + global _lock_fd, _lock_depth + if _lock_depth == 0: + path = _lock_path() + path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(path, os.O_CREAT | os.O_RDWR) + fcntl.flock(fd, fcntl.LOCK_EX) # blocks until any other process's lock is released + _lock_fd = fd + # Cache entries not part of an in-progress deferred batch may be stale — + # another process could have written to disk since we last read them. + # Drop them so the critical section that follows re-reads from disk. + for key in list(_cache): + if key not in _deferred: + del _cache[key] + _lock_depth += 1 + + +def _release_lock() -> None: + global _lock_fd, _lock_depth + if _lock_depth <= 0: + return + _lock_depth -= 1 + if _lock_depth == 0 and _lock_fd is not None: + fcntl.flock(_lock_fd, fcntl.LOCK_UN) + os.close(_lock_fd) + _lock_fd = None + + +@contextmanager +def _locked(): + """Context manager wrapping a single load-mutate-save cycle in the tracker lock.""" + _acquire_lock() + try: + yield + finally: + _release_lock() + + def _load(filename: str) -> dict: path = _tracker_path(filename) key = str(path) @@ -89,6 +152,15 @@ def _save(filename: str, data: dict) -> None: if key in _deferred: _cache[key] = data # accumulate in cache; disk write deferred until flush_batch() _dirty.add(key) + # Flush early every _BATCH_FLUSH_EVERY marks so a crash mid-batch (SIGKILL/OOM/ + # host crash) loses at most a bounded number of marks instead of the whole batch. + # Safe without re-acquiring the lock: begin_batch() already holds it for the + # duration of the batch. + _batch_writes[key] = _batch_writes.get(key, 0) + 1 + if _batch_writes[key] >= _BATCH_FLUSH_EVERY: + _write_to_disk(path, data) + _dirty.discard(key) + _batch_writes[key] = 0 return _write_to_disk(path, data) _cache[key] = data # update cache only after successful write @@ -96,13 +168,19 @@ 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. + in-memory cache until flush_batch() is called (with periodic early flushes — + see _BATCH_FLUSH_EVERY). Use around per-person upload loops to reduce N writes + towards 1. + + Acquires the cross-process tracker lock, held until the matching flush_batch() + releases it, so a concurrent winnow invocation against the same DATA_DIR can't + interleave a read-modify-write with this batch. 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. """ + _acquire_lock() path = _tracker_path(filename) key = str(path) if key in _deferred and key in _dirty: @@ -117,16 +195,26 @@ def begin_batch(filename: str) -> None: _deferred.discard(key) _dirty.discard(key) _deferred.add(key) + _batch_writes[key] = 0 def flush_batch(filename: str) -> None: - """Write the accumulated cache state for filename to disk.""" + """Write the accumulated cache state for filename to disk and release the tracker lock + acquired by the matching begin_batch(). + + The lock release always runs, even if the disk write raises, so a write failure + (e.g. a full disk) can't leave the tracker lock held for the rest of the process. + """ 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) + try: + if key in _dirty and key in _cache: + _write_to_disk(path, _cache[key]) + finally: + _deferred.discard(key) + _dirty.discard(key) + _batch_writes.pop(key, None) + _release_lock() def _flat_key(filename: str) -> str: @@ -165,22 +253,23 @@ def _mark( if not person_name: logger.warning("_mark called with empty person_name for asset %s — asset not recorded", asset_id) return - data = _load(filename) - by_person = dict(data.get("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 - new_data = dict(data) - new_data["by_person"] = by_person - _save(filename, new_data) + with _locked(): + data = _load(filename) + by_person = dict(data.get("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 + 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) @@ -229,14 +318,15 @@ def record_frigate_files_batch(person_name: str, mappings: dict[str, str]) -> No """Record multiple Frigate filename → asset_id mappings in a single load/save.""" if not mappings: return - src = _load(UPLOAD_TRACKER_FILE) - by_person = dict(src.get("by_person", {})) - entry = _migrate_entry(by_person.get(person_name, {})) - entry["frigate_files"].update(mappings) - by_person[person_name] = entry - data = dict(src) - data["by_person"] = by_person - _save(UPLOAD_TRACKER_FILE, data) + with _locked(): + src = _load(UPLOAD_TRACKER_FILE) + by_person = dict(src.get("by_person", {})) + entry = _migrate_entry(by_person.get(person_name, {})) + entry["frigate_files"].update(mappings) + by_person[person_name] = entry + data = dict(src) + data["by_person"] = by_person + _save(UPLOAD_TRACKER_FILE, data) logger.debug(f"Batch-mapped {len(mappings)} Frigate file(s) for {person_name}") @@ -251,20 +341,21 @@ def remove_frigate_file(person_name: str, frigate_filename: str) -> None: def remove_frigate_files_batch(person_name: str, frigate_filenames: list[str]) -> None: """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) - by_person = dict(src.get("by_person", {})) # copy so assignment does not mutate the cache - by_person[person_name] = entry - data = dict(src) - data["by_person"] = by_person - _save(UPLOAD_TRACKER_FILE, data) + with _locked(): + 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) + by_person = dict(src.get("by_person", {})) # copy so assignment does not mutate the cache + by_person[person_name] = entry + data = dict(src) + data["by_person"] = by_person + _save(UPLOAD_TRACKER_FILE, data) logger.debug(f"Removed {len(frigate_filenames)} Frigate file mapping(s) for {person_name}") @@ -378,14 +469,15 @@ def find_by_crop_dimension(size: int) -> list[dict]: 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 = dict(data.get("by_person", {})) - entry = _migrate_entry(by_person.get(person_name, {})) - entry["frigate_count"] = count - by_person[person_name] = entry - new_data = dict(data) - new_data["by_person"] = by_person - _save(UPLOAD_TRACKER_FILE, new_data) + with _locked(): + data = _load(UPLOAD_TRACKER_FILE) + by_person = dict(data.get("by_person", {})) + entry = _migrate_entry(by_person.get(person_name, {})) + entry["frigate_count"] = count + by_person[person_name] = entry + new_data = dict(data) + new_data["by_person"] = by_person + _save(UPLOAD_TRACKER_FILE, new_data) def reset_all_people() -> None: @@ -394,22 +486,25 @@ def reset_all_people() -> None: 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, {}) + with _locked(): + 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") @@ -421,37 +516,38 @@ def reset_person(person_name: str) -> None: files (not in frigate_files) are never touched. Proceeds with tracker reset even if Frigate is unreachable. """ - upload_data = _load(UPLOAD_TRACKER_FILE) - entry = _migrate_entry(upload_data.get("by_person", {}).get(person_name, {})) - frigate_filenames = list(entry.get("frigate_files", {}).keys()) - if frigate_filenames: - 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}") - else: - logger.warning(f"Could not delete Frigate files for {person_name} — tracker reset proceeding anyway") + with _locked(): + upload_data = _load(UPLOAD_TRACKER_FILE) + entry = _migrate_entry(upload_data.get("by_person", {}).get(person_name, {})) + frigate_filenames = list(entry.get("frigate_files", {}).keys()) + if frigate_filenames: + 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}") + else: + logger.warning(f"Could not delete Frigate files for {person_name} — tracker reset proceeding anyway") - changed = False - for filename in (UPLOAD_TRACKER_FILE, REJECT_TRACKER_FILE): - 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 and not isinstance(data[flat_key], list): - logger.warning( - "reset_person: %s has unexpected type for %s (%s) — skipping flat-list cleanup;" - " all persons' legacy IDs in this field are unaffected but unreadable", - filename, flat_key, type(data[flat_key]).__name__, - ) - elif person_ids and flat_key in data: - data[flat_key] = sorted(set(data[flat_key]) - person_ids) - _save(filename, data) - changed = True + changed = False + for filename in (UPLOAD_TRACKER_FILE, REJECT_TRACKER_FILE): + 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 and not isinstance(data[flat_key], list): + logger.warning( + "reset_person: %s has unexpected type for %s (%s) — skipping flat-list cleanup;" + " all persons' legacy IDs in this field are unaffected but unreadable", + filename, flat_key, type(data[flat_key]).__name__, + ) + elif person_ids and flat_key in data: + data[flat_key] = sorted(set(data[flat_key]) - person_ids) + _save(filename, data) + changed = True if changed: logger.info(f"Reset tracking data for {person_name}") else: