Compare commits

...

3 Commits

25 changed files with 3058 additions and 21 deletions

View File

@@ -136,3 +136,36 @@ work_item/scripts/python -m pytest tests/e2e -m phase_d -q
The fault barrier is test-only configuration; without `PHOTO_PIPELINE_FAULT_AFTER`
the apply path has no crash points. Phases AC remain green in the full run above.
### Phase E acceptance gate
Phase E (Epic E05: Immich upload) is the one stage the application cannot take back,
so its gate runs the fake-uploader suites, the black-box upload API journeys, and the
browser suite as a single command:
```bash
work_item/scripts/python -m pytest -m phase_e -q
```
- `tests/integration/test_upload_*.py` drive a **real executable** standing in for
`immich-go` through the real adapter and `subprocess` — argument construction,
output bounding, report parsing, verification, and killing a running process.
- `tests/e2e/test_phase_e_pipeline.py` drives a real server and a real durable worker
over HTTP: credential failure and an unreachable server, preflight blockers and the
explicitly approved partial scope, a new album, an exact duplicate, an upgrade, a
retryable failure and its successful retry, a lost acceptance response, verification
against Immich, an inconclusive answer resolved by an operator with evidence, bytes
edited after upload, cancellation and resume, and an interrupted attempt recovered
across a restart.
- **EXIF precedes upload** is asserted, not assumed: an album without its verified
safety and analysis checkpoints cannot be approved, and the uploader's own argv log
proves it was never executed. Each finished upload re-hashes the files in the folder
the uploader was handed and requires the persisted SHA-256/SHA-1 to match.
- **No secret is retained.** The API key is a sentinel string; after a full upload and
verification it must appear in the uploader's argv and nowhere else — not in the
database, the retained report, or any response the browser can read.
- `tests/e2e/test_uploads_ui.py` covers the browser journeys (preflight preview,
confirmation, progress, stopping a run, verification, manual resolution, stale
bytes, and recovery after a restart).
Phases AD remain green in the full run above.

View File

@@ -18,6 +18,7 @@
<a href="#/analyze" data-nav="analyze">Analyze</a>
<a href="#/albums" data-nav="albums">Albums</a>
<a href="#/renames" data-nav="renames">Renames</a>
<a href="#/uploads" data-nav="uploads">Upload</a>
<a href="#/stats" data-nav="stats">Stats</a>
</nav>
</header>

View File

@@ -114,4 +114,28 @@ export const api = {
request(`/rename-plans/${encodeURIComponent(id)}/rollback`, { method: "POST", ...opts }),
renameRecovery: (opts = {}) => request("/rename-recovery", opts),
resolveRecovery: (opts = {}) => request("/rename-recovery/resolve", { method: "POST", ...opts }),
// ── Uploads: preflight, batches, verification ───────────────────────────
// The API key never travels through here: preflight reports only whether one is
// configured, and every command preview arrives already redacted.
uploadPreflight: (payload = {}, opts = {}) =>
request("/upload-preflight", { method: "POST", body: JSON.stringify(payload), ...opts }),
createUploadBatches: (payload, opts = {}) =>
request("/upload-batches", { method: "POST", body: JSON.stringify(payload), ...opts }),
listUploadBatches: (opts = {}) => request("/upload-batches", opts),
getUploadBatch: (id, opts = {}) => request(`/upload-batches/${encodeURIComponent(id)}`, opts),
startUploadBatch: (id, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/start`, { method: "POST", ...opts }),
cancelUploadBatch: (id, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/cancel`, { method: "POST", ...opts }),
verifyUploadBatch: (id, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/verify`, { method: "POST", ...opts }),
resolveUploadItem: (id, payload, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/resolve`, {
method: "POST",
body: JSON.stringify(payload),
...opts,
}),
uploadVerifications: (id, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/verifications`, opts),
};

View File

@@ -1,6 +1,7 @@
import { api } from "./api.js";
import { navigate, onRouteChange, parseHash } from "./router.js";
import { renderRenames, setRenamesRender } from "./renames.js";
import { renderUploads, setUploadsRender } from "./uploads.js";
import {
renderAlbums,
renderAnalyze,
@@ -364,6 +365,7 @@ function render() {
else if (path === "/analyze") renderAnalyze(root, params);
else if (path === "/albums") renderAlbums(root, params);
else if (path === "/renames") renderRenames(root, params);
else if (path === "/uploads") renderUploads(root, params);
else if (path === "/stats") renderStats(root, params);
else show(errorBanner("Unknown view"));
}
@@ -371,5 +373,6 @@ function render() {
// Let views re-render the current route after a mutation.
setRender(render);
setRenamesRender(render);
setUploadsRender(render);
onRouteChange(render);
render();

704
frontend/js/uploads.js Normal file
View File

@@ -0,0 +1,704 @@
// Upload view (US05-05): preflight a scope, confirm exactly what will be sent,
// watch the batch run, and resolve whatever the uploader left uncertain.
//
// Two rules shape this file. First, nothing here decides what is safe: blockers,
// the preflight token, and the retry policy all come from the server, and an action
// the server would refuse is not offered at all. Second, the API key never reaches
// the browser — preflight reports only that one is configured, and every command
// preview arrives redacted — so nothing in this view may reconstruct, store, or
// route a secret.
import { api } from "./api.js";
import { el, errorBanner, setActiveNav } from "./dom.js";
import { subscribeJob } from "./events.js";
import { navigate } from "./router.js";
// Result of the last command issued from this tab, and the activity of the last
// upload job. Deliberately not persisted: after a reload the page must show what
// the server says happened, not what this page remembers.
let outcome = null;
let activity = [];
let render = () => {};
export function setUploadsRender(fn) {
render = fn;
}
// Report outcomes (services/upload_reports.py) in the words the operator uses.
const OUTCOME_LABEL = [
["uploaded", "new"],
["upgraded", "upgraded"],
["duplicate", "duplicate"],
["skipped", "skipped"],
["failed", "failed"],
["unknown", "uncertain"],
];
export async function renderUploads(root, params = {}) {
setActiveNav("uploads");
const albums = parseAlbums(params.albums);
const allowPartial = params.partial === "1";
let preflight, listed;
try {
[preflight, listed] = await Promise.all([
api.uploadPreflight({ albums, allow_partial: allowPartial }),
api.listUploadBatches(),
]);
} catch (error) {
root.replaceChildren(errorBanner(`Failed to load uploads: ${error.message}`));
return;
}
const batches = listed.batches;
const selectedId = params.batch || (batches.length ? batches[batches.length - 1].id : null);
let batch = null;
let history = [];
if (selectedId) {
try {
const [detail, verifications] = await Promise.all([
api.getUploadBatch(selectedId),
api.uploadVerifications(selectedId),
]);
batch = detail;
history = verifications.verifications;
} catch (error) {
root.replaceChildren(errorBanner(`Failed to load upload batch: ${error.message}`));
return;
}
}
root.replaceChildren(
...[
el("h1", {}, "Upload"),
configurationCard(preflight),
preflightBlockers(preflight),
scopeSection(preflight, params, albums, allowPartial),
confirmBlock(preflight, params, allowPartial),
outcomeBanner(),
activityLog(),
batchList(batches, selectedId),
batch ? batchDetail(batch, history) : null,
].filter(Boolean)
);
}
// Album names can contain commas, so each one is escaped before the URL joins them.
function parseAlbums(value) {
if (!value) return null;
const names = value.split(",").filter(Boolean).map(decodeURIComponent);
return names.length ? names : null;
}
function albumsParam(names) {
return names.map(encodeURIComponent).join(",");
}
// ── configuration ────────────────────────────────────────────────────────────
// What the upload is aimed at, in the only form the browser is ever given: the
// server URL, whether a key exists, and the uploader's version.
function configurationCard(preflight) {
const credentials = preflight.credentials;
const uploader = preflight.uploader;
return el(
"div",
{ class: "card", "data-testid": "upload-config" },
el("h2", {}, "Configuration"),
el(
"dl",
{},
el("dt", {}, "Immich server"),
el("dd", { "data-testid": "config-server" }, credentials.server_url || "not configured"),
el("dt", {}, "API key"),
// Presence, never the value — and never a length or prefix either.
el(
"dd",
{ "data-testid": "config-key" },
credentials.api_key_configured ? "configured (never shown)" : "missing"
),
el("dt", {}, "Reachable"),
el(
"dd",
{ "data-testid": "config-reachable" },
preflight.server.reachable ? "yes" : `no — ${preflight.server.detail || "unknown"}`
),
el("dt", {}, "Uploader"),
el(
"dd",
{ "data-testid": "config-uploader" },
uploader.installed ? uploader.version || "installed" : `${uploader.binary} is not installed`
)
)
);
}
function preflightBlockers(preflight) {
if (!preflight.blockers.length) return null;
return el(
"div",
{ class: "alert", role: "alert", "data-testid": "preflight-blockers" },
el("strong", {}, "This scope cannot be uploaded yet"),
el(
"ul",
{},
...preflight.blockers.map((blocker) =>
el(
"li",
{ "data-testid": "preflight-blocker", "data-code": blocker.code },
`${blocker.code}: ${blocker.message}`
)
)
)
);
}
// ── scope ────────────────────────────────────────────────────────────────────
function scopeSection(preflight, params, albums, allowPartial) {
const selected = new Set(albums || preflight.albums.map((album) => album.album));
function toggle(name, checked) {
const next = new Set(selected);
if (checked) next.add(name);
else next.delete(name);
// An empty selection means "everything" again, which is also what an absent
// parameter means — there is no way to preflight nothing.
navigate("/uploads", { ...params, albums: albumsParam([...next]), batch: params.batch });
}
const rows = preflight.albums.map((album) =>
el(
"tr",
{ "data-testid": "album-row", "data-album": album.album },
el(
"td",
{},
el("input", {
type: "checkbox",
"data-testid": "album-selected",
"aria-label": `Include ${album.album}`,
checked: selected.has(album.album) ? "checked" : false,
onchange: (event) => toggle(album.album, event.target.checked),
})
),
el("td", { "data-testid": "album-name" }, album.album),
// What Immich will call it, which is not always what the folder is called here.
el("td", { "data-testid": "album-immich-name" }, album.album_name),
el("td", { class: "path", "data-testid": "album-folder" }, album.folder),
el("td", { "data-testid": "album-eligible" }, String(album.eligible_count)),
el("td", { "data-testid": "album-blocked" }, String(album.blocked_count)),
el(
"td",
{},
el("span", { class: `badge ${album.state}`, "data-testid": "album-state" }, album.state),
...album.blockers.map((blocker) =>
el(
"div",
{ class: "blocker", "data-testid": "album-blocker", "data-code": blocker.code },
blocker.message
)
)
),
// The exact invocation, as the server built it. Shown so the upload holds no
// surprises; the key is masked at the source, not here.
el("td", { class: "path", "data-testid": "album-command" }, album.command_preview.join(" "))
)
);
const totals = preflight.totals;
return el(
"div",
{ "data-testid": "upload-scope" },
el("h2", {}, "Scope"),
el(
"div",
{ class: "decision-bar" },
el("span", { class: "badge", "data-testid": "total-albums" }, `${totals.albums} album(s)`),
el("span", { class: "badge", "data-testid": "total-eligible" }, `${totals.eligible} ready`),
totals.blocked
? el(
"span",
{ class: "badge attention", "data-testid": "total-blocked" },
`${totals.blocked} blocked`
)
: null,
el(
"label",
{},
el("input", {
type: "checkbox",
"data-testid": "allow-partial",
checked: allowPartial ? "checked" : false,
onchange: (event) =>
navigate("/uploads", { ...params, partial: event.target.checked ? "1" : "" }),
}),
" Upload ready photos and leave the blocked ones behind"
)
),
rows.length
? el(
"table",
{ class: "grid", "data-testid": "albums" },
el(
"thead",
{},
el(
"tr",
{},
...["", "Album", "Immich album", "Folder", "Ready", "Blocked", "State", "Command"].map(
(label) => el("th", { scope: "col" }, label)
)
)
),
el("tbody", {}, ...rows)
)
: el("p", { class: "muted", "data-testid": "no-albums" }, "No album is ready to upload.")
);
}
// ── confirmation ─────────────────────────────────────────────────────────────
function confirmBlock(preflight, params, allowPartial) {
const ready = preflight.state === "ready";
const totals = preflight.totals;
const albums = preflight.albums.map((album) => album.album);
return el(
"div",
{ class: "card", "data-testid": "confirm" },
el("h2", {}, "Confirm"),
// The token is shown, not merely sent: a confirmation the user cannot see is a
// confirmation they cannot check against the preview above.
el(
"p",
{ class: "muted", "data-testid": "confirm-token" },
`Preflight ${preflight.token.slice(0, 20)}… · ${allowPartial ? "partial" : "complete"} scope`
),
el(
"p",
{ "data-testid": "upload-note" },
"Uploading is not reversible from here: Immich decides what to do with each " +
"file, and the app can only record what it reports. One album is sent at a " +
"time and a running album can be stopped."
),
el(
"div",
{ class: "toolbar" },
el(
"button",
{
class: "primary",
"data-testid": "start-upload",
disabled: ready ? false : "disabled",
title: ready ? false : "resolve the blockers above first",
onclick: () =>
run(async () => {
const created = await api.createUploadBatches({
albums,
token: preflight.token,
allow_partial: allowPartial,
});
// One lane: the first batch starts now, the rest wait with their own
// start buttons rather than queueing behind a lock that would reject
// them.
const first = created.batches.find((batch) => !batch.retry_blockers.length);
if (!first) return { batches: created.batches.length, started: null };
const started = await api.startUploadBatch(first.id);
watch(started.job.id, first.id);
return { batches: created.batches.length, started: first.album };
}),
},
`Upload ${totals.albums} album(s) · ${totals.eligible} photo(s)`
)
)
);
}
// ── batches ──────────────────────────────────────────────────────────────────
function batchList(batches, selectedId) {
if (!batches.length) {
return el("p", { class: "muted", "data-testid": "no-batches" }, "No upload has been started yet.");
}
return el(
"div",
{ "data-testid": "upload-batches" },
el("h2", {}, "Batches"),
el(
"table",
{ class: "grid", "data-testid": "batches" },
el(
"thead",
{},
el(
"tr",
{},
...["Album", "State", "Evidence", "Attempts", "Photos"].map((label) =>
el("th", { scope: "col" }, label)
)
)
),
el(
"tbody",
{},
...batches.map((batch) =>
el(
"tr",
{
"data-testid": "batch-row",
"data-album": batch.album,
"aria-current": batch.id === selectedId ? "true" : false,
},
el(
"td",
{},
el("a", { class: "link", href: `#/uploads?batch=${encodeURIComponent(batch.id)}` }, batch.album)
),
el(
"td",
{},
el("span", { class: `badge ${batch.state}`, "data-testid": "batch-state" }, batch.state)
),
el("td", { "data-testid": "batch-outcome-state" }, batch.outcome_state || "not parsed"),
el("td", {}, String(batch.attempt_count)),
el("td", {}, String(batch.asset_count))
)
)
)
)
);
}
function batchDetail(batch, history) {
// Derived from the items, not from the batch's parsed-report summary: verifying
// or resolving an item changes what is true without re-parsing a report, and the
// progress line must show the current answer rather than the uploader's old one.
const counts = {};
for (const item of batch.items) {
const key = item.outcome || "unknown";
counts[key] = (counts[key] || 0) + 1;
}
const blockers = batch.retry_blockers;
const uncertain = batch.outcome_state === "requires_verification" || counts.unknown > 0;
const nodes = [
el("h2", {}, `${batch.album} — attempt ${batch.attempt_count}`),
el(
"div",
{ class: "decision-bar", "data-testid": "batch-progress" },
el("span", { class: `badge ${batch.state}`, "data-testid": "detail-state" }, batch.state),
...OUTCOME_LABEL.map(([key, label]) =>
el(
"span",
{ class: `badge ${key}`, "data-testid": `count-${label}` },
`${label}: ${counts[key] ?? 0}`
)
)
),
];
if (batch.stale_bytes) {
nodes.push(
el(
"div",
{ class: "alert", role: "alert", "data-testid": "stale-bytes" },
"Files in this batch changed after they were uploaded. Immich still holds the " +
"bytes that were sent; re-approve the album through a fresh preflight rather " +
"than uploading the new bytes over it."
)
);
}
if (batch.error_code) {
nodes.push(
el(
"div",
{ class: "alert", role: "alert", "data-testid": "batch-error" },
`${batch.error_code}: ${batch.error_message || ""}`
)
);
}
if (uncertain) {
nodes.push(
el(
"div",
{ class: "alert", role: "alert", "data-testid": "uncertain" },
el("strong", {}, "This upload's outcome is not fully known"),
el(
"p",
{},
"The uploader's report does not account for every file. Retrying could create " +
"a second copy of something Immich already accepted, so verify it first: the " +
"check asks Immich whether it holds the exact bytes that were sent."
)
)
);
}
nodes.push(
el(
"div",
{ class: "toolbar" },
// A start button exists only when the server would accept one. An uncertain
// outcome and changed bytes therefore offer verification, never a retry.
blockers.length
? el(
"div",
{ class: "blocker", "data-testid": "retry-blocked" },
blockers.map((blocker) => `${blocker.code}: ${blocker.message}`).join("; ")
)
: el(
"button",
{
class: "primary",
"data-testid": "start-batch",
onclick: () =>
run(async () => {
const started = await api.startUploadBatch(batch.id);
watch(started.job.id, batch.id);
return { started: batch.album };
}),
},
batch.attempt_count ? "Run this album again" : "Upload this album"
),
["planned", "running"].includes(batch.state)
? el(
"button",
{
"data-testid": "cancel-batch",
onclick: () => run(() => api.cancelUploadBatch(batch.id)),
},
batch.state === "running" ? "Stop after the current file" : "Cancel this album"
)
: null,
// Offered for anything that has run, not only for uncertain outcomes:
// re-checking is read-only and idempotent, and it is how a file edited after
// its upload is discovered.
batch.attempt_count
? el(
"button",
{
"data-testid": "verify-batch",
onclick: () => run(() => api.verifyUploadBatch(batch.id)),
},
"Verify against Immich"
)
: null
),
itemsTable(batch)
);
if (history.length) nodes.push(historyList(history));
return el(
"div",
{ class: "card", "data-testid": "batch-detail", "data-batch": batch.id },
...nodes
);
}
function itemsTable(batch) {
if (!batch.items.length) {
return el("p", { class: "muted", "data-testid": "no-items" }, "This batch has no photos.");
}
return el(
"table",
{ class: "grid", "data-testid": "items" },
el(
"thead",
{},
el(
"tr",
{},
...["Photo", "Outcome", "Evidence", "Verification", "Bytes now", "Resolve"].map((label) =>
el("th", { scope: "col" }, label)
)
)
),
el(
"tbody",
{},
...batch.items.map((item) =>
el(
"tr",
{ "data-testid": "item-row", "data-asset-id": item.asset_id },
el("td", { class: "path", "data-testid": "item-path" }, item.path),
el(
"td",
{},
el(
"span",
{ class: `badge ${item.outcome || ""}`, "data-testid": "item-outcome" },
item.outcome === "uploaded" ? "new" : item.outcome || "pending"
)
),
el("td", { class: "muted", "data-testid": "item-evidence" }, item.evidence || "—"),
el("td", { "data-testid": "item-verification" }, item.verification || "—"),
el(
"td",
{},
item.changed_after_upload
? el("span", { class: "badge attention", "data-testid": "item-changed" }, "changed")
: el("span", { class: "muted" }, "unchanged")
),
el("td", {}, resolveForm(batch, item))
)
)
)
);
}
// Manual resolution is evidence, not permission: the note and the author are
// required by the server, so the form collects both and offers no default.
function resolveForm(batch, item) {
const unresolved = !item.outcome || item.outcome === "unknown" || item.verification === "inconclusive";
if (!unresolved) return el("span", { class: "muted" }, "—");
const outcomeSelect = el(
"select",
{ "data-testid": "resolve-outcome", "aria-label": `Outcome for ${item.path}` },
...OUTCOME_LABEL.map(([key, label]) => el("option", { value: key }, label))
);
const evidence = el("input", {
type: "text",
"data-testid": "resolve-evidence",
"aria-label": `What you checked for ${item.path}`,
placeholder: "What did you check?",
});
const actor = el("input", {
type: "text",
"data-testid": "resolve-actor",
"aria-label": `Who checked ${item.path}`,
placeholder: "Who are you?",
});
return el(
"div",
{ class: "toolbar" },
outcomeSelect,
evidence,
actor,
el(
"button",
{
"data-testid": "resolve-item",
onclick: () =>
run(() =>
api.resolveUploadItem(batch.id, {
asset_id: item.asset_id,
outcome: outcomeSelect.value,
evidence: evidence.value,
actor: actor.value,
})
),
},
"Record"
)
);
}
function historyList(history) {
return el(
"div",
{ "data-testid": "verification-history" },
el("h3", {}, "Verification history"),
el(
"ul",
{},
...history.map((entry) =>
el(
"li",
{ "data-testid": "history-entry", "data-source": entry.source },
`${entry.created_at || ""} · ${entry.action} · ${entry.source} · ${entry.result}` +
`${entry.outcome || "unresolved"}${entry.evidence}` +
(entry.actor ? ` (${entry.actor})` : "")
)
)
)
);
}
// ── running commands ─────────────────────────────────────────────────────────
async function run(action) {
try {
outcome = { kind: "ok", result: await action() };
} catch (error) {
outcome = error.status === 409 ? { kind: "conflict", error } : { kind: "error", error };
}
render();
}
// Live job activity. Events are appended to the log node and the running batch's
// panel is refreshed on its own tick — the uploader reports per file, not per job
// event, so waiting for the next event would leave the panel behind. Only that
// panel is rebuilt: a full re-render would re-run preflight, which re-hashes the
// library, so that happens once when the job ends.
const REFRESH_MS = 1000;
function watch(jobId, batchId) {
activity = [`Started upload job ${jobId}`];
const tick = setInterval(() => refreshBatch(batchId), REFRESH_MS);
subscribeJob(jobId, {
onEvent: (event) => {
activity.push(`${event.type}${event.message ? ": " + event.message : ""}`);
const log = document.querySelector('[data-testid="upload-activity"]');
if (log) log.textContent = activity.join("\n");
},
onDone: () => {
clearInterval(tick);
activity.push("done");
render();
},
});
}
async function refreshBatch(batchId) {
const node = document.querySelector(`[data-testid="batch-detail"][data-batch="${batchId}"]`);
if (!node) return; // the user navigated away from the running batch
try {
const [batch, verifications] = await Promise.all([
api.getUploadBatch(batchId),
api.uploadVerifications(batchId),
]);
node.replaceWith(batchDetail(batch, verifications.verifications));
} catch (_) {
// Transient: the next tick tries again, and the job's end re-renders anyway.
}
}
function activityLog() {
return el(
"pre",
{
class: "activity-log",
role: "status",
"aria-live": "polite",
"data-testid": "upload-activity",
},
activity.join("\n")
);
}
function outcomeBanner() {
if (!outcome) return null;
if (outcome.kind === "conflict") {
return el(
"div",
{ class: "alert", role: "alert", "data-testid": "conflict" },
`The server refused this: ${outcome.error.message}. Nothing was uploaded; the ` +
"state below is the server's current one — review it and decide again."
);
}
if (outcome.kind === "error") {
return el(
"div",
{ class: "alert", role: "alert", "data-testid": "upload-error" },
`Failed: ${outcome.error.message}`
);
}
const result = outcome.result;
if (result && result.batches !== undefined) {
return el(
"div",
{ class: "alert", role: "status", "data-testid": "upload-result" },
`Approved ${result.batches} album batch(es).` +
(result.started ? ` Uploading ${result.started} now.` : " Nothing could be started yet.")
);
}
return el(
"div",
{ class: "alert", role: "status", "data-testid": "upload-result" },
"Done — the state below is the server's."
);
}

