"""photo-pipeline: photo-ingest flow (v2 — batched + checkpointed). Hashes incoming images, checks against the persistent fingerprint DB (exact + near dupes), and registers new ones. Batching: processes in chunks of `batch_size` (default 1000), committing to the DB after each chunk. Checkpointing is DB-native: files already in the DB (by sha256) are skipped on resume — an interrupted run continues where it stopped, never redoing work. For very large trees (e.g. 58K files), scan_directory can be slow to walk; use walk_files for a streaming generator when batch_size is set. """ from pathlib import Path from prefect import flow, task from PIL import Image import photo_db as db EXT_IMAGES = {".jpg", ".jpeg", ".png", ".heic", ".webp", ".gif", ".tif", ".tiff", ".bmp"} def _is_image(p: Path) -> bool: return p.is_file() and p.suffix.lower() in EXT_IMAGES def walk_files(base_dir: str): """Stream image files under base_dir (generator — memory-safe for 50K+ files).""" root = Path(base_dir) if not root.exists(): raise FileNotFoundError(f"{root} does not exist") for p in root.rglob("*"): if _is_image(p): yield str(p) @task def find_unprocessed(base_dir: str, batch_size: int, source: str = None) -> list[str]: """Find the next batch of files NOT yet in the fingerprint DB.""" db.init_db() conn = db.get_db() conn.execute( "CREATE TABLE IF NOT EXISTS known_paths (path TEXT PRIMARY KEY, sha256 TEXT NOT NULL)" ) batch = [] for p_str in walk_files(base_dir): # cheap check: path already seen (image_hashes OR known_paths) known = conn.execute( "SELECT 1 FROM image_hashes WHERE path=? UNION SELECT 1 FROM known_paths WHERE path=?", (p_str, p_str), ).fetchone() if known: continue # content check: sha256 already registered (catches same photo at other paths) sha = db.sha256_file(p_str) sha_known = conn.execute( "SELECT 1 FROM image_hashes WHERE sha256=?", (sha,) ).fetchone() if sha_known: # record this path in known_paths so we don't re-hash it every loop conn.execute( "INSERT OR IGNORE INTO known_paths (path, sha256) VALUES (?,?)", (p_str, sha)) continue batch.append(p_str) if len(batch) >= batch_size: break conn.commit() conn.close() print(f"find_unprocessed: {len(batch)} new files (batch_size={batch_size})") return batch @task def check_and_register(image_paths: list[str], source: str) -> dict: """Hash + dedup-check + register a batch. Returns verdict counts.""" db.init_db() conn = db.get_db() exact_dups = [] near_dups = [] new_images = [] for p_str in image_paths: p = Path(p_str) try: # exact dup by sha sha = db.sha256_file(p) match = conn.execute( "SELECT path, source FROM image_hashes WHERE sha256=?", (sha,) ).fetchone() if match: exact_dups.append((p_str, match[0])) continue # near dup by perceptual hash hx = db.hash_image(p) near, dist = db._near_dup_lookup(conn, hx) if near and dist <= 10: near_dups.append((p_str, near, dist)) continue # new — register with Image.open(p) as im: w, h = im.size conn.execute( "INSERT INTO image_hashes (sha256, phash, dhash, file_size, width, height, path, source) " "VALUES (?,?,?,?,?,?,?,?)", (sha, hx["phash"], hx["dhash"], p.stat().st_size, w, h, str(p), source), ) new_images.append(p_str) except Exception as e: print(f" SKIP {p.name}: {type(e).__name__}: {e}") conn.commit() conn.close() return { "total": len(image_paths), "exact_dups": len(exact_dups), "near_dups": len(near_dups), "new": len(new_images), "exact_dup_list": exact_dups[:20], "near_dup_list": near_dups[:20], "new_list": new_images, } @flow(name="photo-ingest") def photo_ingest(base_dir: str, source: str = "takeout", batch_size: int = 1000, max_batches: int | None = None): """Hash + dedup-check a folder against the persistent library, in batches. Args: base_dir: folder to scan source: label for the batch (e.g. takeout, archive-pictures) batch_size: files per batch/checkpoint (default 1000) max_batches: stop after N batches (useful for testing) — None = all """ db.init_db() processed_batches = 0 totals = {"exact_dups": 0, "near_dups": 0, "new": 0} while True: batch = find_unprocessed(base_dir, batch_size, source) if not batch: print("No more unprocessed files — done.") break result = check_and_register(batch, source) totals["exact_dups"] += result["exact_dups"] totals["near_dups"] += result["near_dups"] totals["new"] += result["new"] processed_batches += 1 print(f"batch {processed_batches} done: {result['total']} files, " f"{result['new']} new, {result['exact_dups']} exact, {result['near_dups']} near") if max_batches and processed_batches >= max_batches: print(f"Stopped after {processed_batches} batches (max_batches={max_batches})") break totals["batches"] = processed_batches return totals if __name__ == "__main__": import sys d = sys.argv[1] if len(sys.argv) > 1 else "/tmp/sample" src = sys.argv[2] if len(sys.argv) > 2 else "cli-test" bs = int(sys.argv[3]) if len(sys.argv) > 3 else 1000 mb = int(sys.argv[4]) if len(sys.argv) > 4 else None r = photo_ingest(d, source=src, batch_size=bs, max_batches=mb) print(r)