Compare commits
5 Commits
us/US07-07
...
us/US08-03
| Author | SHA1 | Date | |
|---|---|---|---|
| 90df4be7f8 | |||
| 888d859e93 | |||
| eea50802af | |||
| ded83178bd | |||
| 74f4b1b640 |
22
.dockerignore
Normal file
22
.dockerignore
Normal file
@@ -0,0 +1,22 @@
|
||||
# Deny-by-default build context (US08-02): the image must contain no secrets, no
|
||||
# photos, no database, no logs, and no .git. An allow list is the only version of this
|
||||
# rule that stays true when a new file appears in the working copy.
|
||||
*
|
||||
|
||||
!pyproject.toml
|
||||
!alembic.ini
|
||||
!README.md
|
||||
!photo_pipeline
|
||||
!migrations
|
||||
!frontend
|
||||
!docker
|
||||
|
||||
# Nothing generated, even under an allowed directory.
|
||||
**/__pycache__
|
||||
**/*.py[cod]
|
||||
**/.DS_Store
|
||||
**/*.env
|
||||
**/*.log
|
||||
**/*.db
|
||||
**/*.db-*
|
||||
**/*.sqlite*
|
||||
85
.env.example
Normal file
85
.env.example
Normal file
@@ -0,0 +1,85 @@
|
||||
# Every setting the application and its composition read, with no values in it.
|
||||
# Copy to `.env`, fill in what you need, and keep that copy out of git (it is
|
||||
# gitignored, and the work-item safety checks refuse to stage it).
|
||||
#
|
||||
# cp .env.example .env
|
||||
#
|
||||
# Empty means "use the default noted beside it". Anything already exported in the
|
||||
# shell wins over this file, both for the app and for `docker compose`.
|
||||
|
||||
# ── the composition (host side; read by docker-compose.yml only) ─────────────
|
||||
# The photo library on this host. Bind-mounted at PHOTO_PIPELINE_LIBRARY_ROOTS.
|
||||
PHOTO_PIPELINE_LIBRARY_HOST_PATH=
|
||||
# Which image to run. Default: photo-pipeline:dev (what `up --build` builds).
|
||||
PHOTO_PIPELINE_IMAGE=
|
||||
# Must be the owner of the library above: what the containers rename and rewrite
|
||||
# keeps this ownership. Default: 1000 / 1000.
|
||||
PHOTO_PIPELINE_UID=
|
||||
PHOTO_PIPELINE_GID=
|
||||
# Host address the API port is published on. Default: 127.0.0.1. Anything else
|
||||
# exposes the app beyond this machine — then ALLOWED_HOSTS and ACCESS_SECRET below
|
||||
# are what stand in for the loopback boundary (US08-01).
|
||||
PHOTO_PIPELINE_PUBLISH_ADDRESS=
|
||||
# This file's own path, if it is not ./.env. Default: .env.
|
||||
PHOTO_PIPELINE_ENV_FILE=
|
||||
|
||||
# ── library and data ────────────────────────────────────────────────────────
|
||||
# os.pathsep-separated. In a container these are the *container-side* mount paths,
|
||||
# and `serve`/`worker` refuse to start when they are not mounted.
|
||||
PHOTO_PIPELINE_LIBRARY_ROOTS=
|
||||
# Database, WAL, thumbnail cache, journals, backups. Default: data (the container
|
||||
# sets /data, which is the persistent volume; a network mount is unsupported).
|
||||
PHOTO_PIPELINE_DATA_DIR=
|
||||
# Database file, if it should not live in the data directory. Default:
|
||||
# <data dir>/photo_pipeline.db.
|
||||
PHOTO_PIPELINE_DB_PATH=
|
||||
|
||||
# ── serving ─────────────────────────────────────────────────────────────────
|
||||
# Bind address. Default: 127.0.0.1. The composition sets 0.0.0.0 inside the
|
||||
# container and publishes to loopback on the host instead.
|
||||
PHOTO_PIPELINE_HOST=
|
||||
# Under the composition this is the *published* host port; the container serves
|
||||
# 8000. Default: 8000.
|
||||
PHOTO_PIPELINE_PORT=
|
||||
# Comma-separated hostnames the app answers to besides loopback. Empty means
|
||||
# loopback only. Naming one makes the access secret mandatory.
|
||||
PHOTO_PIPELINE_ALLOWED_HOSTS=
|
||||
# Traded for the session cookie at GET /api/v1/session via X-Access-Secret.
|
||||
# Required as soon as the app is reachable from anywhere but loopback. Generate
|
||||
# one with: python -c 'import secrets; print(secrets.token_urlsafe(32))'
|
||||
PHOTO_PIPELINE_ACCESS_SECRET=
|
||||
# Comma-separated peer addresses whose X-Forwarded-Proto/-Host may be believed.
|
||||
# Only the reverse proxy's address belongs here. Default: none.
|
||||
PHOTO_PIPELINE_TRUSTED_PROXIES=
|
||||
# Largest request body accepted, in bytes. Default: 1048576.
|
||||
PHOTO_PIPELINE_MAX_REQUEST_BYTES=
|
||||
|
||||
# ── logging ─────────────────────────────────────────────────────────────────
|
||||
# Default: INFO.
|
||||
PHOTO_PIPELINE_LOG_LEVEL=
|
||||
# json or text. Default: json.
|
||||
PHOTO_PIPELINE_LOG_FORMAT=
|
||||
|
||||
# ── limits ──────────────────────────────────────────────────────────────────
|
||||
# Thumbnail cache quota in bytes. Default: 500000000.
|
||||
PHOTO_PIPELINE_THUMBNAIL_CACHE_QUOTA_BYTES=
|
||||
# Refuse to decode images larger than this many pixels. Default: 100000000.
|
||||
PHOTO_PIPELINE_THUMBNAIL_MAX_PIXELS=
|
||||
# Free space an archive destination must keep beyond the transfer. Default:
|
||||
# 1000000000.
|
||||
PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES=
|
||||
|
||||
# ── safety gate ─────────────────────────────────────────────────────────────
|
||||
# true refuses every mutating request until a read-only dry run of this library has
|
||||
# been produced and approved (US07-07). Default: false. Turn it on before pointing
|
||||
# the app at photos that cannot be replaced.
|
||||
PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL=
|
||||
|
||||
# ── external services ───────────────────────────────────────────────────────
|
||||
# Vision provider key for the analysis stage. Without it, analysis cannot run.
|
||||
PHOTO_PIPELINE_VISION_API_KEY=
|
||||
# Immich server and its API key, for the upload stage.
|
||||
PHOTO_PIPELINE_IMMICH_SERVER_URL=
|
||||
PHOTO_PIPELINE_IMMICH_API_KEY=
|
||||
# Uploader binary. Default: immich-go (on PATH; pinned inside the image).
|
||||
PHOTO_PIPELINE_IMMICH_GO_BINARY=
|
||||
9
.gitignore
vendored
9
.gitignore
vendored
@@ -20,3 +20,12 @@ _IGNORE/
|
||||
|
||||
# Test failure evidence (US07-04)
|
||||
.artifacts/
|
||||
|
||||
# Any dotenv, not only the default name.
|
||||
*.env
|
||||
|
||||
# Local virtualenv for running the app.
|
||||
.venv/
|
||||
|
||||
# setuptools editable-install metadata.
|
||||
*.egg-info/
|
||||
|
||||
115
Dockerfile
Normal file
115
Dockerfile
Normal file
@@ -0,0 +1,115 @@
|
||||
# One image, two roles (US08-02).
|
||||
#
|
||||
# The application is not self-contained Python: it shells out to `exiftool` for every
|
||||
# EXIF checkpoint and to `immich-go` for every upload, and it serves the static
|
||||
# frontend from `frontend/`. All three are installed here at pinned versions, because
|
||||
# an image whose external tools drift is an image whose metadata checkpoints and
|
||||
# upload reports drift with them (concept §15, "External integration risks").
|
||||
#
|
||||
# Everything is pinned:
|
||||
# * the base image by tag *and* digest, so a moved tag cannot change the runtime;
|
||||
# * exiftool by its Debian package version, verified against `exiftool -ver`;
|
||||
# * immich-go by release version and per-architecture SHA-256 of the release asset.
|
||||
# The verified versions become image labels and /etc/photo-pipeline/versions.json,
|
||||
# which `python -m photo_pipeline diagnostics` reports — so a running container can
|
||||
# prove what it contains instead of being trusted about it.
|
||||
#
|
||||
# The project is installed editable on purpose: `photo_pipeline.db` resolves
|
||||
# `alembic.ini` and `migrations/`, and the API resolves `frontend/`, relative to the
|
||||
# repository root. An editable install keeps that one layout instead of scattering the
|
||||
# same files across site-packages and a source tree.
|
||||
|
||||
ARG PYTHON_IMAGE=python:3.12.14-slim-trixie@sha256:2c941e860699f878900b0edc2403613c234d4b32eda3cc9fa7036991a2a63c4a
|
||||
|
||||
# ── the uploader, fetched and verified outside the final layer ────────────────
|
||||
FROM ${PYTHON_IMAGE} AS uploader
|
||||
|
||||
ARG IMMICH_GO_VERSION=0.32.0
|
||||
ARG IMMICH_GO_SHA256_AMD64=6e2ad86bafdadb9466d6515de7cb882726c0aea1a21d51164dff361d7d480a97
|
||||
ARG IMMICH_GO_SHA256_ARM64=2c35d9284baae407ef9540bdac5f488971b0bdc7be758a4d7c05ab270af09fdb
|
||||
|
||||
COPY docker/fetch-immich-go.py /tmp/fetch-immich-go.py
|
||||
RUN python /tmp/fetch-immich-go.py \
|
||||
--version "${IMMICH_GO_VERSION}" \
|
||||
--sha256-amd64 "${IMMICH_GO_SHA256_AMD64}" \
|
||||
--sha256-arm64 "${IMMICH_GO_SHA256_ARM64}" \
|
||||
--into /usr/local/bin \
|
||||
&& /usr/local/bin/immich-go version
|
||||
|
||||
# ── the application ──────────────────────────────────────────────────────────
|
||||
FROM ${PYTHON_IMAGE} AS runtime
|
||||
|
||||
ARG EXIFTOOL_VERSION=13.25+dfsg-1
|
||||
ARG IMMICH_GO_VERSION=0.32.0
|
||||
# The library is mounted from the host, so the container's identity must match the
|
||||
# ownership that library already has: everything this application renames, writes
|
||||
# EXIF into, or archives has to stay owned by the host user afterwards.
|
||||
ARG UID=1000
|
||||
ARG GID=1000
|
||||
|
||||
LABEL org.opencontainers.image.title="photo_pipeline" \
|
||||
org.opencontainers.image.source="https://github.com/domverse/photoanalyzer" \
|
||||
io.photoanalyzer.exiftool.version="${EXIFTOOL_VERSION}" \
|
||||
io.photoanalyzer.immich-go.version="${IMMICH_GO_VERSION}"
|
||||
|
||||
ENV PYTHONUNBUFFERED=1 \
|
||||
PYTHONDONTWRITEBYTECODE=1 \
|
||||
PATH=/opt/venv/bin:$PATH \
|
||||
PHOTO_PIPELINE_DATA_DIR=/data
|
||||
|
||||
RUN set -eu; \
|
||||
apt-get update; \
|
||||
DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends \
|
||||
"libimage-exiftool-perl=${EXIFTOOL_VERSION}"; \
|
||||
rm -rf /var/lib/apt/lists/*
|
||||
|
||||
COPY --from=uploader /usr/local/bin/immich-go /usr/local/bin/immich-go
|
||||
|
||||
WORKDIR /app
|
||||
COPY pyproject.toml alembic.ini README.md ./
|
||||
COPY photo_pipeline ./photo_pipeline
|
||||
COPY migrations ./migrations
|
||||
COPY frontend ./frontend
|
||||
COPY docker/entrypoint.sh docker/healthcheck.sh /usr/local/bin/
|
||||
|
||||
# Runtime dependencies only: the `test` extra (pytest, playwright) and the `vision`
|
||||
# extra stay out, and pip's build isolation leaves no build tooling behind.
|
||||
RUN set -eu; \
|
||||
python -m venv /opt/venv; \
|
||||
/opt/venv/bin/pip install --no-cache-dir -e .
|
||||
|
||||
# What is installed must be what was pinned, or the labels and the version record
|
||||
# would be a claim rather than a fact.
|
||||
RUN set -eu; \
|
||||
mkdir -p /etc/photo-pipeline; \
|
||||
exiftool_version="$(exiftool -ver)"; \
|
||||
immich_go_version="$(immich-go version | head -n 1 | tr -d '\r')"; \
|
||||
expected_exiftool="$(printf '%s' "${EXIFTOOL_VERSION}" | cut -d+ -f1 | cut -d- -f1)"; \
|
||||
[ "${exiftool_version}" = "${expected_exiftool}" ] \
|
||||
|| { echo "exiftool ${exiftool_version} is not the pinned ${expected_exiftool}" >&2; exit 1; }; \
|
||||
case "${immich_go_version}" in \
|
||||
*"${IMMICH_GO_VERSION}"*) ;; \
|
||||
*) echo "immich-go '${immich_go_version}' is not pinned ${IMMICH_GO_VERSION}" >&2; exit 1 ;; \
|
||||
esac; \
|
||||
printf '{\n "exiftool": "%s",\n "immich-go": "%s"\n}\n' \
|
||||
"${exiftool_version}" "${IMMICH_GO_VERSION}" > /etc/photo-pipeline/versions.json
|
||||
|
||||
# Non-root, with the host library's ownership. /data is the persistent volume; the
|
||||
# photo library itself is mounted by the deployment (US08-03), never baked in.
|
||||
RUN set -eu; \
|
||||
groupadd --gid "${GID}" --non-unique app; \
|
||||
useradd --uid "${UID}" --gid "${GID}" --non-unique --no-create-home --home-dir /app app; \
|
||||
mkdir -p /data; \
|
||||
chown "${UID}:${GID}" /data
|
||||
USER ${UID}:${GID}
|
||||
|
||||
EXPOSE 8000
|
||||
|
||||
# Readiness, not liveness: an unmigrated or misconfigured database answers
|
||||
# /api/v1/health/ready with 503, and a container that cannot serve must not be
|
||||
# reported healthy. The worker role has no endpoint, so its check is a no-op here.
|
||||
HEALTHCHECK --interval=30s --timeout=10s --start-period=30s --retries=3 \
|
||||
CMD ["/usr/local/bin/healthcheck.sh"]
|
||||
|
||||
ENTRYPOINT ["/usr/local/bin/entrypoint.sh"]
|
||||
CMD ["serve"]
|
||||
209
README.md
209
README.md
@@ -7,16 +7,49 @@ archive workflow. Planning lives in `INTEGRATED_PIPELINE_CONCEPT.md` and
|
||||
## Application (`photo_pipeline`)
|
||||
|
||||
The target application lives in `photo_pipeline/` (FastAPI + SQLAlchemy + Alembic).
|
||||
Run it with:
|
||||
Install it into a virtualenv once:
|
||||
|
||||
```bash
|
||||
python -m photo_pipeline migrate # apply database migrations
|
||||
python -m photo_pipeline serve # start the API + static review UI (127.0.0.1:8000)
|
||||
python3.12 -m venv .venv
|
||||
.venv/bin/pip install -e ".[vision]" # drop [vision] for a review-only install
|
||||
```
|
||||
|
||||
Then run the two processes:
|
||||
|
||||
```bash
|
||||
.venv/bin/python -m photo_pipeline migrate # apply database migrations
|
||||
.venv/bin/python -m photo_pipeline serve # API + review UI at 127.0.0.1:8000/app/
|
||||
.venv/bin/python -m photo_pipeline worker # second terminal: runs the jobs
|
||||
```
|
||||
|
||||
The server enqueues work and serves the UI; nothing actually scans, scores,
|
||||
analyses, uploads, or archives without a worker. `work_item/scripts/python` is the
|
||||
*helper's* launcher — it prefers Conda base and falls back to a bare system
|
||||
interpreter, so it is not how the application is run.
|
||||
|
||||
Configuration comes from `PHOTO_PIPELINE_*` environment variables (see
|
||||
`photo_pipeline/config.py`); secrets are referenced, never logged.
|
||||
|
||||
### Configuration file
|
||||
|
||||
`.env` in the working directory is read at startup, or any path named by
|
||||
`PHOTO_PIPELINE_ENV_FILE`. It is parsed, never executed: `KEY=value` lines,
|
||||
`#` comments, optional quotes — no interpolation and no `export`. **Anything already
|
||||
exported wins**, so the file is the standing configuration and the shell is the
|
||||
override for one run.
|
||||
|
||||
The archived CLI's variable names still work, so an existing `photo_analyzer.env`
|
||||
can be used as-is:
|
||||
|
||||
| in the file | applied as |
|
||||
|---|---|
|
||||
| `LLM_API_KEY` / `GEMINI_API_KEY` | `OPENAI_API_KEY` |
|
||||
| `LLM_BASE_URL` | `OPENAI_BASE_URL` |
|
||||
| `LIBRARY` | `PHOTO_PIPELINE_LIBRARY_ROOTS` |
|
||||
|
||||
`.env` and `*.env` are gitignored and denied by the work-item safety checks: the
|
||||
file holds a real key and must never be committed.
|
||||
|
||||
### API access (US07-02)
|
||||
|
||||
The app listens on loopback, so its attacker is another page in the same browser.
|
||||
@@ -36,6 +69,131 @@ not a loopback name (DNS rebinding), when `Origin` is any other origin, when
|
||||
thumbnail), or when the body exceeds `PHOTO_PIPELINE_MAX_REQUEST_BYTES`. There is no
|
||||
CORS middleware at all, so no other origin can read a response.
|
||||
|
||||
### Reaching it through a hostname or proxy (US08-01)
|
||||
|
||||
| variable | meaning |
|
||||
|---|---|
|
||||
| `PHOTO_PIPELINE_ALLOWED_HOSTS` | comma-separated extra names the app answers to; empty means loopback only |
|
||||
| `PHOTO_PIPELINE_ACCESS_SECRET` | traded for the session cookie at `GET /api/v1/session` via `X-Access-Secret` |
|
||||
| `PHOTO_PIPELINE_TRUSTED_PROXIES` | comma-separated peer addresses whose `X-Forwarded-Proto`/`X-Forwarded-Host` are believed |
|
||||
|
||||
Being reachable *was* the authentication: whoever could open `127.0.0.1:8000` owned
|
||||
the library. So naming any non-loopback host — or binding to one, `0.0.0.0` included
|
||||
— makes the access secret mandatory, and `serve` refuses to start without it rather
|
||||
than publishing the library. Loopback-only deployments need no secret and behave
|
||||
exactly as before.
|
||||
|
||||
```bash
|
||||
curl -sc /tmp/pp.jar -H "X-Access-Secret: $PHOTO_PIPELINE_ACCESS_SECRET" \
|
||||
https://photos.example.com/api/v1/session
|
||||
```
|
||||
|
||||
The browser asks for the secret once per tab and keeps it in `sessionStorage`.
|
||||
Wrong secrets are rate-limited (5 per minute) and logged with the caller's address
|
||||
only. `Host` and `Origin` are judged against the configured names; the *external*
|
||||
scheme and host come from the forwarded headers only when the request arrived from a
|
||||
`PHOTO_PIPELINE_TRUSTED_PROXIES` address, so a client cannot declare its own origin,
|
||||
and the session cookie is marked `Secure` when that external scheme is HTTPS. Health
|
||||
endpoints stay reachable without the secret so an orchestrator can restart the
|
||||
container; nothing else does.
|
||||
|
||||
## Container image (US08-02)
|
||||
|
||||
One image runs either role. It is built from a clean checkout with no arguments:
|
||||
|
||||
```bash
|
||||
docker build -t photo-pipeline:dev .
|
||||
```
|
||||
|
||||
Everything external is pinned, and the build fails rather than drifting: the Python
|
||||
base image by tag *and* digest, `exiftool` by its Debian package version (verified
|
||||
against `exiftool -ver`), and `immich-go` by release version and per-architecture
|
||||
SHA-256 of the release asset. The verified versions become image labels and
|
||||
`/etc/photo-pipeline/versions.json`, which `diagnostics` reports as `tools[].pinned`
|
||||
beside the version actually installed — so a replaced binary shows up as a
|
||||
`tool_version_drift` warning instead of as a misparsed upload report.
|
||||
|
||||
| build argument | default | why change it |
|
||||
|---|---|---|
|
||||
| `UID` / `GID` | `1000` | must match the owner of the mounted photo library |
|
||||
| `PYTHON_IMAGE` | pinned digest | upgrading the base image |
|
||||
| `EXIFTOOL_VERSION` | Debian package version | upgrading exiftool |
|
||||
| `IMMICH_GO_VERSION` + `IMMICH_GO_SHA256_AMD64`/`_ARM64` | pinned release | upgrading the uploader (take the digests from that release's `checksums.txt`) |
|
||||
|
||||
The first argument is the role, and every other management command still works:
|
||||
|
||||
```bash
|
||||
docker run --rm -v /srv/photos:/srv/photos -v pp-data:/data \
|
||||
-e PHOTO_PIPELINE_LIBRARY_ROOTS=/srv/photos photo-pipeline:dev migrate
|
||||
|
||||
docker run -d -p 127.0.0.1:8000:8000 -v /srv/photos:/srv/photos -v pp-data:/data \
|
||||
-e PHOTO_PIPELINE_HOST=0.0.0.0 -e PHOTO_PIPELINE_ACCESS_SECRET=... \
|
||||
-e PHOTO_PIPELINE_LIBRARY_ROOTS=/srv/photos photo-pipeline:dev serve
|
||||
|
||||
docker run -d -v /srv/photos:/srv/photos -v pp-data:/data \
|
||||
-e PHOTO_PIPELINE_LIBRARY_ROOTS=/srv/photos photo-pipeline:dev worker
|
||||
```
|
||||
|
||||
One role per container: `serve` and `worker` each take the library process lock for
|
||||
their role (US07-05), so no supervisor starts both. The container refuses to run as
|
||||
UID 0 — files it renames or writes must keep the ownership the host library expects —
|
||||
and `/data` is the persistent volume holding the database, journals, backups, and
|
||||
thumbnail cache. Binding to `0.0.0.0` makes the access secret mandatory
|
||||
([above](#reaching-it-through-a-hostname-or-proxy-us08-01)); `serve` refuses to start
|
||||
without it. The declared `HEALTHCHECK` polls `/api/v1/health/ready`, so a container
|
||||
whose database is unmigrated or misconfigured is never reported healthy.
|
||||
|
||||
## Composed runtime (US08-03)
|
||||
|
||||
`docker-compose.yml` is the deployment: one `serve` container, one `worker`
|
||||
container, one bind-mounted library, one data volume, and a one-shot `migrate` that
|
||||
both roles wait for.
|
||||
|
||||
```bash
|
||||
cp .env.example .env && $EDITOR .env # nothing has a value in it; fill in yours
|
||||
docker compose up -d --build
|
||||
```
|
||||
|
||||
`.env.example` lists every `PHOTO_PIPELINE_*` variable with its default in a comment
|
||||
and no values at all. The three the composition cannot start without are
|
||||
`PHOTO_PIPELINE_LIBRARY_HOST_PATH` (the library on this host),
|
||||
`PHOTO_PIPELINE_LIBRARY_ROOTS` (where it is mounted *inside* the container), and
|
||||
`PHOTO_PIPELINE_ACCESS_SECRET` — publishing the port means the app is reachable from
|
||||
outside the container, so US08-01 makes the secret mandatory. Set
|
||||
`PHOTO_PIPELINE_UID`/`GID` to the owner of the library: what the containers rename
|
||||
and rewrite keeps that ownership.
|
||||
|
||||
| invariant | how the composition keeps it |
|
||||
|---|---|
|
||||
| one writer | `serve` and `worker` take their role's library lock (US07-05) in the shared `/data` volume, so `--scale worker=2` is refused by the lock, not by convention |
|
||||
| migrations first | `migrate` runs the backup-then-migrate path and must exit 0 before `api` and `worker` start; a failed upgrade leaves the previous database and its pre-migration backup intact |
|
||||
| container paths | `PHOTO_PIPELINE_LIBRARY_ROOTS` is both the mount target and the configured root; a root that is not mounted makes `serve`/`worker` exit 5 at startup instead of writing into the container's throwaway layer |
|
||||
| loopback by default | the port is published to `127.0.0.1` unless `PHOTO_PIPELINE_PUBLISH_ADDRESS` says otherwise, and exposing it needs the hostname in `PHOTO_PIPELINE_ALLOWED_HOSTS` plus the access secret |
|
||||
| restart safety | both roles are `restart: unless-stopped` with a 30 s stop grace period, and a job interrupted by a restart resumes exactly as it does on a host restart |
|
||||
|
||||
The `data` volume holds the database, its write-ahead log, the thumbnail cache,
|
||||
journals, and backups. **It must stay on a local filesystem** — SQLite in WAL mode
|
||||
needs real local locking, so NFS, SMB, and network volume drivers are unsupported
|
||||
there and corrupt the database rather than slow it down. The library bind mount has
|
||||
no such restriction.
|
||||
|
||||
Operating the deployment is operating the same CLI:
|
||||
|
||||
```bash
|
||||
docker compose run --rm --no-deps api diagnostics
|
||||
docker compose run --rm --no-deps api backup --reason pre-upgrade
|
||||
docker compose run --rm --no-deps api verify-backup /data/backups/<name>
|
||||
docker compose run --rm --no-deps api restore /data/backups/<name> --into /data/restored
|
||||
docker compose logs -f worker
|
||||
```
|
||||
|
||||
`--no-deps` keeps a one-off command from starting a second stack; `api` is only the
|
||||
service the command borrows the image and mounts from. `restore` is deliberately not
|
||||
an API call: it replaces the state of an installation and belongs to a stopped one,
|
||||
so stop `api` and `worker` first and restart them against the restored directory.
|
||||
|
||||
Publishing and deploying the image is US08-04.
|
||||
|
||||
## Testing
|
||||
|
||||
One offline command runs the whole suite (unit, integration, and browser
|
||||
@@ -281,6 +439,51 @@ file in the temporary library are copied to `.artifacts/<test id>/` before pytes
|
||||
deletes the directory. Point `PHOTO_PIPELINE_TEST_ARTIFACTS` elsewhere to collect
|
||||
them from CI.
|
||||
|
||||
## Release gate (US07-07)
|
||||
|
||||
One command runs every suite in an isolated stack and keeps the evidence:
|
||||
|
||||
```bash
|
||||
work_item/scripts/python -m photo_pipeline release-gate --output data/release/$(date -u +%Y%m%dT%H%M%SZ)
|
||||
```
|
||||
|
||||
It fails — and exits non-zero — when any stage fails, when a suite skips a test for
|
||||
a reason that is not a documented environment limit (`exiftool not installed`,
|
||||
`root ignores directory permissions`), or when the story matrix has a hole. The
|
||||
evidence directory holds `release-report.json` (revision, per-stage result, timings,
|
||||
summaries), `logs/<stage>.log`, and `CHECKSUMS.sha256` over both.
|
||||
|
||||
**The story matrix** lives in `tests/story_traceability.json`: every story under
|
||||
`delivery_backlog/stories/` is either mapped to test files that exist, or listed in
|
||||
`planned` as an accepted but unimplemented story. A story that is neither, or a
|
||||
mapping to a file that has been deleted, fails the gate.
|
||||
|
||||
**The journey** (`tests/e2e/test_release_journey.py`) takes one fresh library through
|
||||
discovery, duplicate review, safety, analysis, EXIF verification, album proposal,
|
||||
guarded rename, rescan, upload with server-side verification, archive, offline
|
||||
deduplication, and restore — over HTTP against real server and worker processes,
|
||||
with a full restart in the middle and at the end.
|
||||
|
||||
### Real-library dry run and approval
|
||||
|
||||
Before the application is pointed at photos that cannot be replaced:
|
||||
|
||||
```bash
|
||||
work_item/scripts/python -m photo_pipeline dry-run --output dry-run.json
|
||||
work_item/scripts/python -m photo_pipeline approve-dry-run dry-run.json --approver "$(whoami)"
|
||||
```
|
||||
|
||||
The dry run is strictly read-only: it opens no file for writing, writes no database
|
||||
row, and reports what it found — file counts by extension, folders, bytes, unreadable
|
||||
files, excluded directories, and a reconciliation against what the database already
|
||||
knows (already registered, new, recorded but absent). Set
|
||||
`PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL=1` and **every mutating API request is
|
||||
refused with `403 dry_run_not_approved`** until a report for exactly those library
|
||||
roots has been approved. Reading stays open — you have to be able to see what was
|
||||
found in order to approve it — and so does taking a backup. Change the library roots
|
||||
and the approval no longer applies: it approves that reconciliation, not the idea of
|
||||
mutating.
|
||||
|
||||
## Performance budgets (US07-06)
|
||||
|
||||
Budgets are measured, not asserted in prose. `python -m photo_pipeline benchmark`
|
||||
|
||||
105
docker-compose.yml
Normal file
105
docker-compose.yml
Normal file
@@ -0,0 +1,105 @@
|
||||
# The deployed runtime (US08-03): one API, one worker, one library, one volume.
|
||||
#
|
||||
# The same image (US08-02) runs both roles, so what is composed here is process
|
||||
# topology, not a second application. Three invariants shape it:
|
||||
#
|
||||
# * one writer — `serve` and `worker` each take the library process lock for their
|
||||
# role (US07-05), and both containers mount the *same* data volume, which is what
|
||||
# makes the lock file visible to both. Scaling `worker` past 1 is refused by that
|
||||
# lock rather than by anyone remembering not to;
|
||||
# * one local filesystem — the database, its write-ahead log, the thumbnail cache,
|
||||
# and the backups live in the `data` volume, and SQLite in WAL mode requires real
|
||||
# local-filesystem locking. A network mount (NFS, SMB, or a cloud volume driver)
|
||||
# is unsupported for it; that is a corrupted database, not a slow one;
|
||||
# * one library path vocabulary — `PHOTO_PIPELINE_LIBRARY_ROOTS` names the path
|
||||
# *inside* the container, which is also the bind mount's target below. A root
|
||||
# that is not mounted there makes `serve` and `worker` refuse at startup instead
|
||||
# of writing into the container's throwaway layer.
|
||||
#
|
||||
# Configuration and secrets come from the environment only: copy `.env.example` to
|
||||
# `.env` and fill it in. Nothing is baked into the image and nothing with a value in
|
||||
# it is committed.
|
||||
#
|
||||
# cp .env.example .env && $EDITOR .env
|
||||
# docker compose up -d --build
|
||||
#
|
||||
# Operating it is operating the same CLI — `docker compose run --rm --no-deps api
|
||||
# <command>` — see README, "Composed runtime".
|
||||
|
||||
name: photo-pipeline
|
||||
|
||||
x-runtime: &runtime
|
||||
image: ${PHOTO_PIPELINE_IMAGE:-photo-pipeline:dev}
|
||||
build:
|
||||
context: .
|
||||
args:
|
||||
# Everything the app renames or rewrites has to stay owned by the host user
|
||||
# the library already belongs to.
|
||||
UID: ${PHOTO_PIPELINE_UID:-1000}
|
||||
GID: ${PHOTO_PIPELINE_GID:-1000}
|
||||
user: "${PHOTO_PIPELINE_UID:-1000}:${PHOTO_PIPELINE_GID:-1000}"
|
||||
env_file:
|
||||
- ${PHOTO_PIPELINE_ENV_FILE:-.env}
|
||||
volumes:
|
||||
- data:/data
|
||||
- "${PHOTO_PIPELINE_LIBRARY_HOST_PATH:?set PHOTO_PIPELINE_LIBRARY_HOST_PATH to the photo library on this host}:${PHOTO_PIPELINE_LIBRARY_ROOTS:?set PHOTO_PIPELINE_LIBRARY_ROOTS to the container-side library path}"
|
||||
# Jobs check for cancellation between items and leave a resumable record; a
|
||||
# too-short grace period turns an orderly stop into a recovery on next start.
|
||||
stop_grace_period: 30s
|
||||
|
||||
x-environment: &environment
|
||||
# Set here rather than left to the file: these two are what the composition itself
|
||||
# promises, and an `.env` that disagreed would move the database off the volume or
|
||||
# the library off its mount.
|
||||
PHOTO_PIPELINE_DATA_DIR: /data
|
||||
PHOTO_PIPELINE_LIBRARY_ROOTS: ${PHOTO_PIPELINE_LIBRARY_ROOTS}
|
||||
|
||||
services:
|
||||
# Migrations run to completion before either role accepts work, through the same
|
||||
# backup-then-migrate path the roles use (US07-05): a pending upgrade is snapshotted
|
||||
# first, and a failed one exits non-zero with the backup named — so `api` and
|
||||
# `worker` never start, and the previous database is left intact and restorable.
|
||||
migrate:
|
||||
<<: *runtime
|
||||
command: ["migrate"]
|
||||
environment: *environment
|
||||
restart: "no"
|
||||
|
||||
api:
|
||||
<<: *runtime
|
||||
command: ["serve"]
|
||||
environment:
|
||||
<<: *environment
|
||||
# Published to host loopback below. Inside the container the server must bind
|
||||
# the container's own interface for that publish to reach it, which is exactly
|
||||
# what makes the access secret mandatory (US08-01) — `serve` refuses to start
|
||||
# without one. Exposing the port beyond loopback additionally needs
|
||||
# PHOTO_PIPELINE_ALLOWED_HOSTS to name the hostname it is reached under.
|
||||
PHOTO_PIPELINE_HOST: 0.0.0.0
|
||||
PHOTO_PIPELINE_PORT: 8000
|
||||
ports:
|
||||
# Host side only: PHOTO_PIPELINE_PORT in `.env` moves the *published* port, and
|
||||
# the container always serves 8000, which is what the image's health check probes.
|
||||
- "${PHOTO_PIPELINE_PUBLISH_ADDRESS:-127.0.0.1}:${PHOTO_PIPELINE_PORT:-8000}:8000"
|
||||
depends_on:
|
||||
migrate:
|
||||
condition: service_completed_successfully
|
||||
restart: unless-stopped
|
||||
|
||||
# One worker. A second one is refused by the library lock in the shared data
|
||||
# volume, which is the point: `docker compose up --scale worker=2` fails loudly
|
||||
# instead of running two writers against one library.
|
||||
worker:
|
||||
<<: *runtime
|
||||
command: ["worker", "--id", "worker-1"]
|
||||
environment: *environment
|
||||
depends_on:
|
||||
migrate:
|
||||
condition: service_completed_successfully
|
||||
restart: unless-stopped
|
||||
|
||||
volumes:
|
||||
# Local driver on purpose: the database, WAL, thumbnail cache, and backups need a
|
||||
# real local filesystem. Do not point this at NFS, SMB, or a network volume driver.
|
||||
data:
|
||||
driver: local
|
||||
22
docker/entrypoint.sh
Executable file
22
docker/entrypoint.sh
Executable file
@@ -0,0 +1,22 @@
|
||||
#!/bin/sh
|
||||
# One entrypoint, one role per container (US08-02).
|
||||
#
|
||||
# The first argument is the management command the image runs — `serve` and `worker`
|
||||
# are the two roles, and every other `python -m photo_pipeline` command (migrate,
|
||||
# diagnostics, backup, restore, dry-run) is passed through unchanged so operating the
|
||||
# container is operating the same CLI. No supervisor: two roles in one container would
|
||||
# share a process lock they are each meant to hold alone (US07-05).
|
||||
set -eu
|
||||
|
||||
if [ "$(id -u)" = "0" ]; then
|
||||
echo "refusing to run as root: start this image with a non-root UID/GID so files" \
|
||||
"it renames or writes keep the ownership the mounted library expects" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
role="${1:-serve}"
|
||||
# The health check has to know which role it is checking, and only the API has an
|
||||
# endpoint to check. /tmp is writable for the unprivileged user; /run may not be.
|
||||
printf '%s' "${role}" > "${PHOTO_PIPELINE_ROLE_FILE:-/tmp/photo-pipeline-role}" 2>/dev/null || true
|
||||
|
||||
exec python -m photo_pipeline "$@"
|
||||
72
docker/fetch-immich-go.py
Normal file
72
docker/fetch-immich-go.py
Normal file
@@ -0,0 +1,72 @@
|
||||
"""Download one pinned immich-go release and verify it before unpacking (US08-02).
|
||||
|
||||
Run at image build time by the `uploader` stage, with the interpreter that is already
|
||||
in the base image: no curl, no wget, and no download tooling in the layer that ships.
|
||||
The checksum is not advisory — a release asset that does not match the pinned digest
|
||||
is a failed build, not a warning, because the uploader's flags and report format are
|
||||
what the upload parser is written against (concept §15).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import hashlib
|
||||
import platform
|
||||
import tarfile
|
||||
import tempfile
|
||||
import urllib.request
|
||||
from pathlib import Path
|
||||
|
||||
RELEASE_URL = "https://github.com/simulot/immich-go/releases/download/v{version}/{asset}"
|
||||
# Debian/BuildKit architecture as the interpreter sees it → release asset name.
|
||||
ASSETS = {
|
||||
"x86_64": ("immich-go_Linux_x86_64.tar.gz", "amd64"),
|
||||
"amd64": ("immich-go_Linux_x86_64.tar.gz", "amd64"),
|
||||
"aarch64": ("immich-go_Linux_arm64.tar.gz", "arm64"),
|
||||
"arm64": ("immich-go_Linux_arm64.tar.gz", "arm64"),
|
||||
}
|
||||
TIMEOUT_SECONDS = 300
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--version", required=True, help="immich-go release, without the v")
|
||||
parser.add_argument("--sha256-amd64", required=True)
|
||||
parser.add_argument("--sha256-arm64", required=True)
|
||||
parser.add_argument("--into", default="/usr/local/bin")
|
||||
args = parser.parse_args()
|
||||
|
||||
machine = platform.machine().lower()
|
||||
if machine not in ASSETS:
|
||||
raise SystemExit(f"unsupported architecture: {machine}")
|
||||
asset, arch = ASSETS[machine]
|
||||
expected = {"amd64": args.sha256_amd64, "arm64": args.sha256_arm64}[arch]
|
||||
url = RELEASE_URL.format(version=args.version, asset=asset)
|
||||
|
||||
with urllib.request.urlopen(url, timeout=TIMEOUT_SECONDS) as response: # noqa: S310
|
||||
payload = response.read()
|
||||
digest = hashlib.sha256(payload).hexdigest()
|
||||
if digest != expected:
|
||||
raise SystemExit(f"checksum mismatch for {url}: {digest} != {expected}")
|
||||
|
||||
target = Path(args.into)
|
||||
target.mkdir(parents=True, exist_ok=True)
|
||||
with tempfile.TemporaryDirectory() as work:
|
||||
archive = Path(work) / asset
|
||||
archive.write_bytes(payload)
|
||||
with tarfile.open(archive) as tar:
|
||||
member = tar.getmember("immich-go")
|
||||
# Extract exactly the one file this pin is about, by name, so nothing
|
||||
# else in the archive can decide where it lands.
|
||||
extracted = tar.extractfile(member)
|
||||
if extracted is None:
|
||||
raise SystemExit("release archive contains no immich-go binary")
|
||||
binary = target / "immich-go"
|
||||
binary.write_bytes(extracted.read())
|
||||
binary.chmod(0o755)
|
||||
print(f"immich-go {args.version} ({arch}) verified {digest}")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
34
docker/healthcheck.sh
Executable file
34
docker/healthcheck.sh
Executable file
@@ -0,0 +1,34 @@
|
||||
#!/bin/sh
|
||||
# Container health for the `serve` role: readiness, not liveness (US08-02).
|
||||
#
|
||||
# /api/v1/health/ready is 503 until the database is reachable, migrated, and in WAL
|
||||
# mode with foreign keys on, so an unmigrated or misconfigured container never reports
|
||||
# healthy. Health endpoints need no session and no access secret, which is what lets an
|
||||
# orchestrator restart a container it holds no credentials for (US08-01).
|
||||
set -eu
|
||||
|
||||
role="$(cat "${PHOTO_PIPELINE_ROLE_FILE:-/tmp/photo-pipeline-role}" 2>/dev/null || echo unknown)"
|
||||
if [ "${role}" != "serve" ]; then
|
||||
# ponytail: the worker has no endpoint to probe; its liveness is its lease and job
|
||||
# heartbeat in the database. Add a `worker --health` command if a restart policy
|
||||
# ever needs to act on it.
|
||||
exit 0
|
||||
fi
|
||||
|
||||
port="${PHOTO_PIPELINE_PORT:-8000}"
|
||||
exec python - "${port}" <<'PY'
|
||||
import sys
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
url = f"http://127.0.0.1:{sys.argv[1]}/api/v1/health/ready"
|
||||
try:
|
||||
with urllib.request.urlopen(url, timeout=5) as response: # noqa: S310 — loopback
|
||||
sys.exit(0 if response.status == 200 else 1)
|
||||
except urllib.error.HTTPError as error:
|
||||
print(f"not ready: HTTP {error.code}", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
except OSError as error:
|
||||
print(f"not ready: {error}", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
PY
|
||||
@@ -7,9 +7,28 @@ export const BASE = "/api/v1";
|
||||
// what makes it proof that the caller is this app and not another page.
|
||||
let csrfToken = null;
|
||||
|
||||
// A deployment reachable through a proxy trades an operator secret for that cookie.
|
||||
// Kept per tab: sessionStorage dies with the tab, and the secret never enters a URL.
|
||||
const SECRET_KEY = "pp_access_secret";
|
||||
|
||||
async function bootstrap(secret) {
|
||||
return fetch(BASE + "/session", {
|
||||
credentials: "same-origin",
|
||||
headers: secret ? { "X-Access-Secret": secret } : {},
|
||||
});
|
||||
}
|
||||
|
||||
async function session() {
|
||||
if (csrfToken === null) {
|
||||
const response = await fetch(BASE + "/session", { credentials: "same-origin" });
|
||||
let response = await bootstrap(sessionStorage.getItem(SECRET_KEY));
|
||||
if (response.status === 401) {
|
||||
sessionStorage.removeItem(SECRET_KEY);
|
||||
const secret = prompt("Access secret");
|
||||
if (secret) {
|
||||
response = await bootstrap(secret);
|
||||
if (response.ok) sessionStorage.setItem(SECRET_KEY, secret);
|
||||
}
|
||||
}
|
||||
const body = await response.json().catch(() => null);
|
||||
csrfToken = (body && body.csrf_token) || null;
|
||||
}
|
||||
|
||||
@@ -11,8 +11,11 @@ from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Sequence
|
||||
|
||||
from photo_pipeline import path_policy
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.services.app_lock import LegacyProcessActive, LibraryLock, LockHeld
|
||||
from photo_pipeline.services.backup import BackupError, BackupService, migrate_with_backup
|
||||
@@ -63,11 +66,37 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
)
|
||||
bench_cmd.add_argument("--output", help="Write the JSON report here as well as to stdout")
|
||||
|
||||
gate_cmd = commands.add_parser(
|
||||
"release-gate", help="Run every suite in an isolated stack and keep the evidence"
|
||||
)
|
||||
gate_cmd.add_argument("--output", help="Evidence directory (default: data/release/<stamp>)")
|
||||
dry_cmd = commands.add_parser(
|
||||
"dry-run", help="Read-only reconciliation of the configured library (US07-07)"
|
||||
)
|
||||
dry_cmd.add_argument("--output", help="Write the report here as well as to stdout")
|
||||
approve_cmd = commands.add_parser(
|
||||
"approve-dry-run", help="Approve a dry-run report, which is what enables mutation"
|
||||
)
|
||||
approve_cmd.add_argument("report", help="Path to the dry-run report")
|
||||
approve_cmd.add_argument("--approver", required=True, help="Who is accepting this")
|
||||
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
config = Config.from_env()
|
||||
config.database_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# In a container the two roles that touch the library check their boundary before
|
||||
# they take a lock or bind a port: the configured roots must name the mount paths,
|
||||
# and a mismatch is cheaper to refuse than to discover at the first write (US08-03).
|
||||
# Only in a container — on a host an unmounted root is an ordinary Tuesday (an
|
||||
# archive medium that is not plugged in), and refusing to serve would take the
|
||||
# offline half of the library away with it.
|
||||
if args.command in ("serve", "worker") and path_policy.in_container():
|
||||
refusal = path_policy.roots_refusal(config.library_roots, require_mount=True)
|
||||
if refusal is not None:
|
||||
print(refusal, file=sys.stderr)
|
||||
return 5
|
||||
|
||||
if args.command == "migrate":
|
||||
manifest = migrate_with_backup(config)
|
||||
if manifest:
|
||||
@@ -115,6 +144,41 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
# anyone reading the JSON.
|
||||
return 0 if report["ok"] else 1
|
||||
|
||||
if args.command == "release-gate":
|
||||
from photo_pipeline.services import release
|
||||
|
||||
report = release.run_gate(config, output=args.output)
|
||||
print(
|
||||
json.dumps(
|
||||
{k: v for k, v in report.items() if k not in ("stages", "matrix")}, indent=2
|
||||
)
|
||||
)
|
||||
return 0 if report["ok"] else 1
|
||||
|
||||
if args.command == "dry-run":
|
||||
from photo_pipeline.services import release
|
||||
|
||||
try:
|
||||
report = release.dry_run(config)
|
||||
except release.ReleaseError as error:
|
||||
print(str(error))
|
||||
return 1
|
||||
if args.output:
|
||||
Path(args.output).write_text(json.dumps(report, indent=2))
|
||||
print(json.dumps(report, indent=2))
|
||||
return 0
|
||||
|
||||
if args.command == "approve-dry-run":
|
||||
from photo_pipeline.services import release
|
||||
|
||||
try:
|
||||
record = release.approve(config, args.report, approver=args.approver)
|
||||
except (release.ReleaseError, OSError, ValueError) as error:
|
||||
print(str(error))
|
||||
return 1
|
||||
print(json.dumps(record, indent=2))
|
||||
return 0
|
||||
|
||||
if args.command == "diagnostics":
|
||||
from photo_pipeline.services import diagnostics
|
||||
|
||||
@@ -161,13 +225,21 @@ def main(argv: Sequence[str] | None = None) -> int:
|
||||
|
||||
import uvicorn
|
||||
|
||||
from photo_pipeline.api.app import create_app
|
||||
from photo_pipeline.api.app import ConfigurationRefused, create_app
|
||||
|
||||
# An exposed deployment without an access secret must not reach the port at all,
|
||||
# and the operator needs a sentence, not a traceback (US08-01).
|
||||
try:
|
||||
app = create_app(config)
|
||||
except ConfigurationRefused as error:
|
||||
print(str(error), file=sys.stderr)
|
||||
return 4
|
||||
|
||||
lock = LibraryLock(config, "api")
|
||||
if (held := _acquire(lock, allow_legacy=args.allow_legacy)) is not None:
|
||||
return held
|
||||
try:
|
||||
uvicorn.run(create_app(config), host=config.host, port=config.port)
|
||||
uvicorn.run(app, host=config.host, port=config.port)
|
||||
finally:
|
||||
lock.release()
|
||||
return 0
|
||||
@@ -178,8 +250,6 @@ def _acquire(lock: LibraryLock, *, allow_legacy: bool) -> int | None:
|
||||
|
||||
Returns an exit code to return, or ``None`` when the lock was acquired.
|
||||
"""
|
||||
import sys
|
||||
|
||||
try:
|
||||
lock.acquire(allow_legacy=allow_legacy)
|
||||
except LockHeld as error:
|
||||
|
||||
@@ -35,7 +35,13 @@ from photo_pipeline.api.routes import (
|
||||
uploads,
|
||||
workflow,
|
||||
)
|
||||
from photo_pipeline.api.security import DEFAULT_HEADERS, SecurityMiddleware, Session
|
||||
from photo_pipeline.api.security import (
|
||||
DEFAULT_HEADERS,
|
||||
FailureLimiter,
|
||||
SecurityMiddleware,
|
||||
Session,
|
||||
trust_refusal,
|
||||
)
|
||||
|
||||
# Registers the safety_score / analysis job handlers on import.
|
||||
import photo_pipeline.jobs.domain_handlers # noqa: F401
|
||||
@@ -84,9 +90,16 @@ def _install_error_handlers(app: FastAPI) -> None:
|
||||
return _envelope(500, "internal_error", "internal error")
|
||||
|
||||
|
||||
class ConfigurationRefused(RuntimeError):
|
||||
"""The configuration would serve the library to callers it cannot authenticate."""
|
||||
|
||||
|
||||
def create_app(config: Config | None = None) -> FastAPI:
|
||||
config = config or Config.from_env()
|
||||
configure_logging(config.log_level, config.log_format)
|
||||
# Before anything is built, let alone bound to a port (US08-01).
|
||||
if (why := trust_refusal(config)) is not None:
|
||||
raise ConfigurationRefused(why)
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
@@ -114,6 +127,10 @@ def create_app(config: Config | None = None) -> FastAPI:
|
||||
# One session per process: the browser exchanges it for a cookie + CSRF token,
|
||||
# and every other origin is refused before a route ever runs (US07-02).
|
||||
app.state.session = Session.create()
|
||||
app.state.access_limiter = FailureLimiter()
|
||||
# Also set in the lifespan, but the bootstrap route reads it, and a caller can
|
||||
# arrive before anything else has touched app.state.
|
||||
app.state.config = config
|
||||
app.add_middleware(SecurityMiddleware, session=app.state.session, config=config)
|
||||
_install_error_handlers(app)
|
||||
app.include_router(session_routes.router, prefix="/api/v1")
|
||||
|
||||
@@ -3,21 +3,52 @@
|
||||
It sets the ``HttpOnly``/``SameSite=Strict`` session cookie and returns the CSRF
|
||||
token in the body. A foreign page can call this — it just cannot read the answer,
|
||||
because the app sends no CORS headers — and the cookie it received is never attached
|
||||
to a request that foreign page initiates.
|
||||
to a request that foreign page initiated.
|
||||
|
||||
When an access secret is configured (mandatory as soon as the app is reachable from
|
||||
another machine, US08-01) this is also the authentication gate: the secret buys the
|
||||
cookie, and every route behind it keeps asking for exactly the session and CSRF token
|
||||
it asked for before. Wrong secrets are counted, and a burst of them stops being
|
||||
answered — otherwise a proxy-exposed deployment could be guessed at indefinitely.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import secrets
|
||||
|
||||
from fastapi import APIRouter, Request
|
||||
from fastapi.responses import JSONResponse
|
||||
|
||||
from photo_pipeline.api.security import SESSION_COOKIE
|
||||
from photo_pipeline.api.security import ACCESS_SECRET_HEADER, SESSION_COOKIE
|
||||
|
||||
router = APIRouter(tags=["session"])
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _refuse(status: int, code: str, message: str) -> JSONResponse:
|
||||
return JSONResponse(status_code=status, content={"error": {"code": code, "message": message}})
|
||||
|
||||
|
||||
@router.get("/session")
|
||||
def start_session(request: Request) -> JSONResponse:
|
||||
config = request.app.state.config
|
||||
secret = config.access_secret
|
||||
if secret is not None:
|
||||
limiter = request.app.state.access_limiter
|
||||
if limiter.blocked():
|
||||
return _refuse(429, "too_many_attempts", "too many failed attempts; retry later")
|
||||
offered = request.headers.get(ACCESS_SECRET_HEADER, "")
|
||||
if not secrets.compare_digest(offered, secret.get_secret_value()):
|
||||
limiter.record_failure()
|
||||
# The client address is the whole record: the offered secret, the issued
|
||||
# session, and the request body all stay out of the log.
|
||||
log.warning(
|
||||
"access secret rejected", extra={"client": _client(request), "path": "/session"}
|
||||
)
|
||||
return _refuse(401, "access_denied", "a valid access secret is required")
|
||||
|
||||
session = request.app.state.session
|
||||
response = JSONResponse({"csrf_token": session.csrf_token})
|
||||
response.set_cookie(
|
||||
@@ -25,6 +56,13 @@ def start_session(request: Request) -> JSONResponse:
|
||||
session.id,
|
||||
httponly=True,
|
||||
samesite="strict",
|
||||
# HTTPS outside means the cookie must never travel over a plain hop, even one
|
||||
# this process cannot see. Loopback http keeps working unchanged.
|
||||
secure=request.scope.get("state", {}).get("external_scheme") == "https",
|
||||
path="/",
|
||||
)
|
||||
return response
|
||||
|
||||
|
||||
def _client(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
@@ -19,6 +19,14 @@ The defenses stack, because each one alone has a hole:
|
||||
only in the bootstrap response body, which a foreign page cannot read (no CORS) —
|
||||
so possessing it proves the caller is same-origin.
|
||||
|
||||
Behind a reverse proxy (US08-01) the same stack holds with two substitutions: the
|
||||
allowed host set comes from configuration instead of being the loopback names, and
|
||||
the host/scheme the policy judges is the *external* one, which is only read from
|
||||
``X-Forwarded-*`` when the request actually arrived from a configured proxy. The
|
||||
loopback check was standing in for authentication, so naming a non-loopback host
|
||||
also makes an access secret mandatory — ``trust_refusal`` refuses to start without
|
||||
one, and the secret is what the bootstrap endpoint trades for the session cookie.
|
||||
|
||||
``evaluate`` is a pure function over the request metadata: the whole policy is one
|
||||
table that a unit test can enumerate, and the middleware only applies its verdict.
|
||||
"""
|
||||
@@ -26,6 +34,7 @@ table that a unit test can enumerate, and the middleware only applies its verdic
|
||||
from __future__ import annotations
|
||||
|
||||
import secrets
|
||||
import time
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass
|
||||
from urllib.parse import urlsplit
|
||||
@@ -35,6 +44,7 @@ from starlette.responses import JSONResponse
|
||||
|
||||
SESSION_COOKIE = "pp_session"
|
||||
CSRF_HEADER = "x-csrf-token"
|
||||
ACCESS_SECRET_HEADER = "x-access-secret"
|
||||
API_PREFIX = "/api/v1"
|
||||
SAFE_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
|
||||
# Reachable without a session: liveness/readiness (an orchestrator has no cookie)
|
||||
@@ -43,6 +53,9 @@ PUBLIC_PATHS = frozenset(
|
||||
{f"{API_PREFIX}/health/live", f"{API_PREFIX}/health/ready", f"{API_PREFIX}/session"}
|
||||
)
|
||||
LOOPBACK_HOSTS = frozenset({"127.0.0.1", "localhost", "::1", "[::1]"})
|
||||
# Mutating endpoints that must stay reachable while mutation itself is gated: the
|
||||
# backup a careful operator takes first, and its retention (US07-07).
|
||||
MUTATION_EXEMPT_PATHS = frozenset({f"{API_PREFIX}/backups", f"{API_PREFIX}/backups/prune"})
|
||||
|
||||
# Applied to every response. No inline script/style is used by the frontend, so the
|
||||
# policy can stay strict; `frame-ancestors 'none'` and CORP keep other pages from
|
||||
@@ -93,6 +106,28 @@ def split_host(value: str) -> tuple[str, str]:
|
||||
return host, port
|
||||
|
||||
|
||||
def external_view(
|
||||
*,
|
||||
client: str | None,
|
||||
headers: Mapping[str, str],
|
||||
scheme: str,
|
||||
trusted_proxies: frozenset[str],
|
||||
) -> tuple[str, str]:
|
||||
"""The ``(scheme, host)`` the caller used, as opposed to the one this hop saw.
|
||||
|
||||
Forwarded headers are a client-supplied claim. Believing them from anyone lets a
|
||||
request declare its own origin — and origin is half of this module's evidence —
|
||||
so they count only when the connection came from a configured proxy.
|
||||
"""
|
||||
host = headers.get("host", "")
|
||||
if client is None or client not in trusted_proxies:
|
||||
return scheme, host
|
||||
# A chain appends: the first entry is what the original client asked for.
|
||||
forwarded_proto = headers.get("x-forwarded-proto", "").split(",")[0].strip().lower()
|
||||
forwarded_host = headers.get("x-forwarded-host", "").split(",")[0].strip()
|
||||
return forwarded_proto or scheme, forwarded_host or host
|
||||
|
||||
|
||||
def evaluate(
|
||||
*,
|
||||
method: str,
|
||||
@@ -100,20 +135,26 @@ def evaluate(
|
||||
headers: Mapping[str, str],
|
||||
session: Session,
|
||||
allowed_hosts: frozenset[str] = LOOPBACK_HOSTS,
|
||||
scheme: str = "http",
|
||||
max_request_bytes: int,
|
||||
) -> Refusal | None:
|
||||
"""Why this request must be refused, or ``None`` when it may proceed."""
|
||||
"""Why this request must be refused, or ``None`` when it may proceed.
|
||||
|
||||
``headers["host"]`` and ``scheme`` are the external ones (see ``external_view``);
|
||||
the allowed origins are the allowed hosts under that scheme and port, so there is
|
||||
no second list that can drift away from the first.
|
||||
"""
|
||||
host_header = headers.get("host", "")
|
||||
host, port = split_host(host_header)
|
||||
if host.lower() not in allowed_hosts:
|
||||
return Refusal(403, "host_not_allowed", "request host is not a local address")
|
||||
return Refusal(403, "host_not_allowed", "request host is not an allowed address")
|
||||
|
||||
origin = headers.get("origin")
|
||||
if origin is not None and origin != "":
|
||||
parts = urlsplit(origin)
|
||||
origin_host, origin_port = split_host(parts.netloc)
|
||||
if (
|
||||
parts.scheme not in ("http", "https")
|
||||
parts.scheme != scheme
|
||||
or origin_host.lower() not in allowed_hosts
|
||||
or origin_port != port
|
||||
):
|
||||
@@ -139,14 +180,66 @@ def evaluate(
|
||||
return None
|
||||
|
||||
|
||||
def exposed_hosts(config) -> list[str]:
|
||||
"""Configured names by which this application is reachable from another machine."""
|
||||
names = {str(config.host).lower()}
|
||||
names.update(split_host(name)[0].lower() for name in config.allowed_hosts)
|
||||
return sorted(names - LOOPBACK_HOSTS)
|
||||
|
||||
|
||||
def trust_refusal(config) -> str | None:
|
||||
"""Why this configuration must not serve at all, or ``None``.
|
||||
|
||||
Reaching the app used to prove ownership of it. The moment a configuration makes
|
||||
it reachable from elsewhere that stops being true, so serving without a secret
|
||||
would publish the library — refuse at startup rather than at the first request,
|
||||
when the operator is no longer watching (US08-01).
|
||||
"""
|
||||
exposed = exposed_hosts(config)
|
||||
if exposed and config.access_secret is None:
|
||||
return (
|
||||
f"refusing to serve: {', '.join(exposed)} is reachable from outside this "
|
||||
"machine, so PHOTO_PIPELINE_ACCESS_SECRET must be set"
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
class FailureLimiter:
|
||||
"""Bounded failed access-secret attempts, so the secret cannot be guessed online.
|
||||
|
||||
ponytail: one counter for the whole process rather than per client address —
|
||||
behind a proxy every attempt arrives from the same address anyway. Per-caller
|
||||
buckets if the app is ever exposed without one.
|
||||
"""
|
||||
|
||||
def __init__(self, limit: int = 5, window: float = 60.0) -> None:
|
||||
self.limit = limit
|
||||
self.window = window
|
||||
self._failures: list[float] = []
|
||||
|
||||
def blocked(self) -> bool:
|
||||
now = time.monotonic()
|
||||
self._failures = [at for at in self._failures if now - at < self.window]
|
||||
return len(self._failures) >= self.limit
|
||||
|
||||
def record_failure(self) -> None:
|
||||
self._failures.append(time.monotonic())
|
||||
|
||||
|
||||
class SecurityMiddleware:
|
||||
"""Pure-ASGI so the SSE stream keeps streaming (BaseHTTPMiddleware buffers)."""
|
||||
|
||||
def __init__(self, app, *, session: Session, config) -> None:
|
||||
self.app = app
|
||||
self.session = session
|
||||
self.config = config
|
||||
self.max_request_bytes = config.max_request_bytes
|
||||
self.allowed_hosts = frozenset(LOOPBACK_HOSTS | {str(config.host).lower()})
|
||||
self.allowed_hosts = frozenset(
|
||||
LOOPBACK_HOSTS
|
||||
| {str(config.host).lower()}
|
||||
| {split_host(name)[0].lower() for name in config.allowed_hosts}
|
||||
)
|
||||
self.trusted_proxies = frozenset(config.trusted_proxies)
|
||||
|
||||
async def __call__(self, scope, receive, send) -> None:
|
||||
if scope["type"] != "http":
|
||||
@@ -157,14 +250,26 @@ class SecurityMiddleware:
|
||||
# policy never has to parse a Cookie header.
|
||||
lookup = dict(headers)
|
||||
lookup["cookie-session"] = _cookie(headers.get("cookie", ""), SESSION_COOKIE)
|
||||
client = scope.get("client")
|
||||
scheme, lookup["host"] = external_view(
|
||||
client=client[0] if client else None,
|
||||
headers=headers,
|
||||
scheme=scope.get("scheme", "http"),
|
||||
trusted_proxies=self.trusted_proxies,
|
||||
)
|
||||
# What the session cookie's Secure flag is decided from, one hop later.
|
||||
scope.setdefault("state", {})["external_scheme"] = scheme
|
||||
refusal = evaluate(
|
||||
method=scope.get("method", "GET"),
|
||||
path=scope.get("path", "/"),
|
||||
headers=lookup,
|
||||
session=self.session,
|
||||
allowed_hosts=self.allowed_hosts,
|
||||
scheme=scheme,
|
||||
max_request_bytes=self.max_request_bytes,
|
||||
)
|
||||
if refusal is None:
|
||||
refusal = self._mutation_refusal(scope)
|
||||
if refusal is not None:
|
||||
response = JSONResponse(
|
||||
status_code=refusal.status,
|
||||
@@ -183,6 +288,26 @@ class SecurityMiddleware:
|
||||
|
||||
await self.app(scope, receive, send_with_headers)
|
||||
|
||||
def _mutation_refusal(self, scope) -> Refusal | None:
|
||||
"""Refuse every mutating request while the library's dry run is unapproved.
|
||||
|
||||
One choke point for the whole API: every mutation the browser can start is a
|
||||
non-safe method under ``/api/v1``. Reading stays open — an operator has to be
|
||||
able to look at what the application found in order to approve it (US07-07).
|
||||
"""
|
||||
method = scope.get("method", "GET").upper()
|
||||
path = scope.get("path", "/")
|
||||
if method in SAFE_METHODS or not path.startswith(API_PREFIX):
|
||||
return None
|
||||
if path in MUTATION_EXEMPT_PATHS:
|
||||
return None
|
||||
from photo_pipeline.services.release import mutation_blockers
|
||||
|
||||
blockers = mutation_blockers(self.config)
|
||||
if not blockers:
|
||||
return None
|
||||
return Refusal(403, blockers[0]["code"], blockers[0]["message"])
|
||||
|
||||
|
||||
def _cookie(header: str, name: str) -> str:
|
||||
for part in header.split(";"):
|
||||
|
||||
@@ -7,6 +7,9 @@ real external call needs them.
|
||||
|
||||
pydantic-settings would do this too, but a prefix-scan over the declared fields
|
||||
is a few lines and one fewer dependency.
|
||||
|
||||
Tuple-valued settings are lists in one variable: library roots are ``os.pathsep``
|
||||
separated because they are paths, everything else is comma separated.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -18,6 +21,62 @@ from typing import Mapping
|
||||
from pydantic import BaseModel, ConfigDict, SecretStr
|
||||
|
||||
ENV_PREFIX = "PHOTO_PIPELINE_"
|
||||
ENV_FILE_VAR = f"{ENV_PREFIX}ENV_FILE"
|
||||
DEFAULT_ENV_FILE = Path(".env")
|
||||
COMMA_LIST_FIELDS = frozenset({"allowed_hosts", "trusted_proxies"})
|
||||
|
||||
# The archived CLI's variable names, so the configuration file an operator already
|
||||
# has keeps working. The vision provider reads the OpenAI SDK's names, and the
|
||||
# library root is configuration here rather than a bare path (US07-01 donor).
|
||||
LEGACY_ALIASES = {
|
||||
"LLM_API_KEY": "OPENAI_API_KEY",
|
||||
"GEMINI_API_KEY": "OPENAI_API_KEY",
|
||||
"LLM_BASE_URL": "OPENAI_BASE_URL",
|
||||
"LIBRARY": f"{ENV_PREFIX}LIBRARY_ROOTS",
|
||||
}
|
||||
|
||||
|
||||
def parse_env_file(text: str) -> dict[str, str]:
|
||||
"""``KEY=value`` lines into a mapping. Comments, blanks, and quotes handled.
|
||||
|
||||
Deliberately not a shell: no interpolation, no ``export``, no multi-line values.
|
||||
A configuration file that can run code is a configuration file that can be a
|
||||
vulnerability.
|
||||
"""
|
||||
values: dict[str, str] = {}
|
||||
for line in text.splitlines():
|
||||
line = line.strip()
|
||||
if not line or line.startswith("#") or "=" not in line:
|
||||
continue
|
||||
key, _, raw = line.partition("=")
|
||||
key = key.strip()
|
||||
if not key or key.startswith("#"):
|
||||
continue
|
||||
value = raw.strip().strip('"').strip("'")
|
||||
values[key] = value
|
||||
alias = LEGACY_ALIASES.get(key)
|
||||
if alias:
|
||||
values.setdefault(alias, value)
|
||||
return values
|
||||
|
||||
|
||||
def load_env_file(path: Path | str | None = None) -> dict[str, str]:
|
||||
"""Load ``PHOTO_PIPELINE_ENV_FILE`` (or ``./.env``) into the environment.
|
||||
|
||||
Anything already exported wins: a file is the standing configuration, the shell
|
||||
is what you meant *this time*. Returns what it applied, which is what the CLI
|
||||
prints — names only, never values.
|
||||
"""
|
||||
candidate = path or os.environ.get(ENV_FILE_VAR) or DEFAULT_ENV_FILE
|
||||
candidate = Path(candidate)
|
||||
if not candidate.is_file():
|
||||
return {}
|
||||
applied = {}
|
||||
for key, value in parse_env_file(candidate.read_text()).items():
|
||||
if key not in os.environ:
|
||||
os.environ[key] = value
|
||||
applied[key] = value
|
||||
return applied
|
||||
|
||||
|
||||
class Config(BaseModel):
|
||||
@@ -30,6 +89,19 @@ class Config(BaseModel):
|
||||
log_level: str = "INFO"
|
||||
log_format: str = "json" # "json" or "text"
|
||||
|
||||
# Trust boundary (US08-01). Empty means loopback only, which is what the app did
|
||||
# before there was a setting: a request whose Host is not a loopback name is
|
||||
# refused, and no secret is needed because nothing outside this machine can call.
|
||||
# Naming a real hostname here is what makes the app reachable through a reverse
|
||||
# proxy, and it is exactly then that ``access_secret`` becomes mandatory.
|
||||
allowed_hosts: tuple[str, ...] = ()
|
||||
# Addresses whose ``X-Forwarded-Proto``/``X-Forwarded-Host`` may be believed. A
|
||||
# client that is not the proxy can otherwise declare its own origin.
|
||||
trusted_proxies: tuple[str, ...] = ()
|
||||
# Exchanged for the session cookie at the bootstrap endpoint. Once set it is
|
||||
# required even on loopback, so a development setup cannot half-enable it.
|
||||
access_secret: SecretStr | None = None
|
||||
|
||||
# Largest request body the API accepts. Every endpoint takes small JSON commands;
|
||||
# anything larger is a mistake or an attempt to exhaust memory (US07-02).
|
||||
max_request_bytes: int = 1_048_576
|
||||
@@ -42,6 +114,12 @@ class Config(BaseModel):
|
||||
# Free space an archive destination must keep beyond the transfer itself.
|
||||
archive_free_space_reserve_bytes: int = 1_000_000_000
|
||||
|
||||
# Refuse every mutating request until a read-only dry run of the configured
|
||||
# library has been produced and explicitly approved (US07-07). Off by default so
|
||||
# a development setup is unchanged; turn it on before pointing the application at
|
||||
# a library whose photos cannot be replaced.
|
||||
require_dry_run_approval: bool = False
|
||||
|
||||
vision_api_key: SecretStr | None = None
|
||||
immich_api_key: SecretStr | None = None
|
||||
immich_server_url: str = ""
|
||||
@@ -62,10 +140,18 @@ class Config(BaseModel):
|
||||
@classmethod
|
||||
def from_env(cls, environ: Mapping[str, str] | None = None) -> "Config":
|
||||
env = os.environ if environ is None else environ
|
||||
if environ is None:
|
||||
load_env_file() # a file never overrides what the shell already set
|
||||
env = os.environ
|
||||
data: dict = {}
|
||||
for name in cls.model_fields:
|
||||
raw = env.get(ENV_PREFIX + name.upper())
|
||||
if not raw:
|
||||
continue
|
||||
data[name] = raw.split(os.pathsep) if name == "library_roots" else raw
|
||||
if name == "library_roots":
|
||||
data[name] = raw.split(os.pathsep)
|
||||
elif name in COMMA_LIST_FIELDS:
|
||||
data[name] = [part.strip() for part in raw.split(",") if part.strip()]
|
||||
else:
|
||||
data[name] = raw
|
||||
return cls(**data)
|
||||
|
||||
@@ -12,8 +12,10 @@ nt-apply-list, pa-nsfw-filter). No dependency on the archived entry points.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import functools
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
from collections.abc import Iterable
|
||||
|
||||
@@ -32,6 +34,27 @@ def _timeout() -> float:
|
||||
return DEFAULT_TIMEOUT_SECONDS
|
||||
|
||||
|
||||
def find_binary(binary: str = "exiftool") -> str | None:
|
||||
"""Absolute path of exiftool, or ``None`` when it is not installed."""
|
||||
return shutil.which(binary)
|
||||
|
||||
|
||||
@functools.lru_cache(maxsize=1)
|
||||
def version() -> str | None:
|
||||
"""Reported exiftool version, or ``None`` when it is missing or unusable.
|
||||
|
||||
Cached: it cannot change inside a running process, and diagnostics asks for it
|
||||
on every report (US08-02, where a container image pins this version).
|
||||
"""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["exiftool", "-ver"], capture_output=True, text=True, timeout=_timeout()
|
||||
)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return None
|
||||
return (result.stdout or "").strip() or None
|
||||
|
||||
|
||||
def read_keyword_sets(paths: Iterable[str]) -> dict[str, set[str]]:
|
||||
"""Map each path to its lowercased set of ``Keywords`` + ``Subject`` values.
|
||||
|
||||
|
||||
@@ -15,9 +15,7 @@ import os
|
||||
from pathlib import Path
|
||||
from typing import Iterable, Iterator
|
||||
|
||||
SUPPORTED_EXTENSIONS = {
|
||||
".jpg", ".jpeg", ".png", ".webp", ".heic", ".heif", ".tiff", ".tif"
|
||||
}
|
||||
SUPPORTED_EXTENSIONS = {".jpg", ".jpeg", ".png", ".webp", ".heic", ".heif", ".tiff", ".tif"}
|
||||
EXCLUDED_DIR_NAMES = {"_IGNORE", ".@__thumb"}
|
||||
|
||||
|
||||
@@ -72,6 +70,63 @@ def resolve_in_roots(roots: Iterable[os.PathLike | str], path: os.PathLike | str
|
||||
raise PathPolicyError("path is outside the configured library roots")
|
||||
|
||||
|
||||
def in_container() -> bool:
|
||||
"""Whether this process is running inside a container image build of the app."""
|
||||
return Path("/.dockerenv").exists()
|
||||
|
||||
|
||||
def _under_mount(path: Path) -> bool:
|
||||
"""Whether ``path`` or one of its parents below ``/`` is a mounted filesystem."""
|
||||
current = path.resolve()
|
||||
while current != current.parent:
|
||||
if os.path.ismount(current):
|
||||
return True
|
||||
current = current.parent
|
||||
return False
|
||||
|
||||
|
||||
def roots_refusal(roots: Iterable[os.PathLike | str], *, require_mount: bool = False) -> str | None:
|
||||
"""Why the configured library roots cannot be worked with, or ``None`` (US08-03).
|
||||
|
||||
The roots name the directories this installation renames folders in, rewrites
|
||||
EXIF in, and archives from. Container path policy is the same problem as host
|
||||
path policy with one new failure mode: the configured roots must name the
|
||||
*container-side* mount paths. A host path configured inside a container is
|
||||
either absent or an ordinary directory of the image, so the library looks empty
|
||||
and the first write lands in the container's throwaway layer instead of in the
|
||||
library. That is worth refusing at startup, while an operator is still watching,
|
||||
rather than at the first write.
|
||||
|
||||
``require_mount`` is the container-only half of that: inside a container a real
|
||||
library arrives through a bind mount, so a root that is not on (or under) a
|
||||
mount point is not the library the deployment meant.
|
||||
|
||||
Writability is deliberately *not* checked: a bind mount's ownership is
|
||||
virtualised by Docker Desktop and Colima (it arrives as ``root:root``), so
|
||||
``os.access`` there is evidence about the virtio layer rather than about the
|
||||
library. A wrong UID/GID surfaces as a refused rename with the real errno, which
|
||||
is at least true; a refusal here would be false on two supported platforms.
|
||||
"""
|
||||
for root in roots:
|
||||
path = Path(root)
|
||||
if not path.exists():
|
||||
return (
|
||||
f"library root {path} does not exist: PHOTO_PIPELINE_LIBRARY_ROOTS must "
|
||||
"name paths that exist here, and in a container that means the mount path"
|
||||
)
|
||||
if not path.is_dir():
|
||||
return f"library root {path} is not a directory"
|
||||
if not os.access(path, os.R_OK | os.X_OK):
|
||||
return f"library root {path} is not readable by this process"
|
||||
if require_mount and not _under_mount(path):
|
||||
return (
|
||||
f"library root {path} is not on a mounted filesystem in this container: "
|
||||
"the library was not bind-mounted there, so PHOTO_PIPELINE_LIBRARY_ROOTS "
|
||||
"names a directory of the image rather than the library"
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def iter_supported_files(root: os.PathLike | str) -> Iterator[Path]:
|
||||
"""Yield supported, non-excluded files under ``root`` in deterministic order.
|
||||
|
||||
|
||||
@@ -264,8 +264,11 @@ class AnalysisService:
|
||||
asset = session.get(Asset, asset_id)
|
||||
if asset is not None and checkpoint.sha256:
|
||||
# The bytes changed when the container was rewritten; upload must use
|
||||
# the hash of what is actually on disk now (concept §3).
|
||||
# the hash of what is actually on disk now (concept §3), and the
|
||||
# recorded size has to move with it (US07-07).
|
||||
asset.current_sha256 = checkpoint.sha256
|
||||
if checkpoint.byte_size is not None:
|
||||
asset.byte_size = checkpoint.byte_size
|
||||
session.commit()
|
||||
|
||||
def get(self, asset_id: str) -> dict | None:
|
||||
|
||||
@@ -15,9 +15,17 @@ package:
|
||||
"started_at": "...", "library_roots": ["..."]}
|
||||
|
||||
One holder per role: an API and a worker are designed to run together, a second
|
||||
worker is not. A lock whose process is gone is stale and is taken over with the
|
||||
takeover recorded — refusing to start because of a crashed predecessor would turn
|
||||
one outage into two.
|
||||
worker is not. A lock whose process is gone is stale and is taken over — refusing
|
||||
to start because of a crashed predecessor would turn one outage into two.
|
||||
|
||||
Ownership is an advisory ``flock`` on that file, not the record inside it. The
|
||||
record says *who*; the kernel says *whether*. That distinction is what makes the
|
||||
lock work in containers (US08-03), where a PID and a hostname are namespaced: a
|
||||
lock left behind by a container that no longer exists names a pid that still
|
||||
"exists" in the new container and a host that cannot be probed, so believing the
|
||||
file would deadlock every restart. A flock is released when its holder dies however
|
||||
it dies, and is seen by every process that can open the file — which for a local
|
||||
data directory is every container of this deployment.
|
||||
|
||||
Legacy detection is deliberately a heuristic, not a promise: the archived CLI has
|
||||
no lock of its own, so what can be observed is its state files being written right
|
||||
@@ -27,6 +35,7 @@ mutating stage should refuse until it stops.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import fcntl
|
||||
import json
|
||||
import os
|
||||
import socket
|
||||
@@ -109,6 +118,12 @@ def _now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _in_container() -> bool:
|
||||
from photo_pipeline import path_policy
|
||||
|
||||
return path_policy.in_container()
|
||||
|
||||
|
||||
def legacy_activity(config: Config) -> dict:
|
||||
"""Legacy state files written within the activity window, if any."""
|
||||
seen: list[dict] = []
|
||||
@@ -139,6 +154,7 @@ class LibraryLock:
|
||||
self.role = role
|
||||
self.path = Path(config.data_dir) / f"{role}{LOCK_SUFFIX}"
|
||||
self._acquired = False
|
||||
self._handle = None
|
||||
|
||||
# ── inspection ────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -178,12 +194,31 @@ class LibraryLock:
|
||||
"stop it before running the application"
|
||||
)
|
||||
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
# The kernel decides, because the file cannot: a container's PID and hostname
|
||||
# are namespaced, so a lock left by a container that no longer exists names a
|
||||
# pid that "exists" and a host that cannot be probed (US08-03). An advisory
|
||||
# flock is held by a live process or by nobody, is released when that process
|
||||
# dies however it dies, and is shared by every process that can open this
|
||||
# file — which, for a local data directory, is every role in every container
|
||||
# of this deployment.
|
||||
handle = open(self.path, "a+", encoding="utf-8")
|
||||
try:
|
||||
fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
except OSError:
|
||||
handle.close()
|
||||
raise LockHeld(self.holder() or Holder(self.role, -1, "unknown", "unknown")) from None
|
||||
|
||||
current = self.holder()
|
||||
if current is not None:
|
||||
if current.alive:
|
||||
raise LockHeld(current)
|
||||
# Stale: its process is gone. Take over, and say so.
|
||||
self.path.unlink(missing_ok=True)
|
||||
if current is not None and current.host != socket.gethostname() and not _in_container():
|
||||
# We hold the kernel's lock, so nothing on *this* machine holds the file.
|
||||
# On a host that still leaves one case open: a data directory shared with
|
||||
# another machine, whose flock we cannot trust. Believe its record rather
|
||||
# than run two writers. In a container the data volume is local by
|
||||
# construction (US08-03), and a foreign hostname is only a dead container.
|
||||
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
|
||||
handle.close()
|
||||
raise LockHeld(current)
|
||||
|
||||
mine = Holder(
|
||||
role=self.role,
|
||||
@@ -192,15 +227,14 @@ class LibraryLock:
|
||||
started_at=_now().isoformat(),
|
||||
library_roots=tuple(str(root) for root in self._config.library_roots),
|
||||
)
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
payload = {k: v for k, v in mine.as_dict().items() if k != "alive"}
|
||||
# Exclusive create, so two processes racing here cannot both believe they won.
|
||||
try:
|
||||
with open(self.path, "x", encoding="utf-8") as handle:
|
||||
json.dump(payload, handle, indent=2)
|
||||
except FileExistsError:
|
||||
winner = self.holder()
|
||||
raise LockHeld(winner or mine) from None
|
||||
handle.seek(0)
|
||||
handle.truncate()
|
||||
json.dump(payload, handle, indent=2)
|
||||
handle.flush()
|
||||
# Held open on purpose: closing it is what releases the lock, and that must
|
||||
# happen when this process ends, not when this method returns.
|
||||
self._handle = handle
|
||||
self._acquired = True
|
||||
return mine
|
||||
|
||||
@@ -208,9 +242,10 @@ class LibraryLock:
|
||||
"""Give up a lock this process owns. Another holder's lock is left alone."""
|
||||
if not self._acquired:
|
||||
return
|
||||
current = self.holder()
|
||||
if current is not None and current.pid == os.getpid():
|
||||
self.path.unlink(missing_ok=True)
|
||||
self.path.unlink(missing_ok=True)
|
||||
if self._handle is not None:
|
||||
self._handle.close() # closing the descriptor releases the kernel lock
|
||||
self._handle = None
|
||||
self._acquired = False
|
||||
|
||||
def __enter__(self) -> "LibraryLock":
|
||||
|
||||
@@ -15,10 +15,13 @@ application.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import functools
|
||||
import json
|
||||
import shutil
|
||||
from pathlib import Path
|
||||
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.integrations import exiftool, immich_go
|
||||
from photo_pipeline.services import app_lock
|
||||
|
||||
# Below this much free space, mutating stages should stop rather than risk a
|
||||
@@ -26,6 +29,12 @@ from photo_pipeline.services import app_lock
|
||||
LOW_DISK_BYTES = 1_000_000_000
|
||||
CRITICAL_DISK_BYTES = 200_000_000
|
||||
|
||||
# Written into the container image at build time (US08-02). The image pins exiftool
|
||||
# and immich-go, and this file is how a running container reports which versions it
|
||||
# was built with — so a drifted or missing binary is visible here rather than in a
|
||||
# failed EXIF checkpoint or a misparsed upload report.
|
||||
IMAGE_VERSIONS_FILE = Path("/etc/photo-pipeline/versions.json")
|
||||
|
||||
|
||||
def _tree_bytes(path: Path) -> int:
|
||||
if not path.exists():
|
||||
@@ -69,8 +78,47 @@ def disk(path: Path) -> dict:
|
||||
}
|
||||
|
||||
|
||||
def _pinned_versions() -> dict[str, str]:
|
||||
"""The versions this image recorded at build time; empty outside a container."""
|
||||
try:
|
||||
recorded = json.loads(IMAGE_VERSIONS_FILE.read_text())
|
||||
except (OSError, ValueError):
|
||||
return {}
|
||||
if not isinstance(recorded, dict):
|
||||
return {}
|
||||
return {str(name): str(value) for name, value in recorded.items()}
|
||||
|
||||
|
||||
@functools.lru_cache(maxsize=4)
|
||||
def _uploader_version(binary: str) -> str | None:
|
||||
"""Cached: the uploader cannot change version inside one process."""
|
||||
return immich_go.version(binary)
|
||||
|
||||
|
||||
def tools(config: Config) -> list[dict]:
|
||||
"""The external executables the pipeline shells out to, and their versions.
|
||||
|
||||
``pinned`` is what the image was built against, ``version`` is what is actually
|
||||
installed. They differ only when the binary was replaced or mounted over.
|
||||
"""
|
||||
return [
|
||||
{
|
||||
"name": "exiftool",
|
||||
"path": exiftool.find_binary(),
|
||||
"version": exiftool.version(),
|
||||
"pinned": _pinned_versions().get("exiftool"),
|
||||
},
|
||||
{
|
||||
"name": "immich-go",
|
||||
"path": immich_go.find_binary(config.immich_go_binary),
|
||||
"version": _uploader_version(config.immich_go_binary),
|
||||
"pinned": _pinned_versions().get("immich-go"),
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
def report(config: Config) -> dict:
|
||||
"""Sizes, disk headroom, warnings, and who currently holds the library lock."""
|
||||
"""Sizes, disk headroom, tool versions, warnings, and who holds the library lock."""
|
||||
database = config.database_path
|
||||
components = [
|
||||
_component("database", database),
|
||||
@@ -132,6 +180,23 @@ def report(config: Config) -> dict:
|
||||
}
|
||||
)
|
||||
|
||||
installed_tools = tools(config)
|
||||
for tool in installed_tools:
|
||||
# A missing tool is reported as ``version: null`` rather than warned about: on a
|
||||
# development machine the uploader is legitimately absent, and the stages that
|
||||
# need it already refuse to run. A *drifted* tool is different — the image pinned
|
||||
# a version and something replaced it.
|
||||
if tool["version"] and tool["pinned"] and tool["pinned"] not in tool["version"]:
|
||||
warnings.append(
|
||||
{
|
||||
"code": "tool_version_drift",
|
||||
"message": (
|
||||
f"{tool['name']} reports {tool['version']} but this image pinned "
|
||||
f"{tool['pinned']}"
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
locks = {}
|
||||
for role in ("api", "worker"):
|
||||
holder = app_lock.LibraryLock(config, role).holder()
|
||||
@@ -151,6 +216,7 @@ def report(config: Config) -> dict:
|
||||
"components": components,
|
||||
"total_bytes": sum(component["bytes"] for component in components),
|
||||
"disk": space,
|
||||
"tools": installed_tools,
|
||||
"warnings": warnings,
|
||||
"locks": locks,
|
||||
"legacy_activity": legacy,
|
||||
|
||||
@@ -21,6 +21,7 @@ unreadable file, or a write that did not take is not evidence that metadata is f
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
@@ -58,6 +59,11 @@ class CheckpointResult:
|
||||
state: str # verified | divergent | failed
|
||||
changed_fields: tuple[str, ...] = ()
|
||||
sha256: str | None = None
|
||||
# exiftool rewrites the container, so the file's size moves with its hash. Both
|
||||
# are inventory facts about the current bytes and both have to be refreshed
|
||||
# together, or the next stage compares against a size that no longer exists
|
||||
# (US07-07: a rename plan blocked itself forever after any EXIF write).
|
||||
byte_size: int | None = None
|
||||
verified_at: datetime | None = None
|
||||
reason: str | None = None
|
||||
|
||||
@@ -137,9 +143,12 @@ def run(
|
||||
|
||||
changed = compare(before, after)
|
||||
sha256 = hashing.sha256_file(path)
|
||||
byte_size = os.path.getsize(path)
|
||||
if changed:
|
||||
return CheckpointResult(DIVERGENT, changed_fields=changed, sha256=sha256)
|
||||
return CheckpointResult(VERIFIED, sha256=sha256, verified_at=_now())
|
||||
return CheckpointResult(
|
||||
DIVERGENT, changed_fields=changed, sha256=sha256, byte_size=byte_size
|
||||
)
|
||||
return CheckpointResult(VERIFIED, sha256=sha256, byte_size=byte_size, verified_at=_now())
|
||||
|
||||
|
||||
def record(
|
||||
|
||||
399
photo_pipeline/services/release.py
Normal file
399
photo_pipeline/services/release.py
Normal file
@@ -0,0 +1,399 @@
|
||||
"""The release gate, the real-library dry run, and the approval that unlocks
|
||||
mutation (US07-07, concept §18 release gates).
|
||||
|
||||
Three things live here because they are one decision:
|
||||
|
||||
1. **The gate** — one command that provisions an isolated stack, runs every suite in
|
||||
a fixed order, and retains versioned evidence with checksums. A release is not
|
||||
"the tests passed on my machine last Tuesday"; it is a report that says which
|
||||
revision, which suites, how long, and what the artefacts hash to.
|
||||
2. **The dry run** — a strictly read-only pass over the real photo library that
|
||||
answers "what would this application do to it?" before it is allowed to do
|
||||
anything. It opens no file for writing, creates no database rows, and touches no
|
||||
metadata; it counts, classifies, and reconciles against whatever the database
|
||||
already knows.
|
||||
3. **The approval** — a person reads that report and signs it off for exactly the
|
||||
library roots it describes. Until then, with
|
||||
``PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL`` set, every mutating request is
|
||||
refused. Change the roots, or produce a newer report, and the approval no longer
|
||||
matches: it approves *that* reconciliation, not the idea of mutating.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from collections import Counter
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from photo_pipeline import path_policy
|
||||
from photo_pipeline.config import Config
|
||||
|
||||
SCHEMA_VERSION = 1
|
||||
APPROVAL_NAME = "dry-run-approval.json"
|
||||
CHECKSUMS_NAME = "CHECKSUMS.sha256"
|
||||
REPORT_NAME = "release-report.json"
|
||||
|
||||
# The suites, in the order a failure is cheapest to read: units before the stacks
|
||||
# they compose. ``label`` is what the report and the operator see.
|
||||
STAGES: tuple[tuple[str, tuple[str, ...]], ...] = (
|
||||
("unit", ("tests/unit",)),
|
||||
("characterization", ("tests/characterization",)),
|
||||
("integration", ("tests/integration",)),
|
||||
("browser", ("tests/e2e",)),
|
||||
)
|
||||
|
||||
# Skips the gate accepts, because they describe the machine rather than the code. The
|
||||
# container ones (US08-02) belong here for the same reason exiftool does: the image
|
||||
# build needs a Docker daemon and the network, and its definition is still checked
|
||||
# offline in tests/integration/test_container_image.py.
|
||||
ALLOWED_SKIP_REASONS = (
|
||||
"exiftool not installed",
|
||||
"root ignores directory permissions",
|
||||
"no Docker daemon available",
|
||||
"bind-mount ownership is virtualised",
|
||||
)
|
||||
|
||||
|
||||
class ReleaseError(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def sha256_file(path: Path) -> str:
|
||||
digest = hashlib.sha256()
|
||||
with path.open("rb") as handle:
|
||||
for chunk in iter(lambda: handle.read(1024 * 1024), b""):
|
||||
digest.update(chunk)
|
||||
return digest.hexdigest()
|
||||
|
||||
|
||||
def sha256_bytes(payload: bytes) -> str:
|
||||
return hashlib.sha256(payload).hexdigest()
|
||||
|
||||
|
||||
def revision() -> str | None:
|
||||
"""The commit this gate ran against, when the tree is a git checkout."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["git", "rev-parse", "HEAD"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=10,
|
||||
cwd=str(Path(__file__).resolve().parents[2]),
|
||||
)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return None
|
||||
return result.stdout.strip() or None
|
||||
|
||||
|
||||
# ── the story matrix ─────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def story_matrix(repo: Path | None = None) -> dict:
|
||||
"""Every backlog story, and how it is covered.
|
||||
|
||||
A story is ``delivered`` (mapped to test files that exist) or ``planned`` (an
|
||||
accepted, not-yet-implemented story). Anything else — a story file nobody
|
||||
mapped, or a mapping to a file that is gone — is a hole in the matrix, and the
|
||||
gate fails on it rather than reporting a green run over missing coverage.
|
||||
"""
|
||||
repo = repo or Path(__file__).resolve().parents[2]
|
||||
traceability = json.loads((repo / "tests" / "story_traceability.json").read_text())
|
||||
mapped: dict[str, list[str]] = traceability["stories"]
|
||||
planned: list[str] = traceability.get("planned", [])
|
||||
stories = sorted(
|
||||
"-".join(path.stem.split("-")[:2])
|
||||
for path in (repo / "delivery_backlog" / "stories").glob("US*.md")
|
||||
)
|
||||
|
||||
missing_tests = [
|
||||
f"{story}: {rel}"
|
||||
for story, files in mapped.items()
|
||||
for rel in files
|
||||
if not (repo / rel).is_file()
|
||||
]
|
||||
unmapped = [s for s in stories if s not in mapped and s not in planned]
|
||||
unknown = [s for s in list(mapped) + planned if s not in stories]
|
||||
overlap = sorted(set(mapped) & set(planned))
|
||||
return {
|
||||
"stories": len(stories),
|
||||
"delivered": sorted(mapped),
|
||||
"planned": sorted(planned),
|
||||
"problems": [
|
||||
*(f"story with no tests and not planned: {s}" for s in unmapped),
|
||||
*(f"mapped test file is missing — {entry}" for entry in missing_tests),
|
||||
*(f"mapped story is not in the backlog: {s}" for s in unknown),
|
||||
*(f"story is both delivered and planned: {s}" for s in overlap),
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
# ── the gate ─────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@dataclass
|
||||
class StageResult:
|
||||
label: str
|
||||
command: list[str]
|
||||
returncode: int
|
||||
seconds: float
|
||||
summary: str
|
||||
skipped: list[str]
|
||||
|
||||
def as_dict(self) -> dict:
|
||||
return {
|
||||
"stage": self.label,
|
||||
"command": self.command,
|
||||
"returncode": self.returncode,
|
||||
"seconds": round(self.seconds, 2),
|
||||
"summary": self.summary,
|
||||
"skipped": self.skipped,
|
||||
"ok": self.returncode == 0,
|
||||
}
|
||||
|
||||
|
||||
def _run_stage(label: str, paths: tuple[str, ...], *, repo: Path, log_dir: Path) -> StageResult:
|
||||
command = [sys.executable, "-m", "pytest", *paths, "-q", "-rs"]
|
||||
started = time.monotonic()
|
||||
result = subprocess.run(command, cwd=str(repo), capture_output=True, text=True)
|
||||
elapsed = time.monotonic() - started
|
||||
output = result.stdout + result.stderr
|
||||
(log_dir / f"{label}.log").write_text(output)
|
||||
lines = [line for line in output.splitlines() if line.strip()]
|
||||
summary = lines[-1] if lines else ""
|
||||
skipped = [line for line in lines if line.startswith("SKIPPED")]
|
||||
return StageResult(label, command, result.returncode, elapsed, summary, skipped)
|
||||
|
||||
|
||||
def unexpected_skips(results: list[StageResult]) -> list[str]:
|
||||
"""Skips the gate will not accept: everything but the documented environment ones."""
|
||||
return [
|
||||
line
|
||||
for result in results
|
||||
for line in result.skipped
|
||||
if not any(reason in line for reason in ALLOWED_SKIP_REASONS)
|
||||
]
|
||||
|
||||
|
||||
def run_gate(
|
||||
config: Config,
|
||||
*,
|
||||
output: Path | str | None = None,
|
||||
stages: tuple[tuple[str, tuple[str, ...]], ...] = STAGES,
|
||||
repo: Path | None = None,
|
||||
) -> dict:
|
||||
"""Run every suite in an isolated stack and retain checksummed evidence.
|
||||
|
||||
The stack is isolated by construction: each pytest run builds its own temporary
|
||||
data directories and libraries, so the gate never reads or writes the operator's
|
||||
photos. What it keeps afterwards is the report, the per-stage logs, and a
|
||||
checksum file over both.
|
||||
"""
|
||||
repo = repo or Path(__file__).resolve().parents[2]
|
||||
directory = Path(output) if output else Path(config.data_dir) / "release" / _now().strftime(
|
||||
"%Y%m%dT%H%M%SZ"
|
||||
)
|
||||
logs = directory / "logs"
|
||||
logs.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
matrix = story_matrix(repo)
|
||||
results = [_run_stage(label, paths, repo=repo, log_dir=logs) for label, paths in stages]
|
||||
skips = unexpected_skips(results)
|
||||
|
||||
report = {
|
||||
"schema_version": SCHEMA_VERSION,
|
||||
"started_at": _now().isoformat(),
|
||||
"revision": revision(),
|
||||
"python": sys.version.split()[0],
|
||||
"platform": os.uname().sysname,
|
||||
"matrix": matrix,
|
||||
"stages": [result.as_dict() for result in results],
|
||||
"unexpected_skips": skips,
|
||||
"failures": [result.label for result in results if result.returncode != 0],
|
||||
}
|
||||
report["ok"] = not report["failures"] and not matrix["problems"] and not skips
|
||||
report["finished_at"] = _now().isoformat()
|
||||
|
||||
(directory / REPORT_NAME).write_text(json.dumps(report, indent=2))
|
||||
# The evidence is only evidence if it can be shown to be the evidence that was
|
||||
# produced. Checksums are the honest version of "signed" without a key: a real
|
||||
# signature belongs to whatever key management the release actually has.
|
||||
checksums = "\n".join(
|
||||
f"{sha256_file(path)} {path.relative_to(directory)}"
|
||||
for path in sorted(directory.rglob("*"))
|
||||
if path.is_file() and path.name != CHECKSUMS_NAME
|
||||
)
|
||||
(directory / CHECKSUMS_NAME).write_text(checksums + "\n")
|
||||
report["evidence"] = str(directory)
|
||||
return report
|
||||
|
||||
|
||||
# ── the real-library dry run ─────────────────────────────────────────────────
|
||||
|
||||
|
||||
def dry_run(config: Config, *, roots: tuple[Path, ...] | None = None) -> dict:
|
||||
"""Read-only reconciliation of the configured library. Changes nothing.
|
||||
|
||||
Opens no file for writing, writes no database row, and reads only what
|
||||
``os.stat`` and the existing database already say. The point is to be able to
|
||||
look at a real library — the one with the irreplaceable photos in it — and see
|
||||
what the application believes about it before it is allowed to act.
|
||||
"""
|
||||
roots = roots or tuple(Path(root) for root in config.library_roots)
|
||||
if not roots:
|
||||
raise ReleaseError("no library roots are configured")
|
||||
|
||||
by_extension: Counter = Counter()
|
||||
folders: set[str] = set()
|
||||
files: list[str] = []
|
||||
unreadable: list[str] = []
|
||||
excluded = 0
|
||||
total_bytes = 0
|
||||
for root in roots:
|
||||
if not Path(root).is_dir():
|
||||
raise ReleaseError(f"library root {root} is not a directory")
|
||||
for path in sorted(Path(root).rglob("*")):
|
||||
if path.is_dir():
|
||||
# Never traverse into an excluded directory, and never report its
|
||||
# contents: proving exclusion must not require opening it.
|
||||
if path_policy.is_excluded(path):
|
||||
excluded += 1
|
||||
continue
|
||||
if path_policy.is_excluded(path):
|
||||
continue
|
||||
try:
|
||||
stat = path.stat()
|
||||
except OSError:
|
||||
unreadable.append(str(path))
|
||||
continue
|
||||
files.append(str(path))
|
||||
folders.add(str(path.parent))
|
||||
by_extension[path.suffix.lower() or "(none)"] += 1
|
||||
total_bytes += stat.st_size
|
||||
|
||||
known = _known_paths(config)
|
||||
on_disk = set(files)
|
||||
report = {
|
||||
"schema_version": SCHEMA_VERSION,
|
||||
"generated_at": _now().isoformat(),
|
||||
"revision": revision(),
|
||||
"library_roots": [str(root) for root in roots],
|
||||
"files": len(files),
|
||||
"folders": len(folders),
|
||||
"bytes": total_bytes,
|
||||
"excluded_directories": excluded,
|
||||
"unreadable": unreadable,
|
||||
"by_extension": dict(sorted(by_extension.items())),
|
||||
"reconciliation": {
|
||||
"known_to_database": len(known),
|
||||
"already_registered": len(on_disk & known),
|
||||
"new_to_the_application": len(on_disk - known),
|
||||
"recorded_but_absent": sorted(known - on_disk)[:100],
|
||||
"recorded_but_absent_total": len(known - on_disk),
|
||||
},
|
||||
"mutation": "none — this pass is read-only",
|
||||
}
|
||||
report["checksum"] = sha256_bytes(
|
||||
json.dumps(report, sort_keys=True).encode("utf-8")
|
||||
)
|
||||
return report
|
||||
|
||||
|
||||
def _known_paths(config: Config) -> set[str]:
|
||||
"""Current asset paths the database holds, or an empty set if there is none."""
|
||||
if not config.database_path.exists():
|
||||
return set()
|
||||
from sqlalchemy import select
|
||||
|
||||
from photo_pipeline.db import create_db_engine, create_session_factory
|
||||
from photo_pipeline.models import Asset
|
||||
|
||||
engine = create_db_engine(config.database_url)
|
||||
try:
|
||||
with create_session_factory(engine)() as session:
|
||||
return {
|
||||
path
|
||||
for path in session.scalars(select(Asset.current_path))
|
||||
if path is not None
|
||||
}
|
||||
except Exception:
|
||||
return set()
|
||||
finally:
|
||||
engine.dispose()
|
||||
|
||||
|
||||
# ── the approval ─────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def approval_path(config: Config) -> Path:
|
||||
return Path(config.data_dir) / APPROVAL_NAME
|
||||
|
||||
|
||||
def approve(config: Config, report: dict | Path | str, *, approver: str) -> dict:
|
||||
"""Record that a person read this reconciliation and accepts mutation for it."""
|
||||
if isinstance(report, (str, Path)):
|
||||
report = json.loads(Path(report).read_text())
|
||||
if "checksum" not in report:
|
||||
raise ReleaseError("this is not a dry-run report: it has no checksum")
|
||||
record = {
|
||||
"schema_version": SCHEMA_VERSION,
|
||||
"approved_at": _now().isoformat(),
|
||||
"approved_by": approver,
|
||||
"report_checksum": report["checksum"],
|
||||
"library_roots": report["library_roots"],
|
||||
"files": report["files"],
|
||||
"revision": report.get("revision"),
|
||||
}
|
||||
path = approval_path(config)
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps(record, indent=2))
|
||||
return record
|
||||
|
||||
|
||||
def mutation_blockers(config: Config) -> list[dict]:
|
||||
"""Why mutation must stay refused, or an empty list.
|
||||
|
||||
Only enforced when ``require_dry_run_approval`` is configured — the loopback
|
||||
developer setup keeps working unchanged, and an operator turns this on before
|
||||
pointing the application at the library they cannot replace.
|
||||
"""
|
||||
if not config.require_dry_run_approval:
|
||||
return []
|
||||
path = approval_path(config)
|
||||
if not path.exists():
|
||||
return [
|
||||
{
|
||||
"code": "dry_run_not_approved",
|
||||
"message": (
|
||||
"run `python -m photo_pipeline dry-run` and approve its report "
|
||||
"before mutation is enabled"
|
||||
),
|
||||
}
|
||||
]
|
||||
try:
|
||||
record = json.loads(path.read_text())
|
||||
except ValueError:
|
||||
return [{"code": "approval_unreadable", "message": f"{path} is not readable JSON"}]
|
||||
approved_roots = [str(root) for root in record.get("library_roots", [])]
|
||||
configured = [str(root) for root in config.library_roots]
|
||||
if sorted(approved_roots) != sorted(configured):
|
||||
return [
|
||||
{
|
||||
"code": "approval_scope_mismatch",
|
||||
"message": (
|
||||
f"the approval covers {approved_roots}, but the configured library "
|
||||
f"is {configured}; run a new dry run"
|
||||
),
|
||||
}
|
||||
]
|
||||
return []
|
||||
@@ -294,6 +294,7 @@ class SafetyService:
|
||||
|
||||
exif_verified_at = None
|
||||
result_sha256 = None
|
||||
result_byte_size = None
|
||||
if write_exif and decision in (SFW, NSFW) and path:
|
||||
ops = exif_projection(decision)
|
||||
# The full checkpoint: write the owned keyword, read the whole file back,
|
||||
@@ -314,6 +315,7 @@ class SafetyService:
|
||||
if result.verified:
|
||||
exif_verified_at = result.verified_at
|
||||
result_sha256 = result.sha256
|
||||
result_byte_size = result.byte_size
|
||||
|
||||
now = _now()
|
||||
with self._session_factory() as session:
|
||||
@@ -332,6 +334,8 @@ class SafetyService:
|
||||
if result_sha256:
|
||||
asset = session.get(Asset, asset_id)
|
||||
asset.current_sha256 = result_sha256
|
||||
if result_byte_size is not None:
|
||||
asset.byte_size = result_byte_size
|
||||
session.commit()
|
||||
return {
|
||||
"asset_id": asset_id,
|
||||
|
||||
@@ -8,18 +8,32 @@ dependencies = [
|
||||
"sqlalchemy>=2.0",
|
||||
"alembic>=1.13",
|
||||
"pydantic>=2.7",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
test = [
|
||||
"pytest>=8",
|
||||
"httpx>=0.27",
|
||||
# Imaging is runtime, not test-only: thumbnails decode through Pillow, and the
|
||||
# perceptual hash is a DCT over the decoded pixels (services/hashing.py).
|
||||
"pillow>=10",
|
||||
"numpy>=1.26",
|
||||
"scipy>=1.11",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
# The cloud vision provider. Optional because the analysis stage is the only thing
|
||||
# that needs it, and a local review-only install should not pull an API client.
|
||||
vision = ["openai>=1.30"]
|
||||
test = [
|
||||
"pytest>=8",
|
||||
"httpx>=0.27",
|
||||
"playwright>=1.40",
|
||||
"pytest-playwright>=0.4",
|
||||
]
|
||||
|
||||
[build-system]
|
||||
requires = ["setuptools>=68"]
|
||||
build-backend = "setuptools.build_meta"
|
||||
|
||||
[tool.setuptools]
|
||||
# The importable application. ``migrations`` and ``work_item`` live beside it but
|
||||
# are not part of the package; without this, an editable install cannot guess.
|
||||
packages = ["photo_pipeline"]
|
||||
# Browser end-to-end tests also require: python -m playwright install chromium
|
||||
|
||||
[tool.ruff]
|
||||
@@ -35,4 +49,5 @@ markers = [
|
||||
"phase_d: Phase D end-to-end acceptance (US04-06) — guarded rename API, fault, and browser journeys",
|
||||
"phase_e: Phase E end-to-end acceptance (US05-06) — upload preflight, uploader, and browser journeys",
|
||||
"phase_f: Phase F end-to-end acceptance (US06-06) — archive destination, transfer, and restore journeys",
|
||||
"container: builds and runs the container image and its composition (US08-02, US08-03) — needs a Docker daemon, the compose plugin, and the network",
|
||||
]
|
||||
|
||||
337
tests/e2e/test_compose_stack.py
Normal file
337
tests/e2e/test_compose_stack.py
Normal file
@@ -0,0 +1,337 @@
|
||||
"""US08-03: the composition, actually composed.
|
||||
|
||||
Nothing here is faked below the process boundary: Docker builds the image, Compose
|
||||
starts the migrate/API/worker containers against a temporary fixture library on a
|
||||
real bind mount, and every assertion is made over HTTP or against what the stack
|
||||
left in its data volume. The vision provider is the deterministic fake seam the
|
||||
other end-to-end suites use, because an upload of real photos to a real model is not
|
||||
what this story is about — the mount, the lock, the volume, and the restart are.
|
||||
|
||||
The file contract (one API, one worker, migrations first, no committed values) is
|
||||
checked without a daemon in ``tests/integration/test_compose_runtime.py``; only the
|
||||
running proof needs Docker, and CI is where it runs unskipped (US08-04).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from PIL import Image
|
||||
|
||||
REPO = Path(__file__).resolve().parents[2]
|
||||
PROJECT = "photo-pipeline-us0803"
|
||||
IMAGE = "photo-pipeline-test:us08-03"
|
||||
SECRET = "compose-acceptance-secret"
|
||||
CONTAINER_LIBRARY = "/library"
|
||||
READY_TIMEOUT_SECONDS = 180
|
||||
JOB_TIMEOUT_SECONDS = 180
|
||||
UP_TIMEOUT_SECONDS = 30 * 60
|
||||
|
||||
pytestmark = pytest.mark.container
|
||||
|
||||
|
||||
def compose_available() -> bool:
|
||||
try:
|
||||
return (
|
||||
subprocess.run(
|
||||
["docker", "compose", "version"], capture_output=True, timeout=60
|
||||
).returncode
|
||||
== 0
|
||||
)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return False
|
||||
|
||||
|
||||
needs_compose = pytest.mark.skipif(
|
||||
not compose_available(), reason="no Docker daemon with the compose plugin"
|
||||
)
|
||||
|
||||
|
||||
def free_port() -> int:
|
||||
with socket.socket() as sock:
|
||||
sock.bind(("127.0.0.1", 0))
|
||||
return sock.getsockname()[1]
|
||||
|
||||
|
||||
class Stack:
|
||||
"""The composition under test, plus the environment it was started with."""
|
||||
|
||||
def __init__(self, library: Path, env_file: Path) -> None:
|
||||
self.library = library
|
||||
self.port = free_port()
|
||||
self.base = f"http://127.0.0.1:{self.port}"
|
||||
# Compose reads the repository's own .env for substitution; the process
|
||||
# environment wins over it, so the test's values are the ones that apply.
|
||||
self.env = {
|
||||
**os.environ,
|
||||
"PHOTO_PIPELINE_IMAGE": IMAGE,
|
||||
"PHOTO_PIPELINE_ENV_FILE": str(env_file),
|
||||
"PHOTO_PIPELINE_LIBRARY_HOST_PATH": str(library),
|
||||
"PHOTO_PIPELINE_LIBRARY_ROOTS": CONTAINER_LIBRARY,
|
||||
"PHOTO_PIPELINE_PORT": str(self.port),
|
||||
"PHOTO_PIPELINE_UID": str(os.getuid()),
|
||||
"PHOTO_PIPELINE_GID": str(os.getgid()),
|
||||
}
|
||||
|
||||
def compose(self, *args: str, check: bool = True, timeout: int = 300):
|
||||
result = subprocess.run(
|
||||
["docker", "compose", "-p", PROJECT, "-f", str(REPO / "docker-compose.yml"), *args],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env=self.env,
|
||||
cwd=REPO,
|
||||
timeout=timeout,
|
||||
)
|
||||
if check and result.returncode != 0:
|
||||
raise AssertionError(
|
||||
f"docker compose {' '.join(args)} failed:\n{result.stdout}\n{result.stderr}\n"
|
||||
f"{self.compose('logs', '--tail', '80', check=False).stdout}"
|
||||
)
|
||||
return result
|
||||
|
||||
def wait_until_ready(self) -> None:
|
||||
deadline = time.monotonic() + READY_TIMEOUT_SECONDS
|
||||
while time.monotonic() < deadline:
|
||||
try:
|
||||
if httpx.get(f"{self.base}/api/v1/health/ready", timeout=5).status_code == 200:
|
||||
return
|
||||
except httpx.HTTPError:
|
||||
pass
|
||||
time.sleep(0.5)
|
||||
logs = self.compose("logs", "--tail", "120", check=False)
|
||||
raise AssertionError(f"the stack never became ready:\n{logs.stdout}\n{logs.stderr}")
|
||||
|
||||
def client(self) -> httpx.Client:
|
||||
client = httpx.Client(base_url=f"{self.base}/api/v1", timeout=60)
|
||||
bootstrap = client.get("/session", headers={"X-Access-Secret": SECRET})
|
||||
assert bootstrap.status_code == 200, bootstrap.text
|
||||
client.headers["X-CSRF-Token"] = bootstrap.json()["csrf_token"]
|
||||
return client
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def library() -> Path:
|
||||
"""A small fixture library on the host, mounted into both containers.
|
||||
|
||||
Not under pytest's ``tmp_path``: on macOS that is ``/var/folders/...``, which a
|
||||
Docker VM (Colima, Docker Desktop) does not share, so the bind mount would arrive
|
||||
empty and every assertion below would be about nothing. ``$HOME`` is shared by
|
||||
every default configuration.
|
||||
"""
|
||||
base = Path(
|
||||
os.environ.get("PHOTO_PIPELINE_TEST_MOUNT_BASE", Path.home() / ".cache" / "photo-pipeline")
|
||||
)
|
||||
base.mkdir(parents=True, exist_ok=True)
|
||||
root = Path(tempfile.mkdtemp(prefix="library-", dir=base))
|
||||
for album, count in (("01_day", 2), ("02_night", 1)):
|
||||
(root / album).mkdir()
|
||||
for index in range(count):
|
||||
colour = (40 * (index + 1), 90, 160)
|
||||
Image.new("RGB", (64, 48), colour).save(root / album / f"{album}_{index}.jpg")
|
||||
# The exclusion sentinel: it must never be discovered, counted, or analyzed.
|
||||
(root / "_IGNORE").mkdir()
|
||||
Image.new("RGB", (32, 32), (0, 0, 0)).save(root / "_IGNORE" / "sentinel.jpg")
|
||||
try:
|
||||
yield root
|
||||
finally:
|
||||
shutil.rmtree(root, ignore_errors=True)
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def env_file(tmp_path_factory) -> Path:
|
||||
"""Configuration and secrets come from the environment, so the test writes its
|
||||
own file rather than borrowing the operator's."""
|
||||
path = tmp_path_factory.mktemp("config") / "compose.env"
|
||||
path.write_text(
|
||||
"\n".join(
|
||||
[
|
||||
f"PHOTO_PIPELINE_ACCESS_SECRET={SECRET}",
|
||||
"PHOTO_PIPELINE_LOG_FORMAT=text",
|
||||
# The deterministic vision seam, in the data volume so both roles and
|
||||
# the test can see it (concept §18).
|
||||
"PHOTO_PIPELINE_FAKE_VISION_LOG=/data/vision.log",
|
||||
]
|
||||
)
|
||||
+ "\n"
|
||||
)
|
||||
return path
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def stack(library, env_file):
|
||||
if not compose_available():
|
||||
pytest.skip("no Docker daemon with the compose plugin")
|
||||
running = Stack(library, env_file)
|
||||
running.compose("down", "--volumes", "--remove-orphans", check=False)
|
||||
running.compose("up", "--detach", "--build", timeout=UP_TIMEOUT_SECONDS)
|
||||
try:
|
||||
running.wait_until_ready()
|
||||
yield running
|
||||
finally:
|
||||
running.compose("down", "--volumes", "--remove-orphans", check=False, timeout=300)
|
||||
|
||||
|
||||
def await_job(client: httpx.Client, job_id: str, states=("succeeded",)) -> dict:
|
||||
deadline = time.monotonic() + JOB_TIMEOUT_SECONDS
|
||||
snapshot: dict = {}
|
||||
while time.monotonic() < deadline:
|
||||
response = client.get(f"/jobs/{job_id}")
|
||||
if response.status_code == 200:
|
||||
snapshot = response.json()
|
||||
if snapshot["state"] in states:
|
||||
return snapshot
|
||||
time.sleep(0.5)
|
||||
raise AssertionError(f"job {job_id} never reached {states}: {snapshot}")
|
||||
|
||||
|
||||
# ── the mounted library ──────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_the_stack_scans_the_bind_mounted_library_at_its_container_paths(stack):
|
||||
client = stack.client()
|
||||
try:
|
||||
scanned = client.post("/inventory/scan")
|
||||
assert scanned.status_code == 200, scanned.text
|
||||
assets = client.get("/inventory/assets", params={"limit": 200}).json()["items"]
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
assert len(assets) == 3, assets
|
||||
paths = {asset["current_path"] for asset in assets}
|
||||
assert all(path.startswith(CONTAINER_LIBRARY + "/") for path in paths), paths
|
||||
assert not any("_IGNORE" in path or "sentinel" in path for path in paths)
|
||||
# The host paths are what the operator mounted, and they are not what the
|
||||
# application records: the roots are the container's.
|
||||
assert not any(str(stack.library) in path for path in paths)
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_a_library_root_that_is_not_mounted_is_refused_at_startup(stack):
|
||||
"""The container-specific failure: configured roots that name nothing mounted."""
|
||||
refused = stack.compose(
|
||||
"run",
|
||||
"--rm",
|
||||
"--no-deps",
|
||||
"--env",
|
||||
"PHOTO_PIPELINE_LIBRARY_ROOTS=/srv/photos",
|
||||
"api",
|
||||
"serve",
|
||||
check=False,
|
||||
)
|
||||
assert refused.returncode == 5, refused.stdout + refused.stderr
|
||||
assert "library root /srv/photos" in refused.stdout + refused.stderr
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_a_second_worker_is_refused_by_the_library_lock(stack):
|
||||
"""Not by convention: the running worker's lock is in the shared data volume."""
|
||||
refused = stack.compose(
|
||||
"run", "--rm", "--no-deps", "worker", "worker", "--id", "worker-2", check=False
|
||||
)
|
||||
assert refused.returncode == 2, refused.stdout + refused.stderr
|
||||
assert "worker is already running" in refused.stdout + refused.stderr
|
||||
|
||||
|
||||
# ── restart, resume, and the data volume ─────────────────────────────────────
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_a_queued_job_resumes_after_both_containers_restart(stack):
|
||||
client = stack.client()
|
||||
try:
|
||||
client.post("/inventory/scan").raise_for_status()
|
||||
assets = client.get("/inventory/assets", params={"limit": 200}).json()["items"]
|
||||
for asset in assets:
|
||||
decided = client.post(
|
||||
"/safety/decisions", json={"asset_id": asset["id"], "decision": "sfw"}
|
||||
)
|
||||
assert decided.status_code == 200, decided.text
|
||||
|
||||
# Stop the worker first, so the job is provably still queued when the restart
|
||||
# happens: a job that finished before the restart would prove nothing.
|
||||
stack.compose("stop", "worker")
|
||||
job = client.post("/analysis/jobs").json()
|
||||
assert client.get(f"/jobs/{job['id']}").json()["state"] == "queued"
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
stack.compose("restart", "api", "worker")
|
||||
stack.wait_until_ready()
|
||||
|
||||
client = stack.client() # the session is per API process, so it is re-bootstrapped
|
||||
try:
|
||||
finished = await_job(client, job["id"])
|
||||
assert finished["state"] == "succeeded", finished
|
||||
# The database is intact and the work is durable, not merely reported.
|
||||
after = client.get("/inventory/assets", params={"limit": 200}).json()["items"]
|
||||
assert {asset["id"] for asset in after} == {asset["id"] for asset in assets}
|
||||
analysed = client.get(f"/analysis/results/{assets[0]['id']}")
|
||||
assert analysed.status_code == 200, analysed.text
|
||||
assert analysed.json()["description"]
|
||||
# Only the mounted library's own photos were analysed: the sentinel under
|
||||
# _IGNORE is not an asset, so it can never have become an item of this job.
|
||||
assert finished["progress"]["total"] == len(assets)
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_the_data_volume_survives_recreating_the_containers(stack):
|
||||
"""`down` without `--volumes` then `up` is the upgrade path: state stays."""
|
||||
client = stack.client()
|
||||
try:
|
||||
client.post("/inventory/scan").raise_for_status()
|
||||
before = {a["id"] for a in client.get("/inventory/assets").json()["items"]}
|
||||
finally:
|
||||
client.close()
|
||||
|
||||
stack.compose("down", "--remove-orphans", timeout=300)
|
||||
stack.compose("up", "--detach", timeout=UP_TIMEOUT_SECONDS)
|
||||
stack.wait_until_ready()
|
||||
|
||||
client = stack.client()
|
||||
try:
|
||||
after = {a["id"] for a in client.get("/inventory/assets").json()["items"]}
|
||||
finally:
|
||||
client.close()
|
||||
assert after == before, "the same assets, from the same database, on the same volume"
|
||||
|
||||
|
||||
# ── operating it ─────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@needs_compose
|
||||
def test_backup_verify_and_diagnostics_run_as_container_commands(stack):
|
||||
backup = stack.compose("run", "--rm", "--no-deps", "api", "backup", "--reason", "compose")
|
||||
manifest = json.loads(backup.stdout[backup.stdout.index("{") :])
|
||||
assert manifest["name"].startswith("2")
|
||||
|
||||
verified = stack.compose(
|
||||
"run", "--rm", "--no-deps", "api", "verify-backup", f"/data/backups/{manifest['name']}"
|
||||
)
|
||||
assert json.loads(verified.stdout[verified.stdout.index("{") :])["ok"] is True
|
||||
|
||||
report = stack.compose("run", "--rm", "--no-deps", "api", "diagnostics")
|
||||
diagnostics = json.loads(report.stdout[report.stdout.index("{") :])
|
||||
components = {c["name"]: c["path"] for c in diagnostics["components"]}
|
||||
# Database, WAL, thumbnail cache, and backups all live in the mounted volume.
|
||||
for name in ("database", "write_ahead_log", "thumbnail_cache", "backups"):
|
||||
assert components[name].startswith("/data/"), (name, components[name])
|
||||
assert {tool["name"] for tool in diagnostics["tools"]} >= {"exiftool", "immich-go"}
|
||||
# And the worker running beside this command is visible as the lock's holder.
|
||||
assert diagnostics["locks"]["worker"]["role"] == "worker"
|
||||
|
||||
|
||||
if __name__ == "__main__": # a quick way to run just this file
|
||||
raise SystemExit(pytest.main([__file__, "-v", *sys.argv[1:]]))
|
||||
306
tests/e2e/test_container_runtime.py
Normal file
306
tests/e2e/test_container_runtime.py
Normal file
@@ -0,0 +1,306 @@
|
||||
"""US08-02: the built image, actually built and actually run.
|
||||
|
||||
This is the acceptance test for the image itself, so nothing here is faked: Docker
|
||||
builds from a clean context, the container starts under a chosen UID/GID against a
|
||||
mounted data directory, and the assertions are made over HTTP and against the files
|
||||
the container left on the host.
|
||||
|
||||
It is skipped without a Docker daemon — the build also needs the network for the base
|
||||
image, the pinned exiftool package, and the pinned uploader release. The contract the
|
||||
Dockerfile itself has to keep (pins, non-root, health target, build context) is checked
|
||||
offline in ``tests/integration/test_container_image.py``, so a machine without Docker
|
||||
still fails on a broken image definition; only the running proof needs the daemon. CI
|
||||
builds the image on every change (US08-04), which is where this runs unskipped.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import platform
|
||||
import re
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
REPO = Path(__file__).resolve().parents[2]
|
||||
IMAGE = "photo-pipeline-test:us08-02"
|
||||
SECRET = "container-acceptance-secret"
|
||||
HOSTNAME = "photos.test"
|
||||
READY_TIMEOUT_SECONDS = 120
|
||||
BUILD_TIMEOUT_SECONDS = 30 * 60
|
||||
|
||||
pytestmark = pytest.mark.container
|
||||
|
||||
|
||||
def docker_available() -> bool:
|
||||
try:
|
||||
return subprocess.run(["docker", "info"], capture_output=True, timeout=60).returncode == 0
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return False
|
||||
|
||||
|
||||
needs_docker = pytest.mark.skipif(not docker_available(), reason="no Docker daemon available")
|
||||
|
||||
|
||||
def docker(*args: str, check: bool = True, timeout: int = 120) -> subprocess.CompletedProcess:
|
||||
result = subprocess.run(
|
||||
["docker", *args], capture_output=True, text=True, timeout=timeout
|
||||
)
|
||||
if check and result.returncode != 0:
|
||||
raise AssertionError(f"docker {' '.join(args)} failed:\n{result.stdout}\n{result.stderr}")
|
||||
return result
|
||||
|
||||
|
||||
def pins() -> dict[str, str]:
|
||||
"""The pinned versions, read from the Dockerfile that produced the image."""
|
||||
text = (REPO / "Dockerfile").read_text()
|
||||
found = dict(re.findall(r"^ARG\s+([A-Z0-9_]+)=(.+)$", text, re.MULTILINE))
|
||||
return {
|
||||
# The Debian package version carries a packaging suffix; exiftool reports the
|
||||
# upstream version only.
|
||||
"exiftool": found["EXIFTOOL_VERSION"].split("+")[0].split("-")[0],
|
||||
"immich-go": found["IMMICH_GO_VERSION"],
|
||||
}
|
||||
|
||||
|
||||
def free_port() -> int:
|
||||
with socket.socket() as sock:
|
||||
sock.bind(("127.0.0.1", 0))
|
||||
return sock.getsockname()[1]
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def image() -> str:
|
||||
"""Build from a clean checkout: the build context is the repository, unmodified."""
|
||||
if not docker_available():
|
||||
pytest.skip("no Docker daemon available")
|
||||
docker(
|
||||
"build",
|
||||
"--build-arg",
|
||||
f"UID={os.getuid()}",
|
||||
"--build-arg",
|
||||
f"GID={os.getgid()}",
|
||||
"-t",
|
||||
IMAGE,
|
||||
str(REPO),
|
||||
timeout=BUILD_TIMEOUT_SECONDS,
|
||||
)
|
||||
return IMAGE
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def data_dir(tmp_path) -> Path:
|
||||
data = tmp_path / "data"
|
||||
data.mkdir()
|
||||
return data
|
||||
|
||||
|
||||
def run_detached(image: str, port: int, *args: str, data: Path | None = None) -> str:
|
||||
"""Start a container. ``data`` bind-mounts the host's data directory when the test
|
||||
is about the files themselves; otherwise the image's own /data is used, because a
|
||||
macOS bind mount arrives with an ownership the container did not choose."""
|
||||
result = docker(
|
||||
"run",
|
||||
"--detach",
|
||||
"--rm",
|
||||
"--publish",
|
||||
f"127.0.0.1:{port}:8000",
|
||||
*(("--volume", f"{data}:/data") if data is not None else ()),
|
||||
"--env",
|
||||
# Reachable from outside the container means reachable from another machine as
|
||||
# far as the application is concerned, so the access secret is mandatory
|
||||
# (US08-01) — the image must not weaken that.
|
||||
"PHOTO_PIPELINE_HOST=0.0.0.0",
|
||||
"--env",
|
||||
f"PHOTO_PIPELINE_ACCESS_SECRET={SECRET}",
|
||||
"--env",
|
||||
f"PHOTO_PIPELINE_ALLOWED_HOSTS={HOSTNAME}",
|
||||
image,
|
||||
*args,
|
||||
)
|
||||
return result.stdout.strip()
|
||||
|
||||
|
||||
def wait_until_ready(base: str, container: str) -> None:
|
||||
deadline = time.monotonic() + READY_TIMEOUT_SECONDS
|
||||
while time.monotonic() < deadline:
|
||||
try:
|
||||
if httpx.get(f"{base}/api/v1/health/ready", timeout=5).status_code == 200:
|
||||
return
|
||||
except httpx.HTTPError:
|
||||
pass
|
||||
if docker("inspect", "-f", "{{.State.Running}}", container, check=False).stdout.strip() in (
|
||||
"false",
|
||||
"",
|
||||
):
|
||||
break
|
||||
time.sleep(0.5)
|
||||
logs = docker("logs", container, check=False)
|
||||
raise AssertionError(f"container never became ready:\n{logs.stdout}\n{logs.stderr}")
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def serving(image):
|
||||
port = free_port()
|
||||
container = run_detached(image, port, "serve")
|
||||
try:
|
||||
base = f"http://127.0.0.1:{port}"
|
||||
wait_until_ready(base, container)
|
||||
yield base, container
|
||||
finally:
|
||||
docker("rm", "--force", container, check=False)
|
||||
|
||||
|
||||
def session(base: str) -> httpx.Client:
|
||||
client = httpx.Client(base_url=base, timeout=30)
|
||||
bootstrap = client.get("/api/v1/session", headers={"X-Access-Secret": SECRET})
|
||||
assert bootstrap.status_code == 200, bootstrap.text
|
||||
client.headers["X-CSRF-Token"] = bootstrap.json()["csrf_token"]
|
||||
return client
|
||||
|
||||
|
||||
# ── the image serves, and says what it contains ──────────────────────────────
|
||||
|
||||
|
||||
@needs_docker
|
||||
def test_the_container_serves_the_frontend_and_the_pinned_tool_versions(serving):
|
||||
base, container = serving
|
||||
|
||||
index = httpx.get(f"{base}/app/index.html", timeout=30)
|
||||
assert index.status_code == 200
|
||||
assert "<title" in index.text.lower(), "the application shell, not an API error"
|
||||
|
||||
client = session(base)
|
||||
try:
|
||||
tools = {tool["name"]: tool for tool in client.get("/api/v1/diagnostics").json()["tools"]}
|
||||
finally:
|
||||
client.close()
|
||||
for name, pinned in pins().items():
|
||||
assert tools[name]["pinned"] == pinned, name
|
||||
# Recorded *and* installed: the reported version comes from running the binary.
|
||||
assert pinned in tools[name]["version"], (name, tools[name])
|
||||
assert tools[name]["path"], f"{name} is not on PATH inside the image"
|
||||
|
||||
logs = docker("logs", container, check=False)
|
||||
assert SECRET not in logs.stdout + logs.stderr, "the access secret never reaches the log"
|
||||
|
||||
|
||||
@needs_docker
|
||||
def test_the_declared_health_check_reports_readiness(serving):
|
||||
"""The declared HEALTHCHECK is readiness, so Docker's own verdict is the assertion."""
|
||||
_, container = serving
|
||||
deadline = time.monotonic() + READY_TIMEOUT_SECONDS
|
||||
status = ""
|
||||
while time.monotonic() < deadline:
|
||||
status = docker(
|
||||
"inspect", "-f", "{{.State.Health.Status}}", container, check=False
|
||||
).stdout.strip()
|
||||
if status == "healthy":
|
||||
break
|
||||
time.sleep(1)
|
||||
assert status == "healthy"
|
||||
|
||||
probe = docker("exec", container, "/usr/local/bin/healthcheck.sh", check=False)
|
||||
assert probe.returncode == 0
|
||||
# Point the probe at a port nothing serves: the same script must fail, which is
|
||||
# what makes the healthy verdict above evidence rather than a default.
|
||||
unready = docker(
|
||||
"exec",
|
||||
"--env",
|
||||
"PHOTO_PIPELINE_PORT=1",
|
||||
container,
|
||||
"/usr/local/bin/healthcheck.sh",
|
||||
check=False,
|
||||
)
|
||||
assert unready.returncode != 0
|
||||
|
||||
|
||||
# ── identity: never root, always the configured owner ────────────────────────
|
||||
|
||||
|
||||
@needs_docker
|
||||
def test_the_container_refuses_to_run_as_root(image, data_dir):
|
||||
result = docker(
|
||||
"run",
|
||||
"--rm",
|
||||
"--user",
|
||||
"0:0",
|
||||
"--volume",
|
||||
f"{data_dir}:/data",
|
||||
image,
|
||||
"diagnostics",
|
||||
check=False,
|
||||
)
|
||||
assert result.returncode != 0
|
||||
assert "refusing to run as root" in result.stderr + result.stdout
|
||||
assert not list(data_dir.iterdir()), "a refused container writes nothing"
|
||||
|
||||
|
||||
@needs_docker
|
||||
def test_what_the_container_writes_is_owned_by_the_build_arguments(image):
|
||||
"""The identity the image was built with is the identity on disk afterwards.
|
||||
|
||||
Asserted from inside the container so it holds on every host: a macOS bind mount
|
||||
reports an ownership the container never chose. The host-side proof, which is what
|
||||
the mounted library actually needs, is the Linux test below.
|
||||
"""
|
||||
port = free_port()
|
||||
container = run_detached(image, port, "serve")
|
||||
try:
|
||||
wait_until_ready(f"http://127.0.0.1:{port}", container)
|
||||
owner = docker(
|
||||
"exec", container, "stat", "-c", "%u:%g", "/data/photo_pipeline.db"
|
||||
).stdout.strip()
|
||||
assert owner == f"{os.getuid()}:{os.getgid()}"
|
||||
assert docker("exec", container, "id", "-u").stdout.strip() == str(os.getuid())
|
||||
finally:
|
||||
docker("rm", "--force", container, check=False)
|
||||
|
||||
|
||||
@needs_docker
|
||||
@pytest.mark.skipif(
|
||||
platform.system() != "Linux",
|
||||
reason="bind-mount ownership is virtualised by Docker Desktop on macOS/Windows",
|
||||
)
|
||||
def test_files_the_container_writes_keep_the_configured_ownership(image, data_dir):
|
||||
docker("run", "--rm", "--volume", f"{data_dir}:/data", image, "migrate", timeout=300)
|
||||
|
||||
written = sorted(path for path in data_dir.rglob("*") if path.is_file())
|
||||
assert written, "migrate creates the database in the mounted data directory"
|
||||
for path in written:
|
||||
assert (path.stat().st_uid, path.stat().st_gid) == (os.getuid(), os.getgid()), path
|
||||
|
||||
|
||||
@needs_docker
|
||||
def test_the_worker_role_runs_from_the_same_image(image, data_dir):
|
||||
"""One image, two roles: the worker is the same entrypoint with another argument."""
|
||||
port = free_port()
|
||||
container = run_detached(image, port, "worker", "--id", "container-worker")
|
||||
try:
|
||||
# Taking the worker's library lock is the observable proof that it started,
|
||||
# migrated, and reached its job loop — no sleep required (US07-05).
|
||||
deadline = time.monotonic() + READY_TIMEOUT_SECONDS
|
||||
lock = ""
|
||||
while not lock and time.monotonic() < deadline:
|
||||
assert docker("inspect", "-f", "{{.State.Running}}", container).stdout.strip() == (
|
||||
"true"
|
||||
), docker("logs", container, check=False).stdout
|
||||
lock = docker("exec", container, "cat", "/data/worker.lock.json", check=False).stdout
|
||||
time.sleep(0.5)
|
||||
assert lock, docker("logs", container, check=False).stdout
|
||||
assert json.loads(lock)["role"] == "worker"
|
||||
|
||||
role = docker("exec", container, "cat", "/tmp/photo-pipeline-role").stdout.strip()
|
||||
assert role == "worker", "the health check can tell which role this container is"
|
||||
finally:
|
||||
docker("rm", "--force", container, check=False)
|
||||
|
||||
|
||||
if __name__ == "__main__": # a quick way to run just this file
|
||||
raise SystemExit(pytest.main([__file__, "-v", *sys.argv[1:]]))
|
||||
298
tests/e2e/test_release_gate.py
Normal file
298
tests/e2e/test_release_gate.py
Normal file
@@ -0,0 +1,298 @@
|
||||
"""The release gate, the read-only dry run, and the approval that unlocks mutation
|
||||
(US07-07).
|
||||
|
||||
The gate itself is exercised with a tiny stage set — running the whole suite from
|
||||
inside the suite would be a fork bomb with better manners. What is proven here is
|
||||
the machinery a release depends on: the story matrix is complete, a failing stage
|
||||
fails the gate, an unexpected skip fails the gate, and the evidence is written with
|
||||
checksums that match what was written.
|
||||
|
||||
The dry run is proven to be read-only against a real temporary library, and the
|
||||
approval is proven to be what stands between a configured library and any mutation.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import uuid
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from photo_pipeline.api.app import create_app
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.db import create_db_engine, create_session_factory, run_migrations
|
||||
from photo_pipeline.models import Asset
|
||||
from photo_pipeline.services import release
|
||||
|
||||
REPO = Path(__file__).resolve().parents[2]
|
||||
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
|
||||
|
||||
# Two throwaway stages: one that passes, one the test can point at a failure.
|
||||
PASSING = ("tests/e2e/test_traceability.py",)
|
||||
|
||||
|
||||
def _config(tmp_path, **extra) -> Config:
|
||||
data = tmp_path / "data"
|
||||
data.mkdir(parents=True, exist_ok=True)
|
||||
lib = tmp_path / "lib"
|
||||
lib.mkdir(exist_ok=True)
|
||||
return Config.from_env(
|
||||
{
|
||||
"PHOTO_PIPELINE_DATA_DIR": str(data),
|
||||
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib),
|
||||
**extra,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
# ── the story matrix ─────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_every_backlog_story_is_delivered_or_explicitly_planned():
|
||||
matrix = release.story_matrix(REPO)
|
||||
|
||||
assert matrix["problems"] == [], "the story matrix has holes"
|
||||
assert len(matrix["delivered"]) + len(matrix["planned"]) == matrix["stories"]
|
||||
assert "US01-01" in matrix["delivered"] and "US07-07" in matrix["delivered"]
|
||||
|
||||
|
||||
def test_a_story_without_tests_is_a_gate_failure(tmp_path):
|
||||
"""A story file nobody covered must not pass quietly as 'no tests ran'."""
|
||||
fake = tmp_path / "repo"
|
||||
(fake / "delivery_backlog" / "stories").mkdir(parents=True)
|
||||
(fake / "tests").mkdir()
|
||||
(fake / "delivery_backlog" / "stories" / "US99-01-invented.md").write_text("# US99-01")
|
||||
(fake / "tests" / "story_traceability.json").write_text(json.dumps({"stories": {}}))
|
||||
|
||||
matrix = release.story_matrix(fake)
|
||||
|
||||
assert matrix["problems"] == ["story with no tests and not planned: US99-01"]
|
||||
|
||||
|
||||
def test_a_mapping_to_a_deleted_test_file_is_a_gate_failure(tmp_path):
|
||||
fake = tmp_path / "repo"
|
||||
(fake / "delivery_backlog" / "stories").mkdir(parents=True)
|
||||
(fake / "tests").mkdir()
|
||||
(fake / "delivery_backlog" / "stories" / "US99-01-invented.md").write_text("# US99-01")
|
||||
(fake / "tests" / "story_traceability.json").write_text(
|
||||
json.dumps({"stories": {"US99-01": ["tests/gone.py"]}})
|
||||
)
|
||||
|
||||
assert release.story_matrix(fake)["problems"] == [
|
||||
"mapped test file is missing — US99-01: tests/gone.py"
|
||||
]
|
||||
|
||||
|
||||
# ── the gate ─────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_the_gate_runs_its_stages_and_keeps_checksummed_evidence(tmp_path):
|
||||
config = _config(tmp_path)
|
||||
evidence = tmp_path / "evidence"
|
||||
|
||||
report = release.run_gate(config, output=evidence, stages=(("smoke", PASSING),))
|
||||
|
||||
assert report["ok"] is True and report["failures"] == []
|
||||
assert report["stages"][0]["stage"] == "smoke" and report["stages"][0]["ok"] is True
|
||||
assert report["revision"], "the evidence must say which commit it covers"
|
||||
assert report["matrix"]["problems"] == []
|
||||
|
||||
written = json.loads((evidence / release.REPORT_NAME).read_text())
|
||||
assert written["ok"] is True
|
||||
assert (evidence / "logs" / "smoke.log").exists()
|
||||
checksums = (evidence / release.CHECKSUMS_NAME).read_text().splitlines()
|
||||
assert len(checksums) >= 2
|
||||
for line in checksums:
|
||||
digest, name = line.split(" ", 1)
|
||||
assert release.sha256_file(evidence / name) == digest
|
||||
|
||||
|
||||
def test_a_failing_stage_fails_the_gate(tmp_path):
|
||||
config = _config(tmp_path)
|
||||
failing = tmp_path / "failing_test.py"
|
||||
failing.write_text("def test_no():\n assert False\n")
|
||||
|
||||
report = release.run_gate(
|
||||
config, output=tmp_path / "evidence", stages=(("broken", (str(failing),)),)
|
||||
)
|
||||
|
||||
assert report["ok"] is False and report["failures"] == ["broken"]
|
||||
assert report["stages"][0]["returncode"] != 0
|
||||
|
||||
|
||||
def test_an_unexpected_skip_fails_the_gate_but_an_environment_skip_does_not():
|
||||
environment = release.StageResult(
|
||||
"unit", [], 0, 0.1, "1 skipped", ["SKIPPED [1] x.py:1: exiftool not installed"]
|
||||
)
|
||||
silent = release.StageResult(
|
||||
"unit", [], 0, 0.1, "1 skipped", ["SKIPPED [1] x.py:1: flaky, look later"]
|
||||
)
|
||||
|
||||
assert release.unexpected_skips([environment]) == []
|
||||
assert release.unexpected_skips([silent, environment]) == [
|
||||
"SKIPPED [1] x.py:1: flaky, look later"
|
||||
]
|
||||
|
||||
|
||||
# ── the real-library dry run ─────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _library(root: Path) -> None:
|
||||
(root / "album").mkdir(parents=True)
|
||||
(root / "album" / "a.jpg").write_bytes(b"a" * 128)
|
||||
(root / "album" / "b.png").write_bytes(b"b" * 64)
|
||||
(root / "loose.JPG").write_bytes(b"c" * 32)
|
||||
excluded = root / "_IGNORE" / "private"
|
||||
excluded.mkdir(parents=True)
|
||||
(excluded / "secret.jpg").write_bytes(b"never read")
|
||||
|
||||
|
||||
def test_the_dry_run_describes_the_library_without_touching_it(tmp_path):
|
||||
config = _config(tmp_path)
|
||||
root = Path(config.library_roots[0])
|
||||
_library(root)
|
||||
before = {
|
||||
str(p): (p.stat().st_mtime_ns, p.read_bytes()) for p in root.rglob("*") if p.is_file()
|
||||
}
|
||||
|
||||
report = release.dry_run(config)
|
||||
|
||||
assert report["files"] == 3, "the excluded sentinel is not counted"
|
||||
assert report["by_extension"] == {".jpg": 2, ".png": 1}
|
||||
assert report["excluded_directories"] >= 1
|
||||
assert report["mutation"] == "none — this pass is read-only"
|
||||
assert report["checksum"]
|
||||
assert not any("secret" in json.dumps(report) for _ in [0]), "excluded content never appears"
|
||||
after = {
|
||||
str(p): (p.stat().st_mtime_ns, p.read_bytes()) for p in root.rglob("*") if p.is_file()
|
||||
}
|
||||
assert after == before, "a read-only pass changed the library"
|
||||
|
||||
|
||||
def test_the_dry_run_reconciles_against_what_the_database_already_knows(tmp_path):
|
||||
config = _config(tmp_path)
|
||||
root = Path(config.library_roots[0])
|
||||
_library(root)
|
||||
run_migrations(config.database_url)
|
||||
engine = create_db_engine(config.database_url)
|
||||
with create_session_factory(engine)() as session:
|
||||
session.add(
|
||||
Asset(
|
||||
id=str(uuid.uuid4()),
|
||||
original_path=str(root / "album" / "a.jpg"),
|
||||
current_path=str(root / "album" / "a.jpg"),
|
||||
discovered_at=NOW,
|
||||
hash_version=1,
|
||||
byte_size=128,
|
||||
)
|
||||
)
|
||||
session.add(
|
||||
Asset(
|
||||
id=str(uuid.uuid4()),
|
||||
original_path=str(root / "album" / "gone.jpg"),
|
||||
current_path=str(root / "album" / "gone.jpg"),
|
||||
discovered_at=NOW,
|
||||
hash_version=1,
|
||||
byte_size=1,
|
||||
)
|
||||
)
|
||||
session.commit()
|
||||
engine.dispose()
|
||||
|
||||
reconciliation = release.dry_run(config)["reconciliation"]
|
||||
|
||||
assert reconciliation["known_to_database"] == 2
|
||||
assert reconciliation["already_registered"] == 1
|
||||
assert reconciliation["new_to_the_application"] == 2
|
||||
assert reconciliation["recorded_but_absent_total"] == 1
|
||||
assert reconciliation["recorded_but_absent"][0].endswith("gone.jpg")
|
||||
|
||||
|
||||
def test_a_library_root_that_is_not_there_is_refused(tmp_path):
|
||||
config = _config(tmp_path, PHOTO_PIPELINE_LIBRARY_ROOTS=str(tmp_path / "nowhere"))
|
||||
with pytest.raises(release.ReleaseError, match="not a directory"):
|
||||
release.dry_run(config)
|
||||
|
||||
|
||||
# ── the approval ─────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_mutation_is_refused_until_the_dry_run_is_approved(tmp_path):
|
||||
config = _config(tmp_path, PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL="1")
|
||||
_library(Path(config.library_roots[0]))
|
||||
|
||||
blockers = release.mutation_blockers(config)
|
||||
assert [blocker["code"] for blocker in blockers] == ["dry_run_not_approved"]
|
||||
|
||||
record = release.approve(config, release.dry_run(config), approver="domverse")
|
||||
assert record["approved_by"] == "domverse" and record["report_checksum"]
|
||||
assert release.mutation_blockers(config) == []
|
||||
|
||||
|
||||
def test_an_approval_covers_the_library_it_was_written_for(tmp_path):
|
||||
config = _config(tmp_path, PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL="1")
|
||||
_library(Path(config.library_roots[0]))
|
||||
release.approve(config, release.dry_run(config), approver="domverse")
|
||||
|
||||
other = tmp_path / "other-library"
|
||||
other.mkdir()
|
||||
moved = config.model_copy(update={"library_roots": (other,)})
|
||||
|
||||
assert [b["code"] for b in release.mutation_blockers(moved)] == ["approval_scope_mismatch"]
|
||||
|
||||
|
||||
def test_without_the_requirement_nothing_changes(tmp_path):
|
||||
config = _config(tmp_path) # the loopback development default
|
||||
assert release.mutation_blockers(config) == []
|
||||
|
||||
|
||||
def test_the_api_refuses_every_mutation_until_the_report_is_approved(tmp_path):
|
||||
config = _config(tmp_path, PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL="1")
|
||||
_library(Path(config.library_roots[0]))
|
||||
with TestClient(create_app(config)) as client:
|
||||
# Reading stays open: an operator has to see what was found to approve it.
|
||||
assert client.get("/api/v1/workflow").status_code == 200
|
||||
refused = client.post("/api/v1/inventory/scan", json={})
|
||||
assert refused.status_code == 403
|
||||
assert refused.json()["error"]["code"] == "dry_run_not_approved"
|
||||
# A backup is the one mutation a careful operator takes first.
|
||||
assert client.post("/api/v1/backups", json={}).status_code == 201
|
||||
|
||||
release.approve(config, release.dry_run(config), approver="domverse")
|
||||
assert client.post("/api/v1/inventory/scan", json={}).status_code in (200, 201, 202)
|
||||
|
||||
|
||||
def test_approving_something_that_is_not_a_report_is_refused(tmp_path):
|
||||
config = _config(tmp_path)
|
||||
with pytest.raises(release.ReleaseError, match="not a dry-run report"):
|
||||
release.approve(config, {"files": 3}, approver="domverse")
|
||||
|
||||
|
||||
def test_the_cli_runs_the_dry_run_and_the_approval(tmp_path):
|
||||
config = _config(tmp_path, PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL="1")
|
||||
_library(Path(config.library_roots[0]))
|
||||
from photo_pipeline.__main__ import main
|
||||
|
||||
environment = {
|
||||
"PHOTO_PIPELINE_DATA_DIR": str(config.data_dir),
|
||||
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(config.library_roots[0]),
|
||||
"PHOTO_PIPELINE_REQUIRE_DRY_RUN_APPROVAL": "1",
|
||||
}
|
||||
previous = {key: os.environ.get(key) for key in environment}
|
||||
os.environ.update(environment)
|
||||
try:
|
||||
report_path = tmp_path / "dry-run.json"
|
||||
assert main(["dry-run", "--output", str(report_path)]) == 0
|
||||
assert json.loads(report_path.read_text())["files"] == 3
|
||||
assert main(["approve-dry-run", str(report_path), "--approver", "domverse"]) == 0
|
||||
finally:
|
||||
for key, value in previous.items():
|
||||
if value is None:
|
||||
os.environ.pop(key, None)
|
||||
else:
|
||||
os.environ[key] = value
|
||||
assert release.mutation_blockers(config) == []
|
||||
348
tests/e2e/test_release_journey.py
Normal file
348
tests/e2e/test_release_journey.py
Normal file
@@ -0,0 +1,348 @@
|
||||
"""The release journey (US07-07): one library, one fresh environment, every stage.
|
||||
|
||||
This is the acceptance the whole backlog builds up to — discovery, duplicate review,
|
||||
safety, analysis, EXIF verification, album proposal, guarded rename, rescan and
|
||||
reconciliation, upload, archive, offline deduplication, restore — driven over HTTP
|
||||
against real ``photo_pipeline serve`` and worker child processes, with full process
|
||||
restarts in the middle and at the end.
|
||||
|
||||
Nothing is reached into. External services are the deterministic fakes the earlier
|
||||
phases already use, invoked through the real integration layer: a vision fake that
|
||||
records every path it was given, a real fake ``immich-go`` executable, and an
|
||||
archive medium that is an ordinary directory whose marker file is its identity.
|
||||
|
||||
The invariants asserted along the way are the ones the concept calls non-negotiable:
|
||||
|
||||
- an asset's identity survives a rename, an upload, an archive, and a restore;
|
||||
- an ``_IGNORE`` sentinel is never discovered, counted, analysed, or uploaded;
|
||||
- an NSFW asset never reaches the vision provider but still reaches Immich;
|
||||
- no photo's bytes are lost at any point — every hash is still reachable somewhere;
|
||||
- every stage's durable state survives a restart of both processes.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import shutil
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from tests.e2e._pipeline_harness import (
|
||||
SENTINEL_KEY,
|
||||
FakeImmich,
|
||||
Server,
|
||||
fake_uploader,
|
||||
image,
|
||||
seed_library,
|
||||
start_worker,
|
||||
wait_until,
|
||||
)
|
||||
|
||||
TIMEOUT = 30
|
||||
ALBUM = "rome"
|
||||
UPLOADER = 'echo "INFO uploaded $6"\necho "Uploaded 2, duplicates 0"\nexit 0\n'
|
||||
|
||||
|
||||
def _sha256(path: Path) -> str:
|
||||
return hashlib.sha256(path.read_bytes()).hexdigest()
|
||||
|
||||
|
||||
def _hashes(*roots: Path) -> set[str]:
|
||||
return {
|
||||
_sha256(path)
|
||||
for root in roots
|
||||
for path in root.rglob("*.jpg")
|
||||
if path.is_file() and not path.name.startswith(".")
|
||||
}
|
||||
|
||||
|
||||
def _post(base: str, path: str, **kwargs) -> httpx.Response:
|
||||
response = httpx.post(f"{base}/api/v1{path}", timeout=TIMEOUT, **kwargs)
|
||||
response.raise_for_status()
|
||||
return response
|
||||
|
||||
|
||||
def _get(base: str, path: str, **kwargs) -> dict:
|
||||
response = httpx.get(f"{base}/api/v1{path}", timeout=TIMEOUT, **kwargs)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
|
||||
|
||||
def _await_job(base: str, job_id: str, *, states=("succeeded",)) -> dict:
|
||||
return wait_until(
|
||||
lambda: (
|
||||
snapshot
|
||||
if (snapshot := _get(base, f"/jobs/{job_id}"))["state"] in states
|
||||
else None
|
||||
),
|
||||
timeout=90,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def library(tmp_path):
|
||||
"""A fresh library: an album, an exact duplicate, and an excluded sentinel."""
|
||||
seeded = seed_library(tmp_path, {}, {})
|
||||
album = seeded.lib / ALBUM
|
||||
image(album / "a.jpg", 11)
|
||||
image(album / "b.jpg", 12)
|
||||
shutil.copyfile(album / "a.jpg", album / "a-copy.jpg") # exact duplicate
|
||||
ignored = seeded.lib / "_IGNORE" / "private"
|
||||
ignored.mkdir(parents=True)
|
||||
image(ignored / "sentinel-9f3a2b.jpg", 99)
|
||||
return seeded
|
||||
|
||||
|
||||
@pytest.mark.skipif(shutil.which("exiftool") is None, reason="exiftool not installed")
|
||||
def test_the_full_release_journey_survives_every_stage_and_two_restarts(library, tmp_path):
|
||||
immich = FakeImmich()
|
||||
uploader = fake_uploader(tmp_path, UPLOADER)
|
||||
vision_log = tmp_path / "vision.log"
|
||||
archive_root = tmp_path / "medium"
|
||||
archive_root.mkdir()
|
||||
environment = {
|
||||
"PHOTO_PIPELINE_IMMICH_SERVER_URL": immich.url,
|
||||
"PHOTO_PIPELINE_IMMICH_API_KEY": SENTINEL_KEY,
|
||||
"PHOTO_PIPELINE_IMMICH_GO_BINARY": str(uploader),
|
||||
"PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": "0",
|
||||
"PHOTO_PIPELINE_FAKE_VISION_LOG": str(vision_log),
|
||||
}
|
||||
server = Server(library, extra_env=environment).start()
|
||||
worker = start_worker(library, extra_env=environment)
|
||||
base = server.base
|
||||
|
||||
try:
|
||||
# ── 0. discovery ─────────────────────────────────────────────────────
|
||||
_post(base, "/inventory/scan")
|
||||
assets = _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
assert len(assets) == 3, "the sentinel under _IGNORE is not an asset"
|
||||
paths = {asset["current_path"] for asset in assets}
|
||||
assert not any("_IGNORE" in path or "sentinel" in path for path in paths)
|
||||
identity = {asset["id"]: Path(asset["current_path"]).name for asset in assets}
|
||||
|
||||
# ── 1. duplicate review ──────────────────────────────────────────────
|
||||
_post(base, "/duplicates/detect")
|
||||
clusters = _get(base, "/duplicates/clusters")["items"]
|
||||
assert len(clusters) == 1 and clusters[0]["member_total"] == 2
|
||||
cluster = _get(base, f"/duplicates/clusters/{clusters[0]['id']}")
|
||||
canonical = sorted(member["asset_id"] for member in cluster["members"])[0]
|
||||
_post(
|
||||
base,
|
||||
f"/duplicates/clusters/{cluster['id']}/decision",
|
||||
json={
|
||||
"decision": "canonical",
|
||||
"canonical_asset_id": canonical,
|
||||
"expected_version": cluster["version"],
|
||||
},
|
||||
)
|
||||
|
||||
# ── 2. safety, with its EXIF checkpoint ──────────────────────────────
|
||||
queue = _get(base, "/safety/queue", params={"limit": 100})["items"]
|
||||
assert len(queue) == 2, "a non-canonical variant is not reviewed twice"
|
||||
decisions = {}
|
||||
for index, item in enumerate(sorted(queue, key=lambda row: row["current_path"])):
|
||||
decision = "nsfw" if index == 0 else "sfw"
|
||||
decisions[item["asset_id"]] = decision
|
||||
result = _post(
|
||||
base, "/safety/decisions", json={"asset_id": item["asset_id"], "decision": decision}
|
||||
).json()
|
||||
assert result["exif_verified"] is True, "the safety checkpoint must verify"
|
||||
|
||||
# ── restart: everything so far has to be durable ─────────────────────
|
||||
server.stop()
|
||||
server.start()
|
||||
base = server.base
|
||||
assert _get(base, "/safety/counts")["nsfw"] == 1
|
||||
assert {a["id"] for a in _get(base, "/inventory/assets", params={"limit": 200})["items"]} == set(
|
||||
identity
|
||||
)
|
||||
|
||||
# ── 3. analysis, gated to confirmed-SFW assets ───────────────────────
|
||||
job = _post(base, "/analysis/jobs").json()
|
||||
_await_job(base, job["id"])
|
||||
analysed = [
|
||||
name
|
||||
for name, decision in (
|
||||
(identity[asset_id], decision) for asset_id, decision in decisions.items()
|
||||
)
|
||||
if decision == "sfw"
|
||||
]
|
||||
seen = vision_log.read_text().splitlines()
|
||||
assert len(seen) == len(analysed) == 1
|
||||
assert not any("sentinel" in line or "_IGNORE" in line for line in seen)
|
||||
nsfw_id = next(aid for aid, decision in decisions.items() if decision == "nsfw")
|
||||
assert all(identity[nsfw_id] not in line for line in seen), "NSFW reached the provider"
|
||||
|
||||
# From here on no stage may change a photo's bytes: the metadata stages are
|
||||
# done, and moving, uploading, archiving, and restoring only relocate them.
|
||||
stable_hashes = _hashes(library.lib)
|
||||
|
||||
# ── 4. album proposal and guarded rename ─────────────────────────────
|
||||
_post(base, "/albums/proposals", json={})
|
||||
proposal = _get(base, f"/albums/proposals/{ALBUM}")
|
||||
_post(
|
||||
base,
|
||||
f"/albums/proposals/{ALBUM}/edit",
|
||||
json={"name": "2019 Rome", "expected_version": proposal["version"]},
|
||||
)
|
||||
proposal = _get(base, f"/albums/proposals/{ALBUM}")
|
||||
_post(
|
||||
base,
|
||||
f"/albums/proposals/{ALBUM}/approve",
|
||||
json={"expected_version": proposal["version"]},
|
||||
)
|
||||
plan = _post(base, "/rename-plans").json()
|
||||
assert plan["blockers"] == [], [
|
||||
(issue["code"], issue["message"])
|
||||
for op in plan["operations"]
|
||||
for issue in op["issues"]
|
||||
]
|
||||
response = httpx.post(
|
||||
f"{base}/api/v1/rename-plans/{plan['id']}/apply",
|
||||
json={"expected_version": plan["version"], "expected_checksum": plan["checksum"]},
|
||||
timeout=TIMEOUT,
|
||||
)
|
||||
assert response.status_code == 200, response.text
|
||||
applied = response.json()
|
||||
assert applied["failed"] == 0 and applied["applied"] == 1
|
||||
assert (library.lib / "2019 Rome").is_dir() and not (library.lib / ALBUM).exists()
|
||||
|
||||
# ── 5. rescan and reconciliation: identity survives the move ─────────
|
||||
_post(base, "/inventory/scan")
|
||||
after_rename = _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
assert {asset["id"] for asset in after_rename} == set(identity)
|
||||
assert all("2019 Rome" in asset["current_path"] for asset in after_rename)
|
||||
assert _hashes(library.lib) == stable_hashes, "a rename changed a photo's bytes"
|
||||
|
||||
# ── 6. upload ────────────────────────────────────────────────────────
|
||||
report = _post(base, "/upload-preflight", json={"albums": ["2019 Rome"]}).json()
|
||||
assert report["state"] == "ready", report["blockers"]
|
||||
batch = _post(
|
||||
base,
|
||||
"/upload-batches",
|
||||
json={"albums": ["2019 Rome"], "token": report["token"]},
|
||||
).json()["batches"][0]
|
||||
started = _post(base, f"/upload-batches/{batch['id']}/start").json()
|
||||
_await_job(base, started["job"]["id"])
|
||||
uploaded = _get(base, f"/upload-batches/{batch['id']}")
|
||||
assert uploaded["state"] == "succeeded"
|
||||
# The uploader said nothing per file, so the outcome is uncertain until the
|
||||
# server itself is asked whether it holds those exact bytes (US05-04).
|
||||
assert uploaded["outcome_state"] == "requires_verification"
|
||||
verified = _post(base, f"/upload-batches/{batch['id']}/verify").json()
|
||||
assert verified["outcome_state"] == "verified", verified
|
||||
uploaded = _get(base, f"/upload-batches/{batch['id']}")
|
||||
# Reviewed NSFW is uploaded; it simply never reached the analyser.
|
||||
assert {item["asset_id"] for item in uploaded["items"]} >= {nsfw_id}
|
||||
|
||||
# ── 7. archive ───────────────────────────────────────────────────────
|
||||
location = _post(
|
||||
base, "/archive-locations", json={"name": "external", "root": str(archive_root)}
|
||||
).json()
|
||||
preflight = _post(
|
||||
base, "/archive-preflight", json={"location_id": location["id"]}
|
||||
).json()
|
||||
assert preflight["state"] == "ready", [
|
||||
(asset["asset_id"], asset["blockers"])
|
||||
for album in preflight["albums"]
|
||||
for asset in album["assets"]
|
||||
if asset["blockers"]
|
||||
] or preflight
|
||||
archive_plan = _post(
|
||||
base,
|
||||
"/archive-plans",
|
||||
json={"location_id": location["id"], "token": preflight["token"]},
|
||||
).json()
|
||||
uploaded_ids = {item["asset_id"] for item in uploaded["items"]}
|
||||
_post(base, f"/archive-plans/{archive_plan['id']}/apply")
|
||||
wait_until(
|
||||
lambda: all(
|
||||
asset["availability_state"].startswith("archived")
|
||||
for asset in _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
if asset["id"] in uploaded_ids
|
||||
),
|
||||
timeout=90,
|
||||
)
|
||||
assert _hashes(library.lib, archive_root) == stable_hashes, "archiving lost bytes"
|
||||
|
||||
# ── 8. offline deduplication ─────────────────────────────────────────
|
||||
(archive_root / ".photo-pipeline-archive.json").rename(
|
||||
archive_root / ".photo-pipeline-archive.json.away"
|
||||
)
|
||||
# A copy of an archived photo turns up in the library under its own name —
|
||||
# the real shape of "I re-imported an old card" — so nothing occupies the
|
||||
# path the archived original would be restored to.
|
||||
returned = library.lib / "2019 Rome" / "rediscovered.jpg"
|
||||
returned.parent.mkdir(parents=True, exist_ok=True)
|
||||
archived_copy = next(archive_root.rglob("*.jpg"))
|
||||
shutil.copyfile(archived_copy, returned)
|
||||
_post(base, "/inventory/scan")
|
||||
_post(base, "/duplicates/detect")
|
||||
offline = _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
archived = [a for a in offline if a["availability_state"].startswith("archived")]
|
||||
assert archived, "an unmounted medium must not make assets missing"
|
||||
assert all(a["availability_state"] != "missing_unexpected" for a in offline)
|
||||
assert any(
|
||||
cluster["member_total"] >= 2 for cluster in _get(base, "/duplicates/clusters")["items"]
|
||||
), "the rediscovered copy did not meet its archived original"
|
||||
|
||||
# ── 9. restore ───────────────────────────────────────────────────────
|
||||
(archive_root / ".photo-pipeline-archive.json.away").rename(
|
||||
archive_root / ".photo-pipeline-archive.json"
|
||||
)
|
||||
restore_report = _post(
|
||||
base, "/restore-preflight", json={"location_id": location["id"]}
|
||||
).json()
|
||||
restore_plan = _post(
|
||||
base,
|
||||
"/restore-plans",
|
||||
json={"location_id": location["id"], "token": restore_report["token"]},
|
||||
).json()
|
||||
_post(base, f"/restore-plans/{restore_plan['id']}/apply")
|
||||
wait_until(
|
||||
lambda: all(
|
||||
asset["availability_state"] == "active"
|
||||
for asset in _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
if asset["id"] in identity
|
||||
),
|
||||
timeout=90,
|
||||
)
|
||||
|
||||
# ── 10. the final restart proves every stage was durable ─────────────
|
||||
worker.kill()
|
||||
worker.wait(timeout=20)
|
||||
server.stop()
|
||||
server.start()
|
||||
base = server.base
|
||||
final = {
|
||||
asset["id"]: asset
|
||||
for asset in _get(base, "/inventory/assets", params={"limit": 200})["items"]
|
||||
}
|
||||
assert set(identity) <= set(final), "an asset id did not survive the journey"
|
||||
assert _get(base, "/safety/counts")["nsfw"] == 1
|
||||
assert _get(base, "/upload-batches")["batches"][0]["state"] == "succeeded"
|
||||
reachable = {
|
||||
_sha256(path): str(path)
|
||||
for root in (library.lib, archive_root)
|
||||
for path in root.rglob("*.jpg")
|
||||
if path.is_file() and not path.name.startswith(".")
|
||||
}
|
||||
assert stable_hashes <= set(reachable), (
|
||||
"a photo was lost",
|
||||
sorted(stable_hashes - set(reachable)),
|
||||
sorted(reachable.values()),
|
||||
)
|
||||
workflow = _get(base, "/workflow")
|
||||
assert {stage["key"] for stage in workflow["stages"]} >= {
|
||||
"inventory",
|
||||
"duplicates",
|
||||
"safety",
|
||||
"analysis",
|
||||
}
|
||||
finally:
|
||||
worker.kill()
|
||||
worker.wait(timeout=20)
|
||||
server.stop()
|
||||
immich.stop()
|
||||
299
tests/integration/test_compose_runtime.py
Normal file
299
tests/integration/test_compose_runtime.py
Normal file
@@ -0,0 +1,299 @@
|
||||
"""US08-03: the composition's contract, and the library-root check it depends on.
|
||||
|
||||
Bringing the stack up needs a Docker daemon and the network, which is what
|
||||
``tests/e2e/test_compose_stack.py`` does. What can be checked without either is
|
||||
checked here, because the parts that rot silently — a second writer that is only
|
||||
prevented by convention, a data volume that stopped being the same volume for both
|
||||
roles, migrations that stopped running first, a committed value in a file that must
|
||||
carry none — are all readable from the files.
|
||||
|
||||
The startup refusal is the other half: in a container the configured library roots
|
||||
must name the mount paths, and a mismatch has to fail before the lock is taken, not
|
||||
at the first rename.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import stat
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
|
||||
from photo_pipeline import path_policy
|
||||
from photo_pipeline.__main__ import main
|
||||
from photo_pipeline.config import Config
|
||||
|
||||
REPO = Path(__file__).resolve().parents[2]
|
||||
COMPOSE_FILE = REPO / "docker-compose.yml"
|
||||
COMPOSE = yaml.safe_load(COMPOSE_FILE.read_text())
|
||||
ENV_EXAMPLE = REPO / ".env.example"
|
||||
|
||||
SERVICES = COMPOSE["services"]
|
||||
DATA_VOLUME = "data:/data"
|
||||
|
||||
|
||||
def env_example_keys() -> list[str]:
|
||||
"""The variables the example file declares, in file order."""
|
||||
return [
|
||||
line.split("=", 1)[0]
|
||||
for line in ENV_EXAMPLE.read_text().splitlines()
|
||||
if "=" in line and not line.lstrip().startswith("#")
|
||||
]
|
||||
|
||||
|
||||
# ── one API, one worker, one library, one volume ─────────────────────────────
|
||||
|
||||
|
||||
def test_exactly_one_serving_and_one_working_container_from_the_same_image():
|
||||
roles = {name: service["command"][0] for name, service in SERVICES.items()}
|
||||
assert sorted(roles.values()) == ["migrate", "serve", "worker"]
|
||||
assert [name for name, role in roles.items() if role == "serve"] == ["api"]
|
||||
assert [name for name, role in roles.items() if role == "worker"] == ["worker"]
|
||||
|
||||
images = {service["image"] for service in SERVICES.values()}
|
||||
assert len(images) == 1, "both roles must run the same build of the application"
|
||||
assert "latest" not in images.pop()
|
||||
# No `replicas`/`scale` key promising a second worker is fine; the lock decides.
|
||||
assert not any("deploy" in service for service in SERVICES.values())
|
||||
|
||||
|
||||
def test_both_roles_share_the_data_volume_so_the_lock_is_visible_to_both():
|
||||
"""A second worker is refused by the library lock (US07-05) only if it can see it."""
|
||||
for name, service in SERVICES.items():
|
||||
assert DATA_VOLUME in service["volumes"], name
|
||||
assert COMPOSE["volumes"]["data"]["driver"] == "local"
|
||||
text = COMPOSE_FILE.read_text()
|
||||
# The composition has to say why, because the failure is silent corruption.
|
||||
assert "WAL" in text and re.search(r"NFS|SMB|network", text)
|
||||
|
||||
|
||||
def test_the_library_is_a_bind_mount_whose_target_is_the_configured_root():
|
||||
for name, service in SERVICES.items():
|
||||
mounts = [volume for volume in service["volumes"] if volume != DATA_VOLUME]
|
||||
assert len(mounts) == 1, name
|
||||
source, target = re.match(r"^(\$\{.*?\}):(\$\{.*?\})$", mounts[0]).groups()
|
||||
# An unset host path fails the composition rather than mounting something else.
|
||||
assert source.startswith("${PHOTO_PIPELINE_LIBRARY_HOST_PATH:?")
|
||||
# The container-side path and the configured root are one variable, so they
|
||||
# cannot drift apart into a library that is mounted but not configured.
|
||||
assert target.startswith("${PHOTO_PIPELINE_LIBRARY_ROOTS:?")
|
||||
assert service["environment"]["PHOTO_PIPELINE_LIBRARY_ROOTS"] == (
|
||||
"${PHOTO_PIPELINE_LIBRARY_ROOTS}"
|
||||
)
|
||||
assert service["environment"]["PHOTO_PIPELINE_DATA_DIR"] == "/data"
|
||||
|
||||
|
||||
def test_migrations_run_to_completion_before_either_role_accepts_work():
|
||||
"""`migrate` runs the backup-then-migrate path, and a failed upgrade exits
|
||||
non-zero with its pre-migration backup intact — proven in
|
||||
tests/integration/test_backup_recovery.py. What the composition adds is that
|
||||
neither role starts until it succeeded."""
|
||||
assert SERVICES["migrate"]["command"] == ["migrate"]
|
||||
assert SERVICES["migrate"]["restart"] == "no", "a one-shot that retries is not a gate"
|
||||
for role in ("api", "worker"):
|
||||
assert SERVICES[role]["depends_on"] == {
|
||||
"migrate": {"condition": "service_completed_successfully"}
|
||||
}, role
|
||||
|
||||
|
||||
def test_the_api_port_is_published_to_host_loopback_by_default():
|
||||
published = SERVICES["api"]["ports"]
|
||||
assert published == [
|
||||
"${PHOTO_PIPELINE_PUBLISH_ADDRESS:-127.0.0.1}:${PHOTO_PIPELINE_PORT:-8000}:8000"
|
||||
]
|
||||
# Reachable from the host means reachable from elsewhere as far as the app is
|
||||
# concerned, so the access secret stays mandatory (US08-01).
|
||||
assert SERVICES["api"]["environment"]["PHOTO_PIPELINE_HOST"] == "0.0.0.0"
|
||||
assert SERVICES["api"]["environment"]["PHOTO_PIPELINE_PORT"] == 8000
|
||||
assert "PHOTO_PIPELINE_ACCESS_SECRET" not in SERVICES["api"]["environment"]
|
||||
|
||||
|
||||
def test_containers_restart_by_themselves_and_stop_with_time_to_drain():
|
||||
for role in ("api", "worker"):
|
||||
assert SERVICES[role]["restart"] == "unless-stopped", role
|
||||
assert SERVICES[role]["stop_grace_period"] == "30s", role
|
||||
|
||||
|
||||
def test_the_containers_run_as_the_library_owner_and_never_as_root():
|
||||
for name, service in SERVICES.items():
|
||||
assert service["user"] == "${PHOTO_PIPELINE_UID:-1000}:${PHOTO_PIPELINE_GID:-1000}", name
|
||||
assert service["build"]["args"]["UID"] == "${PHOTO_PIPELINE_UID:-1000}", name
|
||||
|
||||
|
||||
# ── configuration comes from the environment, never from a committed file ────
|
||||
|
||||
|
||||
def test_configuration_and_secrets_come_from_the_environment_only():
|
||||
for name, service in SERVICES.items():
|
||||
assert service["env_file"] == ["${PHOTO_PIPELINE_ENV_FILE:-.env}"], name
|
||||
for key, value in service["environment"].items():
|
||||
# Every value is either a variable reference or a property of the
|
||||
# composition itself (the volume path, the container's own port).
|
||||
composed = isinstance(value, int) or value in ("/data", "0.0.0.0")
|
||||
assert composed or value.startswith("${"), (name, key, value)
|
||||
assert not (REPO / ".env").is_file() or ".env" in (REPO / ".gitignore").read_text()
|
||||
|
||||
|
||||
def test_the_example_file_lists_every_setting_and_carries_no_values():
|
||||
declared = env_example_keys()
|
||||
assert declared == sorted(set(declared), key=declared.index), "no variable twice"
|
||||
for line in ENV_EXAMPLE.read_text().splitlines():
|
||||
if "=" in line and not line.lstrip().startswith("#"):
|
||||
assert line.endswith("="), f"a value in the example file: {line}"
|
||||
|
||||
expected = {f"PHOTO_PIPELINE_{name.upper()}" for name in Config.model_fields}
|
||||
assert expected <= set(declared), sorted(expected - set(declared))
|
||||
# And every variable the composition substitutes is documented there too.
|
||||
substituted = set(re.findall(r"\$\{(PHOTO_PIPELINE_[A-Z_]+)", COMPOSE_FILE.read_text()))
|
||||
assert substituted <= set(declared), sorted(substituted - set(declared))
|
||||
|
||||
|
||||
def test_the_example_file_is_not_a_dotenv_that_could_be_loaded_by_accident():
|
||||
"""`.env.example` must not be what `.env` is: no values means nothing to leak."""
|
||||
assert ENV_EXAMPLE.name != ".env"
|
||||
parsed = {k: v for k, v in _parse(ENV_EXAMPLE.read_text()).items() if v}
|
||||
assert parsed == {}
|
||||
|
||||
|
||||
def _parse(text: str) -> dict[str, str]:
|
||||
from photo_pipeline.config import parse_env_file
|
||||
|
||||
return parse_env_file(text)
|
||||
|
||||
|
||||
# ── the lock across container lifetimes ──────────────────────────────────────
|
||||
|
||||
|
||||
def test_a_lock_left_by_a_container_that_is_gone_does_not_block_the_restart(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
"""A restarted container is a new hostname and a recycled pid 1, so the record
|
||||
in the lock file proves nothing; the kernel's flock does (US08-03)."""
|
||||
from photo_pipeline.services.app_lock import LibraryLock
|
||||
|
||||
config = Config(data_dir=tmp_path / "data", library_roots=(tmp_path,))
|
||||
(tmp_path / "data").mkdir()
|
||||
(tmp_path / "data" / "worker.lock.json").write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"lock_version": 1,
|
||||
"role": "worker",
|
||||
"pid": 1, # pid 1 of a container that no longer exists
|
||||
"host": "3f2a1b9c4d5e", # its hostname was its container id
|
||||
"started_at": "2026-01-01T00:00:00+00:00",
|
||||
"library_roots": ["/library"],
|
||||
}
|
||||
)
|
||||
)
|
||||
monkeypatch.setattr(path_policy, "in_container", lambda: True)
|
||||
|
||||
taken = LibraryLock(config, "worker").acquire()
|
||||
|
||||
assert taken.pid == os.getpid(), "the worker must come back after a restart"
|
||||
|
||||
|
||||
def test_a_second_worker_is_still_refused_while_the_first_holds_the_lock(tmp_path, monkeypatch):
|
||||
"""The other half: the same flock refuses a concurrent second writer, whether it
|
||||
is a process or another container of the same composition."""
|
||||
from photo_pipeline.services.app_lock import LibraryLock, LockHeld
|
||||
|
||||
config = Config(data_dir=tmp_path / "data", library_roots=(tmp_path,))
|
||||
monkeypatch.setattr(path_policy, "in_container", lambda: True)
|
||||
first = LibraryLock(config, "worker")
|
||||
first.acquire()
|
||||
|
||||
with pytest.raises(LockHeld, match="worker is already running"):
|
||||
LibraryLock(config, "worker").acquire()
|
||||
|
||||
first.release()
|
||||
LibraryLock(config, "worker").acquire() # free again
|
||||
|
||||
|
||||
# ── the startup check the mount depends on ───────────────────────────────────
|
||||
|
||||
|
||||
def test_configured_roots_that_are_mounted_and_writable_are_accepted(tmp_path):
|
||||
assert path_policy.roots_refusal([tmp_path]) is None
|
||||
assert path_policy.roots_refusal([]) is None, "no roots is a configuration, not a fault"
|
||||
|
||||
|
||||
def test_an_unmounted_library_root_is_refused_by_name(tmp_path):
|
||||
refusal = path_policy.roots_refusal([tmp_path / "srv" / "photos"])
|
||||
assert refusal is not None
|
||||
assert "does not exist" in refusal and "PHOTO_PIPELINE_LIBRARY_ROOTS" in refusal
|
||||
|
||||
|
||||
def test_a_root_that_is_not_a_directory_or_not_readable_is_refused(tmp_path):
|
||||
a_file = tmp_path / "photos.txt"
|
||||
a_file.write_text("not a library")
|
||||
assert "not a directory" in path_policy.roots_refusal([a_file])
|
||||
|
||||
unreadable = tmp_path / "unreadable"
|
||||
unreadable.mkdir()
|
||||
unreadable.chmod(0o000)
|
||||
try:
|
||||
refusal = path_policy.roots_refusal([unreadable])
|
||||
finally:
|
||||
unreadable.chmod(0o755)
|
||||
if os.getuid() != 0: # root ignores the mode, and CI may well be root
|
||||
assert refusal is not None and "not readable" in refusal
|
||||
|
||||
|
||||
def test_an_unwritable_root_is_not_refused_here(tmp_path):
|
||||
"""A bind mount's ownership is virtualised on macOS and Windows, so os.access
|
||||
would refuse a working deployment. The real errno at the first rename is at
|
||||
least true; this check is about the mount, not the mode."""
|
||||
read_only = tmp_path / "read-only"
|
||||
read_only.mkdir()
|
||||
read_only.chmod(stat.S_IRUSR | stat.S_IXUSR)
|
||||
try:
|
||||
assert path_policy.roots_refusal([read_only]) is None
|
||||
finally:
|
||||
read_only.chmod(0o755)
|
||||
|
||||
|
||||
def test_in_a_container_a_root_that_was_never_mounted_is_refused(tmp_path, monkeypatch):
|
||||
"""The container-only failure: the path exists, but it belongs to the image."""
|
||||
unmounted = tmp_path / "library"
|
||||
(unmounted / "album").mkdir(parents=True)
|
||||
refusal = path_policy.roots_refusal([unmounted], require_mount=True)
|
||||
assert refusal is not None and "not on a mounted filesystem" in refusal
|
||||
|
||||
# A bind mount is a mount point, and a root *below* one is mounted too: a
|
||||
# deployment may mount /srv and configure /srv/photos.
|
||||
monkeypatch.setattr(os.path, "ismount", lambda path: Path(path) == unmounted.resolve())
|
||||
assert path_policy.roots_refusal([unmounted], require_mount=True) is None
|
||||
assert path_policy.roots_refusal([unmounted / "album"], require_mount=True) is None
|
||||
|
||||
|
||||
@pytest.mark.parametrize("role", ["serve", "worker"])
|
||||
def test_a_root_mismatch_refuses_at_startup_before_any_lock_is_taken(
|
||||
role, tmp_path, monkeypatch, capsys
|
||||
):
|
||||
data = tmp_path / "data"
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_DATA_DIR", str(data))
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_LIBRARY_ROOTS", str(tmp_path / "not-mounted"))
|
||||
monkeypatch.setattr(path_policy, "in_container", lambda: True)
|
||||
|
||||
assert main([role]) == 5
|
||||
assert "does not exist" in capsys.readouterr().err
|
||||
assert not list(data.glob("*.lock.json")), "nothing started, so nothing is locked"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("role", ["serve", "worker"])
|
||||
def test_an_unmounted_root_on_a_host_is_not_a_reason_to_refuse(role, tmp_path, monkeypatch):
|
||||
"""An archive medium that is not plugged in is a Tuesday, not a misconfiguration:
|
||||
refusing would take the offline half of the library away with it (concept §9)."""
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_DATA_DIR", str(tmp_path / "data"))
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_LIBRARY_ROOTS", str(tmp_path / "not-mounted"))
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_ACCESS_SECRET", "unused-on-loopback")
|
||||
monkeypatch.setattr(path_policy, "in_container", lambda: False)
|
||||
# Reaching the lock is the proof: that is the next thing either role does, and
|
||||
# stopping there keeps the test out of a uvicorn/worker loop.
|
||||
monkeypatch.setattr("photo_pipeline.__main__._acquire", lambda *_, **__: 99)
|
||||
|
||||
assert main([role]) == 99
|
||||
284
tests/integration/test_container_image.py
Normal file
284
tests/integration/test_container_image.py
Normal file
@@ -0,0 +1,284 @@
|
||||
"""US08-02: the image's build contract, its entrypoint, and its health check.
|
||||
|
||||
Building the image needs a Docker daemon and the network, which is what
|
||||
``tests/e2e/test_container_runtime.py`` does. Everything that can be checked without
|
||||
either is checked here, because the parts most likely to rot silently — a pin that
|
||||
stopped being a pin, a build context that started including the library, a health
|
||||
check pointed at liveness instead of readiness — are all readable from the files.
|
||||
|
||||
The entrypoint and health check are shell, so they are exercised as shell: run with a
|
||||
stubbed ``id`` and ``python`` on ``PATH``, which is enough to prove the refusal, the
|
||||
role marker, and the argument pass-through without a container.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import stat
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from photo_pipeline.config import Config
|
||||
from photo_pipeline.services import diagnostics
|
||||
|
||||
REPO = Path(__file__).resolve().parents[2]
|
||||
DOCKERFILE = (REPO / "Dockerfile").read_text()
|
||||
DOCKERIGNORE = (REPO / ".dockerignore").read_text()
|
||||
ENTRYPOINT = REPO / "docker" / "entrypoint.sh"
|
||||
HEALTHCHECK = REPO / "docker" / "healthcheck.sh"
|
||||
|
||||
SHA256 = re.compile(r"^[0-9a-f]{64}$")
|
||||
|
||||
|
||||
def instructions(text: str) -> list[str]:
|
||||
"""The lines that do something: comments explain, they do not build."""
|
||||
return [line.strip() for line in text.splitlines() if line.strip() and not line.startswith("#")]
|
||||
|
||||
|
||||
def build_args() -> dict[str, str]:
|
||||
"""Every ``ARG name=default`` in the Dockerfile — the pins, in other words."""
|
||||
found = {}
|
||||
for match in re.finditer(r"^ARG\s+([A-Z0-9_]+)=(.+)$", DOCKERFILE, re.MULTILINE):
|
||||
found[match.group(1)] = match.group(2).strip()
|
||||
return found
|
||||
|
||||
|
||||
# ── pins ─────────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_the_base_image_is_pinned_by_version_and_digest():
|
||||
base = build_args()["PYTHON_IMAGE"]
|
||||
assert base.startswith("python:3.12.")
|
||||
assert "@sha256:" in base, "a tag can be moved; a digest cannot"
|
||||
assert ":latest" not in DOCKERFILE
|
||||
# Both stages build from the same pinned base, so the tool that was verified in one
|
||||
# is the tool that ships in the other.
|
||||
assert DOCKERFILE.count("FROM ${PYTHON_IMAGE}") == 2
|
||||
|
||||
|
||||
def test_exiftool_and_the_uploader_are_pinned_and_verified():
|
||||
args = build_args()
|
||||
assert re.match(r"^\d+\.\d+", args["EXIFTOOL_VERSION"])
|
||||
assert re.match(r"^\d+\.\d+\.\d+$", args["IMMICH_GO_VERSION"])
|
||||
for arch in ("AMD64", "ARM64"):
|
||||
assert SHA256.match(args[f"IMMICH_GO_SHA256_{arch}"]), arch
|
||||
# The pinned exiftool package is installed by version, not by name alone.
|
||||
assert 'libimage-exiftool-perl=${EXIFTOOL_VERSION}"' in DOCKERFILE
|
||||
# And the build fails if what got installed is not what was pinned.
|
||||
assert "is not the pinned" in DOCKERFILE and "is not pinned" in DOCKERFILE
|
||||
|
||||
|
||||
def test_the_uploader_download_refuses_a_mismatching_checksum(tmp_path):
|
||||
"""The verification is the point of pinning a URL, so it is run, not read."""
|
||||
script = REPO / "docker" / "fetch-immich-go.py"
|
||||
result = subprocess.run(
|
||||
[
|
||||
sys.executable,
|
||||
str(script),
|
||||
"--version",
|
||||
"0.0.0-does-not-exist",
|
||||
"--sha256-amd64",
|
||||
"0" * 64,
|
||||
"--sha256-arm64",
|
||||
"0" * 64,
|
||||
"--into",
|
||||
str(tmp_path),
|
||||
],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
assert result.returncode != 0
|
||||
assert not list(tmp_path.iterdir()), "nothing is written before it is verified"
|
||||
|
||||
|
||||
def test_the_recorded_versions_are_reported_by_diagnostics(tmp_path, monkeypatch):
|
||||
"""What the image records is what `diagnostics` answers with (acceptance criterion 2)."""
|
||||
recorded = tmp_path / "versions.json"
|
||||
recorded.write_text(json.dumps({"exiftool": "13.25", "immich-go": "0.32.0"}))
|
||||
monkeypatch.setattr(diagnostics, "IMAGE_VERSIONS_FILE", recorded)
|
||||
monkeypatch.setattr(diagnostics.exiftool, "version", lambda: "13.25")
|
||||
monkeypatch.setattr(diagnostics, "_uploader_version", lambda _binary: "immich-go 0.32.0")
|
||||
|
||||
config = Config(data_dir=tmp_path / "data")
|
||||
reported = {tool["name"]: tool for tool in diagnostics.tools(config)}
|
||||
|
||||
assert reported["exiftool"]["pinned"] == "13.25"
|
||||
assert reported["immich-go"]["pinned"] == "0.32.0"
|
||||
assert "0.32.0" in reported["immich-go"]["version"]
|
||||
assert diagnostics.report(config)["tools"] == list(reported.values())
|
||||
assert "tool_version_drift" not in {w["code"] for w in diagnostics.report(config)["warnings"]}
|
||||
|
||||
|
||||
def test_a_replaced_tool_is_reported_as_drift(tmp_path, monkeypatch):
|
||||
recorded = tmp_path / "versions.json"
|
||||
recorded.write_text(json.dumps({"exiftool": "13.25"}))
|
||||
monkeypatch.setattr(diagnostics, "IMAGE_VERSIONS_FILE", recorded)
|
||||
monkeypatch.setattr(diagnostics.exiftool, "version", lambda: "12.57")
|
||||
|
||||
report = diagnostics.report(Config(data_dir=tmp_path / "data"))
|
||||
|
||||
drift = [w for w in report["warnings"] if w["code"] == "tool_version_drift"]
|
||||
assert drift and "13.25" in drift[0]["message"] and "12.57" in drift[0]["message"]
|
||||
|
||||
|
||||
def test_versions_are_absent_rather_than_invented_outside_a_container(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(diagnostics, "IMAGE_VERSIONS_FILE", tmp_path / "nothing.json")
|
||||
for tool in diagnostics.tools(Config(data_dir=tmp_path / "data")):
|
||||
assert tool["pinned"] is None
|
||||
|
||||
|
||||
# ── the final layer ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_no_test_or_build_tooling_is_installed_in_the_image():
|
||||
runtime = "\n".join(instructions(DOCKERFILE.split("AS runtime", 1)[1]))
|
||||
for unwanted in ("[test]", "pytest", "playwright", "build-essential", "gcc"):
|
||||
assert unwanted not in runtime, unwanted
|
||||
assert "pip install --no-cache-dir -e ." in runtime
|
||||
|
||||
|
||||
def test_neither_secrets_nor_library_data_can_enter_the_build_context():
|
||||
lines = instructions(DOCKERIGNORE)
|
||||
assert lines[0] == "*", "the context is deny-by-default"
|
||||
allowed = {line[1:] for line in lines if line.startswith("!")}
|
||||
# Everything the Dockerfile copies has to be allowed, and nothing else is.
|
||||
copied = {
|
||||
source
|
||||
for match in re.finditer(r"^COPY (?!--from)(.+)$", DOCKERFILE, re.MULTILINE)
|
||||
for source in match.group(1).split()[:-1]
|
||||
}
|
||||
assert {Path(source).parts[0] for source in copied} <= allowed
|
||||
assert not {"data", ".git", ".env", "tests", ".venv"} & allowed
|
||||
for generated in ("**/*.env", "**/*.db", "**/*.log", "**/__pycache__"):
|
||||
assert generated in lines, generated
|
||||
|
||||
|
||||
def test_the_image_runs_as_a_non_root_user_whose_ids_are_build_arguments():
|
||||
args = build_args()
|
||||
assert args["UID"] == "1000" and args["GID"] == "1000"
|
||||
assert "USER ${UID}:${GID}" in DOCKERFILE
|
||||
assert re.search(r"^USER (root|0)", DOCKERFILE, re.MULTILINE) is None
|
||||
assert 'useradd --uid "${UID}" --gid "${GID}"' in DOCKERFILE
|
||||
|
||||
|
||||
def test_the_health_check_is_readiness_and_the_default_role_is_serve():
|
||||
assert "HEALTHCHECK" in DOCKERFILE
|
||||
assert "/usr/local/bin/healthcheck.sh" in DOCKERFILE
|
||||
assert 'CMD ["serve"]' in DOCKERFILE
|
||||
assert 'ENTRYPOINT ["/usr/local/bin/entrypoint.sh"]' in DOCKERFILE
|
||||
assert "/api/v1/health/ready" in HEALTHCHECK.read_text()
|
||||
assert "/api/v1/health/live" not in HEALTHCHECK.read_text()
|
||||
# No supervisor: one role per container (acceptance criterion 4).
|
||||
for supervisor in ("supervisord", "s6-overlay", "runit"):
|
||||
assert supervisor not in DOCKERFILE
|
||||
|
||||
|
||||
@pytest.mark.parametrize("script", [ENTRYPOINT, HEALTHCHECK])
|
||||
def test_the_scripts_are_executable(script):
|
||||
assert script.stat().st_mode & stat.S_IXUSR, f"{script.name} must be executable in git"
|
||||
|
||||
|
||||
# ── the entrypoint, run as shell ─────────────────────────────────────────────
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def stubs(tmp_path):
|
||||
"""A PATH where ``python`` records its arguments and ``id`` can be told a UID."""
|
||||
bin_dir = tmp_path / "bin"
|
||||
bin_dir.mkdir()
|
||||
recorded = tmp_path / "argv"
|
||||
python = bin_dir / "python"
|
||||
python.write_text(f'#!/bin/sh\nprintf "%s\\n" "$@" > {recorded}\nexit 0\n')
|
||||
python.chmod(0o755)
|
||||
(bin_dir / "id").write_text('#!/bin/sh\nprintf "%s" "${STUB_UID:-1000}"\n')
|
||||
(bin_dir / "id").chmod(0o755)
|
||||
return bin_dir, recorded, tmp_path / "role"
|
||||
|
||||
|
||||
def run_script(script: Path, *args, stubs, env=None):
|
||||
bin_dir, recorded, role_file = stubs
|
||||
result = subprocess.run(
|
||||
["/bin/sh", str(script), *args],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env={
|
||||
"PATH": f"{bin_dir}:{os.environ['PATH']}",
|
||||
"PHOTO_PIPELINE_ROLE_FILE": str(role_file),
|
||||
**(env or {}),
|
||||
},
|
||||
)
|
||||
argv = recorded.read_text().splitlines() if recorded.exists() else []
|
||||
return result, argv
|
||||
|
||||
|
||||
def test_the_container_refuses_to_run_as_root(stubs):
|
||||
result, argv = run_script(ENTRYPOINT, "serve", stubs=stubs, env={"STUB_UID": "0"})
|
||||
assert result.returncode == 1
|
||||
assert "refusing to run as root" in result.stderr
|
||||
assert argv == [], "the application is never started as root"
|
||||
assert not stubs[2].exists(), "not even the role marker is written"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"given,expected",
|
||||
[
|
||||
(["serve"], ["-m", "photo_pipeline", "serve"]),
|
||||
(["worker", "--id", "worker-2"], ["-m", "photo_pipeline", "worker", "--id", "worker-2"]),
|
||||
# Every other management command stays reachable: operating the container is
|
||||
# operating the same CLI.
|
||||
(["diagnostics"], ["-m", "photo_pipeline", "diagnostics"]),
|
||||
([], ["-m", "photo_pipeline"]),
|
||||
],
|
||||
)
|
||||
def test_the_role_selects_the_command_and_arguments_pass_through(given, expected, stubs):
|
||||
result, argv = run_script(ENTRYPOINT, *given, stubs=stubs)
|
||||
assert result.returncode == 0, result.stderr
|
||||
assert argv == expected
|
||||
|
||||
|
||||
def test_the_role_is_recorded_for_the_health_check(stubs):
|
||||
run_script(ENTRYPOINT, "worker", stubs=stubs)
|
||||
assert stubs[2].read_text() == "worker"
|
||||
|
||||
|
||||
def test_the_health_check_only_probes_the_serving_role(stubs):
|
||||
stubs[2].write_text("worker")
|
||||
result, argv = run_script(HEALTHCHECK, stubs=stubs)
|
||||
assert result.returncode == 0 and argv == [], "a worker has no endpoint to probe"
|
||||
|
||||
stubs[2].write_text("serve")
|
||||
result, argv = run_script(HEALTHCHECK, stubs=stubs, env={"PHOTO_PIPELINE_PORT": "9123"})
|
||||
assert result.returncode == 0, result.stderr
|
||||
assert argv == ["-", "9123"], "the configured port is the one probed"
|
||||
|
||||
|
||||
def test_the_health_check_fails_while_the_api_is_not_ready(stubs):
|
||||
"""No python stub: the real interpreter probes a port nothing is listening on."""
|
||||
stubs[2].write_text("serve")
|
||||
result = subprocess.run(
|
||||
["/bin/sh", str(HEALTHCHECK)],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env={
|
||||
"PATH": os.path.dirname(sys.executable) + os.pathsep + os.environ["PATH"],
|
||||
"PHOTO_PIPELINE_ROLE_FILE": str(stubs[2]),
|
||||
"PHOTO_PIPELINE_PORT": "1",
|
||||
},
|
||||
)
|
||||
assert result.returncode == 1
|
||||
assert "not ready" in result.stderr
|
||||
|
||||
|
||||
def test_the_scripts_are_posix_shell():
|
||||
"""They run in the image's /bin/sh, which is dash — not bash."""
|
||||
shells = ["/bin/sh"] + ([dash] if (dash := shutil.which("dash")) else [])
|
||||
for shell in shells:
|
||||
for script in (ENTRYPOINT, HEALTHCHECK):
|
||||
checked = subprocess.run([shell, "-n", str(script)], capture_output=True, text=True)
|
||||
assert checked.returncode == 0, f"{shell} {script.name}: {checked.stderr}"
|
||||
253
tests/integration/test_trusted_hosts.py
Normal file
253
tests/integration/test_trusted_hosts.py
Normal file
@@ -0,0 +1,253 @@
|
||||
"""US08-01: the configurable trust boundary and its authentication gate.
|
||||
|
||||
Until now, reaching the app proved ownership of it: it answered only to loopback
|
||||
names. A container behind a reverse proxy answers to a real hostname, so these tests
|
||||
pin the two halves that replace that proof — the app refuses to start exposed without
|
||||
an access secret, and the secret is the only way to obtain the session every other
|
||||
route already required (US07-02, unchanged and re-asserted here).
|
||||
|
||||
The suite's ``conftest`` bootstraps a session for any ``TestClient`` automatically,
|
||||
which is precisely what an unauthenticated caller does not get; ``raw_client``
|
||||
pre-seeds a placeholder CSRF header to opt out of that convenience.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
from starlette.testclient import TestClient
|
||||
|
||||
from photo_pipeline.api.app import ConfigurationRefused, create_app
|
||||
from photo_pipeline.api.security import ACCESS_SECRET_HEADER, CSRF_HEADER, SESSION_COOKIE
|
||||
from photo_pipeline.config import Config
|
||||
|
||||
SECRET = "operator-secret-value"
|
||||
HOSTNAME = "photos.example.com"
|
||||
# What Starlette reports as the peer address of an in-process request.
|
||||
TESTCLIENT_ADDRESS = "testclient"
|
||||
|
||||
# One of each route class: a read, a mutation, and a media endpoint.
|
||||
PROTECTED = [
|
||||
("GET", "/api/v1/workflow", None),
|
||||
("POST", "/api/v1/albums/proposals", {}),
|
||||
("GET", "/api/v1/assets/unknown-asset/thumbnail?size=256", None),
|
||||
]
|
||||
|
||||
|
||||
def config(tmp_path, **overrides) -> Config:
|
||||
return Config(data_dir=tmp_path / "data", **overrides)
|
||||
|
||||
|
||||
def raw_client(app, base_url="http://127.0.0.1") -> TestClient:
|
||||
client = TestClient(app, base_url=base_url)
|
||||
client.headers[CSRF_HEADER] = "placeholder"
|
||||
return client
|
||||
|
||||
|
||||
def exchange(client, secret=SECRET, headers=None):
|
||||
return client.get("/api/v1/session", headers={ACCESS_SECRET_HEADER: secret, **(headers or {})})
|
||||
|
||||
|
||||
# ── startup: exposure without a secret is refused, loopback is unchanged ──────
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"exposure,exposed",
|
||||
[({"allowed_hosts": (HOSTNAME,)}, HOSTNAME), ({"host": "0.0.0.0"}, "0.0.0.0")],
|
||||
)
|
||||
def test_an_exposed_configuration_refuses_to_serve_without_a_secret(tmp_path, exposure, exposed):
|
||||
with pytest.raises(ConfigurationRefused) as refused:
|
||||
create_app(config(tmp_path, **exposure))
|
||||
assert "PHOTO_PIPELINE_ACCESS_SECRET" in str(refused.value)
|
||||
# The message names what is exposed, so the operator knows which setting did it.
|
||||
assert exposed in str(refused.value)
|
||||
|
||||
|
||||
def test_the_serve_command_reports_the_refusal_instead_of_binding(tmp_path, monkeypatch, capsys):
|
||||
"""Exit before the port, the lock, and the database, with a sentence not a trace."""
|
||||
from photo_pipeline.__main__ import main
|
||||
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_DATA_DIR", str(tmp_path / "data"))
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_ALLOWED_HOSTS", HOSTNAME)
|
||||
monkeypatch.delenv("PHOTO_PIPELINE_ACCESS_SECRET", raising=False)
|
||||
|
||||
assert main(["serve"]) == 4
|
||||
assert "PHOTO_PIPELINE_ACCESS_SECRET" in capsys.readouterr().err
|
||||
|
||||
|
||||
def test_an_exposed_configuration_with_a_secret_starts(tmp_path):
|
||||
app = create_app(config(tmp_path, allowed_hosts=(HOSTNAME,), access_secret=SECRET))
|
||||
with raw_client(app, base_url=f"http://{HOSTNAME}") as client:
|
||||
assert exchange(client).status_code == 200
|
||||
|
||||
|
||||
def test_a_loopback_configuration_still_needs_no_secret(tmp_path):
|
||||
"""An unset trust boundary must behave exactly as it did before this story."""
|
||||
with raw_client(create_app(config(tmp_path))) as client:
|
||||
response = client.get("/api/v1/session")
|
||||
assert response.status_code == 200
|
||||
assert response.json()["csrf_token"]
|
||||
assert "secure" not in response.headers["set-cookie"].lower()
|
||||
|
||||
|
||||
# ── the exchange: secret in, session out ─────────────────────────────────────
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def gated(tmp_path):
|
||||
app = create_app(
|
||||
config(
|
||||
tmp_path,
|
||||
allowed_hosts=(HOSTNAME,),
|
||||
access_secret=SECRET,
|
||||
trusted_proxies=(TESTCLIENT_ADDRESS,),
|
||||
)
|
||||
)
|
||||
with raw_client(app, base_url=f"http://{HOSTNAME}") as client:
|
||||
yield client
|
||||
|
||||
|
||||
def test_the_secret_buys_the_session_and_the_session_buys_the_routes(gated):
|
||||
response = exchange(gated)
|
||||
assert response.status_code == 200
|
||||
cookie = response.headers["set-cookie"].lower()
|
||||
assert "httponly" in cookie and "samesite=strict" in cookie
|
||||
gated.headers[CSRF_HEADER] = response.json()["csrf_token"]
|
||||
|
||||
# The session and CSRF requirements behind the gate are the ones US07-02 set.
|
||||
assert gated.get("/api/v1/workflow").status_code == 200
|
||||
assert gated.post("/api/v1/albums/proposals", json={}).status_code == 200
|
||||
refused = gated.post("/api/v1/albums/proposals", json={}, headers={CSRF_HEADER: "guessed"})
|
||||
assert refused.json()["error"]["code"] == "csrf_failed"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("offered", ["", "wrong-secret", SECRET + "x", SECRET.upper()])
|
||||
def test_a_wrong_secret_buys_nothing(gated, offered):
|
||||
response = exchange(gated, secret=offered)
|
||||
assert response.status_code == 401
|
||||
assert response.json()["error"]["code"] == "access_denied"
|
||||
assert "set-cookie" not in response.headers
|
||||
|
||||
|
||||
def test_a_refusal_never_echoes_the_secret_or_the_session(gated, caplog):
|
||||
with caplog.at_level("WARNING"):
|
||||
response = exchange(gated, secret="wrong-secret")
|
||||
assert SECRET not in response.text and "wrong-secret" not in response.text
|
||||
assert SECRET not in caplog.text
|
||||
# Logged as an event with its caller, without the session it did not get.
|
||||
assert "access secret rejected" in caplog.text
|
||||
|
||||
|
||||
def test_guessing_is_rate_limited(gated):
|
||||
codes = [exchange(gated, secret=f"guess-{n}").status_code for n in range(6)]
|
||||
assert codes.count(401) == 5 and codes[-1] == 429
|
||||
assert gated.get("/api/v1/session").status_code == 429
|
||||
# The right secret is refused too while the limiter holds: that is the point.
|
||||
blocked = exchange(gated)
|
||||
assert blocked.status_code == 429
|
||||
assert SECRET not in blocked.text
|
||||
|
||||
|
||||
def test_every_route_class_is_unreachable_without_the_secret(gated):
|
||||
for method, path, body in PROTECTED:
|
||||
response = gated.request(method, path, json=body)
|
||||
assert response.status_code == 401, path
|
||||
assert response.json()["error"]["code"] == "unauthenticated", path
|
||||
# Health stays open: an orchestrator restarting the container holds no secret.
|
||||
assert gated.get("/api/v1/health/live").status_code == 200
|
||||
assert gated.get("/api/v1/health/ready").status_code == 200
|
||||
|
||||
|
||||
def test_a_session_from_another_process_is_not_replayable(tmp_path):
|
||||
"""Sessions live in the process, so a cookie captured from a previous one — a
|
||||
restarted container, or a second deployment — must not open this one."""
|
||||
settings = dict(allowed_hosts=(HOSTNAME,), access_secret=SECRET)
|
||||
first, second = (create_app(config(tmp_path / str(n), **settings)) for n in (1, 2))
|
||||
with raw_client(first, base_url=f"http://{HOSTNAME}") as client:
|
||||
exchange(client)
|
||||
stolen = client.cookies[SESSION_COOKIE]
|
||||
|
||||
with raw_client(second, base_url=f"http://{HOSTNAME}") as client:
|
||||
client.cookies.set(SESSION_COOKIE, stolen, domain=HOSTNAME)
|
||||
response = client.get("/api/v1/workflow")
|
||||
assert response.status_code == 401
|
||||
assert response.json()["error"]["code"] == "unauthenticated"
|
||||
|
||||
|
||||
def test_a_cross_site_request_is_still_refused_behind_the_gate(gated):
|
||||
gated.headers[CSRF_HEADER] = exchange(gated).json()["csrf_token"]
|
||||
refused = gated.post(
|
||||
"/api/v1/albums/proposals", json={}, headers={"Origin": "https://evil.example"}
|
||||
)
|
||||
assert refused.json()["error"]["code"] == "origin_not_allowed"
|
||||
embedded = gated.get(
|
||||
"/api/v1/assets/unknown-asset/thumbnail?size=256", headers={"Sec-Fetch-Site": "cross-site"}
|
||||
)
|
||||
assert embedded.json()["error"]["code"] == "cross_site_blocked"
|
||||
|
||||
|
||||
def test_an_unconfigured_host_is_refused_even_with_a_valid_session(gated):
|
||||
gated.headers[CSRF_HEADER] = exchange(gated).json()["csrf_token"]
|
||||
for host in ("other.example.com", "192.168.1.10"):
|
||||
response = gated.get("/api/v1/workflow", headers={"Host": host})
|
||||
assert response.status_code == 403, host
|
||||
assert response.json()["error"]["code"] == "host_not_allowed", host
|
||||
|
||||
|
||||
# ── forwarded headers: believed from the proxy, ignored from anyone else ──────
|
||||
|
||||
|
||||
def test_a_trusted_proxys_https_makes_the_cookie_secure(gated):
|
||||
"""The proxy speaks HTTPS outward and HTTP to this app, so only the header knows."""
|
||||
assert "secure" in exchange(gated, headers={"X-Forwarded-Proto": "https"}).headers[
|
||||
"set-cookie"
|
||||
].lower()
|
||||
assert "secure" not in exchange(gated).headers["set-cookie"].lower()
|
||||
|
||||
|
||||
def test_the_external_scheme_is_part_of_the_accepted_origin(gated):
|
||||
gated.headers[CSRF_HEADER] = exchange(gated).json()["csrf_token"]
|
||||
allowed = gated.post(
|
||||
"/api/v1/albums/proposals",
|
||||
json={},
|
||||
headers={"X-Forwarded-Proto": "https", "Origin": f"https://{HOSTNAME}"},
|
||||
)
|
||||
assert allowed.status_code == 200
|
||||
# The scheme is part of the origin: the same name over plain HTTP is not it.
|
||||
refused = gated.post(
|
||||
"/api/v1/albums/proposals",
|
||||
json={},
|
||||
headers={"X-Forwarded-Proto": "https", "Origin": f"http://{HOSTNAME}"},
|
||||
)
|
||||
assert refused.json()["error"]["code"] == "origin_not_allowed"
|
||||
|
||||
|
||||
def test_a_trusted_proxys_forwarded_host_is_the_host_that_is_judged(tmp_path):
|
||||
"""The proxy terminates the operator's hostname and dials this app by address."""
|
||||
app = create_app(
|
||||
config(
|
||||
tmp_path,
|
||||
allowed_hosts=(HOSTNAME,),
|
||||
access_secret=SECRET,
|
||||
trusted_proxies=(TESTCLIENT_ADDRESS,),
|
||||
)
|
||||
)
|
||||
with raw_client(app, base_url="http://10.0.0.5") as client:
|
||||
forwarded = {"X-Forwarded-Host": HOSTNAME}
|
||||
assert exchange(client, headers=forwarded).status_code == 200
|
||||
# Without the header the address it was dialled by is not an allowed name.
|
||||
assert exchange(client).json()["error"]["code"] == "host_not_allowed"
|
||||
|
||||
|
||||
def test_forwarded_headers_from_an_untrusted_client_are_ignored(tmp_path):
|
||||
"""Otherwise any caller could declare the hostname and scheme of its choosing."""
|
||||
app = create_app(config(tmp_path, allowed_hosts=(HOSTNAME,), access_secret=SECRET))
|
||||
with raw_client(app, base_url="http://evil.example") as client:
|
||||
forged = exchange(client, headers={"X-Forwarded-Host": HOSTNAME})
|
||||
assert forged.json()["error"]["code"] == "host_not_allowed"
|
||||
|
||||
with raw_client(app, base_url=f"http://{HOSTNAME}") as client:
|
||||
# A forged scheme would flip the cookie's Secure flag on a plain connection,
|
||||
# which is how a cookie gets set and then never sent again.
|
||||
response = exchange(client, headers={"X-Forwarded-Proto": "https"})
|
||||
assert response.status_code == 200
|
||||
assert "secure" not in response.headers["set-cookie"].lower()
|
||||
@@ -15,9 +15,10 @@
|
||||
"tests/characterization/test_webapp_query.py"
|
||||
],
|
||||
"US01-02": [
|
||||
"tests/unit/test_config.py",
|
||||
"tests/integration/test_app_lifecycle.py",
|
||||
"tests/integration/test_migrations.py",
|
||||
"tests/integration/test_app_lifecycle.py"
|
||||
"tests/unit/test_config.py",
|
||||
"tests/unit/test_env_file.py"
|
||||
],
|
||||
"US01-03": [
|
||||
"tests/unit/test_path_policy.py",
|
||||
@@ -169,6 +170,27 @@
|
||||
],
|
||||
"US07-06": [
|
||||
"tests/integration/test_performance_budgets.py"
|
||||
],
|
||||
"US07-07": [
|
||||
"tests/e2e/test_release_gate.py",
|
||||
"tests/e2e/test_release_journey.py"
|
||||
],
|
||||
"US08-01": [
|
||||
"tests/unit/test_security_policy.py",
|
||||
"tests/integration/test_trusted_hosts.py"
|
||||
],
|
||||
"US08-02": [
|
||||
"tests/integration/test_container_image.py",
|
||||
"tests/e2e/test_container_runtime.py"
|
||||
],
|
||||
"US08-03": [
|
||||
"tests/integration/test_compose_runtime.py",
|
||||
"tests/e2e/test_compose_stack.py"
|
||||
]
|
||||
}
|
||||
},
|
||||
"planned": [
|
||||
"US08-04",
|
||||
"US08-05"
|
||||
],
|
||||
"_planned_comment": "Accepted backlog stories that are not implemented yet. The release gate (US07-07) requires every story file to be either mapped to tests or listed here, so an unimplemented story is a visible decision rather than a hole in the matrix."
|
||||
}
|
||||
|
||||
69
tests/unit/test_env_file.py
Normal file
69
tests/unit/test_env_file.py
Normal file
@@ -0,0 +1,69 @@
|
||||
"""Configuration from a dotenv file, including the archived CLI's variable names.
|
||||
|
||||
An operator who already has a ``photo_analyzer.env`` should not have to rewrite it
|
||||
to run the application it was replaced by. The file is standing configuration; the
|
||||
shell is what you meant this time, so the shell always wins.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
from photo_pipeline.config import Config, load_env_file, parse_env_file
|
||||
|
||||
SAMPLE = """
|
||||
# The archived CLI's shape, comments and all.
|
||||
LLM_API_KEY=not-real
|
||||
LLM_BASE_URL="https://example.invalid/v1beta/openai/"
|
||||
LLM_MODEL='gemini-2.5-flash'
|
||||
LIBRARY=/tmp/pictures
|
||||
|
||||
MAX_WORKERS=4
|
||||
# commented=ignored
|
||||
malformed line without an equals sign
|
||||
"""
|
||||
|
||||
|
||||
def test_the_file_is_parsed_and_never_executed():
|
||||
values = parse_env_file(SAMPLE)
|
||||
|
||||
assert values["LLM_BASE_URL"] == "https://example.invalid/v1beta/openai/" # quotes stripped
|
||||
assert values["LLM_MODEL"] == "gemini-2.5-flash"
|
||||
assert values["MAX_WORKERS"] == "4"
|
||||
assert "commented" not in values and "malformed line without an equals sign" not in values
|
||||
|
||||
|
||||
def test_the_archived_cli_names_still_configure_the_application():
|
||||
values = parse_env_file(SAMPLE)
|
||||
|
||||
assert values["OPENAI_API_KEY"] == "not-real"
|
||||
assert values["OPENAI_BASE_URL"] == "https://example.invalid/v1beta/openai/"
|
||||
assert values["PHOTO_PIPELINE_LIBRARY_ROOTS"] == "/tmp/pictures"
|
||||
|
||||
|
||||
def test_an_explicit_shell_variable_beats_the_file(tmp_path, monkeypatch):
|
||||
path = tmp_path / "photo_analyzer.env"
|
||||
path.write_text(SAMPLE)
|
||||
monkeypatch.setenv("OPENAI_API_KEY", "from-the-shell")
|
||||
monkeypatch.delenv("PHOTO_PIPELINE_LIBRARY_ROOTS", raising=False)
|
||||
|
||||
applied = load_env_file(path)
|
||||
|
||||
assert "OPENAI_API_KEY" not in applied, "the file overrode an exported value"
|
||||
assert os.environ["OPENAI_API_KEY"] == "from-the-shell"
|
||||
assert os.environ["PHOTO_PIPELINE_LIBRARY_ROOTS"] == "/tmp/pictures"
|
||||
assert Config.from_env().library_roots[0].name == "pictures"
|
||||
|
||||
|
||||
def test_the_file_is_found_through_its_variable(tmp_path, monkeypatch):
|
||||
path = tmp_path / "custom.env"
|
||||
path.write_text("PHOTO_PIPELINE_PORT=9123\n")
|
||||
monkeypatch.delenv("PHOTO_PIPELINE_PORT", raising=False)
|
||||
monkeypatch.setenv("PHOTO_PIPELINE_ENV_FILE", str(path))
|
||||
monkeypatch.chdir(tmp_path) # no ./.env here, so only the variable can find it
|
||||
|
||||
assert Config.from_env().port == 9123
|
||||
|
||||
|
||||
def test_a_missing_file_is_not_an_error(tmp_path):
|
||||
assert load_env_file(tmp_path / "nothing-here.env") == {}
|
||||
@@ -13,18 +13,31 @@ import pytest
|
||||
|
||||
from photo_pipeline.api.security import (
|
||||
CSRF_HEADER,
|
||||
LOOPBACK_HOSTS,
|
||||
PUBLIC_PATHS,
|
||||
FailureLimiter,
|
||||
Session,
|
||||
evaluate,
|
||||
exposed_hosts,
|
||||
external_view,
|
||||
split_host,
|
||||
trust_refusal,
|
||||
)
|
||||
from photo_pipeline.config import Config
|
||||
|
||||
SESSION = Session(id="session-id", csrf_token="csrf-token")
|
||||
HOST = "127.0.0.1:8000"
|
||||
LIMIT = 1024
|
||||
HOSTNAME = "photos.example.com"
|
||||
|
||||
|
||||
def check(method="GET", path="/api/v1/workflow", **headers):
|
||||
def check(
|
||||
method="GET",
|
||||
path="/api/v1/workflow",
|
||||
allowed_hosts=LOOPBACK_HOSTS,
|
||||
scheme="http",
|
||||
**headers,
|
||||
):
|
||||
"""Evaluate a request that is authenticated and same-origin unless overridden."""
|
||||
sent = {
|
||||
"host": HOST,
|
||||
@@ -38,6 +51,8 @@ def check(method="GET", path="/api/v1/workflow", **headers):
|
||||
path=path,
|
||||
headers=sent,
|
||||
session=SESSION,
|
||||
allowed_hosts=allowed_hosts,
|
||||
scheme=scheme,
|
||||
max_request_bytes=LIMIT,
|
||||
)
|
||||
|
||||
@@ -155,6 +170,115 @@ def test_refusals_name_no_path_secret_or_internal():
|
||||
assert "/" not in refusal.message
|
||||
|
||||
|
||||
# ── US08-01: the same table with a configured trust boundary ─────────────────
|
||||
|
||||
CONFIGURED = frozenset(LOOPBACK_HOSTS | {HOSTNAME})
|
||||
|
||||
|
||||
def test_a_configured_host_is_accepted_and_its_neighbours_are_not():
|
||||
assert check(host=HOSTNAME, allowed_hosts=CONFIGURED) is None
|
||||
for host in ("other.example.com", f"evil-{HOSTNAME}", "192.168.1.10"):
|
||||
refusal = check(host=host, allowed_hosts=CONFIGURED)
|
||||
assert (refusal.status, refusal.code) == (403, "host_not_allowed"), host
|
||||
|
||||
|
||||
def test_the_loopback_default_refuses_a_host_nobody_configured():
|
||||
"""The default set is what the app enforced before there was a setting."""
|
||||
refusal = check(host=HOSTNAME)
|
||||
assert (refusal.status, refusal.code) == (403, "host_not_allowed")
|
||||
|
||||
|
||||
def test_the_origin_must_match_the_external_scheme():
|
||||
for scheme in ("http", "https"):
|
||||
assert (
|
||||
check(
|
||||
method="POST",
|
||||
host=HOSTNAME,
|
||||
origin=f"{scheme}://{HOSTNAME}",
|
||||
allowed_hosts=CONFIGURED,
|
||||
scheme=scheme,
|
||||
)
|
||||
is None
|
||||
)
|
||||
# An HTTPS deployment whose caller claims plain HTTP is a different origin.
|
||||
refusal = check(
|
||||
method="POST",
|
||||
host=HOSTNAME,
|
||||
origin=f"http://{HOSTNAME}",
|
||||
allowed_hosts=CONFIGURED,
|
||||
scheme="https",
|
||||
)
|
||||
assert (refusal.status, refusal.code) == (403, "origin_not_allowed")
|
||||
|
||||
|
||||
def view(client, *, trusted=(), **headers):
|
||||
sent = {name.replace("_", "-"): value for name, value in headers.items()}
|
||||
return external_view(
|
||||
client=client,
|
||||
headers={"host": HOST, **sent},
|
||||
scheme="http",
|
||||
trusted_proxies=frozenset(trusted),
|
||||
)
|
||||
|
||||
|
||||
def test_forwarded_headers_are_ignored_without_a_trusted_proxy():
|
||||
forged = {"x_forwarded_proto": "https", "x_forwarded_host": HOSTNAME}
|
||||
assert view("10.0.0.9", **forged) == ("http", HOST)
|
||||
assert view(None, **forged) == ("http", HOST)
|
||||
# Configuring *a* proxy does not trust a caller that is not it.
|
||||
assert view("10.0.0.9", trusted=("10.0.0.1",), **forged) == ("http", HOST)
|
||||
|
||||
|
||||
def test_a_trusted_proxy_defines_the_external_scheme_and_host():
|
||||
assert view(
|
||||
"10.0.0.1", trusted=("10.0.0.1",), x_forwarded_proto="https", x_forwarded_host=HOSTNAME
|
||||
) == ("https", HOSTNAME)
|
||||
# A chain: the first entry is what the original client asked for.
|
||||
assert view(
|
||||
"10.0.0.1",
|
||||
trusted=("10.0.0.1",),
|
||||
x_forwarded_proto="https, http",
|
||||
x_forwarded_host=f"{HOSTNAME}, inner.internal",
|
||||
) == ("https", HOSTNAME)
|
||||
# Trusted but silent: this hop's own view stands.
|
||||
assert view("10.0.0.1", trusted=("10.0.0.1",)) == ("http", HOST)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"settings,exposed",
|
||||
[
|
||||
({}, []),
|
||||
({"host": "127.0.0.1"}, []),
|
||||
({"allowed_hosts": ("localhost", "127.0.0.1")}, []),
|
||||
({"allowed_hosts": (f"{HOSTNAME}:8443",)}, [HOSTNAME]),
|
||||
({"host": "0.0.0.0", "allowed_hosts": (HOSTNAME,)}, ["0.0.0.0", HOSTNAME]),
|
||||
],
|
||||
)
|
||||
def test_exposed_hosts_names_only_what_another_machine_can_reach(settings, exposed):
|
||||
assert exposed_hosts(Config(**settings)) == exposed
|
||||
|
||||
|
||||
def test_an_exposed_configuration_without_a_secret_must_not_serve():
|
||||
refusal = trust_refusal(Config(allowed_hosts=(HOSTNAME,)))
|
||||
assert HOSTNAME in refusal and "PHOTO_PIPELINE_ACCESS_SECRET" in refusal
|
||||
assert trust_refusal(Config(allowed_hosts=(HOSTNAME,), access_secret="s")) is None
|
||||
# Loopback-only, with and without a secret, is unchanged.
|
||||
assert trust_refusal(Config()) is None
|
||||
assert trust_refusal(Config(access_secret="s")) is None
|
||||
|
||||
|
||||
def test_failed_attempts_are_bounded_per_window():
|
||||
limiter = FailureLimiter(limit=2, window=60.0)
|
||||
assert not limiter.blocked()
|
||||
limiter.record_failure()
|
||||
assert not limiter.blocked()
|
||||
limiter.record_failure()
|
||||
assert limiter.blocked()
|
||||
# Attempts age out, so a locked-out operator is not locked out forever.
|
||||
limiter._failures = [-120.0, -120.0]
|
||||
assert not limiter.blocked()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"value,expected",
|
||||
[
|
||||
|
||||
Reference in New Issue
Block a user