Files
photo-pipeline/takeout_fetch.py

171 lines
5.9 KiB
Python

"""photo-pipeline: Takeout download + extract flow (v2).
Completes the dump-tracking loop:
1. Download Takeout archives (resume-capable via aria2c)
2. Verify archive sha256 against takeout_archives (no re-download of known archives)
3. Extract to /mnt/data/takeout/<export>/
4. Update status in takeout_archives
Takeout gives each export a set of archive URLs (the "Download" page). Save them
to a manifest file (one URL per line) or pass a list.
"""
import shutil
import subprocess
from pathlib import Path
from prefect import flow, task
import photo_db as db
STAGING = Path("/mnt/data/takeout")
@task
def download_archive(url: str, dest_dir: Path = STAGING) -> Path:
"""Download one Takeout archive with aria2c (resume-capable, 4 connections).
Returns the local archive path. Skips if archive sha256 already known+downloaded.
"""
dest_dir.mkdir(parents=True, exist_ok=True)
archive_name = url.split("/")[-1].split("?")[0]
if not archive_name:
archive_name = f"takeout-{abs(hash(url))}.zip"
dest = dest_dir / archive_name
# Skip if we already have it downloaded (by name)
if dest.exists() and dest.stat().st_size > 0:
print(f" already on disk: {archive_name} ({dest.stat().st_size} bytes) — verifying sha256")
sha = db.sha256_file(dest)
conn = db.get_db()
known = conn.execute(
"SELECT id FROM takeout_archives WHERE archive_name=? AND sha256=?", (archive_name, sha)
).fetchone()
conn.close()
if known:
print(f" verified — already tracked (id={known[0]}). Skipping download.")
return dest
print(f" on disk but untracked sha — re-downloading to be safe")
print(f"Downloading {archive_name}{dest}")
cmd = [
"aria2c",
"-x", "4", "-s", "4", # 4 connections
"-c", # continue/resume
"--auto-file-renaming=false",
"--allow-overwrite=false",
"-d", str(dest_dir),
"-o", archive_name,
url,
]
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode != 0:
raise RuntimeError(f"aria2c failed for {archive_name}:\n{r.stderr[-1000:]}")
return dest
@task
def track_archive(archive_path: Path, export_id: str | None = None) -> int:
"""Record the downloaded archive in takeout_archives (or update it)."""
db.init_db()
conn = db.get_db()
archive_id = db.register_takeout_archive(conn, archive_path.name, archive_path)
if export_id:
conn.execute(
"UPDATE takeout_archives SET export_id=? WHERE id=?", (export_id, archive_id)
)
conn.commit()
conn.close()
print(f"Tracked archive id={archive_id} ({archive_path.name}, {archive_path.stat().st_size} bytes)")
return archive_id
@task
def extract_archive(archive_path: Path, export_dir: str | None = None) -> Path:
"""Extract a Takeout archive to /mnt/data/takeout/<archive_stem>/.
Handles .zip (unzip or python zipfile) and .tgz/.tar.gz (tar).
"""
if export_dir:
out_root = STAGING / export_dir
else:
out_root = STAGING / archive_path.stem.replace(".tar", "")
out_root.mkdir(parents=True, exist_ok=True)
if archive_path.suffix == ".zip" or archive_path.name.endswith(".zip"):
# Try unzip first; fall back to python zipfile (no external dep)
if shutil.which("unzip"):
r = subprocess.run(
["unzip", "-o", "-q", str(archive_path), "-d", str(out_root)],
capture_output=True, text=True,
)
if r.returncode != 0:
raise RuntimeError(f"unzip failed: {r.stderr[-500:]}")
else:
import zipfile
with zipfile.ZipFile(archive_path) as zf:
zf.extractall(out_root)
elif archive_path.name.endswith((".tgz", ".tar.gz", ".tar")):
r = subprocess.run(
["tar", "-xzf", str(archive_path), "-C", str(out_root)],
capture_output=True, text=True,
)
if r.returncode != 0:
raise RuntimeError(f"tar failed: {r.stderr[-500:]}")
else:
raise ValueError(f"Unsupported archive type: {archive_path.name}")
# Count extracted files
n = sum(1 for _ in out_root.rglob("*") if _.is_file())
db.init_db()
conn = db.get_db()
conn.execute(
"UPDATE takeout_archives SET extracted_at=datetime('now'), extract_path=?, file_count=?, status='extracted' "
"WHERE archive_name=?",
(str(out_root), n, archive_path.name),
)
conn.commit()
conn.close()
print(f"Extracted {n} files → {out_root}")
return out_root
@flow(name="takeout-fetch")
def takeout_fetch(
urls: list[str] | None = None,
manifest: str | None = None,
export_id: str | None = None,
):
"""Download + track + extract a set of Takeout archives.
Args:
urls: list of archive download URLs (from takeout.google.com Download page)
manifest: path to a text file with one URL per line (alternative to urls)
export_id: label for this export (e.g. "photos-2026-08")
"""
if manifest:
mp = Path(manifest)
if not mp.exists():
raise FileNotFoundError(f"manifest not found: {mp}")
urls = [l.strip() for l in mp.read_text().splitlines() if l.strip().startswith("http")]
if not urls:
raise ValueError("provide urls or a manifest file")
results = []
for url in urls:
path = download_archive(url)
archive_id = track_archive(path, export_id=export_id)
out = extract_archive(path, export_dir=export_id)
results.append({"archive": path.name, "id": archive_id, "extracted_to": str(out)})
return results
if __name__ == "__main__":
import sys
if len(sys.argv) < 2:
print("usage: python takeout_fetch.py <manifest-file> [export_id]")
sys.exit(1)
man = sys.argv[1]
eid = sys.argv[2] if len(sys.argv) > 2 else None
takeout_fetch(manifest=man, export_id=eid)