View File

@@ -0,0 +1,70 @@
"""Upload verification and manual resolution (US05-04).
Revision ID: 0010_upload_verification
Revises: 0009_upload_report_outcomes
Create Date: 2026-08-16
Per-item server evidence plus the append-only history that produced it, so an
uncertain upload can be resolved without ever guessing that it succeeded.
"""
import sqlalchemy as sa
from alembic import op
revision = "0010_upload_verification"
down_revision = "0009_upload_report_outcomes"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.add_column(
"upload_batches", sa.Column("verified_at", sa.DateTime(timezone=True), nullable=True)
)
# At least one uploaded file was edited afterwards: a visible warning that also
# blocks re-running the batch.
op.add_column(
"upload_batches",
sa.Column("stale_bytes", sa.Boolean(), nullable=False, server_default=sa.false()),
)
# present | absent | inconclusive | manual
op.add_column("upload_items", sa.Column("verification", sa.String(), nullable=True))
op.add_column(
"upload_items", sa.Column("verified_at", sa.DateTime(timezone=True), nullable=True)
)
op.add_column("upload_items", sa.Column("observed_sha256", sa.String(), nullable=True))
op.add_column(
"upload_items",
sa.Column("changed_after_upload", sa.Boolean(), nullable=False, server_default=sa.false()),
)
op.create_table(
"upload_verifications",
sa.Column("id", sa.String(), primary_key=True),
sa.Column(
"batch_id",
sa.String(),
sa.ForeignKey("upload_batches.id", ondelete="CASCADE"),
nullable=False,
index=True,
),
sa.Column("asset_id", sa.String(), nullable=False),
sa.Column("action", sa.String(), nullable=False),
sa.Column("source", sa.String(), nullable=False),
sa.Column("result", sa.String(), nullable=False),
sa.Column("outcome", sa.String(), nullable=True),
sa.Column("evidence", sa.String(), nullable=False),
sa.Column("actor", sa.String(), nullable=True),
sa.Column(
"created_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()
),
)
def downgrade() -> None:
op.drop_table("upload_verifications")
for column in ("changed_after_upload", "observed_sha256", "verified_at", "verification"):
op.drop_column("upload_items", column)
for column in ("stale_bytes", "verified_at"):
op.drop_column("upload_batches", column)

View File

@@ -1,9 +1,12 @@
"""Upload preflight and batch API (US05-01, US05-02).
"""Upload preflight, batch, and verification API (US05-01, US05-02, US05-04).
Preflight is a command, not a resource read: it contacts the Immich server, hashes
the current bytes, and issues a token. Creating a batch requires that token, and
starting one enqueues a durable job on the single ``upload`` lane — the API never
runs the uploader in the request thread.
``verify`` and ``resolve`` are the way out of an uncertain outcome; ``start``
refuses one with ``409`` rather than letting the browser retry it.
"""
from __future__ import annotations
@@ -19,6 +22,11 @@ from photo_pipeline.services.upload_batches import (
BatchError,
UploadBatchService,
)
from photo_pipeline.services.upload_verification import (
UploadVerificationService,
VerificationError,
retry_blockers,
)
from photo_pipeline.services.uploads import UploadError, UploadService
router = APIRouter(tags=["uploads"])
@@ -36,6 +44,16 @@ class CreateBatchRequest(PreflightRequest):
token: str
class ResolveRequest(BaseModel):
asset_id: str
# uploaded | upgraded | duplicate | skipped | failed | unknown
outcome: str
# What the operator actually checked, and who they are — both mandatory so a
# manual resolution can never look like server evidence.
evidence: str
actor: str
def _service(request: Request) -> UploadService:
return UploadService(request.app.state.session_factory, config=request.app.state.config)
@@ -44,6 +62,12 @@ def _batches(request: Request) -> UploadBatchService:
return UploadBatchService(request.app.state.session_factory, config=request.app.state.config)
def _verification(request: Request) -> UploadVerificationService:
return UploadVerificationService(
request.app.state.session_factory, config=request.app.state.config
)
def _error(status: int, code: str, message: str) -> JSONResponse:
return JSONResponse(status_code=status, content={"error": {"code": code, "message": message}})
@@ -61,9 +85,11 @@ def preflight(request: Request, body: PreflightRequest | None = None):
def create_batches(body: CreateBatchRequest, request: Request):
"""Turn an approved preflight into one durable batch per album."""
try:
return {"batches": _batches(request).create(
body.albums, token=body.token, allow_partial=body.allow_partial
)}
return {
"batches": _batches(request).create(
body.albums, token=body.token, allow_partial=body.allow_partial
)
}
except UploadError as error:
return _error(422, "unknown_album", str(error))
except BatchConflict as error:
@@ -92,6 +118,11 @@ def start_batch(batch_id: str, request: Request):
batch = service.get(batch_id)
if batch is None:
return _error(404, "not_found", f"unknown upload batch {batch_id}")
# An uncertain outcome or bytes changed after upload must be resolved first
# (US05-04); the worker refuses them too, but the user is told here.
blocked = retry_blockers(batch)
if blocked:
return _error(409, blocked[0]["code"], blocked[0]["message"])
try:
job = JobService(request.app.state.session_factory).enqueue(
UPLOAD_BATCH,
@@ -105,6 +136,38 @@ def start_batch(batch_id: str, request: Request):
return {"batch_id": batch_id, "job": job}
@router.post("/upload-batches/{batch_id}/verify")
def verify_batch(batch_id: str, request: Request):
"""Check the batch against Immich and the bytes on disk (US05-04)."""
try:
return _verification(request).verify(batch_id)
except VerificationError as error:
return _error(404 if error.code == "not_found" else 422, error.code, str(error))
@router.post("/upload-batches/{batch_id}/resolve")
def resolve_item(batch_id: str, body: ResolveRequest, request: Request):
"""Record an operator's own verification of one item. Evidence is mandatory."""
try:
return _verification(request).resolve(
batch_id,
body.asset_id,
outcome=body.outcome,
evidence=body.evidence,
actor=body.actor,
)
except VerificationError as error:
return _error(404 if error.code == "not_found" else 422, error.code, str(error))
@router.get("/upload-batches/{batch_id}/verifications")
def list_verifications(batch_id: str, request: Request):
batch = _batches(request).get(batch_id)
if batch is None:
return _error(404, "not_found", f"unknown upload batch {batch_id}")
return {"verifications": _verification(request).history(batch_id)}
@router.post("/upload-batches/{batch_id}/cancel")
def cancel_batch(batch_id: str, request: Request):
try:

