diff --git a/dashboard/app.py b/dashboard/app.py index 5d6de60..2238eac 100644 --- a/dashboard/app.py +++ b/dashboard/app.py @@ -96,25 +96,34 @@ def index(request: Request): @app.get("/review", response_class=HTMLResponse) -def review(request: Request, source: str = None, status: str = None): +def review(request: Request, source: str = None, status: str = None, page: int = 1): conn = _conn() + per_page = 200 q = "SELECT sha256, path, status, source FROM image_hashes WHERE 1=1" + count_q = "SELECT COUNT(*) FROM image_hashes WHERE 1=1" params = [] if source: q += " AND source=?" + count_q += " AND source=?" params.append(source) if status: q += " AND status=?" + count_q += " AND status=?" params.append(status) else: q += " AND status IN ('scanned','review')" - q += " ORDER BY added_at DESC LIMIT 300" - rows = conn.execute(q, params).fetchall() + count_q += " AND status IN ('scanned','review')" + total = conn.execute(count_q, params).fetchone()[0] + pages = max(1, (total + per_page - 1) // per_page) + page = max(1, min(page, pages)) + q += " ORDER BY added_at DESC LIMIT ? OFFSET ?" + rows = conn.execute(q, params + [per_page, (page - 1) * per_page]).fetchall() conn.close() items = [_load_item(r) for r in rows] return templates.TemplateResponse( request, "review.html", - {"items": items, "source": source, "status": status}, + {"items": items, "source": source, "status": status, + "page": page, "pages": pages, "total": total}, ) @@ -231,6 +240,47 @@ def thumb_file(name: str): return FileResponse(f) +@app.get("/upload", response_class=HTMLResponse) +def upload_page(request: Request): + """Upload page β€” drop files/archives into the incoming folder.""" + return templates.TemplateResponse(request, "upload.html", {}) + + +@app.post("/upload") +async def upload(request: Request): + """Receive uploaded files β†’ save to /mnt/data/takeout/incoming/.""" + import uuid + + from starlette.datastructures import UploadFile + + form = await request.form() + incoming = STAGING / "takeout" / "incoming" + incoming.mkdir(parents=True, exist_ok=True) + saved = [] + for field in form.values(): + if isinstance(field, UploadFile) and field.filename: + # sanitize: keep name but avoid path traversal + name = Path(field.filename).name + dest = incoming / f"{uuid.uuid4().hex[:8]}_{name}" + with open(dest, "wb") as f: + while chunk := await field.read(1024 * 1024): + f.write(chunk) + saved.append(dest.name) + # notify + try: + import sys + if str(BASE.parent) not in sys.path: + sys.path.insert(0, str(BASE.parent)) + import apprise_helper + apprise_helper.notify( + "πŸ“₯ photo-pipeline: upload received", + f"{len(saved)} file(s) saved to incoming. Watch flow will process them.", + ) + except Exception: + pass + return {"saved": len(saved), "files": saved} + + if __name__ == "__main__": import uvicorn diff --git a/dashboard/templates/review.html b/dashboard/templates/review.html index 1cc25b3..94a6ca3 100644 --- a/dashboard/templates/review.html +++ b/dashboard/templates/review.html @@ -41,6 +41,9 @@ .none { color: #666; font-style: italic; padding: 2rem; text-align: center; } .selall { display: inline-flex; align-items: center; gap: .35rem; background: #222; border: 1px solid #444; border-radius: 6px; padding: .4rem .8rem; cursor: pointer; } .toast { position: fixed; bottom: 1.2rem; right: 1.2rem; background: #1d4; color: #031; padding: .7rem 1rem; border-radius: 8px; display: none; z-index: 50; font-size: .9rem; box-shadow: 0 2px 12px rgba(0,0,0,.5); } + .pager { display: flex; gap: 1rem; align-items: center; justify-content: center; padding: 1.5rem; } + .pager a { color: #6cf; text-decoration: none; padding: .4rem .8rem; background: #222; border: 1px solid #444; border-radius: 6px; } + .pager span { color: #999; } @@ -72,7 +75,7 @@
-

{{ items|length }} images

+

{{ total }} images Β· page {{ page }}/{{ pages }}

{% for item in items %} @@ -81,6 +84,17 @@

No images match the filter.

{% endfor %}
+ {% if pages > 1 %} +
+ {% if page > 1 %} + ← Prev + {% endif %} + page {{ page }} / {{ pages }} + {% if page < pages %} + Next β†’ + {% endif %} +
+ {% endif %}
+ + diff --git a/photo_db.py b/photo_db.py index 8e4cb4c..8b55c3d 100644 --- a/photo_db.py +++ b/photo_db.py @@ -161,3 +161,20 @@ def register_takeout_archive(conn, archive_name: str, path: Path) -> int: if __name__ == "__main__": init_db() print(f"DB ready at {DB_PATH}") + + +def _near_dup_lookup(conn, hx: dict, hamming_threshold: int = 10): + """Find nearest perceptual-hash match, given precomputed hashes (avoids re-hash).""" + rows = conn.execute("SELECT phash, dhash, path, source FROM image_hashes").fetchall() + best, best_dist = None, None + for ph, dh, p, src in rows: + phd = hamming(ph, hx["phash"]) + dhd = hamming(dh, hx["dhash"]) + dist = min(phd, dhd) + if best_dist is None or dist < best_dist: + best, best_dist = (p, src), dist + if best_dist == 0: + break + if best and best_dist <= hamming_threshold: + return best, best_dist + return None, best_dist diff --git a/photo_ingest.py b/photo_ingest.py index 6940cb1..ce8e6c3 100644 --- a/photo_ingest.py +++ b/photo_ingest.py @@ -1,38 +1,69 @@ -"""photo-pipeline: photo-ingest flow (v1). +"""photo-pipeline: photo-ingest flow (v2 β€” batched + checkpointed). -Stage 1 of the pipeline: hash incoming images, check against the persistent -fingerprint DB (exact + near dupes), and register new ones. +Hashes incoming images, checks against the persistent fingerprint DB +(exact + near dupes), and registers new ones. -Run via Prefect deployment on photo-pool (see prefect.yaml). +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. """ -import sqlite3 from pathlib import Path from prefect import flow, task import photo_db as db +EXT_IMAGES = {".jpg", ".jpeg", ".png", ".heic", ".webp", ".gif", ".tif", ".tiff", ".bmp"} -@task -def scan_directory(base_dir: str) -> list[str]: - """Enumerate image files in a directory tree.""" - exts = {".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") - found = [ - str(p) - for p in root.rglob("*") - if p.is_file() and p.suffix.lower() in exts - ] - print(f"Found {len(found)} images under {root}") - return found + for p in root.rglob("*"): + if _is_image(p): + yield str(p) @task -def check_duplicates(image_paths: list[str]) -> dict: - """Check each image against the fingerprint DB. Returns classification.""" +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() + batch = [] + for p_str in walk_files(base_dir): + # skip if already registered for this source (or any source) + row = conn.execute( + "SELECT 1 FROM image_hashes WHERE sha256=?", + (db.sha256_file(p_str),), + ).fetchone() if False else None + # cheap check: path already known? + known = conn.execute( + "SELECT 1 FROM image_hashes WHERE path=?", (p_str,) + ).fetchone() + if known: + continue + batch.append(p_str) + if len(batch) >= batch_size: + break + 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 = [] @@ -41,75 +72,85 @@ def check_duplicates(image_paths: list[str]) -> dict: for p_str in image_paths: p = Path(p_str) try: - match = db.check_exact_dup(conn, p) + # 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)) + exact_dups.append((p_str, match[0])) continue - near, dist = db.check_near_dup(conn, p) - if near: + # 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 __import__("PIL.Image", fromlist=["Image"]).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}: {e}") + print(f" SKIP {p.name}: {type(e).__name__}: {e}") + conn.commit() conn.close() - result = { + return { "total": len(image_paths), "exact_dups": len(exact_dups), "near_dups": len(near_dups), "new": len(new_images), - "exact_dup_list": exact_dups[:50], - "near_dup_list": near_dups[:50], + "exact_dup_list": exact_dups[:20], + "near_dup_list": near_dups[:20], "new_list": new_images, } - print( - f"Check: {result['total']} total, " - f"{result['exact_dups']} exact dups, {result['near_dups']} near dups, " - f"{result['new']} new" - ) - return result - - -@task -def register_new_images(image_paths: list[str], source: str) -> int: - """Add hashes for confirmed-new images into the fingerprint DB.""" - db.init_db() - conn = db.get_db() - registered = 0 - for p_str in image_paths: - p = Path(p_str) - try: - if db.register_image(conn, p, source=source): - registered += 1 - except Exception as e: - print(f" FAIL register {p.name}: {e}") - conn.commit() - conn.close() - print(f"Registered {registered} new images (source={source})") - return registered @flow(name="photo-ingest") -def photo_ingest(base_dir: str, source: str = "takeout", register: bool = True): - """Hash + dedup-check a folder against the persistent library.""" - images = scan_directory(base_dir) - if not images: - print("No images found β€” nothing to do.") - return {"total": 0} +def photo_ingest(base_dir: str, source: str = "takeout", batch_size: int = 1000, + max_batches: int = None): + """Hash + dedup-check a folder against the persistent library, in batches. - result = check_duplicates(images) + 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} - if register and result["new_list"]: - n = register_new_images(result["new_list"], source=source) - result["registered"] = n + 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 - return result + totals["batches"] = processed_batches + return totals if __name__ == "__main__": - # Local run (no deployment) import sys d = sys.argv[1] if len(sys.argv) > 1 else "/tmp/sample" - r = photo_ingest(d) + 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) diff --git a/prefect.yaml b/prefect.yaml index 59c3105..92b1c54 100644 --- a/prefect.yaml +++ b/prefect.yaml @@ -29,14 +29,14 @@ deployments: - name: ingest version: null tags: [photo-pipeline] - description: "Hash + dedup-check a folder against the persistent fingerprint DB" + description: "Hash + dedup-check a folder against the persistent fingerprint DB (batched, checkpointed)" schedule: null flow_name: null entrypoint: photo_ingest.py:photo_ingest parameters: base_dir: /mnt/data/takeout source: takeout - register: true + batch_size: 1000 work_pool: name: photo-pool work_queue_name: null @@ -68,6 +68,7 @@ deployments: parameters: base_dir: /mnt/data/takeout move: false + max_files: 0 work_pool: name: photo-pool work_queue_name: null diff --git a/quality_scan.py b/quality_scan.py index dd66810..e52c9b9 100644 --- a/quality_scan.py +++ b/quality_scan.py @@ -86,9 +86,37 @@ def _move(p: Path, dest_root: Path, src_root: Path): shutil.move(str(p), str(dest)) + + +def _slice_dir(base_dir: str, max_files: int) -> str: + """Copy first N images into a temp dir for CleanVision to audit.""" + import shutil + import tempfile + from pathlib import Path + + src = Path(base_dir) + tmp = Path(tempfile.mkdtemp(prefix="cvslice_")) + exts = {".jpg", ".jpeg", ".png", ".webp", ".gif", ".heic", ".tif", ".bmp"} + n = 0 + for p in src.rglob("*"): + if p.is_file() and p.suffix.lower() in exts: + shutil.copy2(p, tmp / p.name) + n += 1 + if n >= max_files: + break + print(f"_slice_dir: copied {n} files to {tmp}") + return str(tmp) + + @flow(name="photo-quality-scan") -def quality_scan(base_dir: str, move: bool = False, notify: bool = True): - """Audit image quality with CleanVision; classify into keep/review/delete.""" +def quality_scan(base_dir: str, move: bool = False, notify: bool = True, max_files: int = None): + """Audit image quality with CleanVision; classify into keep/review/delete. + + max_files: if set, only audit the first N image files (slices huge folders + into reviewable chunks β€” prevents OOM on 50K-file trees). + """ + if max_files: # 0/None = unlimited + base_dir = _slice_dir(base_dir, max_files) audit = audit_folder(base_dir) print(f"Issue summary: {audit['summary']}") result = classify_and_sort(base_dir, audit["per_image"], move=move)