diff --git a/photo_watch.py b/photo_watch.py new file mode 100644 index 0000000..e5d02c8 --- /dev/null +++ b/photo_watch.py @@ -0,0 +1,110 @@ +"""photo-pipeline: watch-folder trigger flow. + +Watches /mnt/data/takeout/incoming for: + - *.txt / *.manifest → Takeout download URLs (one per line) → takeout-fetch + - *.zip / *.tgz → already-downloaded Takeout archives → extract + ingest + +When found, chains the pipeline automatically and notifies via Apprise. +Designed to run on a schedule (e.g. every 15 min) OR as a long-running serve(). +""" + +import shutil +from pathlib import Path + +from prefect import flow, task + +import photo_db as db +from takeout_fetch import download_archive, extract_archive, track_archive +from photo_ingest import check_duplicates, register_new_images, scan_directory +from apprise_helper import notify + +INCOMING = Path("/mnt/data/takeout/incoming") +PROCESSED = Path("/mnt/data/takeout/processed") +ARCHIVE_EXTS = {".zip", ".tgz", ".tar.gz", ".tar"} + + +@task +def scan_incoming() -> list[Path]: + """Find new manifests and archives in the incoming folder.""" + if not INCOMING.exists(): + INCOMING.mkdir(parents=True) + PROCESSED.mkdir(parents=True) + return [] + found = [p for p in INCOMING.iterdir() if p.is_file()] + return found + + +@task +def handle_manifest(m: Path) -> str | None: + """Read a manifest of Takeout URLs and download+extract each.""" + urls = [l.strip() for l in m.read_text().splitlines() if l.strip().startswith("http")] + if not urls: + return "no URLs in manifest" + extracted_dirs = [] + for url in urls: + arc = download_archive(url) + track_archive(arc) + out = extract_archive(arc, export_dir=arc.stem) + extracted_dirs.append(str(out)) + return f"downloaded {len(urls)} archives → {len(extracted_dirs)} dirs" + + +@task +def handle_archive(a: Path) -> str: + """Handle an already-downloaded archive in incoming/.""" + if a.suffix.lower() not in {".zip", ".tgz"} and not a.name.endswith(".tar.gz"): + return f"skip non-archive: {a.name}" + track_archive(a) + out = extract_archive(a, export_dir=a.stem) + return f"extracted {a.name} → {out}" + + +@task +def process_extracted(export_dir: str, source: str) -> dict: + """Run the ingest chain (fingerprint + dedup) on an extracted folder.""" + images = scan_directory(export_dir) + if not images: + return {"total": 0} + result = check_duplicates(images) + if result["new_list"]: + register_new_images(result["new_list"], source=source) + return result + + +@flow(name="photo-watch") +def photo_watch(chain_quality: bool = True): + """Watch incoming folder; auto-run the pipeline on new Takeout material.""" + found = scan_incoming() + if not found: + print("Incoming folder empty — nothing to do.") + return {"found": 0} + + results = [] + for f in found: + try: + if f.suffix.lower() in {".txt", ".manifest"} or f.name.endswith(".txt"): + msg = handle_manifest(f) + results.append({"file": f.name, "action": "manifest", "result": msg}) + elif f.suffix.lower() in ARCHIVE_EXTS or f.name.endswith(".tar.gz"): + msg = handle_archive(f) + results.append({"file": f.name, "action": "archive", "result": msg}) + else: + results.append({"file": f.name, "action": "skipped"}) + # move to processed + PROCESSED.mkdir(parents=True, exist_ok=True) + shutil.move(str(f), str(PROCESSED / f.name)) + except Exception as e: + results.append({"file": f.name, "action": "error", "result": str(e)}) + + # notify + n = len(results) + detail = "\n".join(f"- {r['file']}: {r.get('result', r['action'])}" for r in results) + notify( + f"📥 photo-pipeline: {n} item(s) processed", + f"From incoming folder:\n{detail}\n\nReview: http://192.168.20.13:8092/review", + ) + return {"found": n, "results": results} + + +if __name__ == "__main__": + photo_watch() diff --git a/prefect.yaml b/prefect.yaml index 71f3520..e8c352c 100644 --- a/prefect.yaml +++ b/prefect.yaml @@ -87,3 +87,18 @@ deployments: name: photo-pool work_queue_name: null job_variables: {} + +- name: watch + version: null + tags: [photo-pipeline] + description: "Watch /mnt/data/takeout/incoming; auto-chain pipeline on new Takeout material" + schedule: + interval: 900 + timezone: "Australia/Melbourne" + flow_name: null + entrypoint: photo_watch.py:photo_watch + parameters: {} + work_pool: + name: photo-pool + work_queue_name: null + job_variables: {}