View File

@@ -8,7 +8,9 @@ activity log. Both come from the same builder so the preview can never drift fro
the command that would actually run.
Server reachability uses ``/api/server/ping`` through stdlib ``urllib`` — the app
has no HTTP client dependency and this is one request.
has no HTTP client dependency and this is one request. :func:`bulk_upload_check`
uses the same client to ask Immich which uploaded bytes it already holds, which is
the authoritative evidence behind upload verification (US05-04).
:func:`run_upload` is the only place the uploader is actually executed. It never
uses a shell (the argument list goes straight to ``execve``, so no path or album
@@ -35,6 +37,11 @@ from pathlib import Path
REDACTED = "***"
PING_PATH = "/api/server/ping"
PING_TIMEOUT_SECONDS = 5.0
# Immich's own deduplication endpoint: the authoritative answer to "do you already
# have these exact bytes?" used to verify uncertain uploads (US05-04).
BULK_CHECK_PATH = "/api/assets/bulk-upload-check"
CHECK_TIMEOUT_SECONDS = 30.0
CHECK_BATCH_SIZE = 500
# Reports are kept in full up to this size; beyond it the tail is dropped and the
# result is flagged truncated rather than growing without bound (concept §17).
MAX_REPORT_BYTES = 4_000_000
@@ -89,6 +96,68 @@ def ping(server_url: str, *, timeout: float = PING_TIMEOUT_SECONDS) -> tuple[boo
return False, "server did not answer with pong"
def bulk_upload_check(
server_url: str,
api_key: str | None,
checksums: dict[str, str],
*,
timeout: float = CHECK_TIMEOUT_SECONDS,
) -> dict:
"""Ask Immich which of these exact bytes it already holds (US05-04).
``checksums`` maps an application key (the asset id) to the SHA-1 of the bytes
that were uploaded — the digest Immich itself deduplicates on. The answer is
``{"reachable", "detail", "present"}`` where ``present`` maps each key to
``True`` (the server rejected it as a duplicate, so it holds those bytes),
``False`` (the server would accept it, so it does not), or ``None`` (the server
answered something this adapter will not interpret).
An unreachable or unparsable server is reported, never guessed at: the caller
must treat it as uncertainty rather than absence.
"""
if not server_url or not api_key:
return {"reachable": False, "detail": "no Immich credentials configured", "present": {}}
keys = list(checksums)
present: dict[str, bool | None] = {}
for start in range(0, len(keys), CHECK_BATCH_SIZE):
# ponytail: fixed chunk size; make it configurable if a server ever rejects it.
chunk = keys[start : start + CHECK_BATCH_SIZE]
payload = {"assets": [{"id": key, "checksum": checksums[key]} for key in chunk]}
request = urllib.request.Request( # noqa: S310 — http(s) URL from configuration
server_url.rstrip("/") + BULK_CHECK_PATH,
data=json.dumps(payload).encode("utf-8"),
headers={"Content-Type": "application/json", "x-api-key": api_key},
method="POST",
)
try:
with urllib.request.urlopen(request, timeout=timeout) as response: # noqa: S310
body = json.loads(response.read().decode("utf-8") or "{}")
except (urllib.error.URLError, OSError, ValueError, TimeoutError) as error:
return {"reachable": False, "detail": f"{type(error).__name__}: {error}", "present": {}}
results = body.get("results")
if not isinstance(results, list):
return {
"reachable": False,
"detail": "unrecognised bulk-upload-check response",
"present": {},
}
for result in results:
if not isinstance(result, dict) or result.get("id") not in checksums:
continue
present[result["id"]] = _holds_bytes(result)
return {"reachable": True, "detail": None, "present": present}
def _holds_bytes(result: dict) -> bool | None:
"""Whether one bulk-upload-check result means the server already has the file."""
action, reason = result.get("action"), result.get("reason")
if action == "reject":
# Only a duplicate proves possession; "unsupported-format" and friends say
# nothing about whether the bytes are there.
return True if reason == "duplicate" else None
return False if action == "accept" else None
def build_command(
*,
binary: str,

View File

@@ -96,6 +96,16 @@ class Worker:
return
try:
if cancelled or snapshot["state"] == JobState.CANCELLING:
# A handler may stop on its own — an upload batch cancelled through
# its own API never touches the job — so the job can still be
# ``running`` here. Record the request before the outcome: a stop is
# always observable as cancelling → cancelled, and ``running ->
# cancelled`` is not a legal jump. Without this hop the transition
# is rejected and the job keeps its lock forever.
if snapshot["state"] == JobState.RUNNING:
self.service.transition(
job_id, JobState.CANCELLING, worker_id=self.worker_id, fencing_token=token
)
self.service.transition(
job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token
)

View File

@@ -14,7 +14,7 @@ from photo_pipeline.models.duplicates import (
from photo_pipeline.models.jobs import Job, JobEvent, JobItem
from photo_pipeline.models.renames import RenameOperation, RenamePlan
from photo_pipeline.models.thumbnails import Thumbnail
from photo_pipeline.models.uploads import UploadBatch, UploadItem
from photo_pipeline.models.uploads import UploadBatch, UploadItem, UploadVerification
from photo_pipeline.models.workflow import AnalysisResult, SafetyReview
__all__ = [
@@ -32,6 +32,7 @@ __all__ = [
"Thumbnail",
"UploadBatch",
"UploadItem",
"UploadVerification",
"SafetyReview",
"AnalysisResult",
]

View File

@@ -69,6 +69,11 @@ class UploadBatch(Base):
outcome_counts: Mapped[str | None] = mapped_column(String) # JSON, from the items
report_counts: Mapped[str | None] = mapped_column(String) # JSON, uploader's own
# Verification (US05-04). ``stale_bytes`` means at least one uploaded file has
# been edited since: the batch carries a visible warning and cannot be re-run.
verified_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
stale_bytes: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
error_code: Mapped[str | None] = mapped_column(String)
error_message: Mapped[str | None] = mapped_column(String)
@@ -102,6 +107,45 @@ class UploadItem(Base):
outcome: Mapped[str | None] = mapped_column(String)
evidence: Mapped[str | None] = mapped_column(String) # the bounded report line
outcome_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
# Verification against the server (US05-04): present | absent | inconclusive |
# manual. NULL until the item has been verified; ``inconclusive`` whenever the
# server could not answer — which is never treated as success.
verification: Mapped[str | None] = mapped_column(String)
verified_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
# The bytes on disk at verification time, and whether they still are the bytes
# this batch uploaded.
observed_sha256: Mapped[str | None] = mapped_column(String)
changed_after_upload: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
)
class UploadVerification(Base):
"""Append-only evidence for every verification and manual resolution (US05-04).
The item row is a projection of the latest answer; this table is the history
that answers "who decided this, on what evidence, and when?". Rows are never
updated or deleted, so a manual resolution can always be told apart from
server evidence.
"""
__tablename__ = "upload_verifications"
id: Mapped[str] = mapped_column(String, primary_key=True)
batch_id: Mapped[str] = mapped_column(
ForeignKey("upload_batches.id", ondelete="CASCADE"), nullable=False, index=True
)
asset_id: Mapped[str] = mapped_column(String, nullable=False)
action: Mapped[str] = mapped_column(String, nullable=False) # verify | resolve
source: Mapped[str] = mapped_column(String, nullable=False) # immich_api | operator
# present | absent | inconclusive for a verify; the recorded outcome for a resolve.
result: Mapped[str] = mapped_column(String, nullable=False)
outcome: Mapped[str | None] = mapped_column(String)
evidence: Mapped[str] = mapped_column(String, nullable=False)
actor: Mapped[str | None] = mapped_column(String)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now()
)

View File

@@ -16,7 +16,8 @@ cannot undo, so the discipline is:
- **an interrupted attempt is uncertain, not failed.** Immich may have accepted
files the app never saw a report for, so ``recover`` marks a batch whose worker
vanished ``unknown_requires_verification`` (concept §15) instead of retrying it
blindly. Resolving that is US05-04.
blindly. :func:`~photo_pipeline.services.upload_verification.retry_blockers`
(US05-04) is what decides whether an attempt may start at all.
Per-asset upload *results* are not interpreted here: after the attempt ends the
report is handed to :class:`~photo_pipeline.services.upload_reports.
@@ -40,6 +41,7 @@ from photo_pipeline.integrations import immich_go
from photo_pipeline.models import UploadBatch, UploadItem
from photo_pipeline.services.hashing import sha1_file
from photo_pipeline.services.upload_reports import UploadReportService
from photo_pipeline.services.upload_verification import retry_blockers
from photo_pipeline.services.uploads import UploadService
@@ -62,6 +64,7 @@ RUNNABLE_STATES = frozenset({BatchState.PLANNED, BatchState.FAILED, BatchState.C
# States whose batch is still the live one for its album.
OPEN_STATES = frozenset({BatchState.PLANNED, BatchState.RUNNING, BatchState.CANCELLING})
class ItemState:
PENDING = "pending"
# ``sent`` means the batch process exited cleanly, not that Immich confirmed the
@@ -162,13 +165,11 @@ class UploadBatchService:
def run(self, batch_id: str, *, worker_id: str = "uploader", cancelled=None) -> dict:
"""Run one batch to completion. Blocks for the duration of the upload."""
batch = self._require(batch_id)
if batch["state"] not in RUNNABLE_STATES:
code = (
"requires_verification"
if batch["state"] == BatchState.UNKNOWN
else "not_runnable"
)
raise BatchError(code, f"batch {batch_id} is {batch['state']}")
# Retry policy (US05-04): a safe failure may run again, an uncertain outcome
# or bytes edited after upload may not.
blocked = retry_blockers(batch)
if blocked:
raise BatchError(blocked[0]["code"], blocked[0]["message"])
busy = [row for row in self.list() if row["id"] != batch_id and row["state"] in LANE_STATES]
if busy:
raise BatchConflict("lane_busy", f"upload batch {busy[0]['id']} is still running")
@@ -189,9 +190,7 @@ class UploadBatchService:
token = self._claim(batch_id, worker_id=worker_id)
attempt = self.get(batch_id)["attempt_count"]
report_path = (
Path(self._config.data_dir) / "uploads" / f"{batch_id}-attempt-{attempt}.log"
)
report_path = Path(self._config.data_dir) / "uploads" / f"{batch_id}-attempt-{attempt}.log"
key = self._config.immich_api_key.get_secret_value() if self._config.immich_api_key else ""
command = immich_go.build_command(
binary=self._config.immich_go_binary,
@@ -398,7 +397,7 @@ class UploadBatchService:
def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict:
return {
batch = {
"id": row.id,
"album": row.album,
"folder": row.folder,
@@ -426,6 +425,10 @@ def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict:
"outcome_state": row.outcome_state,
"outcome_counts": json.loads(row.outcome_counts) if row.outcome_counts else None,
"report_counts": json.loads(row.report_counts) if row.report_counts else None,
# Verification evidence (US05-04). ``stale_bytes`` is the visible warning
# that an uploaded file has since been edited.
"verified_at": row.verified_at.isoformat() if row.verified_at else None,
"stale_bytes": row.stale_bytes,
"started_at": row.started_at.isoformat() if row.started_at else None,
"finished_at": row.finished_at.isoformat() if row.finished_at else None,
"items": [
@@ -438,7 +441,16 @@ def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict:
"outcome": item.outcome,
"evidence": item.evidence,
"outcome_at": item.outcome_at.isoformat() if item.outcome_at else None,
"verification": item.verification,
"verified_at": item.verified_at.isoformat() if item.verified_at else None,
"observed_sha256": item.observed_sha256,
"changed_after_upload": item.changed_after_upload,
}
for item in items
],
}
# Why this batch may not be (re)started, from the one place that decides it
# (US05-04). Carried in the record so the browser can hide an action the server
# would refuse instead of re-implementing the policy (US05-05).
batch["retry_blockers"] = retry_blockers(batch)
return batch

View File

@@ -0,0 +1,349 @@
"""UploadVerificationService — resolving uncertain uploads (US05-04).
US05-03 leaves a batch honest but sometimes uncertain: a killed uploader, an
unpinned report grammar, or a file the report never mentioned all end as
``unknown``. Retrying such a batch is the dangerous move — Immich may already hold
the files, and a blind retry is how a lost response turns into a second server
asset. So this service resolves uncertainty *before* anything is re-run:
- **evidence beats the report.** Verification asks Immich itself whether it holds
the exact SHA-1 the batch recorded before uploading. That answer is authoritative
over the uploader's text, in both directions: present makes an unknown item
``uploaded``, absent makes it ``failed`` and therefore safe to retry.
- **no answer is never "no".** An unreachable server, missing credentials, or a
response this adapter will not interpret leave the item ``inconclusive``. The
batch stays ``unknown_requires_verification`` and stays un-runnable.
- **changed bytes are a stale warning, not a silent re-upload.** Verification
re-hashes what is on disk. A file edited after its upload is flagged, the batch
is marked ``stale_bytes``, and re-running it is refused: uploading again would
create or upgrade a server asset the user never approved (concept §8).
- **manual resolution is evidence, not permission.** An operator may record what
they checked in Immich, but only with a non-empty note and their identity, and
every decision is appended to an immutable history alongside the server answers.
Retry policy lives in :func:`retry_blockers`, which :class:`UploadBatchService`
enforces before every attempt and the API surfaces as ``409``.
"""
from __future__ import annotations
import uuid
from datetime import datetime, timezone
from pathlib import Path
from sqlalchemy import select
from sqlalchemy.orm import sessionmaker
from photo_pipeline.config import Config
from photo_pipeline.integrations import immich_go
from photo_pipeline.integrations import immich_go_report as report_parser
from photo_pipeline.models import UploadBatch, UploadItem, UploadVerification
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.upload_reports import REQUIRES_VERIFICATION, VERIFIED
PRESENT = "present"
ABSENT = "absent"
INCONCLUSIVE = "inconclusive"
MANUAL = "manual"
MAX_EVIDENCE_CHARS = 500
def _now() -> datetime:
return datetime.now(timezone.utc)
class VerificationError(RuntimeError):
"""The request cannot be carried out (unknown batch/item, missing evidence)."""
def __init__(self, code: str, message: str) -> None:
super().__init__(message)
self.code = code
class UploadVerificationService:
def __init__(self, session_factory: sessionmaker, *, config: Config) -> None:
self._session_factory = session_factory
self._config = config
# ── verification ──────────────────────────────────────────────────────────
def verify(self, batch_id: str) -> dict:
"""Check every item of a batch against the server and the bytes on disk.
Returns ``{"batch_id", "state", "outcome_state", "stale_bytes",
"server_reachable", "detail", "counts", "items"}``. Safe to call
repeatedly: the same evidence produces the same rows, and each run appends
its own history entries.
"""
with self._session_factory() as session:
batch = session.get(UploadBatch, batch_id)
if batch is None:
raise VerificationError("not_found", f"unknown upload batch {batch_id!r}")
items = list(
session.scalars(
select(UploadItem)
.where(UploadItem.batch_id == batch_id)
.order_by(UploadItem.path)
)
)
answer = immich_go.bulk_upload_check(
self._config.immich_server_url,
self._config.immich_api_key.get_secret_value()
if self._config.immich_api_key
else None,
{item.asset_id: item.sha1 for item in items if item.sha1},
)
detail = answer["detail"]
for item in items:
observed, changed = _current_bytes(item)
held = answer["present"].get(item.asset_id)
if item.sha1 is None:
result, evidence = INCONCLUSIVE, "no upload hash was recorded for this file"
elif held is True:
result, evidence = PRESENT, f"Immich holds sha1 {item.sha1}"
elif held is False:
result, evidence = ABSENT, f"Immich does not hold sha1 {item.sha1}"
else:
result = INCONCLUSIVE
evidence = detail or "the server did not classify these bytes"
item.verification = result
item.verified_at = _now()
item.observed_sha256 = observed
item.changed_after_upload = changed
# The server is authoritative over the report text — but only when
# it actually answered.
if result == PRESENT:
item.outcome = report_parser.UPLOADED
item.evidence = evidence
item.outcome_at = _now()
elif result == ABSENT:
item.outcome = report_parser.FAILED
item.evidence = evidence
item.outcome_at = _now()
session.add(
_event(
batch_id,
item.asset_id,
action="verify",
source="immich_api",
result=result,
outcome=item.outcome,
evidence=evidence,
)
)
_resolve_batch(batch, items)
session.commit()
return _report(batch, items, server_reachable=answer["reachable"], detail=detail)
# ── manual resolution ─────────────────────────────────────────────────────
def resolve(
self, batch_id: str, asset_id: str, *, outcome: str, evidence: str, actor: str
) -> dict:
"""Record an operator's own verification of one item.
``evidence`` and ``actor`` are mandatory: a manual resolution is only worth
keeping if it says what was checked and who checked it.
"""
if outcome not in report_parser.OUTCOMES:
raise VerificationError("invalid_outcome", f"unknown upload outcome {outcome!r}")
evidence = (evidence or "").strip()
actor = (actor or "").strip()
if not evidence:
raise VerificationError("evidence_required", "a manual resolution must record evidence")
if not actor:
raise VerificationError("actor_required", "a manual resolution must record its author")
with self._session_factory() as session:
batch = session.get(UploadBatch, batch_id)
if batch is None:
raise VerificationError("not_found", f"unknown upload batch {batch_id!r}")
item = session.get(UploadItem, {"batch_id": batch_id, "asset_id": asset_id})
if item is None:
raise VerificationError("not_found", f"{asset_id!r} is not part of this batch")
item.outcome = outcome
item.evidence = evidence[:MAX_EVIDENCE_CHARS]
item.outcome_at = _now()
item.verification = MANUAL
item.verified_at = _now()
session.add(
_event(
batch_id,
asset_id,
action="resolve",
source="operator",
result=MANUAL,
outcome=outcome,
evidence=evidence,
actor=actor,
)
)
items = list(
session.scalars(
select(UploadItem)
.where(UploadItem.batch_id == batch_id)
.order_by(UploadItem.path)
)
)
_resolve_batch(batch, items)
session.commit()
return _report(batch, items, server_reachable=None, detail=None)
# ── history ───────────────────────────────────────────────────────────────
def history(self, batch_id: str) -> list[dict]:
"""Every verification and resolution recorded for this batch, oldest first."""
with self._session_factory() as session:
rows = session.scalars(
select(UploadVerification)
.where(UploadVerification.batch_id == batch_id)
.order_by(UploadVerification.created_at, UploadVerification.id)
)
return [
{
"id": row.id,
"asset_id": row.asset_id,
"action": row.action,
"source": row.source,
"result": row.result,
"outcome": row.outcome,
"evidence": row.evidence,
"actor": row.actor,
"created_at": row.created_at.isoformat() if row.created_at else None,
}
for row in rows
]
# ── retry policy ─────────────────────────────────────────────────────────────
def retry_blockers(batch: dict) -> list[dict]:
"""Why this batch may not be (re)run, in the order the user should fix them.
A plain uploader failure is a *safe* failure: nothing uncertain happened, so it
is retryable. An uncertain outcome and edited bytes are not.
"""
from photo_pipeline.services.upload_batches import BatchState, RUNNABLE_STATES
blockers: list[dict] = []
if batch["state"] == BatchState.UNKNOWN:
blockers.append(
{
"code": "requires_verification",
"message": "this upload's outcome is uncertain; verify it before retrying",
}
)
if batch.get("stale_bytes"):
blockers.append(
{
"code": "changed_after_upload",
"message": "files in this batch changed after they were uploaded; "
"re-approve them through a fresh preflight",
}
)
if not blockers and batch["state"] not in RUNNABLE_STATES:
blockers.append({"code": "not_runnable", "message": f"batch is {batch['state']}"})
return blockers
# ── internals ────────────────────────────────────────────────────────────────
def _current_bytes(item: UploadItem) -> tuple[str | None, bool]:
"""``(hash on disk now, changed since upload)``. A missing file counts as changed."""
path = Path(item.path)
if not path.exists():
return None, item.sha256 is not None
observed = sha256_file(path)
return observed, bool(item.sha256 and observed != item.sha256)
def _event(
batch_id: str,
asset_id: str,
*,
action: str,
source: str,
result: str,
outcome: str | None,
evidence: str,
actor: str | None = None,
) -> UploadVerification:
return UploadVerification(
id=str(uuid.uuid4()),
batch_id=batch_id,
asset_id=asset_id,
action=action,
source=source,
result=result,
outcome=outcome,
evidence=evidence[:MAX_EVIDENCE_CHARS],
actor=actor,
created_at=_now(),
)
def _resolve_batch(batch: UploadBatch, items: list[UploadItem]) -> None:
"""Fold the item evidence back into the batch's own state.
An uncertain batch only leaves that state once every item is accounted for:
all present makes it succeeded, any absent makes it a safe failure to retry,
and a single inconclusive item keeps it uncertain.
"""
from photo_pipeline.services.upload_batches import BatchState
batch.stale_bytes = any(item.changed_after_upload for item in items)
batch.verified_at = _now()
unresolved = [item for item in items if item.outcome in (None, report_parser.UNKNOWN)]
batch.outcome_state = VERIFIED if items and not unresolved else REQUIRES_VERIFICATION
if batch.state != BatchState.UNKNOWN or unresolved:
return
if any(item.outcome == report_parser.FAILED for item in items):
batch.state = BatchState.FAILED
batch.error_code = "verified_incomplete"
batch.error_message = "verification proved some files never reached Immich"
else:
batch.state = BatchState.SUCCEEDED
batch.error_code = batch.error_message = None
def _report(
batch: UploadBatch,
items: list[UploadItem],
*,
server_reachable: bool | None,
detail: str | None,
) -> dict:
counts: dict[str, int] = {}
for item in items:
key = item.verification or "unverified"
counts[key] = counts.get(key, 0) + 1
return {
"batch_id": batch.id,
"state": batch.state,
"outcome_state": batch.outcome_state,
"stale_bytes": batch.stale_bytes,
"server_reachable": server_reachable,
"detail": detail,
"counts": counts,
"items": [
{
"asset_id": item.asset_id,
"path": item.path,
"verification": item.verification,
"outcome": item.outcome,
"evidence": item.evidence,
"changed_after_upload": item.changed_after_upload,
"verified_at": item.verified_at.isoformat() if item.verified_at else None,
}
for item in items
],
}

View File

@@ -31,4 +31,5 @@ markers = [
"phase_b: Phase B end-to-end acceptance (US02-07) — API, worker-recovery, and browser journeys",
"phase_c: Phase C end-to-end acceptance (US03-05) — album proposal API and browser journeys",
"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",
]

View File

@@ -9,14 +9,18 @@ and worker, never mocked inside a test.
from __future__ import annotations
import json
import os
import socket
import stat
import subprocess
import sys
import threading
import time
import uuid
from dataclasses import dataclass
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path
import httpx
@@ -135,15 +139,27 @@ class Server:
self.proc = None
def start_worker(seeded: Seeded, *, fake_vision_log: Path) -> subprocess.Popen:
"""Launch a real durable worker wired to the recording vision fake."""
def start_worker(
seeded: Seeded,
*,
fake_vision_log: Path | None = None,
extra_env: dict[str, str] | None = None,
) -> subprocess.Popen:
"""Launch a real durable worker wired to the recording vision fake.
``extra_env`` carries whatever else the job under test needs — the Immich
credentials and uploader path, for the upload lane.
"""
extra = dict(extra_env or {})
if fake_vision_log is not None:
extra["PHOTO_PIPELINE_FAKE_VISION_LOG"] = str(fake_vision_log)
return subprocess.Popen(
[sys.executable, "-m", "photo_pipeline", "worker", "--id", "e2e-worker"],
cwd=str(REPO),
env=_env(
seeded,
free_port(), # unused by the worker, but keeps the env shape uniform
extra={"PHOTO_PIPELINE_FAKE_VISION_LOG": str(fake_vision_log)},
extra=extra,
),
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
@@ -246,3 +262,182 @@ def approve_album(base: str, *, album: str = "rome", name: str) -> None:
json={**payload, "expected_version": current["version"]},
timeout=10,
).raise_for_status()
# ── Phase E: an upload-ready album, a fake Immich, and a real fake uploader ───
SENTINEL_KEY = "immich-sentinel-9f3a2b"
UPLOADER_VERSION = "immich-go 0.21.0" # a pinned family, so reports are parsable
# Uploader bodies for the pinned ``text-v1`` grammar. ``$6`` is the folder argument
# of ``upload from-folder``.
REPORTING_UPLOADER = (
'echo "INFO uploaded $6/a.jpg"\n'
'echo "INFO server has the same file $6/b.jpg"\n'
'echo "Uploaded 1, duplicates 1"\n'
"exit 0\n"
)
# Exits cleanly but says nothing about any file: the process succeeded, the
# per-file outcome is unknown.
SILENT_UPLOADER = "exit 0\n"
def _immich_handler(state: dict):
class Handler(BaseHTTPRequestHandler):
def do_GET(self): # noqa: N802 (BaseHTTPRequestHandler API)
self._json(200, {"res": "pong"})
def do_POST(self): # noqa: N802
length = int(self.headers.get("Content-Length", 0))
payload = json.loads(self.rfile.read(length) or b"{}")
if state["mode"] == "broken":
self.send_error(500, "bulk-upload-check is unavailable")
return
reject = state["mode"] == "present"
self._json(
200,
{
"results": [
{
"id": asset["id"],
"action": "reject" if reject else "accept",
"reason": "duplicate" if reject else None,
}
for asset in payload.get("assets", [])
]
},
)
def _json(self, code: int, body: dict) -> None:
raw = json.dumps(body).encode()
self.send_response(code)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(raw)))
self.end_headers()
self.wfile.write(raw)
def log_message(self, *args):
pass
return Handler
class FakeImmich:
"""An Immich that answers ping, and says whether it holds the exact bytes.
``mode`` is what the next verification will find: ``present`` (the server
deduplicates them, so it has them), ``absent`` (it would accept them, so it does
not), or ``broken`` (no usable answer at all).
"""
def __init__(self) -> None:
self.state = {"mode": "present"}
self._server = HTTPServer(("127.0.0.1", 0), _immich_handler(self.state))
threading.Thread(target=self._server.serve_forever, daemon=True).start()
self.url = f"http://127.0.0.1:{self._server.server_port}"
self._running = True
def mode(self, mode: str) -> None:
self.state["mode"] = mode
def stop(self) -> None:
"""Idempotent, so a test may take Immich away mid-journey."""
if not self._running:
return
self._running = False
self._server.shutdown()
self._server.server_close()
def fake_uploader(tmp_path: Path, body: str) -> Path:
"""A real executable standing in for immich-go.
``--version`` answers like the real tool; any other invocation appends its
complete argv to ``immich-go.argv`` — which is how a test proves the uploader
ran, what folder it was handed, or that it never ran at all.
"""
path = tmp_path / "immich-go"
path.write_text(
"#!/bin/sh\n"
f'if [ "$1" = "--version" ]; then echo "{UPLOADER_VERSION}"; exit 0; fi\n'
f'printf "%s\\n" "$*" >> "{tmp_path / "immich-go.argv"}"\n'
f"{body}"
)
path.chmod(path.stat().st_mode | stat.S_IEXEC | stat.S_IXGRP | stat.S_IXOTH)
return path
def uploader_argv(tmp_path: Path) -> list[str]:
"""Every upload invocation the fake uploader saw, oldest first."""
log = tmp_path / "immich-go.argv"
return log.read_text().splitlines() if log.exists() else []
def mark_upload_ready(seeded: Seeded, *, unverified: tuple[str, ...] = ()) -> None:
"""Give every seeded photo the verified EXIF checkpoints upload requires.
``unverified`` names stems whose analysis checkpoint stays incomplete, which is
what makes an album partially blocked.
"""
from sqlalchemy import select
from photo_pipeline.models import AnalysisResult, SafetyReview
blocked = {seeded.asset_ids[stem] for stem in unverified}
with session_factory(seeded) as sf:
with sf() as session:
for review in session.scalars(select(SafetyReview)):
review.exif_verified_at = NOW
for analysis in session.scalars(select(AnalysisResult)):
analysis.exif_written_at = None if analysis.asset_id in blocked else NOW
session.commit()
class UploadStack:
"""A seeded, upload-ready library plus the server, worker, and fake Immich."""
def __init__(self, tmp_path: Path, seeded: Seeded) -> None:
self.tmp_path = tmp_path
self.seeded = seeded
self.immich = FakeImmich()
self.server: Server | None = None
self.worker: subprocess.Popen | None = None
self.base = ""
def start(
self,
*,
uploader: str = REPORTING_UPLOADER,
worker: bool = True,
credentials: bool = True,
) -> "UploadStack":
env = {
"PHOTO_PIPELINE_IMMICH_SERVER_URL": self.immich.url if credentials else "",
"PHOTO_PIPELINE_IMMICH_GO_BINARY": str(fake_uploader(self.tmp_path, uploader)),
}
if credentials:
env["PHOTO_PIPELINE_IMMICH_API_KEY"] = SENTINEL_KEY
self.server = Server(self.seeded, extra_env=env).start()
self.base = self.server.base
if worker:
self.worker = start_worker(self.seeded, extra_env=env)
return self
def restart_server(self) -> None:
"""A genuinely fresh process against the same database and library."""
self.server.stop()
self.server.start()
def batches(self) -> list[dict]:
return httpx.get(f"{self.base}/api/v1/upload-batches", timeout=20).json()["batches"]
def argv(self) -> list[str]:
return uploader_argv(self.tmp_path)
def stop(self) -> None:
if self.worker is not None:
self.worker.kill()
self.worker.wait(timeout=10)
if self.server is not None:
self.server.stop()
self.immich.stop()

View File

@@ -0,0 +1,568 @@
"""Phase E end-to-end acceptance (US05-06): Immich upload, black box.
Every journey drives a real ``photo_pipeline serve`` child process and a real durable
worker over HTTP — preflight, approve, upload, duplicate, upgrade, fail, retry, lose
the acceptance response, verify, cancel, crash, restart. Nothing external is mocked
inside the application: ``immich-go`` is a real executable that records the argv it
was handed, and Immich is a real HTTP server answering the same ``ping`` and
``bulk-upload-check`` endpoints the adapter calls in production.
Two invariants are asserted in every relevant journey, because they are what make an
irreversible stage safe:
- **EXIF precedes upload.** An album without its verified safety and analysis
checkpoints cannot be approved, and the uploader's argv log proves it was never
even executed.
- **The persisted hashes are the submitted bytes.** After each upload the recorded
SHA-256/SHA-1 of every item is recomputed from the files in the folder the uploader
was actually given.
The API key is a sentinel string, so the last journey can prove it reached the
uploader and nothing else that was retained.
"""
from __future__ import annotations
import hashlib
from pathlib import Path
import httpx
import pytest
from tests.e2e._pipeline_harness import (
SENTINEL_KEY,
SILENT_UPLOADER,
UploadStack,
mark_upload_ready,
seed_album,
session_factory,
wait_until,
)
pytestmark = pytest.mark.phase_e
TIMEOUT = 20
ALBUM = "rome"
TERMINAL = {"succeeded", "failed", "cancelled", "unknown_requires_verification"}
# Uploader bodies in the pinned ``text-v1`` grammar. ``$6`` is the folder argument of
# ``upload from-folder``, so each line names the real path of a real file.
ALL_NEW = (
'echo "INFO uploaded $6/a.jpg"\n'
'echo "INFO uploaded $6/b.jpg"\n'
'echo "Uploaded 2"\n'
"exit 0\n"
)
EXACT_DUPLICATES = (
'echo "INFO server has the same file $6/a.jpg"\n'
'echo "INFO server has the same file $6/b.jpg"\n'
'echo "Duplicates 2"\n'
"exit 0\n"
)
UPGRADES = (
'echo "INFO server has an older file $6/a.jpg"\n'
'echo "INFO server has an older file $6/b.jpg"\n'
'echo "Upgraded 2"\n'
"exit 0\n"
)
def _once_then(tmp_path: Path, first: str, rest: str) -> str:
"""An uploader that behaves one way on its first attempt and another afterwards.
The flag file is the attempt counter, so retry and resume journeys are
deterministic without any test reaching into the running application.
"""
flag = tmp_path / "first-attempt.flag"
return f'if [ ! -f "{flag}" ]; then\n touch "{flag}"\n{first}fi\n{rest}'
FAILING_FIRST = (' echo "ERROR error uploading $6/a.jpg: connection reset"\n exit 1\n', ALL_NEW)
SLOW_FIRST = (' echo "INFO starting"\n sleep 30\n exit 0\n', ALL_NEW)
@pytest.fixture
def stack(tmp_path):
"""An analysed album whose EXIF checkpoints are already verified."""
seeded = seed_album(tmp_path)
mark_upload_ready(seeded)
running = UploadStack(tmp_path, seeded)
try:
yield running
finally:
running.stop()
@pytest.fixture
def unfinished(tmp_path):
"""The same album *before* its EXIF checkpoints were written."""
seeded = seed_album(tmp_path)
running = UploadStack(tmp_path, seeded)
try:
yield running
finally:
running.stop()
# ── helpers ──────────────────────────────────────────────────────────────────
def _preflight(base: str, **body) -> dict:
response = httpx.post(f"{base}/api/v1/upload-preflight", json=body, timeout=TIMEOUT)
response.raise_for_status()
return response.json()
def _create(base: str, report: dict, **body) -> httpx.Response:
return httpx.post(
f"{base}/api/v1/upload-batches",
json={"token": report["token"], **body},
timeout=TIMEOUT,
)
def _get(base: str, batch_id: str) -> dict:
return httpx.get(f"{base}/api/v1/upload-batches/{batch_id}", timeout=TIMEOUT).json()
def _start(base: str, batch_id: str) -> httpx.Response:
return httpx.post(f"{base}/api/v1/upload-batches/{batch_id}/start", timeout=TIMEOUT)
def _start_accepted(base: str, batch_id: str) -> httpx.Response:
"""Start, waiting out the uploader lane the previous attempt still holds.
A stopped attempt releases its job a moment after the batch itself reaches
``cancelled``; ``lock_held`` is that gap, not a refusal of this batch.
"""
def _attempt():
response = _start(base, batch_id)
if response.status_code == 409 and response.json()["error"]["code"] == "lock_held":
return None
response.raise_for_status()
return response
return wait_until(_attempt)
def _verify(base: str, batch_id: str) -> dict:
response = httpx.post(f"{base}/api/v1/upload-batches/{batch_id}/verify", timeout=TIMEOUT)
response.raise_for_status()
return response.json()
def _await_state(base: str, batch_id: str, states: set[str], *, timeout: float = 60) -> dict:
return wait_until(
lambda: (lambda b: b if b.get("state") in states else None)(_get(base, batch_id)),
timeout=timeout,
)
def _await_report(base: str, batch_id: str, states: set[str] = TERMINAL) -> dict:
"""Wait for a finished attempt *and* the report that explains it.
The batch state is recorded a moment before its report is parsed — the outcome of
the process and the outcome of each file are deliberately separate facts — so a
journey that reads per-item evidence must wait for the second one too.
"""
return wait_until(
lambda: (lambda b: b if b.get("state") in states and b.get("parsed_at") else None)(
_get(base, batch_id)
),
timeout=60,
)
def _approve(stack, **body) -> dict:
"""Preflight, approve exactly that report, and return the created batch."""
report = _preflight(stack.base, **body)
response = _create(stack.base, report, **body)
response.raise_for_status()
return response.json()["batches"][0]
def _upload(stack, **body) -> dict:
"""The whole approved journey, up to whatever terminal state it reaches."""
batch = _approve(stack, **body)
_start(stack.base, batch["id"]).raise_for_status()
return _await_report(stack.base, batch["id"])
def _outcomes(batch: dict) -> dict[str, str]:
return {Path(item["path"]).name: item["outcome"] for item in batch["items"]}
def _uploaded_folder(stack) -> Path:
"""The folder the uploader was actually handed, from its own argv log."""
invocations = stack.argv()
assert invocations, "the uploader was never executed"
return Path(invocations[-1].split()[-1])
def _assert_hashes_match_submitted_bytes(stack, batch: dict) -> None:
folder = _uploaded_folder(stack)
for item in batch["items"]:
submitted = folder / Path(item["path"]).name
raw = submitted.read_bytes()
assert item["sha256"] == hashlib.sha256(raw).hexdigest(), submitted
assert item["sha1"] == hashlib.sha1(raw).hexdigest(), submitted # noqa: S324 — Immich's
# ── credentials ──────────────────────────────────────────────────────────────
def test_missing_credentials_block_the_preflight_and_no_upload_is_attempted(stack):
stack.start(worker=False, credentials=False)
report = _preflight(stack.base)
assert report["state"] == "blocked"
assert [issue["code"] for issue in report["blockers"]] == ["credentials_missing"]
assert report["credentials"]["api_key_configured"] is False
# A blocked scope still issues a token; approving it is what is refused.
refused = _create(stack.base, report)
assert refused.status_code == 422
assert refused.json()["error"]["code"] == "not_ready"
assert stack.batches() == []
assert stack.argv() == [], "the uploader must not run without credentials"
def test_a_server_that_stops_answering_blocks_the_preflight(stack):
stack.start(worker=False)
assert _preflight(stack.base)["state"] == "ready"
stack.immich.stop() # Immich goes away between one preview and the next
report = _preflight(stack.base)
assert report["state"] == "blocked"
assert [issue["code"] for issue in report["blockers"]] == ["server_unreachable"]
assert report["server"]["reachable"] is False
assert _create(stack.base, report).status_code == 422
assert stack.argv() == []
# ── EXIF precedes upload ─────────────────────────────────────────────────────
def test_exif_checkpoints_must_be_written_before_anything_is_uploaded(unfinished):
stack = unfinished
stack.start()
report = _preflight(stack.base)
assert report["state"] == "blocked"
codes = {
issue["code"]
for album in report["albums"]
for asset in album["assets"]
for issue in asset["blockers"]
}
assert codes == {"safety_exif_unverified"}
assert _create(stack.base, report).status_code == 422
assert stack.argv() == [], "the uploader must not run before the EXIF checkpoints"
mark_upload_ready(stack.seeded) # the checkpoints are written
batch = _upload(stack)
assert batch["state"] == "succeeded"
assert stack.argv(), "the same scope uploads once its checkpoints exist"
# The recorded checkpoints predate the attempt that was allowed to run.
from sqlalchemy import select
from photo_pipeline.models import AnalysisResult, SafetyReview
with session_factory(stack.seeded) as sf, sf() as session:
checkpoints = [
*[row.exif_verified_at for row in session.scalars(select(SafetyReview))],
*[row.exif_written_at for row in session.scalars(select(AnalysisResult))],
]
started = batch["started_at"]
assert checkpoints and all(stamp.isoformat() < started for stamp in checkpoints)
def test_an_unfinished_photo_blocks_its_album_until_a_partial_upload_is_approved(tmp_path):
seeded = seed_album(tmp_path)
mark_upload_ready(seeded, unverified=("b",))
stack = UploadStack(tmp_path, seeded)
try:
stack.start()
report = _preflight(stack.base)
assert report["state"] == "blocked"
assert [issue["code"] for issue in report["albums"][0]["blockers"]] == ["partial_scope"]
assert report["totals"] == {
"albums": 1,
"ready_albums": 0,
"assets": 2,
"eligible": 1,
"blocked": 1,
}
assert _create(stack.base, report).status_code == 422
partial = _preflight(stack.base, allow_partial=True)
assert partial["state"] == "ready"
assert partial["token"] != report["token"]
# A full-scope approval can never be replayed as a partial one.
replayed = _create(stack.base, report, allow_partial=True)
assert replayed.status_code == 409
assert replayed.json()["error"]["code"] == "stale_preflight"
batch = _upload(stack, allow_partial=True)
assert batch["state"] == "succeeded"
assert batch["allow_partial"] is True
assert [Path(item["path"]).name for item in batch["items"]] == ["a.jpg"]
finally:
stack.stop()
# ── outcomes ─────────────────────────────────────────────────────────────────
def test_a_new_album_uploads_and_persists_the_hashes_that_were_submitted(stack):
stack.start(uploader=ALL_NEW)
batch = _upload(stack)
assert batch["state"] == "succeeded"
assert batch["outcome_state"] == "verified"
assert _outcomes(batch) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
assert batch["outcome_counts"]["uploaded"] == 2
assert batch["report_counts"] == {"uploaded": 2}
assert batch["parser"] == "text-v1"
_assert_hashes_match_submitted_bytes(stack, batch)
# One album, one invocation, scoped to that album's own folder.
assert len(stack.argv()) == 1
assert f"--album-name={ALBUM}" in stack.argv()[0]
assert _uploaded_folder(stack) == stack.seeded.lib / ALBUM
def test_an_exact_duplicate_is_recorded_as_a_duplicate_not_a_new_asset(stack):
stack.start(uploader=EXACT_DUPLICATES)
batch = _upload(stack)
assert batch["state"] == "succeeded"
assert batch["outcome_state"] == "verified"
assert set(_outcomes(batch).values()) == {"duplicate"}
assert batch["outcome_counts"]["uploaded"] == 0
_assert_hashes_match_submitted_bytes(stack, batch)
def test_a_better_copy_is_recorded_as_an_upgrade(stack):
stack.start(uploader=UPGRADES)
batch = _upload(stack)
assert batch["state"] == "succeeded"
assert batch["outcome_state"] == "verified"
assert set(_outcomes(batch).values()) == {"upgraded"}
assert batch["report_counts"] == {"upgraded": 2}
# ── failure and retry ────────────────────────────────────────────────────────
def test_a_plain_uploader_failure_is_retryable_and_the_retry_succeeds(stack, tmp_path):
stack.start(uploader=_once_then(tmp_path, *FAILING_FIRST))
failed = _upload(stack)
assert failed["state"] == "failed"
assert failed["error_code"] == "uploader_failed"
assert failed["exit_code"] == 1
assert _outcomes(failed)["a.jpg"] == "failed"
# Nothing uncertain happened, so the batch is offered again rather than blocked.
assert failed["retry_blockers"] == []
_start_accepted(stack.base, failed["id"])
retried = _await_report(stack.base, failed["id"], {"succeeded"})
assert retried["attempt_count"] == 2
assert _outcomes(retried) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
assert len(stack.argv()) == 2
# ── uncertainty ──────────────────────────────────────────────────────────────
def test_a_lost_acceptance_response_stays_uncertain_and_is_never_retried_blindly(stack):
stack.start(uploader=SILENT_UPLOADER)
batch = _upload(stack)
# The process succeeded; what happened to each file is simply not known.
assert batch["state"] == "succeeded"
assert batch["outcome_state"] == "requires_verification"
assert set(_outcomes(batch).values()) == {"unknown"}
refused = _start(stack.base, batch["id"])
assert refused.status_code == 409
assert refused.json()["error"]["code"] == "not_runnable"
assert len(stack.argv()) == 1, "a blind retry must not reach the uploader"
def test_verification_asks_immich_and_resolves_every_uncertain_item(stack):
stack.start(uploader=SILENT_UPLOADER)
stack.immich.mode("present") # Immich holds exactly the bytes that were sent
batch = _upload(stack)
result = _verify(stack.base, batch["id"])
assert result["counts"] == {"present": 2}
assert result["outcome_state"] == "verified"
verified = _get(stack.base, batch["id"])
assert set(_outcomes(verified).values()) == {"uploaded"}
history = httpx.get(
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications", timeout=TIMEOUT
).json()["verifications"]
assert {entry["source"] for entry in history} == {"immich_api"}
def test_an_unusable_answer_stays_uncertain_until_an_operator_records_evidence(stack):
stack.start(uploader=SILENT_UPLOADER)
stack.immich.mode("broken") # answers, but nothing this adapter will interpret
batch = _upload(stack)
result = _verify(stack.base, batch["id"])
# No answer is never "no": the items stay uncertain rather than being called failed.
assert result["counts"] == {"inconclusive": 2}
assert result["outcome_state"] == "requires_verification"
assert set(_outcomes(_get(stack.base, batch["id"])).values()) == {"unknown"}
asset_id = batch["items"][0]["asset_id"]
incomplete = httpx.post(
f"{stack.base}/api/v1/upload-batches/{batch['id']}/resolve",
json={"asset_id": asset_id, "outcome": "uploaded", "evidence": "", "actor": "dom"},
timeout=TIMEOUT,
)
assert incomplete.status_code == 422
assert incomplete.json()["error"]["code"] == "evidence_required"
resolved = httpx.post(
f"{stack.base}/api/v1/upload-batches/{batch['id']}/resolve",
json={
"asset_id": asset_id,
"outcome": "uploaded",
"evidence": "found it in Immich by checksum",
"actor": "dom",
},
timeout=TIMEOUT,
)
resolved.raise_for_status()
assert _outcomes(_get(stack.base, batch["id"]))["a.jpg"] == "uploaded"
manual = httpx.get(
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications", timeout=TIMEOUT
).json()["verifications"][-1]
assert manual["source"] == "operator"
assert manual["actor"] == "dom"
def test_bytes_edited_after_the_upload_are_flagged_and_block_another_run(stack):
stack.start(uploader=ALL_NEW)
batch = _upload(stack)
assert batch["state"] == "succeeded"
(stack.seeded.lib / ALBUM / "a.jpg").write_bytes(b"edited after the upload")
result = _verify(stack.base, batch["id"])
assert result["stale_bytes"] is True
changed = [item for item in result["items"] if item["changed_after_upload"]]
assert [Path(item["path"]).name for item in changed] == ["a.jpg"]
refused = _start(stack.base, batch["id"])
assert refused.status_code == 409
assert refused.json()["error"]["code"] == "changed_after_upload"
# The preflight agrees: those bytes are no longer approved for any new upload.
report = _preflight(stack.base)
assert report["state"] == "blocked"
assert "bytes_changed" in {
issue["code"]
for album in report["albums"]
for asset in album["assets"]
for issue in asset["blockers"]
}
# ── cancellation and resume ──────────────────────────────────────────────────
def test_a_running_album_can_be_stopped_and_run_again_from_that_boundary(stack, tmp_path):
stack.start(uploader=_once_then(tmp_path, *SLOW_FIRST))
batch = _approve(stack)
_start(stack.base, batch["id"]).raise_for_status()
_await_state(stack.base, batch["id"], {"running"})
httpx.post(
f"{stack.base}/api/v1/upload-batches/{batch['id']}/cancel", timeout=TIMEOUT
).raise_for_status()
cancelled = _await_state(stack.base, batch["id"], {"cancelled"})
# A stopped album is a clean boundary, not an uncertain one.
assert cancelled["retry_blockers"] == []
_start_accepted(stack.base, batch["id"])
resumed = _await_report(stack.base, batch["id"], {"succeeded"})
assert resumed["attempt_count"] == 2
assert _outcomes(resumed) == {"a.jpg": "uploaded", "b.jpg": "uploaded"}
def test_an_interrupted_attempt_is_uncertain_after_a_restart_and_stays_blocked(stack, tmp_path):
"""The worker vanishes mid-upload: Immich may hold the files, so the outcome is
unknown. Startup recovery must say so, refuse a retry, and survive the restart."""
stack.start(uploader=_once_then(tmp_path, *SLOW_FIRST))
batch = _approve(stack)
_start(stack.base, batch["id"]).raise_for_status()
_await_state(stack.base, batch["id"], {"running"})
stack.worker.kill() # no chance to record any outcome
stack.worker.wait(timeout=10)
stack.restart_server()
recovered = _get(stack.base, batch["id"])
assert recovered["state"] == "unknown_requires_verification"
assert recovered["error_code"] == "interrupted"
assert recovered["attempt_count"] == 1
refused = _start(stack.base, batch["id"])
assert refused.status_code == 409
assert refused.json()["error"]["code"] == "requires_verification"
# Verification is the only way out, and it is what makes the batch certain again.
stack.immich.mode("present")
result = _verify(stack.base, batch["id"])
assert result["state"] == "succeeded"
assert result["outcome_state"] == "verified"
stack.restart_server() # the resolution is durable, not in-process memory
after = _get(stack.base, batch["id"])
assert after["state"] == "succeeded"
assert set(_outcomes(after).values()) == {"uploaded"}
# ── privacy ──────────────────────────────────────────────────────────────────
def test_no_secret_appears_in_any_retained_artifact(stack):
stack.start(uploader=ALL_NEW)
planned = _approve(stack)
job = _start(stack.base, planned["id"]).json()["job"]
batch = _await_report(stack.base, planned["id"])
stack.immich.mode("present")
_verify(stack.base, batch["id"])
# The uploader really was given the key…
assert f"--api-key={SENTINEL_KEY}" in stack.argv()[0]
# …and it is in nothing that was kept: not the database, not the retained report,
# not any response the browser can read.
retained = [path for path in stack.seeded.data.rglob("*") if path.is_file()]
assert any(path.suffix == ".log" for path in retained), "the report was not retained"
for path in retained:
assert SENTINEL_KEY.encode() not in path.read_bytes(), path
for url in (
f"{stack.base}/api/v1/upload-batches",
f"{stack.base}/api/v1/upload-batches/{batch['id']}",
f"{stack.base}/api/v1/upload-batches/{batch['id']}/verifications",
f"{stack.base}/api/v1/jobs/{job['id']}",
f"{stack.base}/api/v1/jobs/{job['id']}/events",
):
assert SENTINEL_KEY not in httpx.get(url, timeout=TIMEOUT).text, url
assert SENTINEL_KEY not in httpx.post(
f"{stack.base}/api/v1/upload-preflight", json={}, timeout=TIMEOUT
).text
assert "--api-key=***" in " ".join(_get(stack.base, batch["id"])["command"])

View File

@@ -13,6 +13,7 @@ REPO = Path(__file__).resolve().parents[2]
MAP = json.loads((REPO / "tests" / "story_traceability.json").read_text())["stories"]
PHASE_A_STORIES = {f"US01-0{n}" for n in range(1, 8)}
PHASE_D_STORIES = {f"US04-0{n}" for n in range(1, 7)}
PHASE_E_STORIES = {f"US05-0{n}" for n in range(1, 7)}
def test_all_phase_a_stories_are_mapped():
@@ -25,6 +26,12 @@ def test_all_phase_d_stories_are_mapped():
assert PHASE_D_STORIES <= set(MAP)
def test_all_phase_e_stories_are_mapped():
"""US05-06 acceptance: every upload story, preflight through browser, is tied to
automated tests — upload is the one stage the app cannot take back."""
assert PHASE_E_STORIES <= set(MAP)
def test_every_mapped_test_file_exists_and_is_nonempty():
for story, files in MAP.items():
assert files, f"{story} maps to no tests"

View File

@@ -0,0 +1,283 @@
"""Browser journeys for the upload view (US05-05).
Covers the preflight preview (scope, redacted configuration, blockers, exact
confirmation), a real upload through the real worker with per-outcome progress,
stopping a running album, and the two ways out of an uncertain outcome —
verification against Immich and a manual resolution that records its evidence.
Nothing external is mocked inside the browser: the uploader is a real executable
driven by the real worker process, and Immich is a real HTTP server answering the
same ``ping``/``bulk-upload-check`` endpoints the adapter calls in production. The
API key is a sentinel string, so the last test can prove it never reached the page.
"""
from __future__ import annotations
import httpx
import pytest
from playwright.sync_api import expect
from tests.e2e._pipeline_harness import (
SENTINEL_KEY,
SILENT_UPLOADER,
UPLOADER_VERSION,
UploadStack,
mark_upload_ready,
seed_album,
session_factory,
wait_until,
)
pytestmark = pytest.mark.phase_e
TIMEOUT = 10
@pytest.fixture
def stack(tmp_path):
seeded = seed_album(tmp_path)
mark_upload_ready(seeded)
running = UploadStack(tmp_path, seeded)
try:
yield running
finally:
running.stop()
def _open(page, stack) -> None:
page.goto(f"{stack.base}/app/#/uploads")
page.get_by_test_id("upload-scope").wait_for()
def _upload(page, stack) -> None:
"""Confirm the upload and wait for the worker to finish the album."""
_open(page, stack)
page.get_by_test_id("start-upload").click()
expect(page.get_by_test_id("detail-state")).not_to_have_text("planned", timeout=30_000)
# ── preflight ────────────────────────────────────────────────────────────────
def test_the_preview_shows_scope_configuration_and_an_exact_confirmation(page, stack):
errors = []
page.on("console", lambda m: errors.append(m.text) if m.type == "error" else None)
stack.start(worker=False)
_open(page, stack)
expect(page.get_by_test_id("config-server")).to_have_text(stack.immich.url)
expect(page.get_by_test_id("config-key")).to_have_text("configured (never shown)")
expect(page.get_by_test_id("config-reachable")).to_have_text("yes")
expect(page.get_by_test_id("config-uploader")).to_have_text(UPLOADER_VERSION)
row = page.get_by_test_id("album-row").first
expect(row.get_by_test_id("album-name")).to_have_text("rome")
expect(row.get_by_test_id("album-immich-name")).to_have_text("rome")
expect(row.get_by_test_id("album-eligible")).to_have_text("2")
expect(row.get_by_test_id("album-state")).to_have_text("ready")
# The exact invocation is previewed, with the key masked at the source.
command = row.get_by_test_id("album-command").inner_text()
assert "upload from-folder" in command and "--album-name=rome" in command
assert "--api-key=***" in command
# The confirmation names the scope it is about to send, not just "Upload".
expect(page.get_by_test_id("start-upload")).to_have_text("Upload 1 album(s) · 2 photo(s)")
expect(page.get_by_test_id("start-upload")).to_be_enabled()
assert errors == [], f"console errors: {errors}"
def test_an_unfinished_photo_blocks_its_album_and_the_confirmation(page, stack):
mark_upload_ready(stack.seeded, unverified=("b",))
stack.start(worker=False)
_open(page, stack)
expect(page.get_by_test_id("album-state")).to_have_text("blocked")
expect(page.get_by_test_id("album-blocker")).to_have_attribute("data-code", "partial_scope")
expect(page.get_by_test_id("album-eligible")).to_have_text("1")
expect(page.get_by_test_id("album-blocked")).to_have_text("1")
expect(page.get_by_test_id("start-upload")).to_be_disabled()
# Partial upload exists, but only as a deliberate act: ticking it re-runs the
# preflight under that policy and the confirmation then names the smaller scope.
page.get_by_test_id("allow-partial").check()
expect(page.get_by_test_id("album-state")).to_have_text("ready")
expect(page.get_by_test_id("start-upload")).to_have_text("Upload 1 album(s) · 1 photo(s)")
assert stack.batches() == [], "nothing may be created by previewing"
# ── uploading ────────────────────────────────────────────────────────────────
def test_a_confirmed_upload_runs_and_reports_each_outcome(page, stack):
stack.start()
_upload(page, stack)
expect(page.get_by_test_id("detail-state")).to_have_text("succeeded")
expect(page.get_by_test_id("count-new")).to_have_text("new: 1")
expect(page.get_by_test_id("count-duplicate")).to_have_text("duplicate: 1")
expect(page.get_by_test_id("count-uncertain")).to_have_text("uncertain: 0")
expect(page.get_by_test_id("count-failed")).to_have_text("failed: 0")
expect(page.get_by_test_id("batch-outcome-state")).to_have_text("verified")
outcomes = sorted(page.get_by_test_id("item-outcome").all_inner_texts())
assert outcomes == ["duplicate", "new"]
# A finished album offers no restart: the server would refuse one.
expect(page.get_by_test_id("retry-blocked")).to_contain_text("not_runnable")
expect(page.get_by_test_id("start-batch")).to_have_count(0)
def test_a_running_album_can_be_stopped(page, stack):
stack.start(uploader='echo "INFO starting"; sleep 20; exit 0\n')
_open(page, stack)
page.get_by_test_id("start-upload").click()
stop = page.get_by_test_id("cancel-batch")
expect(stop).to_have_text("Stop after the current file", timeout=30_000)
stop.click()
expect(page.get_by_test_id("detail-state")).to_have_text("cancelled", timeout=30_000)
# A stopped album is a clean boundary, not an uncertain one: it can run again.
expect(page.get_by_test_id("start-batch")).to_be_visible()
def test_the_finished_upload_survives_a_reload(page, stack):
stack.start()
_upload(page, stack)
expect(page.get_by_test_id("detail-state")).to_have_text("succeeded")
page.reload()
expect(page.get_by_test_id("detail-state")).to_have_text("succeeded")
expect(page.get_by_test_id("count-new")).to_have_text("new: 1")
# The result banner is this tab's memory, not server state, so it stays gone.
expect(page.get_by_test_id("upload-result")).to_have_count(0)
# ── uncertainty ──────────────────────────────────────────────────────────────
def test_an_uncertain_outcome_offers_verification_and_no_retry(page, stack):
stack.start(uploader=SILENT_UPLOADER)
_upload(page, stack)
expect(page.get_by_test_id("detail-state")).to_have_text("succeeded")
expect(page.get_by_test_id("batch-outcome-state")).to_have_text("requires_verification")
expect(page.get_by_test_id("count-uncertain")).to_have_text("uncertain: 2")
expect(page.get_by_test_id("uncertain")).to_be_visible()
# The point of the story: verification is offered, a retry is not.
expect(page.get_by_test_id("verify-batch")).to_be_visible()
expect(page.get_by_test_id("start-batch")).to_have_count(0)
def test_verification_asks_immich_and_resolves_the_uncertain_items(page, stack):
stack.start(uploader=SILENT_UPLOADER)
stack.immich.mode("present") # Immich holds exactly the bytes that were sent
_upload(page, stack)
page.get_by_test_id("verify-batch").click()
expect(page.get_by_test_id("batch-outcome-state")).to_have_text("verified")
expect(page.get_by_test_id("count-new")).to_have_text("new: 2")
expect(page.get_by_test_id("count-uncertain")).to_have_text("uncertain: 0")
expect(page.get_by_test_id("item-verification").first).to_have_text("present")
expect(page.get_by_test_id("history-entry").first).to_have_attribute(
"data-source", "immich_api"
)
expect(page.get_by_test_id("uncertain")).to_have_count(0)
def test_an_unusable_answer_stays_uncertain_until_someone_records_evidence(page, stack):
stack.start(uploader=SILENT_UPLOADER)
stack.immich.mode("broken") # answers, but nothing this adapter will interpret
_upload(page, stack)
page.get_by_test_id("verify-batch").click()
# No answer is never "no": the items stay uncertain rather than being called failed.
expect(page.get_by_test_id("item-verification").first).to_have_text("inconclusive")
expect(page.get_by_test_id("uncertain")).to_be_visible()
expect(page.get_by_test_id("start-batch")).to_have_count(0)
row = page.get_by_test_id("item-row").first
row.get_by_test_id("resolve-outcome").select_option("uploaded")
row.get_by_test_id("resolve-evidence").fill("found it in Immich by checksum")
row.get_by_test_id("resolve-actor").fill("dom")
row.get_by_test_id("resolve-item").click()
expect(page.get_by_test_id("item-row").first.get_by_test_id("item-outcome")).to_have_text("new")
manual = page.get_by_test_id("history-entry").last
expect(manual).to_have_attribute("data-source", "operator")
expect(manual).to_contain_text("found it in Immich by checksum")
expect(manual).to_contain_text("dom")
def test_bytes_changed_after_upload_are_flagged_and_block_another_run(page, stack):
stack.start()
_upload(page, stack)
expect(page.get_by_test_id("detail-state")).to_have_text("succeeded")
# The user edits a photo after it was uploaded; Immich still holds the old bytes.
(stack.seeded.lib / "rome" / "a.jpg").write_bytes(b"edited after the upload")
page.get_by_test_id("verify-batch").click()
expect(page.get_by_test_id("stale-bytes")).to_be_visible()
expect(page.get_by_test_id("item-changed")).to_have_count(1)
expect(page.get_by_test_id("retry-blocked")).to_contain_text("changed_after_upload")
expect(page.get_by_test_id("start-batch")).to_have_count(0)
# The preflight agrees: those bytes are no longer approved for any new upload.
expect(page.get_by_test_id("album-state")).to_have_text("blocked")
expect(page.get_by_test_id("album-blocked")).to_have_text("1")
# ── privacy ──────────────────────────────────────────────────────────────────
def test_the_api_key_never_reaches_the_browser(page, stack):
logs = []
page.on("console", lambda message: logs.append(message.text))
stack.start()
_upload(page, stack)
page.get_by_test_id("verify-batch").click()
expect(page.get_by_test_id("item-verification").first).to_have_text("present")
storage = page.evaluate(
"() => JSON.stringify([{...localStorage}, {...sessionStorage}, document.cookie])"
)
assert SENTINEL_KEY not in page.content()
assert SENTINEL_KEY not in page.url
assert SENTINEL_KEY not in storage
assert SENTINEL_KEY not in "\n".join(logs)
# The uploader was given the real key even though nothing on the page shows it.
report = wait_until(lambda: sorted((stack.seeded.data / "uploads").glob("*.log")))[0]
assert SENTINEL_KEY not in report.read_text()
# ── recovery ─────────────────────────────────────────────────────────────────
def test_an_interrupted_attempt_is_shown_as_uncertain_after_a_restart(page, stack):
"""A worker that vanished mid-upload leaves a batch whose outcome nobody knows.
Startup recovery marks it uncertain, and the view must not offer to retry it."""
from photo_pipeline.models import UploadBatch
stack.start(worker=False)
_open(page, stack)
token = httpx.post(f"{stack.base}/api/v1/upload-preflight", json={}, timeout=TIMEOUT).json()[
"token"
]
created = httpx.post(
f"{stack.base}/api/v1/upload-batches", json={"token": token}, timeout=TIMEOUT
).json()["batches"][0]
with session_factory(stack.seeded) as sf: # what a killed worker leaves behind
with sf() as session:
session.get(UploadBatch, created["id"]).state = "running"
session.commit()
stack.restart_server() # the same port, so recovery runs in a genuinely fresh process
page.reload()
expect(page.get_by_test_id("detail-state")).to_have_text("unknown_requires_verification")
expect(page.get_by_test_id("batch-error")).to_contain_text("interrupted")
expect(page.get_by_test_id("uncertain")).to_be_visible()
expect(page.get_by_test_id("start-batch")).to_have_count(0)
expect(page.get_by_test_id("retry-blocked")).to_contain_text("requires_verification")

View File

@@ -38,6 +38,8 @@ from photo_pipeline.services.upload_batches import (
)
from photo_pipeline.services.uploads import UploadService
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SENTINEL_KEY = "immich-sentinel-9f3a2b"
UPLOADER_VERSION = "immich-go 0.21.0"

View File

@@ -29,6 +29,8 @@ from photo_pipeline.models import (
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.uploads import UploadError, UploadService
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
# A sentinel credential: every assertion below proves it never leaves configuration.
SENTINEL_KEY = "immich-sentinel-9f3a2b"

View File

@@ -32,6 +32,8 @@ from photo_pipeline.services.upload_reports import (
)
from photo_pipeline.services.uploads import UploadService
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SUPPORTED_VERSION = "immich-go 0.21.0"

View File

@@ -0,0 +1,537 @@
"""Verifying, retrying, and resolving uncertain uploads (US05-04).
The Immich boundary is a real HTTP server here: it answers ``/api/server/ping``
and ``/api/assets/bulk-upload-check`` exactly as the app's own adapter parses
them, and the set of checksums it "holds" is what each fault scenario controls.
The uploader stays a real executable driven through the real batch service, so
every state under test is reached the way production reaches it.
The invariant these tests defend is one-directional: uncertainty may only become
success when something authoritative said so — the server, or an operator who
recorded what they checked.
"""
import json
import stat
import threading
import uuid
from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, HTTPServer
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 AnalysisResult, Asset, SafetyReview, UploadBatch
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.upload_batches import BatchError, BatchState, UploadBatchService
from photo_pipeline.services.upload_reports import REQUIRES_VERIFICATION, VERIFIED
from photo_pipeline.services.upload_verification import (
ABSENT,
INCONCLUSIVE,
MANUAL,
PRESENT,
UploadVerificationService,
VerificationError,
retry_blockers,
)
from photo_pipeline.services.uploads import UploadService
pytestmark = pytest.mark.phase_e # part of the Phase E acceptance gate (US05-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SENTINEL_KEY = "immich-sentinel-9f3a2b"
UPLOADER_VERSION = "immich-go 0.21.0"
# ── fake Immich ──────────────────────────────────────────────────────────────
class _FakeImmich:
"""The server's view of the world: which checksums it holds, and whether it
is willing to answer at all."""
def __init__(self) -> None:
self.held: set[str] = set()
self.available = True
self.checked: list[str] = []
def _handler(state: _FakeImmich):
class Handler(BaseHTTPRequestHandler):
def do_GET(self): # noqa: N802 (BaseHTTPRequestHandler API)
self._json(200, {"res": "pong"})
def do_POST(self): # noqa: N802
body = json.loads(self.rfile.read(int(self.headers["Content-Length"] or 0)) or "{}")
if not state.available:
self._json(503, {"error": "unavailable"})
return
if self.headers.get("x-api-key") != SENTINEL_KEY:
self._json(401, {"error": "unauthorized"})
return
results = []
for asset in body.get("assets", []):
state.checked.append(asset["checksum"])
results.append(
{"id": asset["id"], "action": "reject", "reason": "duplicate"}
if asset["checksum"] in state.held
else {"id": asset["id"], "action": "accept"}
)
self._json(200, {"results": results})
def _json(self, status: int, payload: dict) -> None:
body = json.dumps(payload).encode()
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def log_message(self, *args):
pass
return Handler
@pytest.fixture
def immich():
state = _FakeImmich()
server = HTTPServer(("127.0.0.1", 0), _handler(state))
threading.Thread(target=server.serve_forever, daemon=True).start()
state.url = f"http://127.0.0.1:{server.server_port}"
yield state
server.shutdown()
server.server_close()
# ── environment ──────────────────────────────────────────────────────────────
def _uploader(tmp_path, report_body: str = "", *, exit_code: int = 0, version=UPLOADER_VERSION):
path = tmp_path / "immich-go"
path.write_text(
"#!/bin/sh\n"
f'if [ "$1" = "--version" ]; then echo "{version}"; exit 0; fi\n'
f"cat <<'REPORT'\n{report_body}\nREPORT\n"
f"exit {exit_code}\n"
)
path.chmod(path.stat().st_mode | stat.S_IEXEC | stat.S_IXGRP | stat.S_IXOTH)
return path
def _env(tmp_path, immich, uploader=None):
(tmp_path / "data").mkdir(exist_ok=True)
lib = tmp_path / "lib"
lib.mkdir(exist_ok=True)
config = Config.from_env(
{
"PHOTO_PIPELINE_DATA_DIR": str(tmp_path / "data"),
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib),
"PHOTO_PIPELINE_IMMICH_SERVER_URL": immich.url,
"PHOTO_PIPELINE_IMMICH_API_KEY": SENTINEL_KEY,
"PHOTO_PIPELINE_IMMICH_GO_BINARY": str(
uploader if uploader is not None else _uploader(tmp_path)
),
}
)
run_migrations(config.database_url)
return config, create_session_factory(create_db_engine(config.database_url)), lib
def _album(sf, lib, album="rome", names=("a.jpg", "b.jpg")):
folder = lib / album
folder.mkdir(parents=True, exist_ok=True)
with sf() as session:
for name in names:
path = folder / name
path.write_bytes(f"{album}/{name}".encode() * 16)
asset_id = str(uuid.uuid4())
session.add(
Asset(
id=asset_id,
original_path=str(path),
current_path=str(path),
discovered_at=NOW,
hash_version=1,
byte_size=path.stat().st_size,
current_sha256=sha256_file(path),
)
)
session.add(
SafetyReview(
id=str(uuid.uuid4()), asset_id=asset_id, decision="sfw", exif_verified_at=NOW
)
)
session.add(AnalysisResult(asset_id=asset_id, status="analyzed", exif_written_at=NOW))
session.commit()
return folder
def _batch(sf, config, albums=None):
report = UploadService(sf, config=config).preflight(albums)
assert report["state"] == "ready", report["blockers"]
(batch,) = UploadBatchService(sf, config=config).create(albums, token=report["token"])
return batch
def _interrupt(sf, config, batch_id):
"""Leave behind exactly what a killed worker leaves: a mid-flight attempt."""
with sf() as session:
session.get(UploadBatch, batch_id).state = BatchState.RUNNING
session.commit()
UploadBatchService(sf, config=config).recover()
def _server_holds(immich, batch):
immich.held.update(item["sha1"] for item in batch["items"])
def _by_name(batch) -> dict:
return {Path(item["path"]).name: item for item in batch["items"]}
def _uncertain_batch(tmp_path, immich, *, hold: bool):
"""A batch whose attempt died mid-upload, with the server holding the bytes or not."""
config, sf, lib = _env(tmp_path, immich)
_album(sf, lib)
batch = _batch(sf, config)
if hold:
_server_holds(immich, batch)
_interrupt(sf, config, batch["id"])
return config, sf, lib, batch
# ── accepted-but-unknown ─────────────────────────────────────────────────────
def test_acceptance_then_lost_response_is_confirmed_by_the_server(tmp_path, immich):
"""The classic lost response: Immich took the files, the app never saw a report."""
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["state"] == BatchState.SUCCEEDED, "server evidence resolves the uncertainty"
assert result["outcome_state"] == VERIFIED
assert result["counts"] == {PRESENT: 2}
verified = UploadBatchService(sf, config=config).get(batch["id"])
for item in verified["items"]:
assert item["outcome"] == "uploaded"
assert item["verification"] == PRESENT
assert item["sha1"] in item["evidence"], "the exact bytes are named in the evidence"
assert immich.checked, "verification actually asked the server"
def test_timeout_before_acceptance_leaves_a_safe_retry(tmp_path, immich):
"""Nothing arrived, so the batch becomes a plain failure and may run again."""
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=False)
service = UploadBatchService(sf, config=config)
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["state"] == BatchState.FAILED
assert result["counts"] == {ABSENT: 2}
assert {item["outcome"] for item in service.get(batch["id"])["items"]} == {"failed"}
assert retry_blockers(service.get(batch["id"])) == []
assert service.run(batch["id"])["state"] == BatchState.SUCCEEDED
def test_an_uncertain_batch_cannot_be_retried_before_verification(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
service = UploadBatchService(sf, config=config)
with pytest.raises(BatchError) as error:
service.run(batch["id"])
assert error.value.code == "requires_verification"
with TestClient(create_app(config)) as client:
response = client.post(f"/api/v1/upload-batches/{batch['id']}/start")
assert response.status_code == 409
assert response.json()["error"]["code"] == "requires_verification"
def test_a_partly_arrived_batch_is_a_failure_not_a_success(tmp_path, immich):
config, sf, lib = _env(tmp_path, immich)
_album(sf, lib)
batch = _batch(sf, config)
immich.held.add(batch["items"][0]["sha1"]) # only one file made it
_interrupt(sf, config, batch["id"])
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["state"] == BatchState.FAILED
assert result["counts"] == {PRESENT: 1, ABSENT: 1}
assert result["outcome_state"] == VERIFIED, "every file is accounted for"
# ── parser uncertainty ───────────────────────────────────────────────────────
def test_parser_uncertainty_is_settled_by_server_evidence(tmp_path, immich):
"""An unpinned uploader version leaves every item unknown; the server decides."""
config, sf, lib = _env(tmp_path, immich, _uploader(tmp_path, "done", version="immich-go 9.9.9"))
_album(sf, lib)
batch = _batch(sf, config)
service = UploadBatchService(sf, config=config)
ran = service.run(batch["id"])
assert ran["state"] == BatchState.SUCCEEDED and ran["outcome_state"] == REQUIRES_VERIFICATION
_server_holds(immich, ran)
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["outcome_state"] == VERIFIED
assert {item["outcome"] for item in service.get(batch["id"])["items"]} == {"uploaded"}
def test_an_unreachable_server_never_turns_uncertainty_into_success(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
immich.available = False
service = UploadBatchService(sf, config=config)
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["server_reachable"] is False
assert result["detail"], "the reason the server could not answer is reported"
assert result["counts"] == {INCONCLUSIVE: 2}
assert result["state"] == BatchState.UNKNOWN, "still uncertain, not succeeded"
assert result["outcome_state"] == REQUIRES_VERIFICATION
assert [b["code"] for b in retry_blockers(service.get(batch["id"]))] == [
"requires_verification"
]
# ── safe failures ────────────────────────────────────────────────────────────
def test_a_plain_uploader_failure_is_retryable_without_verification(tmp_path, immich):
"""Nothing uncertain happened: the process failed before/while reporting an error."""
config, sf, lib = _env(tmp_path, immich, _uploader(tmp_path, "boom", exit_code=1))
_album(sf, lib)
batch = _batch(sf, config)
service = UploadBatchService(sf, config=config)
assert service.run(batch["id"])["state"] == BatchState.FAILED
assert retry_blockers(service.get(batch["id"])) == []
with TestClient(create_app(config)) as client:
assert client.post(f"/api/v1/upload-batches/{batch['id']}/start").status_code == 200
# ── changed bytes ────────────────────────────────────────────────────────────
def test_bytes_changed_after_upload_warn_and_block_a_rerun(tmp_path, immich):
config, sf, lib = _env(tmp_path, immich)
folder = _album(sf, lib)
batch = _batch(sf, config)
service = UploadBatchService(sf, config=config)
ran = service.run(batch["id"])
_server_holds(immich, ran)
(folder / "a.jpg").write_bytes(b"edited after the upload")
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert result["stale_bytes"] is True
changed = _by_name(result)["a.jpg"]
assert changed["changed_after_upload"] is True
assert _by_name(result)["b.jpg"]["changed_after_upload"] is False
stored = service.get(batch["id"])
assert stored["stale_bytes"] is True, "the warning is durable, not only in the response"
assert _by_name(stored)["a.jpg"]["observed_sha256"] != _by_name(stored)["a.jpg"]["sha256"]
with pytest.raises(BatchError) as error:
service.run(batch["id"])
assert error.value.code == "changed_after_upload"
with TestClient(create_app(config)) as client:
response = client.post(f"/api/v1/upload-batches/{batch['id']}/start")
assert response.status_code == 409
assert response.json()["error"]["code"] == "changed_after_upload"
def test_a_deleted_file_counts_as_changed_after_upload(tmp_path, immich):
config, sf, lib = _env(tmp_path, immich)
folder = _album(sf, lib)
batch = _batch(sf, config)
ran = UploadBatchService(sf, config=config).run(batch["id"])
_server_holds(immich, ran)
(folder / "a.jpg").unlink()
result = UploadVerificationService(sf, config=config).verify(batch["id"])
assert _by_name(result)["a.jpg"]["changed_after_upload"] is True
assert result["stale_bytes"] is True
# The bytes are still on the server: verification is about the upload, not the
# local file's continued existence.
assert _by_name(result)["a.jpg"]["verification"] == PRESENT
# ── repeated verification and history ────────────────────────────────────────
def test_repeated_verification_converges_and_keeps_every_answer(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
service = UploadVerificationService(sf, config=config)
first = service.verify(batch["id"])
second = service.verify(batch["id"])
assert first["counts"] == second["counts"] == {PRESENT: 2}
assert first["state"] == second["state"] == BatchState.SUCCEEDED
history = service.history(batch["id"])
assert len(history) == 4, "the audit trail appends, it never overwrites"
assert {entry["source"] for entry in history} == {"immich_api"}
assert all(entry["action"] == "verify" for entry in history)
def test_verification_that_changes_its_mind_keeps_both_answers(tmp_path, immich):
"""A file that was absent and later present must show both, in order."""
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=False)
service = UploadVerificationService(sf, config=config)
service.verify(batch["id"])
_server_holds(immich, batch)
service.verify(batch["id"])
asset_id = batch["items"][0]["asset_id"]
results = [e["result"] for e in service.history(batch["id"]) if e["asset_id"] == asset_id]
assert results == [ABSENT, PRESENT]
# ── manual resolution ────────────────────────────────────────────────────────
def test_manual_resolution_requires_evidence_and_an_author(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
service = UploadVerificationService(sf, config=config)
asset_id = batch["items"][0]["asset_id"]
for kwargs, code in (
({"evidence": " ", "actor": "dom"}, "evidence_required"),
({"evidence": "checked in Immich", "actor": ""}, "actor_required"),
({"evidence": "checked", "actor": "dom", "outcome": "definitely-fine"}, "invalid_outcome"),
):
with pytest.raises(VerificationError) as error:
service.resolve(batch["id"], asset_id, **{"outcome": "uploaded", **kwargs})
assert error.value.code == code
assert service.history(batch["id"]) == [], "a refused resolution records nothing"
assert UploadBatchService(sf, config=config).get(batch["id"])["state"] == BatchState.UNKNOWN
def test_manual_resolution_is_recorded_as_operator_evidence(tmp_path, immich):
"""An operator may settle what the server cannot — but never anonymously."""
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
immich.available = False
service = UploadVerificationService(sf, config=config)
service.verify(batch["id"]) # inconclusive: the server is down
for item in batch["items"]:
result = service.resolve(
batch["id"],
item["asset_id"],
outcome="uploaded",
evidence="found in Immich by checksum in the web UI",
actor="dom",
)
assert result["state"] == BatchState.SUCCEEDED
assert result["outcome_state"] == VERIFIED
resolutions = [e for e in service.history(batch["id"]) if e["action"] == "resolve"]
assert len(resolutions) == 2
assert {e["actor"] for e in resolutions} == {"dom"}
assert {e["source"] for e in resolutions} == {"operator"}
assert all("web UI" in e["evidence"] for e in resolutions)
stored = UploadBatchService(sf, config=config).get(batch["id"])
assert {item["verification"] for item in stored["items"]} == {MANUAL}, (
"a manual answer stays distinguishable from server evidence"
)
def test_resolving_one_item_leaves_the_batch_uncertain(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
immich.available = False
service = UploadVerificationService(sf, config=config)
service.verify(batch["id"])
result = service.resolve(
batch["id"],
batch["items"][0]["asset_id"],
outcome="uploaded",
evidence="visible in Immich",
actor="dom",
)
assert result["state"] == BatchState.UNKNOWN
assert result["outcome_state"] == REQUIRES_VERIFICATION
def test_resolving_an_unknown_item_or_batch_is_refused(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
service = UploadVerificationService(sf, config=config)
for batch_id, asset_id in ((batch["id"], "not-in-this-batch"), ("no-such-batch", "x")):
with pytest.raises(VerificationError) as error:
service.resolve(batch_id, asset_id, outcome="uploaded", evidence="checked", actor="dom")
assert error.value.code == "not_found"
def test_verifying_an_unknown_batch_is_refused(tmp_path, immich):
config, sf, lib = _env(tmp_path, immich)
with pytest.raises(VerificationError) as error:
UploadVerificationService(sf, config=config).verify("does-not-exist")
assert error.value.code == "not_found"
# ── API surface and durability ───────────────────────────────────────────────
def test_the_api_verifies_resolves_and_lists_history_without_secrets(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=False)
with TestClient(create_app(config)) as client:
verified = client.post(f"/api/v1/upload-batches/{batch['id']}/verify")
resolved = client.post(
f"/api/v1/upload-batches/{batch['id']}/resolve",
json={
"asset_id": batch["items"][0]["asset_id"],
"outcome": "skipped",
"evidence": "the file was withdrawn from the album",
"actor": "dom",
},
)
refused = client.post(
f"/api/v1/upload-batches/{batch['id']}/resolve",
json={
"asset_id": batch["items"][0]["asset_id"],
"outcome": "uploaded",
"evidence": "",
"actor": "dom",
},
)
history = client.get(f"/api/v1/upload-batches/{batch['id']}/verifications")
missing = client.post("/api/v1/upload-batches/does-not-exist/verify")
assert verified.status_code == 200 and verified.json()["counts"] == {ABSENT: 2}
assert resolved.status_code == 200
assert refused.status_code == 422 and refused.json()["error"]["code"] == "evidence_required"
assert len(history.json()["verifications"]) == 3
assert missing.status_code == 404
assert SENTINEL_KEY not in verified.text + resolved.text + history.text
def test_verification_survives_a_restart(tmp_path, immich):
config, sf, lib, batch = _uncertain_batch(tmp_path, immich, hold=True)
UploadVerificationService(sf, config=config).verify(batch["id"])
with TestClient(create_app(config)) as client: # a fresh application process
fetched = client.get(f"/api/v1/upload-batches/{batch['id']}").json()
history = client.get(f"/api/v1/upload-batches/{batch['id']}/verifications").json()
assert fetched["state"] == BatchState.SUCCEEDED
assert fetched["verified_at"]
assert {item["verification"] for item in fetched["items"]} == {PRESENT}
assert len(history["verifications"]) == 2
assert json.dumps(fetched) # the record stays JSON-serialisable for the UI

View File

@@ -90,6 +90,26 @@ def test_cooperative_cancellation_leaves_items_resumable(sf, jobs):
assert by_state.get(ItemState.QUEUED) == 1 # "b" left resumable
def test_a_handler_that_stops_itself_releases_the_lock(sf, jobs):
"""A handler may stop without anyone cancelling the *job* — an upload batch
cancelled through its own API does exactly that. The job is still ``running``
when it raises, so it has to reach ``cancelled`` through ``cancelling``; if that
hop is skipped the transition is rejected and the lock is held forever."""
from photo_pipeline.jobs.handlers import Cancelled
def handler(item, ctx):
raise Cancelled("the work this job wraps was stopped elsewhere")
worker = Worker(sf, {"scan": handler}, "w1")
job = jobs.enqueue("scan", lock="library_write", items=["a"])
worker.run_once()
assert jobs.get(job["id"])["state"] == JobState.CANCELLED
assert jobs.progress(job["id"])["by_state"] == {ItemState.QUEUED: 1} # resumable
# The lane is free: the next job may be enqueued under the same lock.
assert jobs.enqueue("scan", lock="library_write", items=["b"])["state"] == JobState.QUEUED
def test_fencing_rejects_superseded_worker(sf, jobs):
job = jobs.enqueue("scan", items=["a"])
stale = jobs.claim(["scan"], "old")

View File

@@ -113,6 +113,16 @@
"US05-03": [
"tests/unit/test_immich_go_report.py",
"tests/integration/test_upload_reports.py"
],
"US05-04": [
"tests/unit/test_immich_bulk_check.py",
"tests/integration/test_upload_verification.py"
],
"US05-05": [
"tests/e2e/test_uploads_ui.py"
],
"US05-06": [
"tests/e2e/test_phase_e_pipeline.py"
]
}
}

View File

@@ -0,0 +1,27 @@
"""How one bulk-upload-check result is read (US05-04).
Only a duplicate rejection proves Immich holds the bytes. Every other answer — a
rejection for another reason, an action this adapter does not know — must stay
uncertain, because "the server did not say yes" is not "the file is not there".
"""
from photo_pipeline.integrations.immich_go import _holds_bytes, bulk_upload_check
def test_a_duplicate_rejection_is_the_only_proof_of_possession():
assert _holds_bytes({"action": "reject", "reason": "duplicate"}) is True
assert _holds_bytes({"action": "accept"}) is False
def test_any_other_answer_is_uncertain_never_absent():
assert _holds_bytes({"action": "reject", "reason": "unsupported-format"}) is None
assert _holds_bytes({"action": "quarantine"}) is None
assert _holds_bytes({}) is None
def test_missing_credentials_are_reported_not_silently_treated_as_absence():
result = bulk_upload_check("", None, {"asset": "abc"})
assert result["reachable"] is False
assert result["present"] == {}
assert "credentials" in result["detail"]