Compare commits
3 Commits
us/US05-04
...
us/US05-06
| Author | SHA1 | Date | |
|---|---|---|---|
| fe027eada9 | |||
| ac9943884d | |||
| 3c440f3d43 |
33
README.md
33
README.md
@@ -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 A–C 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 A–D remain green in the full run above.
|
||||
|
||||
@@ -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>
|
||||
|
||||
@@ -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),
|
||||
};
|
||||
|
||||
@@ -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
704
frontend/js/uploads.js
Normal 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."
|
||||
);
|
||||
}
|
||||
70
migrations/versions/0010_upload_verification.py
Normal file
70
migrations/versions/0010_upload_verification.py
Normal 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)
|
||||
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
)
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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()
|
||||
)
|
||||
|
||||
@@ -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
|
||||
|
||||
349
photo_pipeline/services/upload_verification.py
Normal file
349
photo_pipeline/services/upload_verification.py
Normal 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
|
||||
],
|
||||
}
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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()
|
||||
|
||||
568
tests/e2e/test_phase_e_pipeline.py
Normal file
568
tests/e2e/test_phase_e_pipeline.py
Normal 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"])
|
||||
@@ -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"
|
||||
|
||||
283
tests/e2e/test_uploads_ui.py
Normal file
283
tests/e2e/test_uploads_ui.py
Normal 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")
|
||||
@@ -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"
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
537
tests/integration/test_upload_verification.py
Normal file
537
tests/integration/test_upload_verification.py
Normal 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
|
||||
@@ -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")
|
||||
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
27
tests/unit/test_immich_bulk_check.py
Normal file
27
tests/unit/test_immich_bulk_check.py
Normal 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"]
|
||||
Reference in New Issue
Block a user