Files
immich-photo-flow/docs/superpowers/plans/2026-06-27-immich-photo-flow-m1.md
T
2026-06-27 19:25:15 +02:00

3732 lines
141 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# immich-photo-flow M1 — Foundation + trip-cluster Implementation Plan
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** Stand up the monorepo + shared packages (`immich`, `core`, `ui`) and ship the `trip-cluster` app end-to-end (ingest → cluster → review UI → idempotent write-back to Immich), the POC that proves the shared foundation.
**Architecture:** A Python 3.12 monorepo. `shared/` is a single editable distribution `photoflow` with three sub-packages: `photoflow.immich` (the one Immich REST client + `_pipeline/` tag conventions), `photoflow.core` (SQLite store + domain dataclasses), `photoflow.ui` (base.html / DaisyUI+Tailwind+Alpine+HTMX / shared Jinja macros + grid-lightbox JS). `apps/trip-cluster/` is a Flask factory + argparse CLI that consumes those packages. Immich is the source of truth; SQLite is a rebuildable working/review layer; all durable results are written back as tags.
**Tech Stack:** Python 3.12, Flask 3.1, requests 2.32, sqlite3 (stdlib), Pillow, DaisyUI 4 + Tailwind + Alpine.js + HTMX (CDN, no build step), pytest 8.3 + pytest-httpserver 1.1 + pytest-playwright 0.6.
## Global Constraints
- **Python 3.12.** Stack mirrors the sibling apps (`/home/mischa/Projects/claude-image-rater`, `claude-travel-memories`): Flask factory `create_app()` + argparse CLI sharing one package; no-build frontend; pytest + pytest-httpserver + Playwright; TDD (failing test first); frequent commits.
- **One module talks to Immich** (`photoflow.immich.client.ImmichClient`); **one place owns SQLite** (`photoflow.core.store.Store`). No other module issues HTTP to Immich or SQL.
- **trip-cluster serves on port 8084** (8082 = travel-memories, 8083 = image-rater). Reserve it; never reuse 8082/8083.
- **Pin versions** in every `requirements`/`dependencies` list: `flask==3.1.0`, `requests==2.32.3`, `Pillow==11.0.0`, `pytest==8.3.4`, `pytest-httpserver==1.1.0`, `pytest-playwright==0.6.2`. `anthropic` is NOT a dependency in M1 (`shared/ai` deferred to M3); `ANTHROPIC_API_KEY` plumbing stays optional.
- **Env vars** (same names as siblings): `IMMICH_URL`, `IMMICH_API_KEY` (required), `ANTHROPIC_API_KEY` (optional in M1), `DATA_DIR` (SQLite DB + thumbnails). Loaded from `.env`.
- **Tag conventions:** content/trip/location tags are user-facing and **never namespaced**; pipeline meta-tags nest under a single parent **`_pipeline/`** — `_pipeline/processed`, `_pipeline/non-trip`. Defined once in `photoflow.immich.pipeline`. Always **reconcile against the live Immich tag list** before creating a tag (use `upsert_tag`, never blind-create duplicates).
- **Idempotent write-back:** every change pushed to Immich is recorded in `writeback_log`; re-runs skip what's already applied. **Explicit confirmation before any write** (CLI prompt / UI button).
- **Docker:** one container for trip-cluster, port 8084, `user: ${UID}:${GID}`, DB + thumbs on a mounted volume.
- **Import names:** `from photoflow.immich import ImmichClient`, `from photoflow.immich import pipeline`, `from photoflow.core import Store`, `from photoflow.core.models import Asset, Cluster, ...`, `from photoflow.ui import register_shared_ui`. App code lives in package `app` under `apps/trip-cluster/`.
---
## File Structure
```
immich-photo-flow/
pyproject.toml # repo-root: dev tooling + pytest config
docker-compose.yml # trip-cluster service (port 8084)
.env.example
.gitignore
README.md
CLAUDE.md
shared/
pyproject.toml # dist "photoflow" (editable)
photoflow/
__init__.py
immich/
__init__.py # exports ImmichClient
client.py # the ONLY Immich REST client
pipeline.py # _pipeline/ tag conventions
core/
__init__.py # exports Store
models.py # Asset, Tag, AssetTag, Cluster, ClusterMember, WritebackLog, Meta
store.py # SQLite schema + data-access (the ONLY SQL)
ui/
__init__.py # TEMPLATE_DIR, STATIC_DIR, register_shared_ui()
templates/
base.html
macros.html # image_grid, lightbox, badges
static/
shared.js # photoGrid() Alpine factory (grid + lightbox)
tests/
test_immich_client.py
test_pipeline.py
test_store_assets.py
test_store_clusters.py
test_ui.py
apps/
trip-cluster/
pyproject.toml # dist "trip-cluster", depends on photoflow
Dockerfile
categorize.py # CLI entry point -> app.cli.main()
app/
__init__.py # create_app()
config.py
cli.py # argparse: ingest | cluster | serve | apply
ingest.py
clustering.py # pure: assets -> candidate clusters
coverage.py # pure: completeness flags + outliers
cluster_run.py # orchestration: store <-> clustering
review.py # cluster review operations over Store
writeback.py # apply decisions to Immich (idempotent)
routes/
__init__.py
nav.py # "/" master/detail page + empty states
review.py # cluster actions + apply endpoints
proxy.py # /thumb/<asset_id>
templates/
review.html # master/detail layout
_detail.html # cluster detail partial (HTMX)
static/
app.js # clusterReview() Alpine component
tests/
conftest.py
test_config.py
test_ingest.py
test_clustering.py
test_coverage.py
test_cluster_run.py
test_review.py
test_routes.py
test_writeback.py
ui/
conftest.py
test_smoke_ui.py
test_review_ui.py
```
---
## Task 1: Monorepo scaffolding + editable install
**Files:**
- Create: `shared/pyproject.toml`, `shared/photoflow/__init__.py`, `shared/photoflow/immich/__init__.py`, `shared/photoflow/core/__init__.py`, `shared/photoflow/ui/__init__.py`
- Create: `apps/trip-cluster/pyproject.toml`, `apps/trip-cluster/app/__init__.py`, `apps/trip-cluster/categorize.py`
- Create: `pyproject.toml` (root), `.gitignore`, `.env.example`, `apps/trip-cluster/tests/conftest.py`, `shared/tests/test_imports.py`
**Interfaces:**
- Produces: importable empty packages `photoflow`, `photoflow.immich`, `photoflow.core`, `photoflow.ui`; an editable dev environment where `import photoflow.immich` and `import app` both work.
- [ ] **Step 1: Write the failing import test**
`shared/tests/test_imports.py`:
```python
def test_shared_packages_import():
import photoflow.immich
import photoflow.core
import photoflow.ui
```
- [ ] **Step 2: Run it to verify it fails**
Run: `python -m pytest shared/tests/test_imports.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'photoflow'`
- [ ] **Step 3: Create the shared distribution**
`shared/pyproject.toml`:
```toml
[build-system]
requires = ["setuptools>=68"]
build-backend = "setuptools.build_meta"
[project]
name = "photoflow"
version = "0.1.0"
requires-python = ">=3.12"
dependencies = ["requests==2.32.3", "flask==3.1.0", "Pillow==11.0.0"]
[tool.setuptools.packages.find]
where = ["."]
include = ["photoflow*"]
[tool.setuptools.package-data]
"photoflow.ui" = ["templates/*.html", "static/*.js"]
```
Create the package files (each `__init__.py` empty for now except where noted):
- `shared/photoflow/__init__.py` → empty
- `shared/photoflow/immich/__init__.py` → empty
- `shared/photoflow/core/__init__.py` → empty
- `shared/photoflow/ui/__init__.py` → empty
- [ ] **Step 4: Create the trip-cluster distribution**
`apps/trip-cluster/pyproject.toml`:
```toml
[build-system]
requires = ["setuptools>=68"]
build-backend = "setuptools.build_meta"
[project]
name = "trip-cluster"
version = "0.1.0"
requires-python = ">=3.12"
dependencies = ["photoflow", "flask==3.1.0", "requests==2.32.3", "Pillow==11.0.0"]
[tool.setuptools.packages.find]
where = ["."]
include = ["app*"]
```
`apps/trip-cluster/app/__init__.py` → empty for now.
`apps/trip-cluster/categorize.py`:
```python
import sys
from app.cli import main
if __name__ == "__main__":
sys.exit(main())
```
- [ ] **Step 5: Root tooling, gitignore, env example**
`pyproject.toml` (root):
```toml
[tool.pytest.ini_options]
testpaths = ["shared/tests", "apps/trip-cluster/tests"]
addopts = "-q"
```
`.gitignore`:
```
__pycache__/
*.pyc
.venv/
data/
*.db
.env
.pytest_cache/
*.egg-info/
```
`.env.example`:
```
IMMICH_URL=http://your-immich-host:2283
IMMICH_API_KEY=your-immich-api-key
ANTHROPIC_API_KEY=
DATA_DIR=./data
UID=1000
GID=1000
```
`apps/trip-cluster/tests/conftest.py`:
```python
# Editable installs put `app` and `photoflow` on sys.path; nothing extra needed.
```
- [ ] **Step 6: Create the venv and editable-install both distributions**
Run:
```bash
python3.12 -m venv .venv
.venv/bin/pip install -e ./shared -e ./apps/trip-cluster pytest==8.3.4 pytest-httpserver==1.1.0 pytest-playwright==0.6.2
```
Expected: both packages install; `Successfully installed photoflow-0.1.0 trip-cluster-0.1.0 ...`
- [ ] **Step 7: Run the import test to verify it passes**
Run: `.venv/bin/python -m pytest shared/tests/test_imports.py -q`
Expected: PASS (1 passed)
- [ ] **Step 8: Commit**
```bash
git add -A
git commit -m "feat: monorepo scaffold — photoflow shared dist + trip-cluster app skeleton"
```
---
## Task 2: ImmichClient — read (search, tags, thumbnail)
**Files:**
- Create: `shared/photoflow/immich/client.py`
- Modify: `shared/photoflow/immich/__init__.py`
- Test: `shared/tests/test_immich_client.py`
**Interfaces:**
- Consumes: nothing.
- Produces:
- `ImmichClient(base_url: str, api_key: str, timeout: int = 30)`
- `.list_tags() -> list[dict]` — raw Immich tag dicts (`id`, `name`, `value`)
- `.resolve_tag_id(name: str) -> str | None` — matches on `value` (full nested path) or `name`
- `.search_assets(*, taken_after: str | None = None, taken_before: str | None = None, tag_ids: list[str] | None = None, updated_after: str | None = None) -> list[dict]` — each dict has keys: `id`, `original_filename`, `taken_at`, `gps_lat` (float|None), `gps_lon` (float|None), `place_city` (str|None), `place_country` (str|None), `type` (str), `tags` (list[str]), `rating` (int), `updated_at` (str). Paginates `POST /api/search/metadata` with `withExif: true`.
- `.download_thumbnail(asset_id: str) -> bytes`
- [ ] **Step 1: Write the failing tests**
`shared/tests/test_immich_client.py`:
```python
from werkzeug.wrappers import Response
from photoflow.immich import ImmichClient
def test_resolve_tag_id_matches_name_or_value(httpserver):
httpserver.expect_request("/api/tags").respond_with_json([
{"id": "t1", "name": "Italy 2019", "value": "Italy 2019"},
{"id": "t2", "name": "processed", "value": "_pipeline/processed"},
])
c = ImmichClient(httpserver.url_for(""), "k")
assert c.resolve_tag_id("Italy 2019") == "t1"
assert c.resolve_tag_id("_pipeline/processed") == "t2"
assert c.resolve_tag_id("nope") is None
def test_search_assets_normalizes_and_paginates(httpserver):
def handler(request):
page = request.json.get("page", 1)
assert request.json.get("withExif") is True
if page == 1:
return Response(
'{"assets": {"items": [{"id": "a", "originalFileName": "a.jpg",'
' "localDateTime": "2019-06-01T10:00:00.000Z", "type": "IMAGE",'
' "updatedAt": "2026-01-01T00:00:00Z",'
' "exifInfo": {"latitude": 45.4, "longitude": 12.3, "city": "Venezia",'
' "country": "Italy", "rating": 4},'
' "tags": [{"name": "Italy 2019"}]}], "nextPage": 2}}',
content_type="application/json")
return Response(
'{"assets": {"items": [{"id": "b", "originalFileName": "b.jpg",'
' "localDateTime": "2019-06-02T11:00:00.000Z", "type": "IMAGE",'
' "updatedAt": "2026-01-02T00:00:00Z", "exifInfo": {}, "tags": []}],'
' "nextPage": null}}', content_type="application/json")
httpserver.expect_request("/api/search/metadata", method="POST").respond_with_handler(handler)
c = ImmichClient(httpserver.url_for(""), "k")
assets = c.search_assets(taken_after="2019-01-01", taken_before="2020-01-01")
assert [a["id"] for a in assets] == ["a", "b"]
a = assets[0]
assert a["gps_lat"] == 45.4 and a["place_city"] == "Venezia"
assert a["tags"] == ["Italy 2019"] and a["rating"] == 4
assert a["taken_at"] == "2019-06-01T10:00:00.000Z"
assert assets[1]["gps_lat"] is None and assets[1]["tags"] == []
def test_download_thumbnail(httpserver):
httpserver.expect_request("/api/assets/a/thumbnail").respond_with_data(
b"\xff\xd8\xffjpegbytes", content_type="image/jpeg")
c = ImmichClient(httpserver.url_for(""), "k")
assert c.download_thumbnail("a").startswith(b"\xff\xd8\xff")
```
- [ ] **Step 2: Run to verify it fails**
Run: `.venv/bin/python -m pytest shared/tests/test_immich_client.py -q`
Expected: FAIL — `ImportError: cannot import name 'ImmichClient'`
- [ ] **Step 3: Implement the client**
`shared/photoflow/immich/client.py`:
```python
import requests
def _normalize(item: dict) -> dict:
exif = item.get("exifInfo") or {}
lat = exif.get("latitude")
lon = exif.get("longitude")
return {
"id": item["id"],
"original_filename": item.get("originalFileName", ""),
"taken_at": item.get("localDateTime", ""),
"gps_lat": float(lat) if lat is not None else None,
"gps_lon": float(lon) if lon is not None else None,
"place_city": exif.get("city"),
"place_country": exif.get("country"),
"type": item.get("type", ""),
"tags": [t.get("name", "") for t in (item.get("tags") or [])],
"rating": int(exif.get("rating") or 0),
"updated_at": item.get("updatedAt", ""),
}
class ImmichClient:
def __init__(self, base_url: str, api_key: str, timeout: int = 30):
self.base_url = base_url.rstrip("/")
self.timeout = timeout
self.session = requests.Session()
self.session.headers.update({"x-api-key": api_key, "Accept": "application/json"})
def _url(self, path: str) -> str:
return f"{self.base_url}{path}"
def list_tags(self) -> list[dict]:
r = self.session.get(self._url("/api/tags"), timeout=self.timeout)
r.raise_for_status()
return r.json()
def resolve_tag_id(self, name: str):
for tag in self.list_tags():
if tag.get("value") == name or tag.get("name") == name:
return tag["id"]
return None
def search_assets(self, *, taken_after=None, taken_before=None,
tag_ids=None, updated_after=None) -> list[dict]:
body = {"withExif": True}
if taken_after:
body["takenAfter"] = taken_after
if taken_before:
body["takenBefore"] = taken_before
if tag_ids:
body["tagIds"] = tag_ids
if updated_after:
body["updatedAfter"] = updated_after
out = []
page = 1
while True:
payload = dict(body, size=1000, page=page)
r = self.session.post(self._url("/api/search/metadata"),
json=payload, timeout=self.timeout)
r.raise_for_status()
block = r.json().get("assets", {})
out.extend(_normalize(it) for it in block.get("items", []))
nxt = block.get("nextPage")
if not nxt:
break
page = int(nxt)
return out
def download_thumbnail(self, asset_id: str) -> bytes:
r = self.session.get(self._url(f"/api/assets/{asset_id}/thumbnail?size=preview"),
timeout=self.timeout)
r.raise_for_status()
return r.content
```
`shared/photoflow/immich/__init__.py`:
```python
from photoflow.immich.client import ImmichClient
__all__ = ["ImmichClient"]
```
- [ ] **Step 4: Run to verify it passes**
Run: `.venv/bin/python -m pytest shared/tests/test_immich_client.py -q`
Expected: PASS (3 passed)
- [ ] **Step 5: Commit**
```bash
git add shared/photoflow/immich/ shared/tests/test_immich_client.py
git commit -m "feat(immich): read client — search_assets, tags, thumbnail"
```
---
## Task 3: ImmichClient — write + `_pipeline/` tag conventions
**Files:**
- Modify: `shared/photoflow/immich/client.py`
- Create: `shared/photoflow/immich/pipeline.py`
- Modify: `shared/photoflow/immich/__init__.py`
- Test: `shared/tests/test_pipeline.py`, append to `shared/tests/test_immich_client.py`
**Interfaces:**
- Consumes: `ImmichClient` (Task 2).
- Produces:
- `ImmichClient.upsert_tag(name: str) -> str` — idempotent create-or-resolve of a (possibly nested) tag via `PUT /api/tags` `{"tags": [name]}`; returns the tag id.
- `ImmichClient.tag_assets(tag_id: str, asset_ids: list[str]) -> None``PUT /api/tags/{tag_id}/assets` `{"ids": [...]}`.
- `photoflow.immich.pipeline`: constants `ROOT="_pipeline"`, `PROCESSED="_pipeline/processed"`, `NON_TRIP="_pipeline/non-trip"`; `ai_rating(n: int) -> str`; `is_pipeline_tag(name: str) -> bool`.
- [ ] **Step 1: Write the failing tests**
`shared/tests/test_pipeline.py`:
```python
from photoflow.immich import pipeline
def test_constants_and_helpers():
assert pipeline.ROOT == "_pipeline"
assert pipeline.PROCESSED == "_pipeline/processed"
assert pipeline.NON_TRIP == "_pipeline/non-trip"
assert pipeline.ai_rating(4) == "_pipeline/ai-rating/4"
def test_is_pipeline_tag():
assert pipeline.is_pipeline_tag("_pipeline/processed") is True
assert pipeline.is_pipeline_tag("_pipeline") is True
assert pipeline.is_pipeline_tag("Italy 2019") is False
assert pipeline.is_pipeline_tag("Kiev") is False
```
Append to `shared/tests/test_immich_client.py`:
```python
def test_upsert_tag_returns_id_by_value(httpserver):
httpserver.expect_request("/api/tags", method="PUT").respond_with_json(
[{"id": "p1", "name": "non-trip", "value": "_pipeline/non-trip"}])
c = ImmichClient(httpserver.url_for(""), "k")
assert c.upsert_tag("_pipeline/non-trip") == "p1"
def test_tag_assets_posts_ids(httpserver):
httpserver.expect_request("/api/tags/p1/assets", method="PUT").respond_with_json({"ok": True})
c = ImmichClient(httpserver.url_for(""), "k")
c.tag_assets("p1", ["a", "b"]) # should not raise
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest shared/tests/test_pipeline.py shared/tests/test_immich_client.py -q`
Expected: FAIL — `ModuleNotFoundError: ... pipeline` and `AttributeError: ... upsert_tag`
- [ ] **Step 3: Implement pipeline conventions**
`shared/photoflow/immich/pipeline.py`:
```python
ROOT = "_pipeline"
PROCESSED = f"{ROOT}/processed"
NON_TRIP = f"{ROOT}/non-trip"
def ai_rating(n: int) -> str:
return f"{ROOT}/ai-rating/{n}"
def is_pipeline_tag(name: str) -> bool:
return name == ROOT or name.startswith(ROOT + "/")
```
- [ ] **Step 4: Implement the write methods**
Append to `shared/photoflow/immich/client.py` (inside `ImmichClient`):
```python
def upsert_tag(self, name: str) -> str:
r = self.session.put(self._url("/api/tags"),
json={"tags": [name]}, timeout=self.timeout)
r.raise_for_status()
for tag in r.json():
if tag.get("value") == name or tag.get("name") == name:
return tag["id"]
resolved = self.resolve_tag_id(name)
if resolved is None:
raise RuntimeError(f"upsert_tag: could not resolve id for {name!r}")
return resolved
def tag_assets(self, tag_id: str, asset_ids: list[str]) -> None:
if not asset_ids:
return
r = self.session.put(self._url(f"/api/tags/{tag_id}/assets"),
json={"ids": asset_ids}, timeout=self.timeout)
r.raise_for_status()
```
Update `shared/photoflow/immich/__init__.py`:
```python
from photoflow.immich.client import ImmichClient
from photoflow.immich import pipeline
__all__ = ["ImmichClient", "pipeline"]
```
- [ ] **Step 5: Run to verify all pass**
Run: `.venv/bin/python -m pytest shared/tests/test_pipeline.py shared/tests/test_immich_client.py -q`
Expected: PASS (7 passed)
- [ ] **Step 6: Commit**
```bash
git add shared/photoflow/immich/ shared/tests/test_pipeline.py shared/tests/test_immich_client.py
git commit -m "feat(immich): write-back (upsert_tag, tag_assets) + _pipeline tag conventions"
```
---
## Task 4: core — domain models + Store (schema, connection, meta)
**Files:**
- Create: `shared/photoflow/core/models.py`, `shared/photoflow/core/store.py`
- Modify: `shared/photoflow/core/__init__.py`
- Test: `shared/tests/test_store_assets.py` (meta portion)
**Interfaces:**
- Produces:
- Dataclasses in `photoflow.core.models`: `Asset`, `Tag`, `AssetTag`, `Cluster`, `ClusterMember`, `WritebackLog`. (Field lists below — later tasks rely on exact names.)
- `Store(db_path: str)` with `.connect() -> Store` (creates schema, idempotent), `.close()`, `.conn` (sqlite3.Connection, `row_factory=Row`), `.set_meta(key, value)`, `.get_meta(key, default=None) -> str | None`. `SCHEMA_VERSION = 1`, written to `meta` on first connect.
- [ ] **Step 1: Write the failing test**
`shared/tests/test_store_assets.py`:
```python
from photoflow.core import Store
def test_connect_creates_schema_and_version(tmp_path):
db = str(tmp_path / "t.db")
s = Store(db).connect()
assert s.get_meta("schema_version") == "1"
# tables exist
names = {r["name"] for r in s.conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'")}
assert {"assets", "tags", "asset_tags", "clusters",
"cluster_members", "writeback_log", "meta"} <= names
s.close()
def test_meta_roundtrip_and_default(tmp_path):
s = Store(str(tmp_path / "t.db")).connect()
assert s.get_meta("missing") is None
assert s.get_meta("missing", "x") == "x"
s.set_meta("last_ingest_at", "2026-06-27T00:00:00Z")
assert s.get_meta("last_ingest_at") == "2026-06-27T00:00:00Z"
s.set_meta("last_ingest_at", "newer") # upsert
assert s.get_meta("last_ingest_at") == "newer"
s.close()
```
- [ ] **Step 2: Run to verify it fails**
Run: `.venv/bin/python -m pytest shared/tests/test_store_assets.py -q`
Expected: FAIL — `ImportError: cannot import name 'Store'`
- [ ] **Step 3: Write the domain models**
`shared/photoflow/core/models.py`:
```python
from dataclasses import dataclass, field
from typing import Optional
@dataclass
class Asset:
immich_id: str
taken_at: str = ""
gps_lat: Optional[float] = None
gps_lon: Optional[float] = None
place_city: Optional[str] = None
place_country: Optional[str] = None
type: str = ""
has_gps: bool = False
thumb_path: Optional[str] = None
processed: bool = False
ingested_at: str = ""
updated_at: str = ""
@dataclass
class Tag:
name: str
immich_tag_id: Optional[str] = None
count: int = 0
@dataclass
class AssetTag:
immich_id: str
tag_name: str
@dataclass
class Cluster:
id: Optional[int] = None
start_at: str = ""
end_at: str = ""
count: int = 0
suggested_name: str = ""
confidence: float = 0.0
kind_guess: str = "trip" # trip | everyday
status: str = "pending" # pending|approved|non_trip|merged|split|skipped
decided_name: Optional[str] = None
reviewed_at: Optional[str] = None
notes: Optional[str] = None
@dataclass
class ClusterMember:
cluster_id: int
immich_id: str
member_confidence: float = 1.0
is_outlier: bool = False
included: bool = True
flagged_coverage: bool = False
@dataclass
class WritebackLog:
id: Optional[int]
immich_id: str
action: str # trip | non-trip | processed
tag: Optional[str]
result: str # ok | error:<reason>
applied_at: str
```
- [ ] **Step 4: Write the Store base (schema + connect + meta)**
`shared/photoflow/core/store.py`:
```python
import sqlite3
SCHEMA_VERSION = 1
SCHEMA = """
CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT);
CREATE TABLE IF NOT EXISTS assets (
immich_id TEXT PRIMARY KEY,
taken_at TEXT,
gps_lat REAL, gps_lon REAL,
place_city TEXT, place_country TEXT,
type TEXT,
has_gps INTEGER NOT NULL DEFAULT 0,
thumb_path TEXT,
processed INTEGER NOT NULL DEFAULT 0,
ingested_at TEXT,
updated_at TEXT
);
CREATE TABLE IF NOT EXISTS tags (
name TEXT PRIMARY KEY,
immich_tag_id TEXT,
count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS asset_tags (
immich_id TEXT NOT NULL,
tag_name TEXT NOT NULL,
PRIMARY KEY (immich_id, tag_name)
);
CREATE TABLE IF NOT EXISTS clusters (
id INTEGER PRIMARY KEY AUTOINCREMENT,
start_at TEXT, end_at TEXT,
count INTEGER NOT NULL DEFAULT 0,
suggested_name TEXT,
confidence REAL NOT NULL DEFAULT 0,
kind_guess TEXT NOT NULL DEFAULT 'trip',
status TEXT NOT NULL DEFAULT 'pending',
decided_name TEXT,
reviewed_at TEXT,
notes TEXT
);
CREATE TABLE IF NOT EXISTS cluster_members (
cluster_id INTEGER NOT NULL,
immich_id TEXT NOT NULL,
member_confidence REAL NOT NULL DEFAULT 1.0,
is_outlier INTEGER NOT NULL DEFAULT 0,
included INTEGER NOT NULL DEFAULT 1,
flagged_coverage INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (cluster_id, immich_id)
);
CREATE TABLE IF NOT EXISTS writeback_log (
id INTEGER PRIMARY KEY AUTOINCREMENT,
immich_id TEXT NOT NULL,
action TEXT NOT NULL,
tag TEXT,
result TEXT NOT NULL,
applied_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_assets_taken_at ON assets(taken_at);
CREATE INDEX IF NOT EXISTS idx_members_cluster ON cluster_members(cluster_id);
CREATE INDEX IF NOT EXISTS idx_members_asset ON cluster_members(immich_id);
"""
class Store:
def __init__(self, db_path: str):
self.db_path = db_path
self._conn = None
def connect(self) -> "Store":
self._conn = sqlite3.connect(self.db_path)
self._conn.row_factory = sqlite3.Row
self._conn.execute("PRAGMA foreign_keys=ON")
self.migrate()
return self
@property
def conn(self) -> sqlite3.Connection:
if self._conn is None:
raise RuntimeError("Store not connected; call connect()")
return self._conn
def migrate(self) -> None:
self.conn.executescript(SCHEMA)
if self.get_meta("schema_version") is None:
self.set_meta("schema_version", str(SCHEMA_VERSION))
self.conn.commit()
def set_meta(self, key: str, value: str) -> None:
self.conn.execute(
"INSERT INTO meta(key, value) VALUES(?, ?) "
"ON CONFLICT(key) DO UPDATE SET value=excluded.value", (key, value))
self.conn.commit()
def get_meta(self, key: str, default=None):
row = self.conn.execute("SELECT value FROM meta WHERE key=?", (key,)).fetchone()
return row["value"] if row else default
def close(self) -> None:
if self._conn is not None:
self._conn.close()
self._conn = None
```
`shared/photoflow/core/__init__.py`:
```python
from photoflow.core.store import Store, SCHEMA_VERSION
__all__ = ["Store", "SCHEMA_VERSION"]
```
- [ ] **Step 5: Run to verify it passes**
Run: `.venv/bin/python -m pytest shared/tests/test_store_assets.py -q`
Expected: PASS (2 passed)
- [ ] **Step 6: Commit**
```bash
git add shared/photoflow/core/ shared/tests/test_store_assets.py
git commit -m "feat(core): domain models + Store (schema, connection, meta)"
```
---
## Task 5: core — asset & tag data-access
**Files:**
- Modify: `shared/photoflow/core/store.py`
- Test: append to `shared/tests/test_store_assets.py`
**Interfaces:**
- Consumes: `Store`, `Asset`, `Tag` (Task 4).
- Produces (methods on `Store`):
- `upsert_asset(a: Asset) -> None` — by `immich_id`; sets `has_gps` from gps presence; preserves `processed=1` if already set (a re-ingest never un-processes).
- `get_asset(immich_id: str) -> Asset | None`
- `all_assets(include_processed: bool = True) -> list[Asset]` — ordered by `taken_at`.
- `assets_in_range(start_at: str, end_at: str) -> list[Asset]` — inclusive, ordered by `taken_at`.
- `mark_processed(immich_id: str) -> None`
- `set_asset_tags(immich_id: str, tag_names: list[str]) -> None` — replaces the asset's tag rows.
- `asset_tags(immich_id: str) -> list[str]`
- `upsert_tag(name: str, immich_tag_id: str | None = None, count: int = 0) -> None`
- `all_tags() -> list[Tag]`
- [ ] **Step 1: Write the failing tests**
Append to `shared/tests/test_store_assets.py`:
```python
from photoflow.core.models import Asset
def _store(tmp_path):
return Store(str(tmp_path / "t.db")).connect()
def test_upsert_asset_roundtrip_and_has_gps(tmp_path):
s = _store(tmp_path)
s.upsert_asset(Asset(immich_id="a", taken_at="2019-06-01T10:00:00",
gps_lat=45.4, gps_lon=12.3, place_city="Venezia", type="IMAGE"))
got = s.get_asset("a")
assert got.place_city == "Venezia" and got.has_gps is True
s.upsert_asset(Asset(immich_id="b", taken_at="2019-06-02T10:00:00"))
assert s.get_asset("b").has_gps is False
s.close()
def test_upsert_preserves_processed(tmp_path):
s = _store(tmp_path)
s.upsert_asset(Asset(immich_id="a", taken_at="t"))
s.mark_processed("a")
s.upsert_asset(Asset(immich_id="a", taken_at="t2")) # re-ingest
got = s.get_asset("a")
assert got.processed is True and got.taken_at == "t2"
s.close()
def test_all_assets_ordered_and_range(tmp_path):
s = _store(tmp_path)
for i, t in [("c", "2019-06-03"), ("a", "2019-06-01"), ("b", "2019-06-02")]:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
assert [a.immich_id for a in s.all_assets()] == ["a", "b", "c"]
rng = s.assets_in_range("2019-06-02", "2019-06-03")
assert [a.immich_id for a in rng] == ["b", "c"]
s.close()
def test_asset_tags_replace(tmp_path):
s = _store(tmp_path)
s.upsert_asset(Asset(immich_id="a", taken_at="t"))
s.set_asset_tags("a", ["Italy 2019", "Kiev"])
assert sorted(s.asset_tags("a")) == ["Italy 2019", "Kiev"]
s.set_asset_tags("a", ["Italy 2019"]) # replace
assert s.asset_tags("a") == ["Italy 2019"]
s.close()
def test_tags_inventory(tmp_path):
s = _store(tmp_path)
s.upsert_tag("Italy 2019", immich_tag_id="t1", count=42)
s.upsert_tag("Italy 2019", immich_tag_id="t1", count=43) # upsert
names = {t.name: t for t in s.all_tags()}
assert names["Italy 2019"].count == 43 and names["Italy 2019"].immich_tag_id == "t1"
s.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest shared/tests/test_store_assets.py -q`
Expected: FAIL — `AttributeError: 'Store' object has no attribute 'upsert_asset'`
- [ ] **Step 3: Implement the data-access methods**
Append to `shared/photoflow/core/store.py` (inside `Store`; add `from photoflow.core.models import Asset, Tag` at top of file):
```python
def upsert_asset(self, a: Asset) -> None:
has_gps = 1 if (a.gps_lat is not None and a.gps_lon is not None) else 0
self.conn.execute(
"""INSERT INTO assets
(immich_id, taken_at, gps_lat, gps_lon, place_city, place_country,
type, has_gps, thumb_path, processed, ingested_at, updated_at)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?)
ON CONFLICT(immich_id) DO UPDATE SET
taken_at=excluded.taken_at, gps_lat=excluded.gps_lat,
gps_lon=excluded.gps_lon, place_city=excluded.place_city,
place_country=excluded.place_country, type=excluded.type,
has_gps=excluded.has_gps, thumb_path=excluded.thumb_path,
ingested_at=excluded.ingested_at, updated_at=excluded.updated_at""",
(a.immich_id, a.taken_at, a.gps_lat, a.gps_lon, a.place_city,
a.place_country, a.type, has_gps, a.thumb_path,
1 if a.processed else 0, a.ingested_at, a.updated_at))
self.conn.commit()
def _asset_from_row(self, r) -> Asset:
return Asset(
immich_id=r["immich_id"], taken_at=r["taken_at"],
gps_lat=r["gps_lat"], gps_lon=r["gps_lon"],
place_city=r["place_city"], place_country=r["place_country"],
type=r["type"], has_gps=bool(r["has_gps"]), thumb_path=r["thumb_path"],
processed=bool(r["processed"]), ingested_at=r["ingested_at"],
updated_at=r["updated_at"])
def get_asset(self, immich_id: str):
r = self.conn.execute("SELECT * FROM assets WHERE immich_id=?", (immich_id,)).fetchone()
return self._asset_from_row(r) if r else None
def all_assets(self, include_processed: bool = True) -> list:
sql = "SELECT * FROM assets"
if not include_processed:
sql += " WHERE processed=0"
sql += " ORDER BY taken_at"
return [self._asset_from_row(r) for r in self.conn.execute(sql)]
def assets_in_range(self, start_at: str, end_at: str) -> list:
return [self._asset_from_row(r) for r in self.conn.execute(
"SELECT * FROM assets WHERE taken_at>=? AND taken_at<=? ORDER BY taken_at",
(start_at, end_at))]
def mark_processed(self, immich_id: str) -> None:
self.conn.execute("UPDATE assets SET processed=1 WHERE immich_id=?", (immich_id,))
self.conn.commit()
def set_asset_tags(self, immich_id: str, tag_names: list) -> None:
self.conn.execute("DELETE FROM asset_tags WHERE immich_id=?", (immich_id,))
self.conn.executemany(
"INSERT OR IGNORE INTO asset_tags(immich_id, tag_name) VALUES(?, ?)",
[(immich_id, n) for n in tag_names])
self.conn.commit()
def asset_tags(self, immich_id: str) -> list:
return [r["tag_name"] for r in self.conn.execute(
"SELECT tag_name FROM asset_tags WHERE immich_id=? ORDER BY tag_name",
(immich_id,))]
def upsert_tag(self, name: str, immich_tag_id=None, count: int = 0) -> None:
self.conn.execute(
"""INSERT INTO tags(name, immich_tag_id, count) VALUES(?,?,?)
ON CONFLICT(name) DO UPDATE SET
immich_tag_id=excluded.immich_tag_id, count=excluded.count""",
(name, immich_tag_id, count))
self.conn.commit()
def all_tags(self) -> list:
return [Tag(name=r["name"], immich_tag_id=r["immich_tag_id"], count=r["count"])
for r in self.conn.execute("SELECT * FROM tags ORDER BY name")]
```
- [ ] **Step 4: Run to verify they pass**
Run: `.venv/bin/python -m pytest shared/tests/test_store_assets.py -q`
Expected: PASS (7 passed)
- [ ] **Step 5: Commit**
```bash
git add shared/photoflow/core/store.py shared/tests/test_store_assets.py
git commit -m "feat(core): asset & tag data-access"
```
---
## Task 6: core — cluster data-access (insert, attention sort, neighbors, split/merge, writeback log)
**Files:**
- Modify: `shared/photoflow/core/store.py`
- Test: `shared/tests/test_store_clusters.py`
**Interfaces:**
- Consumes: `Store`, `Cluster`, `ClusterMember` (Task 4/5).
- Produces (methods on `Store`):
- `insert_cluster(c: Cluster, members: list[ClusterMember]) -> int` — inserts cluster (ignoring `c.id`), assigns each member's `cluster_id`, returns new id.
- `clear_clusters() -> None` — wipe `clusters` + `cluster_members` (full re-cluster).
- `get_cluster(cluster_id: int) -> Cluster | None`
- `all_clusters() -> list[Cluster]` — ordered by `start_at`.
- `clusters_by_attention() -> list[Cluster]` — pending first, then `confidence ASC`, then longer date-span first.
- `cluster_members(cluster_id: int) -> list[tuple[Asset, ClusterMember]]` — ordered by asset `taken_at`.
- `chronological_neighbors(cluster_id: int) -> tuple[int | None, int | None]` — prev/next cluster id by `start_at` across all clusters.
- `update_cluster(cluster_id, *, status=None, decided_name=None, suggested_name=None, notes=None, reviewed_at=None) -> None`
- `set_member_inclusion(cluster_id, immich_id, included: bool) -> None`
- `split_cluster(cluster_id, boundary_immich_id) -> tuple[int, int]` — partition members so `boundary` starts the second cluster; original → `status='split'`; two new pending clusters; returns their ids.
- `merge_clusters(cluster_id_a, cluster_id_b) -> int` — union members into a new pending cluster; both originals → `status='merged'`; returns new id.
- `log_writeback(immich_id, action, tag, result) -> None`
- `already_applied(immich_id, action, tag) -> bool` — true if a prior `result='ok'` row exists for that (immich_id, action, tag).
- [ ] **Step 1: Write the failing tests**
`shared/tests/test_store_clusters.py`:
```python
from photoflow.core import Store
from photoflow.core.models import Asset, Cluster, ClusterMember
def _seed(tmp_path):
s = Store(str(tmp_path / "t.db")).connect()
for i, t in [("a", "2019-06-01"), ("b", "2019-06-02"), ("c", "2019-06-03"),
("d", "2019-07-01"), ("e", "2019-07-02")]:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
return s
def _members(ids):
return [ClusterMember(cluster_id=0, immich_id=i) for i in ids]
def test_insert_and_members_ordered(tmp_path):
s = _seed(tmp_path)
cid = s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-03",
count=3, suggested_name="Trip 1", confidence=0.9),
_members(["c", "a", "b"]))
pairs = s.cluster_members(cid)
assert [a.immich_id for a, m in pairs] == ["a", "b", "c"]
assert s.get_cluster(cid).suggested_name == "Trip 1"
s.close()
def test_attention_sort(tmp_path):
s = _seed(tmp_path)
s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-10",
confidence=0.9, status="pending"), _members(["a"]))
low = s.insert_cluster(Cluster(start_at="2019-07-01", end_at="2019-07-02",
confidence=0.2, status="pending"), _members(["d"]))
s.insert_cluster(Cluster(start_at="2019-06-02", end_at="2019-06-03",
confidence=0.1, status="approved"), _members(["b"]))
order = [c.id for c in s.clusters_by_attention()]
assert order[0] == low # lowest-confidence pending first
assert order[-1] != low # approved sinks to the bottom
s.close()
def test_chronological_neighbors(tmp_path):
s = _seed(tmp_path)
c1 = s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-03"), _members(["a"]))
c2 = s.insert_cluster(Cluster(start_at="2019-07-01", end_at="2019-07-02"), _members(["d"]))
assert s.chronological_neighbors(c1) == (None, c2)
assert s.chronological_neighbors(c2) == (c1, None)
s.close()
def test_update_and_member_inclusion(tmp_path):
s = _seed(tmp_path)
cid = s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-03"),
_members(["a", "b"]))
s.update_cluster(cid, status="approved", decided_name="Venice", reviewed_at="now")
c = s.get_cluster(cid)
assert c.status == "approved" and c.decided_name == "Venice"
s.set_member_inclusion(cid, "b", False)
inc = {m.immich_id: m.included for _, m in s.cluster_members(cid)}
assert inc == {"a": True, "b": False}
s.close()
def test_split_cluster(tmp_path):
s = _seed(tmp_path)
cid = s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-03",
suggested_name="Trip"), _members(["a", "b", "c"]))
id1, id2 = s.split_cluster(cid, "c") # boundary "c" begins the second cluster
assert s.get_cluster(cid).status == "split"
left = [a.immich_id for a, _ in s.cluster_members(id1)]
right = [a.immich_id for a, _ in s.cluster_members(id2)]
assert left == ["a", "b"] and right == ["c"]
assert s.get_cluster(id1).status == "pending"
s.close()
def test_merge_clusters(tmp_path):
s = _seed(tmp_path)
a = s.insert_cluster(Cluster(start_at="2019-06-01", end_at="2019-06-03",
suggested_name="A"), _members(["a", "b"]))
b = s.insert_cluster(Cluster(start_at="2019-07-01", end_at="2019-07-02",
suggested_name="B"), _members(["d", "e"]))
new = s.merge_clusters(a, b)
assert s.get_cluster(a).status == "merged" and s.get_cluster(b).status == "merged"
ids = [x.immich_id for x, _ in s.cluster_members(new)]
assert ids == ["a", "b", "d", "e"]
c = s.get_cluster(new)
assert c.start_at == "2019-06-01" and c.end_at == "2019-07-02" and c.status == "pending"
s.close()
def test_writeback_log_idempotency(tmp_path):
s = _seed(tmp_path)
assert s.already_applied("a", "trip", "Venice") is False
s.log_writeback("a", "trip", "Venice", "ok")
assert s.already_applied("a", "trip", "Venice") is True
s.log_writeback("b", "trip", "Venice", "error:boom")
assert s.already_applied("b", "trip", "Venice") is False # only ok counts
s.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest shared/tests/test_store_clusters.py -q`
Expected: FAIL — `AttributeError: 'Store' object has no attribute 'insert_cluster'`
- [ ] **Step 3: Implement the cluster data-access**
Append to `shared/photoflow/core/store.py` (inside `Store`; add `from photoflow.core.models import Cluster, ClusterMember` to the top import):
```python
def insert_cluster(self, c: Cluster, members: list) -> int:
cur = self.conn.execute(
"""INSERT INTO clusters
(start_at, end_at, count, suggested_name, confidence, kind_guess,
status, decided_name, reviewed_at, notes)
VALUES (?,?,?,?,?,?,?,?,?,?)""",
(c.start_at, c.end_at, c.count or len(members), c.suggested_name,
c.confidence, c.kind_guess, c.status, c.decided_name, c.reviewed_at, c.notes))
cid = cur.lastrowid
self.conn.executemany(
"""INSERT OR REPLACE INTO cluster_members
(cluster_id, immich_id, member_confidence, is_outlier, included, flagged_coverage)
VALUES (?,?,?,?,?,?)""",
[(cid, m.immich_id, m.member_confidence, 1 if m.is_outlier else 0,
1 if m.included else 0, 1 if m.flagged_coverage else 0) for m in members])
self.conn.commit()
return cid
def clear_clusters(self) -> None:
self.conn.execute("DELETE FROM cluster_members")
self.conn.execute("DELETE FROM clusters")
self.conn.commit()
def _cluster_from_row(self, r) -> Cluster:
return Cluster(
id=r["id"], start_at=r["start_at"], end_at=r["end_at"], count=r["count"],
suggested_name=r["suggested_name"], confidence=r["confidence"],
kind_guess=r["kind_guess"], status=r["status"], decided_name=r["decided_name"],
reviewed_at=r["reviewed_at"], notes=r["notes"])
def get_cluster(self, cluster_id: int):
r = self.conn.execute("SELECT * FROM clusters WHERE id=?", (cluster_id,)).fetchone()
return self._cluster_from_row(r) if r else None
def all_clusters(self) -> list:
return [self._cluster_from_row(r)
for r in self.conn.execute("SELECT * FROM clusters ORDER BY start_at, id")]
def clusters_by_attention(self) -> list:
return [self._cluster_from_row(r) for r in self.conn.execute(
"""SELECT * FROM clusters
ORDER BY CASE status WHEN 'pending' THEN 0 ELSE 1 END,
confidence ASC,
(julianday(end_at) - julianday(start_at)) DESC,
id""")]
def cluster_members(self, cluster_id: int) -> list:
rows = self.conn.execute(
"""SELECT a.*, m.member_confidence AS m_conf, m.is_outlier AS m_out,
m.included AS m_inc, m.flagged_coverage AS m_cov
FROM cluster_members m JOIN assets a ON a.immich_id = m.immich_id
WHERE m.cluster_id=? ORDER BY a.taken_at, a.immich_id""", (cluster_id,))
out = []
for r in rows:
asset = self._asset_from_row(r)
member = ClusterMember(
cluster_id=cluster_id, immich_id=r["immich_id"],
member_confidence=r["m_conf"], is_outlier=bool(r["m_out"]),
included=bool(r["m_inc"]), flagged_coverage=bool(r["m_cov"]))
out.append((asset, member))
return out
def chronological_neighbors(self, cluster_id: int):
ids = [r["id"] for r in self.conn.execute(
"SELECT id FROM clusters ORDER BY start_at, id")]
if cluster_id not in ids:
return (None, None)
i = ids.index(cluster_id)
prev_id = ids[i - 1] if i > 0 else None
next_id = ids[i + 1] if i < len(ids) - 1 else None
return (prev_id, next_id)
def update_cluster(self, cluster_id: int, *, status=None, decided_name=None,
suggested_name=None, notes=None, reviewed_at=None) -> None:
sets, vals = [], []
for col, val in [("status", status), ("decided_name", decided_name),
("suggested_name", suggested_name), ("notes", notes),
("reviewed_at", reviewed_at)]:
if val is not None:
sets.append(f"{col}=?")
vals.append(val)
if not sets:
return
vals.append(cluster_id)
self.conn.execute(f"UPDATE clusters SET {', '.join(sets)} WHERE id=?", vals)
self.conn.commit()
def set_member_inclusion(self, cluster_id: int, immich_id: str, included: bool) -> None:
self.conn.execute(
"UPDATE cluster_members SET included=? WHERE cluster_id=? AND immich_id=?",
(1 if included else 0, cluster_id, immich_id))
self.conn.commit()
def _recompute_span(self, members: list) -> tuple:
pairs = members
starts = [a.taken_at for a, _ in pairs]
return (min(starts), max(starts)) if starts else ("", "")
def split_cluster(self, cluster_id: int, boundary_immich_id: str):
pairs = self.cluster_members(cluster_id)
ids = [a.immich_id for a, _ in pairs]
if boundary_immich_id not in ids:
raise ValueError(f"boundary {boundary_immich_id!r} not in cluster")
idx = ids.index(boundary_immich_id)
base = self.get_cluster(cluster_id)
left, right = pairs[:idx], pairs[idx:]
def _new(part, suffix):
ms = [ClusterMember(cluster_id=0, immich_id=a.immich_id,
member_confidence=m.member_confidence,
is_outlier=m.is_outlier, included=m.included,
flagged_coverage=m.flagged_coverage) for a, m in part]
start, end = self._recompute_span(part)
return self.insert_cluster(Cluster(
start_at=start, end_at=end, count=len(ms),
suggested_name=f"{base.suggested_name} ({suffix})",
confidence=base.confidence, kind_guess=base.kind_guess,
status="pending"), ms)
id1, id2 = _new(left, 1), _new(right, 2)
self.update_cluster(cluster_id, status="split")
return (id1, id2)
def merge_clusters(self, cluster_id_a: int, cluster_id_b: int):
pairs = self.cluster_members(cluster_id_a) + self.cluster_members(cluster_id_b)
pairs.sort(key=lambda p: (p[0].taken_at, p[0].immich_id))
base = self.get_cluster(cluster_id_a)
ms = [ClusterMember(cluster_id=0, immich_id=a.immich_id,
member_confidence=m.member_confidence, is_outlier=m.is_outlier,
included=m.included, flagged_coverage=m.flagged_coverage)
for a, m in pairs]
start, end = self._recompute_span(pairs)
new = self.insert_cluster(Cluster(
start_at=start, end_at=end, count=len(ms),
suggested_name=base.suggested_name, confidence=base.confidence,
kind_guess=base.kind_guess, status="pending"), ms)
self.update_cluster(cluster_id_a, status="merged")
self.update_cluster(cluster_id_b, status="merged")
return new
def log_writeback(self, immich_id: str, action: str, tag, result: str) -> None:
import datetime
self.conn.execute(
"INSERT INTO writeback_log(immich_id, action, tag, result, applied_at) "
"VALUES (?,?,?,?,?)",
(immich_id, action, tag, result,
datetime.datetime.now(datetime.timezone.utc).isoformat()))
self.conn.commit()
def already_applied(self, immich_id: str, action: str, tag) -> bool:
row = self.conn.execute(
"SELECT 1 FROM writeback_log WHERE immich_id=? AND action=? "
"AND IFNULL(tag,'')=IFNULL(?, '') AND result='ok' LIMIT 1",
(immich_id, action, tag)).fetchone()
return row is not None
```
- [ ] **Step 4: Run to verify they pass**
Run: `.venv/bin/python -m pytest shared/tests/test_store_clusters.py -q`
Expected: PASS (7 passed)
- [ ] **Step 5: Commit**
```bash
git add shared/photoflow/core/store.py shared/tests/test_store_clusters.py
git commit -m "feat(core): cluster data-access — attention sort, neighbors, split/merge, writeback log"
```
---
## Task 7: shared/ui — base.html, macros, grid-lightbox JS, Flask wiring
**Files:**
- Create: `shared/photoflow/ui/__init__.py`, `shared/photoflow/ui/templates/base.html`, `shared/photoflow/ui/templates/macros.html`, `shared/photoflow/ui/static/shared.js`
- Test: `shared/tests/test_ui.py`
**Interfaces:**
- Consumes: nothing (Flask only).
- Produces:
- `photoflow.ui.TEMPLATE_DIR: str`, `photoflow.ui.STATIC_DIR: str`
- `photoflow.ui.register_shared_ui(app: Flask) -> None` — adds `TEMPLATE_DIR` to the app's Jinja search path (via `ChoiceLoader`) and serves shared static at `/shared-static/<file>` (so `base.html` can `<script src="/shared-static/shared.js">`).
- `base.html` — DaisyUI 4 + Tailwind + Alpine + HTMX (CDN); blocks: `title`, `navbar_title`, `content`, `extra_scripts`. Loads `/shared-static/shared.js`.
- `macros.html``lightbox()`, `status_badge(status)`, `confidence_badge(confidence)`. (`thumb_url_prefix` like `/thumb/`.)
- `shared.js` — global `photoGrid()` Alpine factory: `focused`, `select(el)`, `cards()`, `navigate(dir)`, `openLightbox(el)`, `closeLightbox()`, `updateLightbox()`. Apps compose it.
- [ ] **Step 1: Write the failing test**
`shared/tests/test_ui.py`:
```python
import os
from flask import Flask, render_template
from photoflow.ui import TEMPLATE_DIR, STATIC_DIR, register_shared_ui
def test_dirs_exist():
assert os.path.isfile(os.path.join(TEMPLATE_DIR, "base.html"))
assert os.path.isfile(os.path.join(STATIC_DIR, "shared.js"))
def test_register_serves_shared_static_and_template():
app = Flask(__name__)
register_shared_ui(app)
@app.route("/page")
def page():
return render_template("base.html")
client = app.test_client()
r = client.get("/page")
assert r.status_code == 200
assert b"/shared-static/shared.js" in r.data
js = client.get("/shared-static/shared.js")
assert js.status_code == 200 and b"photoGrid" in js.data
```
- [ ] **Step 2: Run to verify it fails**
Run: `.venv/bin/python -m pytest shared/tests/test_ui.py -q`
Expected: FAIL — `ImportError: cannot import name 'TEMPLATE_DIR'`
- [ ] **Step 3: Implement the wiring**
`shared/photoflow/ui/__init__.py`:
```python
import os
from flask import Blueprint, Flask
from jinja2 import ChoiceLoader, FileSystemLoader
TEMPLATE_DIR = os.path.join(os.path.dirname(__file__), "templates")
STATIC_DIR = os.path.join(os.path.dirname(__file__), "static")
def register_shared_ui(app: Flask) -> None:
app.jinja_loader = ChoiceLoader([app.jinja_loader, FileSystemLoader(TEMPLATE_DIR)])
bp = Blueprint("shared_ui", __name__, static_folder=STATIC_DIR,
static_url_path="/shared-static")
app.register_blueprint(bp)
__all__ = ["TEMPLATE_DIR", "STATIC_DIR", "register_shared_ui"]
```
- [ ] **Step 4: Implement base.html**
`shared/photoflow/ui/templates/base.html`:
```html
<!DOCTYPE html>
<html data-theme="forest" lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>{% block title %}photoflow{% endblock %}</title>
<link href="https://cdn.jsdelivr.net/npm/daisyui@4/dist/full.min.css" rel="stylesheet">
<script src="https://cdn.tailwindcss.com"></script>
<script src="https://unpkg.com/htmx.org@1.9.12"></script>
<script defer src="https://cdn.jsdelivr.net/npm/alpinejs@3/dist/cdn.min.js"></script>
</head>
<body class="min-h-screen bg-base-200">
<div class="navbar bg-base-100 shadow-sm sticky top-0 z-40">
<div class="navbar-start px-4 font-bold text-lg">
<a href="/">{% block navbar_title %}photoflow{% endblock %}</a>
</div>
</div>
<div class="p-4">
{% block content %}{% endblock %}
</div>
<script src="/shared-static/shared.js"></script>
{% block extra_scripts %}{% endblock %}
</body>
</html>
```
- [ ] **Step 5: Implement macros.html**
`shared/photoflow/ui/templates/macros.html`:
```html
{% macro status_badge(status) %}
{% set cls = {'pending':'badge-ghost','approved':'badge-success','non_trip':'badge-neutral',
'skipped':'badge-warning','merged':'badge-info','split':'badge-info'} %}
<span class="badge badge-sm {{ cls.get(status, 'badge-ghost') }}">{{ status }}</span>
{% endmacro %}
{% macro confidence_badge(confidence) %}
{% if confidence < 0.4 %}
<span class="badge badge-sm badge-warning" title="needs your eye">low</span>
{% elif confidence < 0.75 %}
<span class="badge badge-sm badge-info">med</span>
{% else %}
<span class="badge badge-sm badge-success">high</span>
{% endif %}
{% endmacro %}
{% macro lightbox() %}
<div id="lb" class="fixed inset-0 z-50 bg-black/95 flex items-center justify-center" style="display:none">
<button class="absolute top-4 right-4 btn btn-circle btn-sm btn-ghost text-white"
@click="closeLightbox()">&#10005;</button>
<button class="absolute left-3 top-1/2 -translate-y-1/2 btn btn-circle btn-ghost text-white text-4xl"
@click="navigate(-1)">&#8249;</button>
<button class="absolute right-3 top-1/2 -translate-y-1/2 btn btn-circle btn-ghost text-white text-4xl"
@click="navigate(1)">&#8250;</button>
<div class="flex flex-col items-center gap-3 px-16 max-w-full">
<img id="lb-img" src="" class="max-h-[80vh] max-w-[88vw] object-contain rounded-lg" alt="">
<div class="text-white/40 text-xs">&larr; &rarr; navigate · Esc close</div>
</div>
</div>
{% endmacro %}
```
- [ ] **Step 6: Implement shared.js**
`shared/photoflow/ui/static/shared.js`:
```javascript
// Generic grid + lightbox behavior shared across photoflow apps.
// Apps spread this into their own Alpine component: { ...photoGrid(), ...appLogic }
function photoGrid() {
return {
focused: null,
lightboxOpen: false,
cards() {
return [...document.querySelectorAll('.photo-card')]
.filter(c => c.style.display !== 'none');
},
select(el) {
if (this.focused) this.focused.classList.remove('ring-4', 'ring-white', 'z-10');
this.focused = el;
if (el) {
el.classList.add('ring-4', 'ring-white', 'z-10');
el.scrollIntoView({ behavior: 'smooth', block: 'nearest' });
if (this.lightboxOpen) this.updateLightbox();
}
},
selectFirst() {
const cards = this.cards();
if (cards.length) this.select(cards[0]);
},
navigate(dir) {
const cards = this.cards();
if (!cards.length) return;
const idx = this.focused ? cards.indexOf(this.focused) : -1;
const next = cards[Math.max(0, Math.min(cards.length - 1, idx + dir))];
if (next) this.select(next);
},
openLightbox(el) {
this.select(el);
this.lightboxOpen = true;
document.getElementById('lb').style.display = '';
this.updateLightbox();
},
closeLightbox() {
this.lightboxOpen = false;
const lb = document.getElementById('lb');
if (lb) lb.style.display = 'none';
},
updateLightbox() {
const el = this.focused;
if (!el) return;
const img = el.querySelector('img');
const lbImg = document.getElementById('lb-img');
if (img && lbImg) lbImg.src = img.src;
},
};
}
```
- [ ] **Step 7: Run to verify it passes**
Run: `.venv/bin/python -m pytest shared/tests/test_ui.py -q`
Expected: PASS (2 passed)
- [ ] **Step 8: Commit**
```bash
git add shared/photoflow/ui/ shared/tests/test_ui.py
git commit -m "feat(ui): base.html + macros + shared grid-lightbox JS + Flask wiring"
```
> **Review revision (doc-review 2026-06-27):** Drop the `image_grid` macro — it has no consumer in M1 (`_detail.html` inlines its own grid). Restore it when a template actually calls it.
---
## Task 8: trip-cluster — config, Flask factory, CLI skeleton, health + proxy
**Files:**
- Create: `apps/trip-cluster/app/config.py`, `apps/trip-cluster/app/__init__.py` (replace empty), `apps/trip-cluster/app/cli.py`, `apps/trip-cluster/app/routes/__init__.py`, `apps/trip-cluster/app/routes/nav.py`, `apps/trip-cluster/app/routes/proxy.py`
- Test: `apps/trip-cluster/tests/test_config.py`, `apps/trip-cluster/tests/test_routes.py` (health + proxy portion)
**Interfaces:**
- Consumes: `photoflow.ui.register_shared_ui`, `photoflow.core.Store`.
- Produces:
- `app.config.Config` dataclass: `immich_url`, `immich_api_key`, `anthropic_api_key` (may be ""), `data_dir`; properties `db_path` (`<data_dir>/trip-cluster.db`), `thumbs_dir` (`<data_dir>/thumbs`). `load_config(env=None) -> Config`; requires only `IMMICH_URL`, `IMMICH_API_KEY`; raises `ConfigError(missing: list[str])`.
- `app.create_app(config=None) -> Flask` — registers shared UI + `nav`, `review`, `proxy` blueprints; stores `app.config["APP_CONFIG"]` and `app.config["DATA_DIR"]`.
- `app.cli.main(argv=None) -> int` — argparse subcommands `ingest|cluster|serve|apply` (only `serve` implemented here; others added in later tasks but registered now as stubs returning 0 with a "not yet" message removed as they're built).
- Routes: `GET /health -> "ok"`; `GET /thumb/<asset_id>` serves `<thumbs_dir>/<asset_id>.jpg`.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_config.py`:
```python
import os
import pytest
from app.config import load_config, ConfigError
def test_missing_required_raises():
with pytest.raises(ConfigError) as e:
load_config({"IMMICH_URL": "http://x"})
assert "IMMICH_API_KEY" in e.value.missing
def test_anthropic_optional_and_paths(tmp_path):
cfg = load_config({"IMMICH_URL": "http://x/", "IMMICH_API_KEY": "k",
"DATA_DIR": str(tmp_path)})
assert cfg.immich_url == "http://x" # trailing slash stripped
assert cfg.anthropic_api_key == "" # optional in M1
assert cfg.db_path == os.path.join(str(tmp_path), "trip-cluster.db")
assert cfg.thumbs_dir == os.path.join(str(tmp_path), "thumbs")
```
`apps/trip-cluster/tests/test_routes.py`:
```python
import os
from app import create_app
from app.config import Config
def _app(tmp_path):
cfg = Config(immich_url="http://x", immich_api_key="k",
anthropic_api_key="", data_dir=str(tmp_path))
app = create_app(cfg)
app.config.update(TESTING=True)
return app
def test_health(tmp_path):
assert _app(tmp_path).test_client().get("/health").data == b"ok"
def test_thumb_served(tmp_path):
thumbs = os.path.join(str(tmp_path), "thumbs")
os.makedirs(thumbs, exist_ok=True)
with open(os.path.join(thumbs, "a.jpg"), "wb") as f:
f.write(b"\xff\xd8\xffjpeg")
client = _app(tmp_path).test_client()
r = client.get("/thumb/a")
assert r.status_code == 200 and r.mimetype == "image/jpeg"
assert _app(tmp_path).test_client().get("/thumb/missing").status_code == 404
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_config.py apps/trip-cluster/tests/test_routes.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.config'`
- [ ] **Step 3: Implement config**
`apps/trip-cluster/app/config.py`:
```python
import os
from dataclasses import dataclass
from typing import Mapping, Optional
REQUIRED = ["IMMICH_URL", "IMMICH_API_KEY"]
class ConfigError(Exception):
def __init__(self, missing: list):
self.missing = missing
super().__init__(f"Missing required environment variables: {', '.join(missing)}")
@dataclass
class Config:
immich_url: str
immich_api_key: str
anthropic_api_key: str
data_dir: str
@property
def db_path(self) -> str:
return os.path.join(self.data_dir, "trip-cluster.db")
@property
def thumbs_dir(self) -> str:
return os.path.join(self.data_dir, "thumbs")
def load_config(env: Optional[Mapping] = None) -> Config:
env = env if env is not None else os.environ
missing = [k for k in REQUIRED if not (env.get(k) or "").strip()]
if missing:
raise ConfigError(missing)
data_dir = (env.get("DATA_DIR") or "").strip() or os.path.join(os.getcwd(), "data")
return Config(
immich_url=env["IMMICH_URL"].strip().rstrip("/"),
immich_api_key=env["IMMICH_API_KEY"].strip(),
anthropic_api_key=(env.get("ANTHROPIC_API_KEY") or "").strip(),
data_dir=data_dir,
)
```
- [ ] **Step 4: Implement the factory + routes**
`apps/trip-cluster/app/__init__.py`:
```python
from flask import Flask
from app.config import load_config
from photoflow.ui import register_shared_ui
def create_app(config=None) -> Flask:
app = Flask(__name__)
cfg = config or load_config()
app.config["APP_CONFIG"] = cfg
app.config["DATA_DIR"] = cfg.data_dir
register_shared_ui(app)
from app.routes.nav import bp as nav_bp
from app.routes.review import bp as review_bp
from app.routes.proxy import bp as proxy_bp
app.register_blueprint(nav_bp)
app.register_blueprint(review_bp)
app.register_blueprint(proxy_bp)
return app
```
`apps/trip-cluster/app/routes/__init__.py` → empty file.
`apps/trip-cluster/app/routes/proxy.py`:
```python
import os
from flask import Blueprint, current_app, send_file, abort
bp = Blueprint("proxy", __name__)
@bp.route("/thumb/<asset_id>")
def thumb(asset_id):
cfg = current_app.config["APP_CONFIG"]
safe = os.path.basename(asset_id)
path = os.path.join(cfg.thumbs_dir, f"{safe}.jpg")
if not os.path.exists(path):
abort(404)
return send_file(path, mimetype="image/jpeg")
```
`apps/trip-cluster/app/routes/nav.py` (minimal for now; full master/detail in Task 15):
```python
from flask import Blueprint
bp = Blueprint("nav", __name__)
@bp.route("/health")
def health():
return "ok"
@bp.route("/")
def index():
return "trip-cluster"
```
> Note: `create_app` imports `app.routes.review`. Create a stub now so the factory imports cleanly; Task 13 fills it in.
`apps/trip-cluster/app/routes/review.py` (stub):
```python
from flask import Blueprint
bp = Blueprint("review", __name__)
```
- [ ] **Step 5: Implement the CLI skeleton**
`apps/trip-cluster/app/cli.py`:
```python
import argparse
import sys
from app.config import load_config
def _store(cfg):
import os
from photoflow.core import Store
os.makedirs(cfg.data_dir, exist_ok=True)
return Store(cfg.db_path).connect()
def _immich(deps, cfg):
if "immich" in deps:
return deps["immich"]
from photoflow.immich import ImmichClient
return ImmichClient(cfg.immich_url, cfg.immich_api_key)
def cmd_serve(deps) -> int:
from app import create_app
create_app(deps["config"]).run(host="0.0.0.0", port=8084)
return 0
def build_parser() -> argparse.ArgumentParser:
p = argparse.ArgumentParser(prog="categorize")
sub = p.add_subparsers(dest="command", required=True)
ing = sub.add_parser("ingest")
ing.add_argument("--from", dest="date_from")
ing.add_argument("--to", dest="date_to")
ing.add_argument("--tag")
ing.add_argument("--subset", type=int)
ing.add_argument("--full", action="store_true", help="ignore incremental updatedAfter")
cl = sub.add_parser("cluster")
cl.add_argument("--gap-factor", type=float, default=6.0)
sub.add_parser("serve")
ap = sub.add_parser("apply")
ap.add_argument("--yes", action="store_true", help="skip the confirmation prompt")
return p
def main(argv=None) -> int:
args = build_parser().parse_args(argv)
deps = {"config": load_config()}
if args.command == "serve":
return cmd_serve(deps)
# ingest / cluster / apply are wired in later tasks.
print(f"Command '{args.command}' is not implemented yet.")
return 1
if __name__ == "__main__":
sys.exit(main())
```
- [ ] **Step 6: Run to verify all pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_config.py apps/trip-cluster/tests/test_routes.py -q`
Expected: PASS (4 passed)
- [ ] **Step 7: Commit**
```bash
git add apps/trip-cluster/app apps/trip-cluster/tests/test_config.py apps/trip-cluster/tests/test_routes.py
git commit -m "feat(trip-cluster): config, factory, CLI skeleton, health + thumb proxy"
```
---
## Task 9: trip-cluster — ingest (scopeable, incremental, processed read-back)
**Files:**
- Create: `apps/trip-cluster/app/ingest.py`
- Modify: `apps/trip-cluster/app/cli.py`
- Test: `apps/trip-cluster/tests/test_ingest.py`
**Interfaces:**
- Consumes: `ImmichClient` (read methods), `Store` (asset/tag/meta methods), `photoflow.immich.pipeline.PROCESSED`.
- Produces:
- `app.ingest.run_ingest(client, store, thumbs_dir, *, date_from=None, date_to=None, tag=None, subset=None, full=False) -> dict` returning `{"fetched": int, "processed_marked": int}`. Resolves `--tag` to a tag id; uses `store.get_meta("last_ingest_at")` as `updated_after` unless `full`; downloads each thumbnail to `<thumbs_dir>/<id>.jpg` (skips if present); upserts assets + asset_tags + tag inventory; marks `processed` for assets carrying `pipeline.PROCESSED`; advances `last_ingest_at` to the max `updated_at` seen.
- CLI `ingest` subcommand wired to it.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_ingest.py`:
```python
import os
from photoflow.core import Store
from photoflow.immich import pipeline
from app.ingest import run_ingest
class FakeImmich:
def __init__(self, assets, tag_map=None, thumb=b"\xff\xd8\xffjpeg"):
self._assets = assets
self._tag_map = tag_map or {}
self._thumb = thumb
self.searches = []
def resolve_tag_id(self, name):
return self._tag_map.get(name)
def search_assets(self, **kwargs):
self.searches.append(kwargs)
return list(self._assets)
def download_thumbnail(self, asset_id):
return self._thumb
def _asset(i, taken, tags=None, updated="2026-01-01T00:00:00Z"):
return {"id": i, "original_filename": f"{i}.jpg", "taken_at": taken,
"gps_lat": None, "gps_lon": None, "place_city": None,
"place_country": None, "type": "IMAGE", "tags": tags or [],
"rating": 0, "updated_at": updated}
def test_ingest_upserts_assets_tags_and_thumbs(tmp_path):
store = Store(str(tmp_path / "t.db")).connect()
thumbs = str(tmp_path / "thumbs")
os.makedirs(thumbs, exist_ok=True)
client = FakeImmich([_asset("a", "2019-06-01", tags=["Italy 2019"]),
_asset("b", "2019-06-02")])
res = run_ingest(client, store, thumbs, date_from="2019-01-01", date_to="2020-01-01")
assert res["fetched"] == 2
assert store.get_asset("a").taken_at == "2019-06-01"
assert store.asset_tags("a") == ["Italy 2019"]
assert os.path.exists(os.path.join(thumbs, "a.jpg"))
assert store.get_meta("last_ingest_at") == "2026-01-01T00:00:00Z"
store.close()
def test_ingest_marks_processed_from_pipeline_tag(tmp_path):
store = Store(str(tmp_path / "t.db")).connect()
thumbs = str(tmp_path / "thumbs"); os.makedirs(thumbs, exist_ok=True)
client = FakeImmich([_asset("a", "2019-06-01", tags=[pipeline.PROCESSED])])
res = run_ingest(client, store, thumbs)
assert res["processed_marked"] == 1
assert store.get_asset("a").processed is True
store.close()
def test_ingest_incremental_passes_updated_after(tmp_path):
store = Store(str(tmp_path / "t.db")).connect()
thumbs = str(tmp_path / "thumbs"); os.makedirs(thumbs, exist_ok=True)
store.set_meta("last_ingest_at", "2026-05-01T00:00:00Z")
client = FakeImmich([_asset("a", "2019-06-01")])
run_ingest(client, store, thumbs)
assert client.searches[0].get("updated_after") == "2026-05-01T00:00:00Z"
# --full ignores it
client2 = FakeImmich([_asset("a", "2019-06-01")])
run_ingest(client2, store, thumbs, full=True)
assert client2.searches[0].get("updated_after") is None
store.close()
def test_ingest_subset_and_tag_resolution(tmp_path):
store = Store(str(tmp_path / "t.db")).connect()
thumbs = str(tmp_path / "thumbs"); os.makedirs(thumbs, exist_ok=True)
client = FakeImmich([_asset("a", "2019-06-01"), _asset("b", "2019-06-02")],
tag_map={"Italy 2019": "t1"})
res = run_ingest(client, store, thumbs, tag="Italy 2019", subset=1)
assert res["fetched"] == 1
assert client.searches[0].get("tag_ids") == ["t1"]
store.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_ingest.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.ingest'`
- [ ] **Step 3: Implement ingest**
`apps/trip-cluster/app/ingest.py`:
```python
import datetime
import os
from photoflow.core.models import Asset
from photoflow.immich import pipeline
def _now() -> str:
return datetime.datetime.now(datetime.timezone.utc).isoformat()
def run_ingest(client, store, thumbs_dir, *, date_from=None, date_to=None,
tag=None, subset=None, full=False) -> dict:
os.makedirs(thumbs_dir, exist_ok=True)
tag_ids = None
if tag:
tid = client.resolve_tag_id(tag)
if not tid:
raise ValueError(f"Tag {tag!r} not found in Immich")
tag_ids = [tid]
updated_after = None if full else store.get_meta("last_ingest_at")
assets = client.search_assets(taken_after=date_from, taken_before=date_to,
tag_ids=tag_ids, updated_after=updated_after)
if subset is not None:
assets = assets[:subset]
processed_marked = 0
max_updated = updated_after or ""
tag_counts: dict = {}
for a in assets:
thumb_path = os.path.join(thumbs_dir, f"{a['id']}.jpg")
if not os.path.exists(thumb_path):
with open(thumb_path, "wb") as f:
f.write(client.download_thumbnail(a["id"]))
store.upsert_asset(Asset(
immich_id=a["id"], taken_at=a["taken_at"], gps_lat=a["gps_lat"],
gps_lon=a["gps_lon"], place_city=a["place_city"],
place_country=a["place_country"], type=a["type"],
thumb_path=thumb_path, ingested_at=_now(), updated_at=a["updated_at"]))
store.set_asset_tags(a["id"], a["tags"])
for t in a["tags"]:
tag_counts[t] = tag_counts.get(t, 0) + 1
if pipeline.PROCESSED in a["tags"]:
store.mark_processed(a["id"])
processed_marked += 1
if a["updated_at"] and a["updated_at"] > max_updated:
max_updated = a["updated_at"]
for name, count in tag_counts.items():
store.upsert_tag(name, count=count)
if max_updated:
store.set_meta("last_ingest_at", max_updated)
return {"fetched": len(assets), "processed_marked": processed_marked}
```
- [ ] **Step 4: Wire the CLI `ingest` subcommand**
In `apps/trip-cluster/app/cli.py`, add a handler and dispatch. Add this function:
```python
def cmd_ingest(deps, *, date_from, date_to, tag, subset, full) -> int:
cfg = deps["config"]
store = _store(cfg)
client = _immich(deps, cfg)
from app.ingest import run_ingest
res = run_ingest(client, store, cfg.thumbs_dir, date_from=date_from,
date_to=date_to, tag=tag, subset=subset, full=full)
store.close()
print(f"Ingested {res['fetched']} asset(s); "
f"{res['processed_marked']} already-processed.")
return 0
```
In `main()`, replace the fall-through for `ingest`:
```python
if args.command == "ingest":
return cmd_ingest(deps, date_from=args.date_from, date_to=args.date_to,
tag=args.tag, subset=args.subset, full=args.full)
```
- [ ] **Step 5: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_ingest.py -q`
Expected: PASS (4 passed)
- [ ] **Step 6: Commit**
```bash
git add apps/trip-cluster/app/ingest.py apps/trip-cluster/app/cli.py apps/trip-cluster/tests/test_ingest.py
git commit -m "feat(trip-cluster): ingest — scopeable, incremental, processed read-back"
```
> **Review revisions (doc-review 2026-06-27):** (1) Make `run_ingest` resilient to thumbnail-download failures — wrap each `download_thumbnail` in try/except, log and `continue` instead of aborting the run; write to a temp path and `os.replace` into place so a failure never leaves a 0-byte `<id>.jpg`, and treat an existing 0-byte file as missing in the skip check. (2) Sanitize the Immich-provided id before using it as a path component: `safe_id = os.path.basename(a["id"])` (mirror the `/thumb` proxy), guarding against path traversal from a crafted or MITM'd response.
---
## Task 10: trip-cluster — clustering algorithm (pure, density-adaptive)
**Files:**
- Create: `apps/trip-cluster/app/clustering.py`
- Test: `apps/trip-cluster/tests/test_clustering.py`
**Interfaces:**
- Consumes: `photoflow.core.models.Asset`.
- Produces:
- `app.clustering.CandidateCluster` dataclass: `member_ids: list[str]`, `start_at: str`, `end_at: str`, `suggested_name: str`, `confidence: float`, `kind_guess: str`, `seed_tag: Optional[str]`.
- `app.clustering.cluster_assets(assets: list[Asset], tags_by_asset: dict[str, list[str]], seed_tags: set[str], *, gap_factor: float = 6.0, hard_split_days: int = 14, min_floor_seconds: int = 3600) -> list[CandidateCluster]`. Existing-tag seeds first (confidence 0.95, not gap-split); remaining assets gap-clustered density-adaptively (rolling-median baseline with a 2-gap bootstrap + a hard-cap safety net); free clusters named from location anchors and scored by GPS fraction + size.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_clustering.py`:
```python
from photoflow.core.models import Asset
from app.clustering import cluster_assets
def _a(i, taken, gps=False, city=None):
return Asset(immich_id=i, taken_at=taken,
gps_lat=45.0 if gps else None, gps_lon=12.0 if gps else None,
place_city=city)
def test_seed_tag_forms_one_cluster_not_gap_split():
# Two assets months apart but sharing a trip tag -> ONE seeded cluster.
assets = [_a("a", "2019-06-01T10:00:00"), _a("b", "2019-09-01T10:00:00")]
tags = {"a": ["Italy 2019"], "b": ["Italy 2019"]}
clusters = cluster_assets(assets, tags, {"Italy 2019"})
assert len(clusters) == 1
c = clusters[0]
assert c.seed_tag == "Italy 2019" and sorted(c.member_ids) == ["a", "b"]
assert c.confidence >= 0.9 and c.suggested_name == "Italy 2019"
def test_sparse_old_regime_splits_on_adaptive_threshold():
# ~1 day intra-trip gaps; trips separated by 10 days (< 14d hard cap),
# so only the adaptive rule can split them.
a = [_a(f"a{i}", f"2008-06-0{i+1}T12:00:00") for i in range(5)] # Jun 1..5
b = [_a(f"b{i}", f"2008-06-1{i+5}T12:00:00") for i in range(3)] # Jun 15..17
assets = a + b
tags = {x.immich_id: [] for x in assets}
clusters = cluster_assets(assets, tags, set())
assert len(clusters) == 2
assert sorted(clusters[0].member_ids) == ["a0", "a1", "a2", "a3", "a4"]
def test_dense_recent_regime_splits_on_adaptive_threshold():
# Hourly bursts within a day; 2-day gap between days.
day1 = [_a(f"d{i}", f"2024-03-10T{10+i:02d}:00:00") for i in range(4)]
day3 = [_a(f"e{i}", f"2024-03-12T{10+i:02d}:00:00") for i in range(4)]
assets = day1 + day3
tags = {x.immich_id: [] for x in assets}
clusters = cluster_assets(assets, tags, set())
assert len(clusters) == 2
assert sorted(clusters[0].member_ids) == ["d0", "d1", "d2", "d3"]
def test_location_anchor_names_and_gps_confidence():
assets = [_a("a", "2020-05-01T10:00:00", gps=True, city="Kiev"),
_a("b", "2020-05-01T12:00:00", gps=True, city="Kiev"),
_a("c", "2020-05-01T14:00:00", gps=True, city="Kiev"),
_a("d", "2020-05-01T16:00:00", gps=True, city="Kiev"),
_a("e", "2020-05-01T18:00:00", gps=True, city="Kiev")]
tags = {x.immich_id: [] for x in assets}
c = cluster_assets(assets, tags, set())[0]
assert c.suggested_name == "Kiev"
assert c.confidence > 0.6 # full GPS lifts confidence
assert c.kind_guess == "trip"
def test_small_scattered_cluster_marked_everyday():
assets = [_a("a", "2015-01-01T10:00:00"), _a("b", "2015-01-01T11:00:00")]
tags = {"a": [], "b": []}
c = cluster_assets(assets, tags, set())[0]
assert c.kind_guess == "everyday"
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_clustering.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.clustering'`
- [ ] **Step 3: Implement the algorithm**
`apps/trip-cluster/app/clustering.py`:
```python
import datetime
from collections import Counter
from dataclasses import dataclass
from typing import Optional
EVERYDAY_MAX_COUNT = 4
SEED_CONFIDENCE = 0.95
@dataclass
class CandidateCluster:
member_ids: list
start_at: str
end_at: str
suggested_name: str
confidence: float
kind_guess: str
seed_tag: Optional[str] = None
def _epoch(taken_at: str) -> float:
s = (taken_at or "").strip()
if not s:
return 0.0
s = s.replace("Z", "")
if "." in s:
s = s.split(".", 1)[0]
try:
if "T" in s:
return datetime.datetime.fromisoformat(s).timestamp()
return datetime.datetime.fromisoformat(s + "T00:00:00").timestamp()
except ValueError:
return 0.0
def _median(values: list) -> float:
if not values:
return 0.0
xs = sorted(values)
n = len(xs)
mid = n // 2
return xs[mid] if n % 2 else (xs[mid - 1] + xs[mid]) / 2
def _span(members: list) -> tuple:
ts = [m.taken_at for m in members]
return (min(ts), max(ts)) if ts else ("", "")
def _name(members: list, start_at: str) -> str:
cities = Counter(m.place_city for m in members if m.place_city)
if cities:
return cities.most_common(1)[0][0]
countries = Counter(m.place_country for m in members if m.place_country)
if countries:
return countries.most_common(1)[0][0]
return f"Trip {start_at[:10]}"
def _confidence(members: list) -> float:
count = len(members)
gps_frac = sum(1 for m in members if m.gps_lat is not None) / count if count else 0
conf = 0.30 + 0.40 * gps_frac + 0.20 * min(count / 20, 1)
return round(min(conf, 0.85), 2)
def _free_cluster(members: list) -> CandidateCluster:
start, end = _span(members)
return CandidateCluster(
member_ids=[m.immich_id for m in members],
start_at=start, end_at=end,
suggested_name=_name(members, start),
confidence=_confidence(members),
kind_guess="everyday" if len(members) <= EVERYDAY_MAX_COUNT else "trip")
def _gap_cluster(assets: list, *, gap_factor, hard_split_days, min_floor_seconds) -> list:
ordered = sorted(assets, key=lambda a: a.taken_at)
if not ordered:
return []
hard_cap = hard_split_days * 86400
groups = []
group = [ordered[0]]
group_gaps: list = []
for prev, cur in zip(ordered, ordered[1:]):
gap = _epoch(cur.taken_at) - _epoch(prev.taken_at)
if gap > hard_cap:
split = True
elif len(group_gaps) < 2: # bootstrap: accept first 2 gaps
split = False
else:
threshold = max(gap_factor * _median(group_gaps), min_floor_seconds)
split = gap > threshold
if split:
groups.append(group)
group = [cur]
group_gaps = []
else:
group.append(cur)
group_gaps.append(gap)
groups.append(group)
return [_free_cluster(g) for g in groups]
def cluster_assets(assets, tags_by_asset, seed_tags, *, gap_factor=6.0,
hard_split_days=14, min_floor_seconds=3600) -> list:
by_id = {a.immich_id: a for a in assets}
used = set()
clusters = []
# 1. Seed clusters from existing trip tags (authoritative; never gap-split).
for tag in sorted(seed_tags):
members = [by_id[aid] for aid in by_id
if aid not in used and tag in tags_by_asset.get(aid, [])]
if not members:
continue
members.sort(key=lambda a: a.taken_at)
used.update(m.immich_id for m in members)
start, end = _span(members)
clusters.append(CandidateCluster(
member_ids=[m.immich_id for m in members], start_at=start, end_at=end,
suggested_name=tag, confidence=SEED_CONFIDENCE, kind_guess="trip",
seed_tag=tag))
# 2. Gap-cluster the remaining (free) assets.
free = [a for a in assets if a.immich_id not in used]
clusters.extend(_gap_cluster(free, gap_factor=gap_factor,
hard_split_days=hard_split_days,
min_floor_seconds=min_floor_seconds))
clusters.sort(key=lambda c: c.start_at)
return clusters
```
- [ ] **Step 4: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_clustering.py -q`
Expected: PASS (5 passed)
- [ ] **Step 5: Commit**
```bash
git add apps/trip-cluster/app/clustering.py apps/trip-cluster/tests/test_clustering.py
git commit -m "feat(trip-cluster): density-adaptive timestamp clustering + tag seeds + anchors"
```
> **Review revisions (doc-review 2026-06-27):** (1) Document the rationale for the `_confidence` constants (`0.30`/`0.40`/`0.20`, cap `0.85`) and add a temporal-tightness term so a GPS-poor but densely-packed cluster can still earn high confidence — as written, `gps_frac=0` caps confidence at `0.50`, below the `0.75` bulk-approve gate, so "approve high-confidence" never fires on GPS-poor data (see Task 13). (2) Expose `_epoch`/`_median` through a public interface (drop the underscore, or extract to e.g. `app/timeutil.py`) since `coverage.py` consumes them across the module boundary.
---
## Task 11: trip-cluster — coverage detection (completeness flags + outliers)
**Files:**
- Create: `apps/trip-cluster/app/coverage.py`
- Test: `apps/trip-cluster/tests/test_coverage.py`
**Interfaces:**
- Consumes: `app.clustering.CandidateCluster`, `app.clustering._epoch`, `app.clustering._median`, `photoflow.core.models.Asset`, `photoflow.core.models.ClusterMember`.
- Produces:
- `app.coverage.coverage_members(candidate: CandidateCluster, assets_by_id: dict[str, Asset], *, outlier_factor: float = 8.0) -> list[ClusterMember]` — returns the cluster's real members (with `is_outlier` set for tagged assets sitting far from a **seeded** cluster's bulk) plus, for **seeded** clusters only, coverage candidates: assets whose `taken_at` falls inside `[start_at, end_at]` but which are not members, returned as `ClusterMember(included=False, flagged_coverage=True)`.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_coverage.py`:
```python
from photoflow.core.models import Asset
from app.clustering import CandidateCluster
from app.coverage import coverage_members
def _a(i, taken):
return Asset(immich_id=i, taken_at=taken)
def test_coverage_candidate_inside_seeded_window():
members = [_a("a", "2019-06-01T10:00:00"), _a("b", "2019-06-01T12:00:00"),
_a("c", "2019-06-02T10:00:00")]
intruder = _a("x", "2019-06-01T13:00:00") # in window, untagged
outside = _a("y", "2019-07-01T10:00:00") # out of window
by_id = {m.immich_id: m for m in members + [intruder, outside]}
cand = CandidateCluster(member_ids=["a", "b", "c"], start_at="2019-06-01T10:00:00",
end_at="2019-06-02T10:00:00", suggested_name="Italy 2019",
confidence=0.95, kind_guess="trip", seed_tag="Italy 2019")
out = coverage_members(cand, by_id)
flagged = {m.immich_id for m in out if m.flagged_coverage}
assert flagged == {"x"} # only the in-window intruder
assert all(not m.included for m in out if m.flagged_coverage)
def test_outlier_member_far_from_bulk():
members = [_a("a", "2019-06-01T10:00:00"), _a("a2", "2019-06-01T11:00:00"),
_a("a3", "2019-06-01T12:00:00"),
_a("z", "2019-09-01T10:00:00")] # tagged but months away
by_id = {m.immich_id: m for m in members}
cand = CandidateCluster(member_ids=["a", "a2", "a3", "z"],
start_at="2019-06-01T10:00:00", end_at="2019-09-01T10:00:00",
suggested_name="Italy 2019", confidence=0.95,
kind_guess="trip", seed_tag="Italy 2019")
out = coverage_members(cand, by_id)
outliers = {m.immich_id for m in out if m.is_outlier}
assert outliers == {"z"}
def test_non_seeded_cluster_has_no_coverage_or_outliers():
members = [_a("a", "2019-06-01T10:00:00"), _a("b", "2019-06-01T12:00:00")]
intruder = _a("x", "2019-06-01T11:00:00")
by_id = {m.immich_id: m for m in members + [intruder]}
cand = CandidateCluster(member_ids=["a", "b"], start_at="2019-06-01T10:00:00",
end_at="2019-06-01T12:00:00", suggested_name="Trip",
confidence=0.4, kind_guess="everyday", seed_tag=None)
out = coverage_members(cand, by_id)
assert {m.immich_id for m in out} == {"a", "b"}
assert not any(m.flagged_coverage or m.is_outlier for m in out)
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_coverage.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.coverage'`
- [ ] **Step 3: Implement coverage**
`apps/trip-cluster/app/coverage.py`:
```python
from photoflow.core.models import ClusterMember
from app.clustering import _epoch, _median
def coverage_members(candidate, assets_by_id, *, outlier_factor: float = 8.0) -> list:
member_ids = [i for i in candidate.member_ids if i in assets_by_id]
members = sorted((assets_by_id[i] for i in member_ids), key=lambda a: a.taken_at)
epochs = [_epoch(a.taken_at) for a in members]
gaps = [b - a for a, b in zip(epochs, epochs[1:])]
base = _median(gaps) if gaps else 0.0
out = []
for idx, a in enumerate(members):
is_outlier = False
if candidate.seed_tag and base > 0:
left = epochs[idx] - epochs[idx - 1] if idx > 0 else 0
right = epochs[idx + 1] - epochs[idx] if idx < len(members) - 1 else 0
nearest = min([g for g in (left, right) if g > 0], default=0)
if nearest > outlier_factor * base:
is_outlier = True
out.append(ClusterMember(
cluster_id=0, immich_id=a.immich_id,
member_confidence=candidate.confidence, is_outlier=is_outlier,
included=True, flagged_coverage=False))
if candidate.seed_tag and candidate.start_at and candidate.end_at:
mset = set(member_ids)
for a in sorted(assets_by_id.values(), key=lambda x: x.taken_at):
if a.immich_id in mset:
continue
if candidate.start_at <= a.taken_at <= candidate.end_at:
out.append(ClusterMember(
cluster_id=0, immich_id=a.immich_id, member_confidence=0.0,
is_outlier=False, included=False, flagged_coverage=True))
return out
```
- [ ] **Step 4: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_coverage.py -q`
Expected: PASS (3 passed)
- [ ] **Step 5: Commit**
```bash
git add apps/trip-cluster/app/coverage.py apps/trip-cluster/tests/test_coverage.py
git commit -m "feat(trip-cluster): coverage detection — completeness flags + outliers"
```
---
## Task 12: trip-cluster — cluster orchestration + CLI `cluster`
**Files:**
- Create: `apps/trip-cluster/app/cluster_run.py`
- Modify: `apps/trip-cluster/app/cli.py`
- Test: `apps/trip-cluster/tests/test_cluster_run.py`
**Interfaces:**
- Consumes: `Store` (asset/tag reads + `clear_clusters`/`insert_cluster`), `cluster_assets`, `coverage_members`, `photoflow.immich.pipeline`, `app.clustering._epoch`.
- Produces:
- `app.cluster_run.run_cluster(store, *, gap_factor: float = 6.0, seed_max_span_days: int = 60) -> dict` returning `{"clusters": int}`. Reads non-processed assets + their tags, derives `seed_tags` (non-pipeline tags whose tagged assets span ≤ `seed_max_span_days` and tag ≥ 2 assets), runs clustering + coverage, replaces all clusters in the store.
- CLI `cluster` subcommand wired to it.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_cluster_run.py`:
```python
from photoflow.core import Store
from photoflow.core.models import Asset
from photoflow.immich import pipeline
from app.cluster_run import run_cluster
def _store(tmp_path):
return Store(str(tmp_path / "t.db")).connect()
def test_run_cluster_seeds_bounded_tag_not_people_tag(tmp_path):
s = _store(tmp_path)
# "Italy 2019" spans 2 days -> seeds; "Mum" spans years -> not a seed.
rows = [("a", "2019-06-01T10:00:00", ["Italy 2019", "Mum"]),
("b", "2019-06-02T10:00:00", ["Italy 2019"]),
("c", "2010-01-01T10:00:00", ["Mum"]),
("d", "2022-01-01T10:00:00", ["Mum"])]
for i, t, tags in rows:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
s.set_asset_tags(i, tags)
res = run_cluster(s)
seeded = [c for c in s.all_clusters() if c.suggested_name == "Italy 2019"]
assert len(seeded) == 1
assert sorted(x.immich_id for x, _ in s.cluster_members(seeded[0].id)) == ["a", "b"]
assert res["clusters"] >= 1
s.close()
def test_run_cluster_excludes_processed_and_replaces(tmp_path):
s = _store(tmp_path)
s.upsert_asset(Asset(immich_id="a", taken_at="2019-06-01T10:00:00"))
s.upsert_asset(Asset(immich_id="p", taken_at="2019-06-01T11:00:00"))
s.set_asset_tags("a", [])
s.set_asset_tags("p", [pipeline.PROCESSED])
s.mark_processed("p")
run_cluster(s)
run_cluster(s) # idempotent replace — not doubled
all_member_ids = [x.immich_id for c in s.all_clusters()
for x, _ in s.cluster_members(c.id)]
assert "p" not in all_member_ids
assert all_member_ids.count("a") == 1
s.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_cluster_run.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.cluster_run'`
- [ ] **Step 3: Implement orchestration**
`apps/trip-cluster/app/cluster_run.py`:
```python
from collections import defaultdict
from photoflow.core.models import Cluster
from photoflow.immich import pipeline
from app.clustering import cluster_assets, _epoch
from app.coverage import coverage_members
def _seed_tags(assets, tags_by_asset, max_span_days: int) -> set:
times = defaultdict(list)
for a in assets:
for t in tags_by_asset.get(a.immich_id, []):
if not pipeline.is_pipeline_tag(t):
times[t].append(a.taken_at)
seeds = set()
for tag, ts in times.items():
if len(ts) < 2:
continue
if _epoch(max(ts)) - _epoch(min(ts)) <= max_span_days * 86400:
seeds.add(tag)
return seeds
def run_cluster(store, *, gap_factor: float = 6.0, seed_max_span_days: int = 60) -> dict:
assets = store.all_assets(include_processed=False)
tags_by_asset = {a.immich_id: store.asset_tags(a.immich_id) for a in assets}
seed_tags = _seed_tags(assets, tags_by_asset, seed_max_span_days)
candidates = cluster_assets(assets, tags_by_asset, seed_tags, gap_factor=gap_factor)
assets_by_id = {a.immich_id: a for a in assets}
store.clear_clusters()
n = 0
for cand in candidates:
members = coverage_members(cand, assets_by_id)
store.insert_cluster(Cluster(
start_at=cand.start_at, end_at=cand.end_at,
count=sum(1 for m in members if m.included),
suggested_name=cand.suggested_name, confidence=cand.confidence,
kind_guess=cand.kind_guess, status="pending"), members)
n += 1
return {"clusters": n}
```
- [ ] **Step 4: Wire the CLI `cluster` subcommand**
In `apps/trip-cluster/app/cli.py`, add:
```python
def cmd_cluster(deps, *, gap_factor) -> int:
cfg = deps["config"]
store = _store(cfg)
from app.cluster_run import run_cluster
res = run_cluster(store, gap_factor=gap_factor)
store.close()
print(f"Built {res['clusters']} candidate cluster(s). Run 'categorize serve' to review.")
return 0
```
In `main()`, add the dispatch:
```python
if args.command == "cluster":
return cmd_cluster(deps, gap_factor=args.gap_factor)
```
- [ ] **Step 5: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_cluster_run.py -q`
Expected: PASS (2 passed)
- [ ] **Step 6: Commit**
```bash
git add apps/trip-cluster/app/cluster_run.py apps/trip-cluster/app/cli.py apps/trip-cluster/tests/test_cluster_run.py
git commit -m "feat(trip-cluster): cluster orchestration (seed derivation + coverage) + CLI"
```
> **Review revision (doc-review 2026-06-27):** Don't gate seed derivation on raw span alone. The `seed_max_span_days = 60` default both mis-seeds short-span non-trip tags (a one-weekend "Birthday" or name burst becomes an un-splittable trip) **and** wrongly rejects genuine long trips (this library has a 2011 trip longer than 60 days). Gate seeding on temporal contiguity (do the tag's photos dominate their window?) and raise the span default well above 60; add a unit test asserting a short-span non-trip tag does NOT seed.
---
## Task 13: trip-cluster — review operations
**Files:**
- Create: `apps/trip-cluster/app/review.py`
- Test: `apps/trip-cluster/tests/test_review.py`
**Interfaces:**
- Consumes: `Store` cluster methods (Task 6).
- Produces (all take `store` first, stamp `reviewed_at` where they finalize a decision):
- `app.review.approve(store, cluster_id, name=None) -> None` — status `approved`; `decided_name` = `name` or the cluster's `suggested_name`.
- `app.review.mark_non_trip(store, cluster_id) -> None` — status `non_trip`.
- `app.review.skip(store, cluster_id) -> None` — status `skipped`.
- `app.review.split(store, cluster_id, boundary_immich_id) -> tuple[int, int]`
- `app.review.merge(store, cluster_id_a, cluster_id_b) -> int`
- `app.review.set_member(store, cluster_id, immich_id, included: bool) -> None`
- `app.review.approve_high_confidence(store, threshold: float = 0.75) -> int` — approve every `pending`, `kind_guess='trip'` cluster with `confidence >= threshold`; returns the count.
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_review.py`:
```python
from photoflow.core import Store
from photoflow.core.models import Asset, Cluster, ClusterMember
from app import review
def _store(tmp_path):
s = Store(str(tmp_path / "t.db")).connect()
for i, t in [("a", "2019-06-01"), ("b", "2019-06-02"), ("c", "2019-06-03")]:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
return s
def _cluster(s, name="Trip", conf=0.9, kind="trip", ids=("a", "b")):
return s.insert_cluster(
Cluster(start_at="2019-06-01", end_at="2019-06-03", suggested_name=name,
confidence=conf, kind_guess=kind, status="pending"),
[ClusterMember(cluster_id=0, immich_id=i) for i in ids])
def test_approve_uses_suggested_when_no_name(tmp_path):
s = _store(tmp_path)
cid = _cluster(s, name="Venice")
review.approve(s, cid)
c = s.get_cluster(cid)
assert c.status == "approved" and c.decided_name == "Venice" and c.reviewed_at
review.approve(s, cid, name="Venezia")
assert s.get_cluster(cid).decided_name == "Venezia"
s.close()
def test_non_trip_skip(tmp_path):
s = _store(tmp_path)
cid = _cluster(s)
review.mark_non_trip(s, cid)
assert s.get_cluster(cid).status == "non_trip"
c2 = _cluster(s)
review.skip(s, c2)
assert s.get_cluster(c2).status == "skipped"
s.close()
def test_set_member(tmp_path):
s = _store(tmp_path)
cid = _cluster(s, ids=("a", "b"))
review.set_member(s, cid, "b", False)
inc = {m.immich_id: m.included for _, m in s.cluster_members(cid)}
assert inc["b"] is False
s.close()
def test_approve_high_confidence_only_trips(tmp_path):
s = _store(tmp_path)
hi = _cluster(s, conf=0.9, kind="trip")
_cluster(s, conf=0.3, kind="trip") # too low
_cluster(s, conf=0.9, kind="everyday") # everyday excluded
n = review.approve_high_confidence(s, threshold=0.75)
assert n == 1
assert s.get_cluster(hi).status == "approved"
s.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_review.py -q`
Expected: FAIL — `ImportError: cannot import name 'review'` / `AttributeError`
- [ ] **Step 3: Implement review operations**
`apps/trip-cluster/app/review.py`:
```python
import datetime
def _now() -> str:
return datetime.datetime.now(datetime.timezone.utc).isoformat()
def approve(store, cluster_id, name=None) -> None:
c = store.get_cluster(cluster_id)
decided = (name or "").strip() or c.suggested_name
store.update_cluster(cluster_id, status="approved",
decided_name=decided, reviewed_at=_now())
def mark_non_trip(store, cluster_id) -> None:
store.update_cluster(cluster_id, status="non_trip", reviewed_at=_now())
def skip(store, cluster_id) -> None:
store.update_cluster(cluster_id, status="skipped", reviewed_at=_now())
def split(store, cluster_id, boundary_immich_id):
return store.split_cluster(cluster_id, boundary_immich_id)
def merge(store, cluster_id_a, cluster_id_b):
return store.merge_clusters(cluster_id_a, cluster_id_b)
def set_member(store, cluster_id, immich_id, included: bool) -> None:
store.set_member_inclusion(cluster_id, immich_id, included)
def approve_high_confidence(store, threshold: float = 0.75) -> int:
count = 0
for c in store.all_clusters():
if c.status == "pending" and c.kind_guess == "trip" and c.confidence >= threshold:
approve(store, c.id)
count += 1
return count
```
- [ ] **Step 4: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_review.py -q`
Expected: PASS (4 passed)
- [ ] **Step 5: Commit**
```bash
git add apps/trip-cluster/app/review.py apps/trip-cluster/tests/test_review.py
git commit -m "feat(trip-cluster): cluster review operations (approve/non-trip/skip/split/merge/bulk)"
```
> **Review revision (doc-review 2026-06-27):** Tie `approve_high_confidence`'s `0.75` threshold to observed validation-gate results rather than a fixed constant (see Task 10) — otherwise it never fires on GPS-poor clusters whose confidence caps at `0.50`.
---
## Task 14: trip-cluster — write-back (idempotent) + CLI `apply`
**Files:**
- Create: `apps/trip-cluster/app/writeback.py`
- Modify: `apps/trip-cluster/app/cli.py`
- Test: `apps/trip-cluster/tests/test_writeback.py`
**Interfaces:**
- Consumes: `ImmichClient.upsert_tag`/`tag_assets`, `Store` (cluster members, `already_applied`, `log_writeback`, `mark_processed`), `photoflow.immich.pipeline.{NON_TRIP, PROCESSED}`.
- Produces:
- `app.writeback.apply_cluster(client, store, cluster_id) -> dict` returning `{"cluster_id", "status", "succeeded": list[str], "failed": list[tuple[str, str]]}`. Writes only **included** members; `approved` → content trip tag (`decided_name` or `suggested_name`); `non_trip``_pipeline/non-trip`; all three applyable states then mark `_pipeline/processed`. Idempotent (skips ids already `ok` in `writeback_log`); per-asset failures recorded, cluster left in place for retry.
- `app.writeback.apply_all(client, store) -> list[dict]` — applies every cluster in an applyable state (`approved`/`non_trip`/`skipped`), returns per-cluster results.
- CLI `apply` subcommand with a confirmation prompt (skipped by `--yes`).
- [ ] **Step 1: Write the failing tests**
`apps/trip-cluster/tests/test_writeback.py`:
```python
from photoflow.core import Store
from photoflow.core.models import Asset, Cluster, ClusterMember
from photoflow.immich import pipeline
from app.writeback import apply_cluster, apply_all
class FakeImmich:
def __init__(self):
self.tagged = []
self.fail_tag_id = None
self._ids = {}
def upsert_tag(self, name):
self._ids.setdefault(name, f"id:{name}")
return self._ids[name]
def tag_assets(self, tag_id, ids):
if self.fail_tag_id is not None and tag_id == self.fail_tag_id:
raise RuntimeError("boom")
self.tagged.append((tag_id, list(ids)))
def _store(tmp_path):
s = Store(str(tmp_path / "t.db")).connect()
for i in ("a", "b", "x"):
s.upsert_asset(Asset(immich_id=i, taken_at="2019-06-01"))
return s
def _approved(s, name="Venice"):
cid = s.insert_cluster(
Cluster(start_at="2019-06-01", end_at="2019-06-02", suggested_name=name,
status="approved", decided_name=name),
[ClusterMember(cluster_id=0, immich_id="a"),
ClusterMember(cluster_id=0, immich_id="b"),
ClusterMember(cluster_id=0, immich_id="x", included=False, flagged_coverage=True)])
return cid
def test_apply_approved_tags_included_then_processed(tmp_path):
s = _store(tmp_path)
cid = _approved(s)
client = FakeImmich()
res = apply_cluster(client, s, cid)
assert sorted(res["succeeded"]) == ["a", "b"] and res["failed"] == []
# trip tag on a,b ; processed on a,b ; x (excluded) never tagged
assert ("id:Venice", ["a", "b"]) in client.tagged
assert ("id:_pipeline/processed", ["a", "b"]) in client.tagged
assert all("x" not in ids for _, ids in client.tagged)
assert s.get_asset("a").processed is True
s.close()
def test_apply_is_idempotent(tmp_path):
s = _store(tmp_path)
cid = _approved(s)
client = FakeImmich()
apply_cluster(client, s, cid)
before = len(client.tagged)
apply_cluster(client, s, cid) # second run writes nothing new
assert len(client.tagged) == before
s.close()
def test_partial_failure_leaves_retryable(tmp_path):
s = _store(tmp_path)
cid = _approved(s)
client = FakeImmich()
client.fail_tag_id = "id:Venice" # trip tag write fails
res = apply_cluster(client, s, cid)
assert res["succeeded"] == [] and sorted(i for i, _ in res["failed"]) == ["a", "b"]
assert s.get_asset("a").processed is False # not marked processed on failure
assert s.already_applied("a", "trip", "Venice") is False # retryable
s.close()
def test_apply_non_trip(tmp_path):
s = _store(tmp_path)
cid = s.insert_cluster(
Cluster(start_at="2019-06-01", end_at="2019-06-02", status="non_trip"),
[ClusterMember(cluster_id=0, immich_id="a")])
client = FakeImmich()
apply_cluster(client, s, cid)
assert ("id:_pipeline/non-trip", ["a"]) in client.tagged
assert ("id:_pipeline/processed", ["a"]) in client.tagged
s.close()
def test_apply_all_reports_per_cluster(tmp_path):
s = _store(tmp_path)
_approved(s, name="Venice")
client = FakeImmich()
results = apply_all(client, s)
assert len(results) == 1 and results[0]["status"] == "approved"
s.close()
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_writeback.py -q`
Expected: FAIL — `ModuleNotFoundError: No module named 'app.writeback'`
- [ ] **Step 3: Implement write-back**
`apps/trip-cluster/app/writeback.py`:
```python
from photoflow.immich import pipeline
APPLYABLE = ("approved", "non_trip", "skipped")
def _apply_tag(client, store, asset_ids, action, tag, tag_id):
todo = [a for a in asset_ids if not store.already_applied(a, action, tag)]
if not todo:
return [], []
try:
client.tag_assets(tag_id, todo)
except Exception as e: # noqa: BLE001 — recorded, surfaced
for a in todo:
store.log_writeback(a, action, tag, f"error:{e}")
return [], [(a, str(e)) for a in todo]
for a in todo:
store.log_writeback(a, action, tag, "ok")
return todo, []
def apply_cluster(client, store, cluster_id) -> dict:
c = store.get_cluster(cluster_id)
if c is None or c.status not in APPLYABLE:
return {"cluster_id": cluster_id, "status": c.status if c else None,
"succeeded": [], "failed": []}
included = [a.immich_id for a, m in store.cluster_members(cluster_id) if m.included]
succeeded, failed = [], []
if c.status == "approved":
tag = c.decided_name or c.suggested_name
ok, fail = _apply_tag(client, store, included, "trip", tag, client.upsert_tag(tag))
succeeded += ok
failed += fail
elif c.status == "non_trip":
ok, fail = _apply_tag(client, store, included, "non-trip", pipeline.NON_TRIP,
client.upsert_tag(pipeline.NON_TRIP))
succeeded += ok
failed += fail
# 'skipped': no content/non-trip tag, only processed below.
failed_ids = {i for i, _ in failed}
proc_targets = [a for a in included if a not in failed_ids]
if proc_targets:
_apply_tag(client, store, proc_targets, "processed", pipeline.PROCESSED,
client.upsert_tag(pipeline.PROCESSED))
for a in proc_targets:
if store.already_applied(a, "processed", pipeline.PROCESSED):
store.mark_processed(a)
return {"cluster_id": cluster_id, "status": c.status,
"succeeded": succeeded, "failed": failed}
def apply_all(client, store) -> list:
return [apply_cluster(client, store, c.id)
for c in store.all_clusters() if c.status in APPLYABLE]
```
- [ ] **Step 4: Wire the CLI `apply` subcommand**
In `apps/trip-cluster/app/cli.py`, add:
```python
def cmd_apply(deps, *, yes) -> int:
cfg = deps["config"]
store = _store(cfg)
from app.writeback import apply_all, APPLYABLE
pending = [c for c in store.all_clusters() if c.status in APPLYABLE]
if not pending:
print("Nothing to apply.")
store.close()
return 0
if not yes:
ans = input(f"Apply {len(pending)} cluster decision(s) to Immich? [y/N] ")
if ans.strip().lower() not in ("y", "yes"):
print("Aborted.")
store.close()
return 1
client = _immich(deps, cfg)
results = apply_all(client, store)
store.close()
ok = sum(len(r["succeeded"]) for r in results)
bad = sum(len(r["failed"]) for r in results)
print(f"Applied {ok} tag write(s); {bad} failure(s) across {len(results)} cluster(s).")
return 0 if bad == 0 else 1
```
In `main()`, add:
```python
if args.command == "apply":
return cmd_apply(deps, yes=args.yes)
```
- [ ] **Step 5: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_writeback.py -q`
Expected: PASS (5 passed)
- [ ] **Step 6: Commit**
```bash
git add apps/trip-cluster/app/writeback.py apps/trip-cluster/app/cli.py apps/trip-cluster/tests/test_writeback.py
git commit -m "feat(trip-cluster): idempotent write-back + CLI apply with confirmation"
```
> **Review revision (doc-review 2026-06-27):** Add `test_apply_skipped`: verify a `skipped` cluster tags only `_pipeline/processed` (no content/non-trip tag), marks its assets processed, and returns the expected `succeeded`/`failed` shape. The `skipped` branch is in `APPLYABLE` but has no coverage, and its `succeeded=[]`-on-success shape makes a successful apply indistinguishable from an idempotent no-op in the CLI summary.
---
## Task 15: trip-cluster — review routes + master/detail templates
**Files:**
- Modify: `apps/trip-cluster/app/routes/nav.py`, `apps/trip-cluster/app/routes/review.py`
- Create: `apps/trip-cluster/app/templates/review.html`, `apps/trip-cluster/app/templates/_detail.html`, `apps/trip-cluster/app/static/app.js`
- Test: append to `apps/trip-cluster/tests/test_routes.py`
**Interfaces:**
- Consumes: `Store`, `app.review`, `app.writeback`, shared `macros.html`/`base.html`, `photoGrid()`.
- Produces:
- `GET /` — master/detail page (left rail = `clusters_by_attention`, detail HTMX-loaded), with three empty states keyed on `stats.assets` / `stats.clusters` / `stats.pending`.
- `GET /cluster/<int:cid>``_detail.html` partial (chronological grid + coverage/outlier toggles + merge-neighbour buttons + lightbox).
- `POST /cluster/<int:cid>/approve` `{name?}`, `/non-trip`, `/skip`, `/split` `{boundary_asset_id}`, `/merge` `{other_id}`, `/member` `{asset_id, included}`, `/apply` — JSON.
- `POST /approve-high-confidence` `{threshold?}``{approved}`; `POST /apply-all``{results}`.
- Test seam: routes build the Immich client from `current_app.config["IMMICH_FACTORY"]()` when set, else from config.
- [ ] **Step 1: Write the failing route tests**
Append to `apps/trip-cluster/tests/test_routes.py`:
```python
from photoflow.core import Store
from photoflow.core.models import Asset, Cluster, ClusterMember
def _seed_store(cfg):
s = Store(cfg.db_path).connect()
for i, t in [("a", "2019-06-01"), ("b", "2019-06-02"), ("c", "2019-07-10")]:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
cid = s.insert_cluster(
Cluster(start_at="2019-06-01", end_at="2019-06-02", suggested_name="Venice",
confidence=0.9, kind_guess="trip", status="pending"),
[ClusterMember(cluster_id=0, immich_id="a"),
ClusterMember(cluster_id=0, immich_id="b")])
other = s.insert_cluster(
Cluster(start_at="2019-07-10", end_at="2019-07-10", suggested_name="Rome",
confidence=0.3, kind_guess="trip", status="pending"),
[ClusterMember(cluster_id=0, immich_id="c")])
s.close()
return cid, other
def test_index_lists_clusters(tmp_path):
app = _app(tmp_path)
_seed_store(app.config["APP_CONFIG"])
r = app.test_client().get("/")
assert r.status_code == 200 and b"Venice" in r.data and b"Rome" in r.data
def test_index_empty_state(tmp_path):
r = _app(tmp_path).test_client().get("/")
assert r.status_code == 200 and b"No clusters yet" in r.data
def test_detail_and_approve(tmp_path):
app = _app(tmp_path)
cid, _ = _seed_store(app.config["APP_CONFIG"])
client = app.test_client()
d = client.get(f"/cluster/{cid}")
assert d.status_code == 200 and b"/thumb/a" in d.data
r = client.post(f"/cluster/{cid}/approve", json={"name": "Venezia"})
assert r.get_json()["status"] == "approved"
s = Store(app.config["APP_CONFIG"].db_path).connect()
assert s.get_cluster(cid).decided_name == "Venezia"
s.close()
def test_member_toggle_and_high_confidence(tmp_path):
app = _app(tmp_path)
cid, _ = _seed_store(app.config["APP_CONFIG"])
client = app.test_client()
client.post(f"/cluster/{cid}/member", json={"asset_id": "b", "included": False})
s = Store(app.config["APP_CONFIG"].db_path).connect()
assert {m.immich_id: m.included for _, m in s.cluster_members(cid)}["b"] is False
s.close()
r = client.post("/approve-high-confidence", json={"threshold": 0.75})
assert r.get_json()["approved"] == 1 # only the 0.9 cluster
def test_apply_all_with_injected_client(tmp_path):
app = _app(tmp_path)
cid, _ = _seed_store(app.config["APP_CONFIG"])
class FakeImmich:
def __init__(self): self.tagged = []
def upsert_tag(self, name): return f"id:{name}"
def tag_assets(self, tid, ids): self.tagged.append((tid, list(ids)))
fake = FakeImmich()
app.config["IMMICH_FACTORY"] = lambda: fake
client = app.test_client()
client.post(f"/cluster/{cid}/approve", json={"name": "Venice"})
r = client.post("/apply-all", json={})
assert r.status_code == 200
assert any(tid == "id:Venice" for tid, _ in fake.tagged)
```
- [ ] **Step 2: Run to verify they fail**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_routes.py -q`
Expected: FAIL — index 500/`TemplateNotFound: review.html` and missing endpoints (404).
- [ ] **Step 3: Implement nav.index**
Replace `apps/trip-cluster/app/routes/nav.py`:
```python
from flask import Blueprint, current_app, render_template
from photoflow.core import Store
bp = Blueprint("nav", __name__)
def _store():
cfg = current_app.config["APP_CONFIG"]
return Store(cfg.db_path).connect()
@bp.route("/health")
def health():
return "ok"
@bp.route("/")
def index():
s = _store()
clusters = s.clusters_by_attention()
stats = {"assets": len(s.all_assets()), "clusters": len(clusters),
"pending": sum(1 for c in clusters if c.status == "pending")}
s.close()
return render_template("review.html", clusters=clusters, stats=stats)
```
- [ ] **Step 4: Implement the review routes**
Replace `apps/trip-cluster/app/routes/review.py`:
```python
from flask import Blueprint, current_app, render_template, request, jsonify, abort
from photoflow.core import Store
from app import review
from app.writeback import apply_cluster, apply_all
bp = Blueprint("review", __name__)
def _store():
cfg = current_app.config["APP_CONFIG"]
return Store(cfg.db_path).connect()
def _client():
factory = current_app.config.get("IMMICH_FACTORY")
if factory:
return factory()
from photoflow.immich import ImmichClient
cfg = current_app.config["APP_CONFIG"]
return ImmichClient(cfg.immich_url, cfg.immich_api_key)
def _body():
return request.get_json(silent=True) or {}
@bp.route("/cluster/<int:cid>")
def detail(cid):
s = _store()
c = s.get_cluster(cid)
if c is None:
s.close()
abort(404)
members = s.cluster_members(cid)
prev_id, next_id = s.chronological_neighbors(cid)
s.close()
return render_template("_detail.html", c=c, members=members,
prev_id=prev_id, next_id=next_id)
@bp.route("/cluster/<int:cid>/approve", methods=["POST"])
def approve(cid):
s = _store()
review.approve(s, cid, _body().get("name"))
c = s.get_cluster(cid)
s.close()
return jsonify({"ok": True, "status": c.status, "name": c.decided_name})
@bp.route("/cluster/<int:cid>/non-trip", methods=["POST"])
def non_trip(cid):
s = _store()
review.mark_non_trip(s, cid)
s.close()
return jsonify({"ok": True, "status": "non_trip"})
@bp.route("/cluster/<int:cid>/skip", methods=["POST"])
def skip(cid):
s = _store()
review.skip(s, cid)
s.close()
return jsonify({"ok": True, "status": "skipped"})
@bp.route("/cluster/<int:cid>/split", methods=["POST"])
def split(cid):
s = _store()
boundary = _body().get("boundary_asset_id")
if not boundary:
s.close()
return jsonify({"error": "boundary_asset_id required"}), 400
try:
id1, id2 = review.split(s, cid, boundary)
except ValueError as e:
s.close()
return jsonify({"error": str(e)}), 400
s.close()
return jsonify({"ok": True, "ids": [id1, id2]})
@bp.route("/cluster/<int:cid>/merge", methods=["POST"])
def merge(cid):
s = _store()
other = _body().get("other_id")
if other is None:
s.close()
return jsonify({"error": "other_id required"}), 400
new_id = review.merge(s, cid, int(other))
s.close()
return jsonify({"ok": True, "id": new_id})
@bp.route("/cluster/<int:cid>/member", methods=["POST"])
def member(cid):
s = _store()
body = _body()
review.set_member(s, cid, body["asset_id"], bool(body.get("included")))
s.close()
return jsonify({"ok": True})
@bp.route("/approve-high-confidence", methods=["POST"])
def approve_high_confidence():
s = _store()
threshold = float(_body().get("threshold", 0.75))
n = review.approve_high_confidence(s, threshold)
s.close()
return jsonify({"approved": n})
@bp.route("/cluster/<int:cid>/apply", methods=["POST"])
def apply_one(cid):
s = _store()
res = apply_cluster(_client(), s, cid)
s.close()
return jsonify(res)
@bp.route("/apply-all", methods=["POST"])
def apply_everything():
s = _store()
results = apply_all(_client(), s)
s.close()
return jsonify({"results": results})
```
- [ ] **Step 5: Implement review.html**
`apps/trip-cluster/app/templates/review.html`:
```html
{% extends "base.html" %}
{% from "macros.html" import confidence_badge, status_badge %}
{% block title %}trip-cluster{% endblock %}
{% block navbar_title %}trip-cluster{% endblock %}
{% block content %}
{% if stats.clusters == 0 %}
<div class="hero py-20"><div class="hero-content text-center"><div>
<h1 class="text-2xl font-bold">No clusters yet</h1>
{% if stats.assets == 0 %}
<p class="opacity-70 mt-2">Nothing ingested. Run
<code>categorize ingest</code> then <code>categorize cluster</code>.</p>
{% else %}
<p class="opacity-70 mt-2">{{ stats.assets }} assets ingested, but this scope
produced no clusters. Try a wider <code>categorize ingest</code> scope.</p>
{% endif %}
</div></div></div>
{% else %}
<div x-data="clusterReview()" @keydown.window="onKey($event)" class="flex gap-4">
<div class="w-72 shrink-0 max-h-[85vh] overflow-y-auto">
<div class="flex items-center justify-between mb-2">
<span class="font-bold">Clusters</span>
<button class="btn btn-xs btn-primary" @click="approveHighConfidence()">Approve high-conf</button>
</div>
{% if stats.pending == 0 %}
<div class="alert alert-success text-xs mb-2">All reviewed — “Apply all” or <code>categorize apply</code>.</div>
{% endif %}
{% for c in clusters %}
<div class="cluster-row block p-2 rounded cursor-pointer hover:bg-base-300"
:class="selected == {{ c.id }} && 'bg-base-300'"
data-cluster-id="{{ c.id }}" @click="selectCluster({{ c.id }})">
<div class="flex items-center justify-between gap-1">
<span class="truncate text-sm font-medium">{{ c.decided_name or c.suggested_name }}</span>
{{ confidence_badge(c.confidence) }}
</div>
<div class="flex items-center gap-1 mt-1">
{{ status_badge(c.status) }}
<span class="badge badge-xs">{{ c.count }}</span>
<span class="text-xs opacity-50">{{ c.start_at[:10] }}</span>
</div>
</div>
{% endfor %}
<button class="btn btn-xs btn-block mt-3" @click="applyAll()">Apply all approved</button>
</div>
<div id="detail" class="flex-1"
hx-get="/cluster/{{ clusters[0].id }}" hx-trigger="load" hx-target="#detail"></div>
</div>
<div class="fixed bottom-2 right-2 text-xs opacity-50">
[ ] cluster · ← → grid · Enter open · A approve · N non-trip · S split · X skip · Esc close
</div>
{% endif %}
{% endblock %}
{% block extra_scripts %}<script src="/static/app.js"></script>{% endblock %}
```
- [ ] **Step 6: Implement _detail.html**
`apps/trip-cluster/app/templates/_detail.html`:
```html
{% from "macros.html" import lightbox %}
<div data-cluster-id="{{ c.id }}">
<div class="flex flex-wrap items-center gap-2 mb-3">
<input id="cluster-name" class="input input-bordered input-sm"
value="{{ c.decided_name or c.suggested_name }}">
<button class="btn btn-sm btn-success" @click="approve()">Approve (A)</button>
<button class="btn btn-sm" @click="nonTrip()">Non-trip (N)</button>
<button class="btn btn-sm" @click="split()">Split (S)</button>
<button class="btn btn-sm" @click="skip()">Skip (X)</button>
<button class="btn btn-sm btn-ghost" @click="apply()">Apply</button>
{% if prev_id %}<button class="btn btn-xs" @click="merge({{ prev_id }})">⤺ merge prev</button>{% endif %}
{% if next_id %}<button class="btn btn-xs" @click="merge({{ next_id }})">merge next ⤻</button>{% endif %}
</div>
<div class="text-sm opacity-60 mb-2">
{{ c.start_at[:10] }} → {{ c.end_at[:10] }} · {{ members | length }} assets · status {{ c.status }}
</div>
<div class="grid grid-cols-3 sm:grid-cols-4 md:grid-cols-6 gap-2">
{% for a, m in members %}
<div class="photo-card relative cursor-pointer rounded-lg overflow-hidden border-2
{% if m.flagged_coverage %}border-info{% elif m.is_outlier %}border-warning{% else %}border-transparent{% endif %}"
data-asset-id="{{ a.immich_id }}" data-included="{{ '1' if m.included else '' }}"
tabindex="0" @click="openLightbox($el)" @focus="select($el)">
<img src="/thumb/{{ a.immich_id }}" class="w-full aspect-square object-cover" loading="lazy" alt="">
{% if m.flagged_coverage %}
<button class="absolute bottom-1 left-1 badge badge-xs badge-info"
@click.stop="setMember('{{ a.immich_id }}', true)">+ include</button>
{% elif m.is_outlier %}
<button class="absolute bottom-1 left-1 badge badge-xs badge-warning"
@click.stop="setMember('{{ a.immich_id }}', false)"> exclude</button>
{% endif %}
{% if not m.included %}<div class="absolute inset-0 bg-black/50 pointer-events-none"></div>{% endif %}
</div>
{% endfor %}
</div>
{{ lightbox() }}
</div>
```
- [ ] **Step 7: Implement app.js (clusterReview)**
`apps/trip-cluster/app/static/app.js`:
```javascript
function clusterReview() {
return {
...photoGrid(),
selected: null,
clusterIds: [],
init() {
this.clusterIds = [...document.querySelectorAll('.cluster-row')]
.map(el => parseInt(el.dataset.clusterId, 10));
this.selected = this.clusterIds.length ? this.clusterIds[0] : null;
document.body.addEventListener('htmx:afterSwap', (e) => {
if (e.target.id === 'detail') this.selectFirst();
});
},
selectCluster(id) {
this.selected = id;
htmx.ajax('GET', `/cluster/${id}`, { target: '#detail' });
},
moveCluster(dir) {
const i = this.clusterIds.indexOf(this.selected);
const j = Math.max(0, Math.min(this.clusterIds.length - 1, i + dir));
if (this.clusterIds[j] != null) this.selectCluster(this.clusterIds[j]);
},
nameValue() {
const el = document.getElementById('cluster-name');
return el ? el.value : '';
},
async post(path, body) {
const res = await fetch(path, {
method: 'POST', headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(body || {}),
});
return res.json();
},
async approve() { await this.post(`/cluster/${this.selected}/approve`, { name: this.nameValue() }); location.reload(); },
async nonTrip() { await this.post(`/cluster/${this.selected}/non-trip`); location.reload(); },
async skip() { await this.post(`/cluster/${this.selected}/skip`); location.reload(); },
async split() {
if (!this.focused) { alert('Focus the first photo of the second trip, then Split.'); return; }
await this.post(`/cluster/${this.selected}/split`, { boundary_asset_id: this.focused.dataset.assetId });
location.reload();
},
async merge(otherId) { await this.post(`/cluster/${this.selected}/merge`, { other_id: otherId }); location.reload(); },
async setMember(assetId, included) { await this.post(`/cluster/${this.selected}/member`, { asset_id: assetId, included }); this.selectCluster(this.selected); },
async approveHighConfidence() { const r = await this.post('/approve-high-confidence', {}); alert(`Approved ${r.approved} cluster(s).`); location.reload(); },
async apply() { const r = await this.post(`/cluster/${this.selected}/apply`, {}); alert(`Applied ${r.succeeded.length}, failed ${r.failed.length}.`); },
async applyAll() {
if (!confirm('Apply all approved decisions to Immich?')) return;
const r = await this.post('/apply-all', {});
const ok = r.results.reduce((n, x) => n + x.succeeded.length, 0);
const bad = r.results.reduce((n, x) => n + x.failed.length, 0);
alert(`Applied ${ok} write(s), ${bad} failure(s).`);
location.reload();
},
onKey(e) {
if (e.target.tagName === 'INPUT') return;
const k = e.key.toLowerCase();
if (e.key === '[') { e.preventDefault(); this.moveCluster(-1); }
else if (e.key === ']') { e.preventDefault(); this.moveCluster(1); }
else if (e.key === 'ArrowLeft') { e.preventDefault(); this.navigate(-1); }
else if (e.key === 'ArrowRight') { e.preventDefault(); this.navigate(1); }
else if (e.key === 'Enter') { if (this.focused) this.openLightbox(this.focused); }
else if (e.key === 'Escape') { this.closeLightbox(); }
else if (k === 'a') { this.approve(); }
else if (k === 'n') { this.nonTrip(); }
else if (k === 's') { this.split(); }
else if (k === 'x') { this.skip(); }
},
};
}
```
- [ ] **Step 8: Run to verify they pass**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/test_routes.py -q`
Expected: PASS (7 passed)
- [ ] **Step 9: Commit**
```bash
git add apps/trip-cluster/app/routes apps/trip-cluster/app/templates apps/trip-cluster/app/static apps/trip-cluster/tests/test_routes.py
git commit -m "feat(trip-cluster): master/detail review routes + templates + clusterReview JS"
```
> **Review revisions (doc-review 2026-06-27):**
> - **F3 — preserve position.** Replace the blanket `location.reload()` after each action with a targeted HTMX swap that advances to the next *pending* cluster, so acting on a cluster doesn't bounce the reviewer back to the attention-sorted queue head.
> - **F4 — surface apply failures.** After `apply-all`, distinguish write-failed clusters in the rail (e.g. a `badge-error`) and list them inline so the per-cluster retry the write-back layer supports is actionable.
> - **F7 — split affordance.** Disable the Split button / suppress the `S` shortcut until a boundary photo is focused, with a visible instruction and an updated hint bar.
> - **F8 — explain the sort.** Label the cluster rail ("needs attention first — lowest-confidence pending").
> - **F16 — loading state.** Add an `hx-indicator` loading state on `#detail` for both declarative and programmatic swaps.
> - **F17 — confirm merge.** Confirm before `merge` (it terminally marks both clusters `merged` with no un-merge in the UI).
> - **F19 — confirm bulk approve.** Show a pre-action count and confirm before `approve-high-confidence`.
---
## Task 16: trip-cluster — Playwright UI tests (grid, lightbox, cluster switching, approve)
**Files:**
- Create: `apps/trip-cluster/tests/ui/conftest.py`, `apps/trip-cluster/tests/ui/test_smoke_ui.py`, `apps/trip-cluster/tests/ui/test_review_ui.py`
- Test: themselves (run under `tests/ui`)
**Interfaces:**
- Consumes: `create_app`, `Config`, `Store`. Mirrors the sibling apps' UI harness (real Flask server thread + Playwright `page`). Port **8095** for tests (never 8084/8083/8082).
- [ ] **Step 1: Write the UI harness + smoke test**
`apps/trip-cluster/tests/ui/conftest.py`:
```python
import io
import os
import threading
import time
import pytest
from PIL import Image
from werkzeug.serving import make_server
from app import create_app
from app.config import Config
from photoflow.core import Store
from photoflow.core.models import Asset, Cluster, ClusterMember
def _jpeg(color):
buf = io.BytesIO()
Image.new("RGB", (48, 48), color).save(buf, format="JPEG")
return buf.getvalue()
@pytest.fixture(scope="session")
def data_dir(tmp_path_factory):
d = str(tmp_path_factory.mktemp("data"))
thumbs = os.path.join(d, "thumbs")
os.makedirs(thumbs, exist_ok=True)
s = Store(os.path.join(d, "trip-cluster.db")).connect()
palette = {"a": (10, 20, 30), "b": (40, 80, 120), "c": (200, 50, 90)}
for i, t in [("a", "2019-06-01"), ("b", "2019-06-02"), ("c", "2019-07-10")]:
s.upsert_asset(Asset(immich_id=i, taken_at=t))
with open(os.path.join(thumbs, f"{i}.jpg"), "wb") as f:
f.write(_jpeg(palette[i]))
s.insert_cluster(
Cluster(start_at="2019-06-01", end_at="2019-06-02", suggested_name="Venice",
confidence=0.9, kind_guess="trip", status="pending"),
[ClusterMember(cluster_id=0, immich_id="a"),
ClusterMember(cluster_id=0, immich_id="b")])
s.insert_cluster(
Cluster(start_at="2019-07-10", end_at="2019-07-10", suggested_name="Rome",
confidence=0.3, kind_guess="trip", status="pending"),
[ClusterMember(cluster_id=0, immich_id="c")])
s.close()
return d
@pytest.fixture(scope="session")
def flask_app(data_dir):
cfg = Config(immich_url="http://127.0.0.1:1", immich_api_key="k",
anthropic_api_key="", data_dir=data_dir)
return create_app(cfg)
@pytest.fixture(scope="session")
def base_url(flask_app):
server = make_server("127.0.0.1", 8095, flask_app)
t = threading.Thread(target=server.serve_forever, daemon=True)
t.start()
time.sleep(0.2)
yield "http://127.0.0.1:8095"
server.shutdown()
```
`apps/trip-cluster/tests/ui/test_smoke_ui.py`:
```python
def test_health(base_url, page):
page.goto(f"{base_url}/health")
assert "ok" in page.content()
def test_index_lists_clusters(base_url, page):
page.goto(base_url)
assert "Venice" in page.inner_text("body")
assert "Rome" in page.inner_text("body")
```
- [ ] **Step 2: Install the browser (one-time) and run the smoke test**
Run:
```bash
.venv/bin/python -m playwright install chromium
.venv/bin/python -m pytest apps/trip-cluster/tests/ui/test_smoke_ui.py -q
```
Expected: PASS (2 passed). (The lowest-confidence cluster "Rome" sorts first in the rail; both names appear.)
- [ ] **Step 3: Write the interaction tests**
`apps/trip-cluster/tests/ui/test_review_ui.py`:
```python
def _open(page, base_url):
page.goto(base_url)
page.wait_for_selector(".photo-card")
def test_grid_arrow_navigation_moves_focus_ring(base_url, page):
# Venice (2 photos) is high-confidence; click it to load its detail.
_open(page, base_url)
page.click("text=Venice")
page.wait_for_selector(".photo-card")
cards = page.query_selector_all(".photo-card")
assert len(cards) == 2
page.keyboard.press("ArrowRight")
# second card gains the focus ring class
assert "ring-4" in (page.query_selector_all(".photo-card")[1].get_attribute("class"))
def test_lightbox_opens_and_closes(base_url, page):
_open(page, base_url)
page.click("text=Venice")
page.wait_for_selector(".photo-card")
page.keyboard.press("ArrowRight")
page.keyboard.press("Enter")
assert page.is_visible("#lb")
page.keyboard.press("Escape")
page.wait_for_selector("#lb", state="hidden")
def test_cluster_switching_with_brackets(base_url, page):
_open(page, base_url)
page.click("text=Venice")
page.wait_for_selector(".photo-card")
page.keyboard.press("]") # move to the next cluster
page.wait_for_function(
"document.querySelectorAll('.photo-card').length === 1") # Rome has 1 photo
def test_approve_updates_status(base_url, page):
_open(page, base_url)
page.click("text=Venice")
page.wait_for_selector("#cluster-name")
page.once("dialog", lambda d: d.accept()) # no dialog expected, but be safe
page.click("text=Approve (A)")
page.wait_for_selector("text=approved")
assert "approved" in page.inner_text("body")
```
- [ ] **Step 4: Run the interaction tests**
Run: `.venv/bin/python -m pytest apps/trip-cluster/tests/ui -q`
Expected: PASS (6 passed total).
- [ ] **Step 5: Commit**
```bash
git add apps/trip-cluster/tests/ui
git commit -m "test(trip-cluster): Playwright UI — grid nav, lightbox, cluster switching, approve"
```
---
## Task 17: Docker, docs, validation gate, full-suite green
**Files:**
- Create: `apps/trip-cluster/Dockerfile`, `docker-compose.yml`, `README.md`, `CLAUDE.md`, `docs/M1-validation-gate.md`
- Modify: `apps/trip-cluster/pyproject.toml` (package-data)
**Interfaces:**
- Consumes: everything above. Produces the runnable container + the validation protocol that gates widening past the hard sample.
- [ ] **Step 1: Add package-data so templates/static ship in non-editable installs**
Append to `apps/trip-cluster/pyproject.toml`:
```toml
[tool.setuptools.package-data]
"app" = ["templates/*.html", "static/*.js"]
```
- [ ] **Step 2: Dockerfile (build context = repo root)**
`apps/trip-cluster/Dockerfile`:
```dockerfile
FROM python:3.12-slim
WORKDIR /srv
COPY shared/ ./shared/
COPY apps/trip-cluster/ ./apps/trip-cluster/
RUN pip install --no-cache-dir ./shared ./apps/trip-cluster
ENV DATA_DIR=/data
EXPOSE 8084
WORKDIR /srv/apps/trip-cluster
CMD ["python", "categorize.py", "serve"]
```
- [ ] **Step 3: docker-compose.yml (root)**
`docker-compose.yml`:
```yaml
services:
trip-cluster:
build:
context: .
dockerfile: apps/trip-cluster/Dockerfile
ports:
- "8084:8084"
volumes:
- ./data:/data
env_file: .env
environment:
- DATA_DIR=/data
user: "${UID}:${GID}"
```
- [ ] **Step 4: README.md**
`README.md`:
```markdown
# immich-photo-flow
A monorepo of small tools that clean up and structure a large [Immich](https://immich.app)
library into travel "memories". **M1** ships the shared foundation
(`shared/photoflow/{immich,core,ui}`) and the `trip-cluster` app.
## Layout
- `shared/photoflow/immich` — the one Immich REST client + `_pipeline/` tag conventions.
- `shared/photoflow/core` — SQLite store + domain models (the only SQL).
- `shared/photoflow/ui` — base template, DaisyUI/Tailwind/Alpine/HTMX, shared grid+lightbox.
- `apps/trip-cluster` — CLI + Flask review UI (port 8084).
## Dev setup
```bash
python3.12 -m venv .venv
.venv/bin/pip install -e ./shared -e ./apps/trip-cluster \
pytest==8.3.4 pytest-httpserver==1.1.0 pytest-playwright==0.6.2
.venv/bin/python -m pytest # unit + route tests
.venv/bin/python -m playwright install chromium
.venv/bin/python -m pytest apps/trip-cluster/tests/ui # Playwright UI
```
## trip-cluster workflow
```bash
cp .env.example .env # fill IMMICH_URL + IMMICH_API_KEY
cd apps/trip-cluster
python categorize.py ingest --tag "Italy 2019" # or --from/--to/--subset
python categorize.py cluster
python categorize.py serve # review at http://localhost:8084
python categorize.py apply # write approved tags back (asks to confirm)
```
Immich is the source of truth; SQLite is a rebuildable working layer. Content/trip tags
are never namespaced; pipeline meta-tags live under `_pipeline/`.
```
- [ ] **Step 5: CLAUDE.md**
`CLAUDE.md`:
```markdown
# CLAUDE.md
## What this is
A monorepo (M1) of an Immich photo-flow foundation + the `trip-cluster` app. See
`docs/ROADMAP.md` for milestones and `docs/superpowers/specs/2026-06-27-immich-photo-flow-design.md`
for the M1 design.
## Architecture map
- `shared/photoflow/immich` — **the only** Immich client (`client.py`) + `_pipeline/` tag
conventions (`pipeline.py`).
- `shared/photoflow/core` — **the only** SQLite owner: `store.py` (data-access) + `models.py`.
- `shared/photoflow/ui` — `base.html`, shared Jinja macros, `shared.js` grid+lightbox.
- `apps/trip-cluster/app` — `config.py`, `cli.py` (ingest|cluster|serve|apply),
`ingest.py`, `clustering.py` (pure), `coverage.py` (pure), `cluster_run.py`,
`review.py`, `writeback.py`, `routes/`, `templates/`, `static/app.js`.
## Key invariants
- **Immich is the source of truth**; SQLite is rebuildable. Only **applied** decisions
survive loss of SQLite (via `_pipeline/processed` + `writeback_log`).
- Content/trip tags are **never namespaced**; pipeline meta-tags nest under `_pipeline/`.
- Trip detection is **timestamp-first, density-adaptive**, anchored by existing trip tags,
refined by GPS. Existing trip tags are authoritative seeds.
- Write-back is **idempotent** (`writeback_log`) and needs **explicit confirmation**.
- Ingest is **scopeable** (`--from/--to`, `--tag`, `--subset`) and incremental (`updatedAfter`).
- trip-cluster serves on **8084**.
## Dev commands
```bash
.venv/bin/python -m pytest # unit + route (Immich mocked)
.venv/bin/python -m pytest apps/trip-cluster/tests/ui # Playwright
docker compose up # UI on :8084
```
Write a failing test first (TDD). Mirrors the sibling apps in `/home/mischa/Projects/`.
```
- [ ] **Step 6: docs/M1-validation-gate.md**
`docs/M1-validation-gate.md`:
```markdown
# M1 Validation Gate
Before widening past the hard sample, trip-cluster must clear a quantified bar on a
**deliberately hard slice** — a GPS-poor, multi-year, low-density window (a well-remembered
old trip plus its surrounding everyday photos), **not** the easy phone-era last trip.
## Protocol
1. `categorize ingest --from <hard-window-start> --to <hard-window-end>` (or `--tag <old-trip>`).
2. `categorize cluster`, then `categorize serve`.
3. Hand-label the slice: the true trip boundaries + which surrounding photos are non-trip.
4. Compare candidate clusters against the labels.
## Acceptance bar (record actual numbers per run)
| Metric | Definition | Target |
|--------|------------|--------|
| Trip-boundary precision | proposed boundaries that are real | ≥ 0.8 |
| Trip-boundary recall | real boundaries proposed | ≥ 0.8 |
| Coverage-flag recall | in-window missing-tag assets surfaced | ≥ 0.9 |
| Over-split rate | extra clusters per real trip | ≤ 0.5 |
| False-cluster rate | clusters that are pure noise | ≤ 0.1 |
If unmet, tune `--gap-factor` / seed-span / confidence thresholds — do **not** widen the
backlog. See the F4 open question (everyday cluster blow-up): observe the real cluster
count on a representative subset first, then choose the surfacing/collapsing strategy.
```
- [ ] **Step 7: Reinstall (package-data changed) and run the FULL suite**
Run:
```bash
.venv/bin/pip install -e ./shared -e ./apps/trip-cluster
.venv/bin/python -m pytest
.venv/bin/python -m pytest apps/trip-cluster/tests/ui
```
Expected: all green across `shared/tests` + `apps/trip-cluster/tests` (+ UI).
- [ ] **Step 8: Commit**
```bash
git add apps/trip-cluster/Dockerfile apps/trip-cluster/pyproject.toml docker-compose.yml \
README.md CLAUDE.md docs/M1-validation-gate.md
git commit -m "docs+docker: README, CLAUDE.md, Dockerfile/compose (8084), M1 validation gate"
```
---
## Self-Review (against the M1 spec)
**Spec coverage**
- Monorepo layout + editable shared packages → Task 1. Per-app container/port 8084 → Tasks 8, 17.
- `shared/immich` (search w/ GPS·city·country·type·tags·updatedAt, thumbnail, upsert/tag_assets, `_pipeline/` namespace) → Tasks 23.
- `shared/core` (SQLite store, all seven tables, domain models, attention sort, neighbors, split/merge, writeback log) → Tasks 46.
- `shared/ai` → correctly **absent** (deferred to M3); `ANTHROPIC_API_KEY` optional → Tasks 1, 8.
- `shared/ui` (base.html, macros, grid+lightbox JS) → Task 7.
- Ingest: scopeable, incremental, `_pipeline/processed` read-back → Task 9.
- Cluster: existing-tag seeds, density-adaptive gap (sparse-old + dense-recent fixtures), location anchors, coverage flags + outliers, confidence/kind_guess → Tasks 1012.
- Review: master/detail, attention sort, approve/non-trip/split/merge/skip, coverage include/exclude, bulk approve-high-confidence, empty states, keyboard-first → Tasks 13, 15, 16.
- Write-back: idempotent, trip/non-trip/processed, partial-failure per-cluster, confirmation, reconcile via upsert → Task 14.
- Tag conventions (content never namespaced; `_pipeline/`) → Tasks 3, 14.
- Testing: pytest + pytest-httpserver + Playwright, TDD → every task.
- Validation gate (hard sample, quantified bar) → Task 17.
- pgvector spike → **out of this plan by decision** (tracked separately; M1.5 dependency).
**Type/name consistency** — verified: `Asset`, `Cluster`, `ClusterMember` field names match across `models.py`, `store.py`, `clustering.py`, `coverage.py`, `writeback.py`; `cluster_members()` returns `list[(Asset, ClusterMember)]` everywhere it's consumed; `apply_cluster` result shape (`succeeded`/`failed`) matches the JS in Task 15.
**Note for the implementer:** the existing-tag **seed derivation** (a non-pipeline tag seeds a cluster only when its tagged assets span ≤ 60 days) is a heuristic to separate time-bounded trip tags from year-spanning people/location tags. If a real library has long multi-leg trips tagged as one, raise `seed_max_span_days` (Task 12) — surfaced as a tunable, not hard-coded policy.
---
## Deferred / Open Questions
### From 2026-06-27 doc review
- **F15 — Everyday-cluster blow-up (links to F4).** Cluster-by-cluster review assumes a few dozen clusters, but timestamp-clustering the everyday/non-trip backlog could mint hundredsthousands of low-confidence clusters that "approve all high-confidence" won't relieve. Before committing to the review surfacing/collapsing strategy, observe the real cluster count on a representative subset.
- **F18 — Image `alt` text for accessibility.** Photos in `_detail.html` render with `alt=""`, so screen-reader users get nothing per image in a tool whose entire purpose is reviewing those images. Consider a composite `alt` from `place_city` / `place_country` / `taken_at`.
### Reviewed and accepted as-is (2026-06-27 doc review)
Findings the review raised and we deliberately declined — recorded so they are not re-litigated as oversights:
- **F1 / F5 / F6 — clustering validated late · no real-API smoke test · gate may not measure the algorithm.** Accepted: the clustering algorithm is treated as best-effort and tuned by hand against real data; the validation gate stays in Task 17 as written.
- **F2 — unauthenticated mutation + thumbnail endpoints on `0.0.0.0`.** Accepted: single-user, self-hosted localhost trust model.
- **F12 — `IMMICH_URL` defaults to `http://`.** Accepted: same localhost trust model.