Compare commits

...

2 Commits

Author SHA1 Message Date
fe027eada9 US05-06: Automate Phase E End-to-End Acceptance 2026-08-16 17:04:00 +02:00
ac9943884d US05-05: Operate Uploads in the Browser (#75) 2026-08-16 16:08:32 +02:00
18 changed files with 1872 additions and 4 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 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. 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="#/analyze" data-nav="analyze">Analyze</a>
<a href="#/albums" data-nav="albums">Albums</a> <a href="#/albums" data-nav="albums">Albums</a>
<a href="#/renames" data-nav="renames">Renames</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> <a href="#/stats" data-nav="stats">Stats</a>
</nav> </nav>
</header> </header>

View File

@@ -114,4 +114,28 @@ export const api = {
request(`/rename-plans/${encodeURIComponent(id)}/rollback`, { method: "POST", ...opts }), request(`/rename-plans/${encodeURIComponent(id)}/rollback`, { method: "POST", ...opts }),
renameRecovery: (opts = {}) => request("/rename-recovery", opts), renameRecovery: (opts = {}) => request("/rename-recovery", opts),
resolveRecovery: (opts = {}) => request("/rename-recovery/resolve", { method: "POST", ...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 { api } from "./api.js";
import { navigate, onRouteChange, parseHash } from "./router.js"; import { navigate, onRouteChange, parseHash } from "./router.js";
import { renderRenames, setRenamesRender } from "./renames.js"; import { renderRenames, setRenamesRender } from "./renames.js";
import { renderUploads, setUploadsRender } from "./uploads.js";
import { import {
renderAlbums, renderAlbums,
renderAnalyze, renderAnalyze,
@@ -364,6 +365,7 @@ function render() {
else if (path === "/analyze") renderAnalyze(root, params); else if (path === "/analyze") renderAnalyze(root, params);
else if (path === "/albums") renderAlbums(root, params); else if (path === "/albums") renderAlbums(root, params);
else if (path === "/renames") renderRenames(root, params); else if (path === "/renames") renderRenames(root, params);
else if (path === "/uploads") renderUploads(root, params);
else if (path === "/stats") renderStats(root, params); else if (path === "/stats") renderStats(root, params);
else show(errorBanner("Unknown view")); else show(errorBanner("Unknown view"));
} }
@@ -371,5 +373,6 @@ function render() {
// Let views re-render the current route after a mutation. // Let views re-render the current route after a mutation.
setRender(render); setRender(render);
setRenamesRender(render); setRenamesRender(render);
setUploadsRender(render);
onRouteChange(render); onRouteChange(render);
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

@@ -96,6 +96,16 @@ class Worker:
return return
try: try:
if cancelled or snapshot["state"] == JobState.CANCELLING: 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( self.service.transition(
job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token job_id, JobState.CANCELLED, worker_id=self.worker_id, fencing_token=token
) )

View File

@@ -397,7 +397,7 @@ class UploadBatchService:
def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict: def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict:
return { batch = {
"id": row.id, "id": row.id,
"album": row.album, "album": row.album,
"folder": row.folder, "folder": row.folder,
@@ -449,3 +449,8 @@ def _batch_dict(row: UploadBatch, items: list[UploadItem]) -> dict:
for item in items 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

@@ -31,4 +31,5 @@ markers = [
"phase_b: Phase B end-to-end acceptance (US02-07) — API, worker-recovery, and browser journeys", "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_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_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 from __future__ import annotations
import json
import os import os
import socket import socket
import stat
import subprocess import subprocess
import sys import sys
import threading
import time import time
import uuid import uuid
from dataclasses import dataclass from dataclasses import dataclass
from datetime import datetime, timezone from datetime import datetime, timezone
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path from pathlib import Path
import httpx import httpx
@@ -135,15 +139,27 @@ class Server:
self.proc = None self.proc = None
def start_worker(seeded: Seeded, *, fake_vision_log: Path) -> subprocess.Popen: def start_worker(
"""Launch a real durable worker wired to the recording vision fake.""" 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( return subprocess.Popen(
[sys.executable, "-m", "photo_pipeline", "worker", "--id", "e2e-worker"], [sys.executable, "-m", "photo_pipeline", "worker", "--id", "e2e-worker"],
cwd=str(REPO), cwd=str(REPO),
env=_env( env=_env(
seeded, seeded,
free_port(), # unused by the worker, but keeps the env shape uniform 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, stdout=subprocess.PIPE,
stderr=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"]}, json={**payload, "expected_version": current["version"]},
timeout=10, timeout=10,
).raise_for_status() ).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"] 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_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_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(): 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) 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(): def test_every_mapped_test_file_exists_and_is_nonempty():
for story, files in MAP.items(): for story, files in MAP.items():
assert files, f"{story} maps to no tests" 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 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) NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SENTINEL_KEY = "immich-sentinel-9f3a2b" SENTINEL_KEY = "immich-sentinel-9f3a2b"
UPLOADER_VERSION = "immich-go 0.21.0" 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.hashing import sha256_file
from photo_pipeline.services.uploads import UploadError, UploadService 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) NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
# A sentinel credential: every assertion below proves it never leaves configuration. # A sentinel credential: every assertion below proves it never leaves configuration.
SENTINEL_KEY = "immich-sentinel-9f3a2b" 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 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) NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SUPPORTED_VERSION = "immich-go 0.21.0" SUPPORTED_VERSION = "immich-go 0.21.0"

View File

@@ -40,6 +40,8 @@ from photo_pipeline.services.upload_verification import (
) )
from photo_pipeline.services.uploads import UploadService 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) NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
SENTINEL_KEY = "immich-sentinel-9f3a2b" SENTINEL_KEY = "immich-sentinel-9f3a2b"
UPLOADER_VERSION = "immich-go 0.21.0" UPLOADER_VERSION = "immich-go 0.21.0"

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 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): def test_fencing_rejects_superseded_worker(sf, jobs):
job = jobs.enqueue("scan", items=["a"]) job = jobs.enqueue("scan", items=["a"])
stale = jobs.claim(["scan"], "old") stale = jobs.claim(["scan"], "old")

View File

@@ -117,6 +117,12 @@
"US05-04": [ "US05-04": [
"tests/unit/test_immich_bulk_check.py", "tests/unit/test_immich_bulk_check.py",
"tests/integration/test_upload_verification.py" "tests/integration/test_upload_verification.py"
],
"US05-05": [
"tests/e2e/test_uploads_ui.py"
],
"US05-06": [
"tests/e2e/test_phase_e_pipeline.py"
] ]
} }
} }