Add watch-folder trigger flow (auto-chain on incoming manifests/archives)

This commit is contained in:
2026-08-08 09:19:58 +10:00
parent 10b4b0b2f1
commit 4209a1153e
2 changed files with 125 additions and 0 deletions

110
photo_watch.py Normal file
View File

@@ -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()

View File

@@ -87,3 +87,18 @@ deployments:
name: photo-pool name: photo-pool
work_queue_name: null work_queue_name: null
job_variables: {} 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: {}