Compare commits

..

3 Commits

16 changed files with 2609 additions and 43 deletions

View File

@@ -19,6 +19,7 @@
<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="#/archive" data-nav="archive">Archive</a>
<a href="#/stats" data-nav="stats">Stats</a>
</nav>
</header>

View File

@@ -138,4 +138,31 @@ export const api = {
}),
uploadVerifications: (id, opts = {}) =>
request(`/upload-batches/${encodeURIComponent(id)}/verifications`, opts),
// ── Archive and restore: destinations, preflight, plans, recovery ────────
archiveLocations: (opts = {}) => request("/archive-locations", opts),
registerArchiveLocation: (payload, opts = {}) =>
request("/archive-locations", { method: "POST", body: JSON.stringify(payload), ...opts }),
archivePreflight: (payload, opts = {}) =>
request("/archive-preflight", { method: "POST", body: JSON.stringify(payload), ...opts }),
createArchivePlan: (payload, opts = {}) =>
request("/archive-plans", { method: "POST", body: JSON.stringify(payload), ...opts }),
listArchivePlans: (opts = {}) => request("/archive-plans", opts),
getArchivePlan: (id, opts = {}) => request(`/archive-plans/${encodeURIComponent(id)}`, opts),
applyArchivePlan: (id, opts = {}) =>
request(`/archive-plans/${encodeURIComponent(id)}/apply`, { method: "POST", ...opts }),
archiveRecovery: (opts = {}) => request("/archive-recovery", opts),
resolveArchiveRecovery: (opts = {}) =>
request("/archive-recovery/resolve", { method: "POST", ...opts }),
restorePreflight: (payload, opts = {}) =>
request("/restore-preflight", { method: "POST", body: JSON.stringify(payload), ...opts }),
createRestorePlan: (payload, opts = {}) =>
request("/restore-plans", { method: "POST", body: JSON.stringify(payload), ...opts }),
listRestorePlans: (opts = {}) => request("/restore-plans", opts),
getRestorePlan: (id, opts = {}) => request(`/restore-plans/${encodeURIComponent(id)}`, opts),
applyRestorePlan: (id, opts = {}) =>
request(`/restore-plans/${encodeURIComponent(id)}/apply`, { method: "POST", ...opts }),
restoreRecovery: (opts = {}) => request("/restore-recovery", opts),
resolveRestoreRecovery: (opts = {}) =>
request("/restore-recovery/resolve", { method: "POST", ...opts }),
};

View File

@@ -1,4 +1,5 @@
import { api } from "./api.js";
import { renderArchive, setArchiveRender } from "./archive.js";
import { navigate, onRouteChange, parseHash } from "./router.js";
import { renderRenames, setRenamesRender } from "./renames.js";
import { renderUploads, setUploadsRender } from "./uploads.js";
@@ -366,6 +367,7 @@ function render() {
else if (path === "/albums") renderAlbums(root, params);
else if (path === "/renames") renderRenames(root, params);
else if (path === "/uploads") renderUploads(root, params);
else if (path === "/archive") renderArchive(root, params);
else if (path === "/stats") renderStats(root, params);
else show(errorBanner("Unknown view"));
}
@@ -374,5 +376,6 @@ function render() {
setRender(render);
setRenamesRender(render);
setUploadsRender(render);
setArchiveRender(render);
onRouteChange(render);
render();

870
frontend/js/archive.js Normal file
View File

@@ -0,0 +1,870 @@
// Archive view (US06-05): preview what would leave active storage, confirm it
// exactly, watch the transfer, recover an interrupted one, browse what is already
// archived, and bring it back.
//
// Archiving is the only stage that removes originals, so this view never decides
// anything itself: the destination's identity, every blocker, the confirmation
// token, and what recovery may do all come from the server, and an action the
// server would refuse is not offered. Two things follow from that. A medium that
// is not mounted produces an instruction naming it rather than a disabled mystery,
// and an operation whose evidence is ambiguous offers no button at all.
import { api } from "./api.js";
import { el, errorBanner, setActiveNav } from "./dom.js";
import { subscribeJob } from "./events.js";
import { navigate } from "./router.js";
let outcome = null;
let activity = [];
let render = () => {};
export function setArchiveRender(fn) {
render = fn;
}
// The per-file journal states, split into the three things an operator actually
// wants told apart: bytes moving, bytes proven, original removed (concept §9).
const PHASES = [
["planned", "planned", "waiting"],
["transferring", "transfer", "copying to the medium"],
["verified", "verified", "archive copy hashed and manifested"],
["removing", "removing", "removing the active original"],
["complete", "complete", "archived and removed"],
["failed", "failed", "left alone for a decision"],
];
const AVAILABILITY_LABEL = {
active: "in the library",
archived_online: "archived · medium mounted",
archived_offline: "archived · medium away",
missing_unexpected: "missing — unexplained",
};
export async function renderArchive(root, params = {}) {
setActiveNav("archive");
let locations;
try {
locations = (await api.archiveLocations()).locations;
} catch (error) {
root.replaceChildren(errorBanner(`Failed to load archive locations: ${error.message}`));
return;
}
const location = locations.find((l) => l.id === params.location) || locations[0] || null;
const nodes = [el("h1", {}, "Archive"), locationsCard(locations, location, params)];
if (!location) {
nodes.push(
el(
"p",
{ class: "muted", "data-testid": "no-locations" },
"Register the disk, NAS share, or removable medium that will hold archived originals."
),
outcomeBanner()
);
root.replaceChildren(...nodes.filter(Boolean));
return;
}
const [preflight, archivePlans, restorePlans, recovery, restoreRecovery, archived, restore] =
await Promise.all([
load(() => api.archivePreflight({ location_id: location.id })),
load(() => api.listArchivePlans()),
load(() => api.listRestorePlans()),
load(() => api.archiveRecovery()),
load(() => api.restoreRecovery()),
load(() => api.listAssets({ limit: 200 })),
load(() => api.restorePreflight({ location_id: location.id })),
]);
// Both directions share the runs table: one lane moves these originals, so one
// history is what an operator has to reason about.
const runs = [...plans(archivePlans), ...plans(restorePlans)].sort((a, b) =>
(a.created_at || "").localeCompare(b.created_at || "")
);
const chosen = runs.find((run) => run.id === params.plan) || runs[runs.length - 1] || null;
const selectedPlan = chosen ? await load(() => planDetailOf(chosen)) : null;
nodes.push(
preflight ? previewSection(preflight, location) : null,
preflight ? confirmBlock(preflight, location) : null,
outcomeBanner(),
activityLog(),
recoverySection(mergeRecovery(recovery, restoreRecovery)),
planList(runs, chosen && chosen.id),
selectedPlan ? planDetail(selectedPlan) : null,
archived ? archivedSection(archived.items, locations) : null,
restore ? restoreSection(restore, location) : null
);
root.replaceChildren(...nodes.filter(Boolean));
}
function plans(listed) {
return listed ? listed.plans : [];
}
function planDetailOf(run) {
return run.direction === "restore" ? api.getRestorePlan(run.id) : api.getArchivePlan(run.id);
}
// Archive and restore recovery answer the same question about the same lane, so
// they are one list; an unresolved item of either kind blocks the other.
function mergeRecovery(archive, restore) {
if (!archive && !restore) return null;
return {
operations: [...((archive || {}).operations || []), ...((restore || {}).operations || [])],
manual: [...((archive || {}).manual || []), ...((restore || {}).manual || [])],
};
}
// A section whose data failed to load must not take the rest of the view with it:
// the medium being away is exactly when the archived-asset list matters most.
async function load(call) {
try {
return await call();
} catch (_) {
return null;
}
}
// ── locations ────────────────────────────────────────────────────────────────
// A location is a medium, not a path: the marker's ``media_id`` is what proves the
// right disk is mounted, so it is shown next to the state it produced.
function locationsCard(locations, selected, params) {
const name = el("input", {
type: "text",
"data-testid": "location-name",
"aria-label": "Archive location name",
placeholder: "External disk",
});
const root = el("input", {
type: "text",
"data-testid": "location-root",
"aria-label": "Archive location path",
placeholder: "/Volumes/archive",
});
return el(
"div",
{ class: "card", "data-testid": "archive-locations" },
el("h2", {}, "Destinations"),
locations.length
? el(
"table",
{ class: "grid", "data-testid": "locations" },
el(
"thead",
{},
el(
"tr",
{},
...["", "Name", "Root", "Medium", "State", "Last seen"].map((label) =>
el("th", { scope: "col" }, label)
)
)
),
el(
"tbody",
{},
...locations.map((location) =>
el(
"tr",
{
"data-testid": "location-row",
"data-name": location.name,
"aria-current": selected && location.id === selected.id ? "true" : false,
},
el(
"td",
{},
el("input", {
type: "radio",
name: "archive-location",
"data-testid": "select-location",
"aria-label": `Use ${location.name}`,
checked: selected && location.id === selected.id ? "checked" : false,
onchange: () => navigate("/archive", { ...params, location: location.id }),
})
),
el("td", { "data-testid": "location-label" }, location.name),
el("td", { class: "path", "data-testid": "location-path" }, location.root),
el("td", { class: "path", "data-testid": "location-media" }, location.media_id),
el(
"td",
{},
el(
"span",
{
class: `badge ${location.state === "online" ? "complete" : "attention"}`,
"data-testid": "location-state",
},
location.state
)
),
el("td", { class: "muted" }, location.last_seen_at || "never")
)
)
)
)
: null,
selected && selected.state !== "online" ? mountInstruction(selected) : null,
el(
"div",
{ class: "toolbar" },
name,
root,
el(
"button",
{
"data-testid": "register-location",
onclick: () =>
run(() => api.registerArchiveLocation({ name: name.value, root: root.value })),
},
"Register destination"
)
)
);
}
// The one thing the app cannot do for the user: name the medium to connect.
function mountInstruction(location) {
const detail =
location.state === "wrong_volume"
? `A different medium is mounted at ${location.root}.`
: `Nothing is mounted at ${location.root}.`;
return el(
"div",
{ class: "confirm", role: "status", "data-testid": "mount-instruction" },
`${detail} Connect “${location.name}” (medium ${location.media_id}) and mount it there, ` +
"then reload this view. Archived photos stay listed and searchable meanwhile."
);
}
// ── preview ──────────────────────────────────────────────────────────────────
function previewSection(preflight, location) {
const totals = preflight.totals;
const capacity = preflight.capacity;
const rows = preflight.albums.map((album) =>
el(
"tr",
{ "data-testid": "archive-album-row", "data-album": album.album },
el("td", { "data-testid": "album-name" }, album.album),
el("td", { class: "path", "data-testid": "album-folder" }, album.folder),
el("td", { class: "path", "data-testid": "album-destination" }, album.destination),
el(
"td",
{ "data-testid": "album-method" },
album.transfer_method === "move" ? "move (same filesystem)" : "copy · verify · remove"
),
el("td", { "data-testid": "album-assets" }, String(album.asset_count)),
el("td", { "data-testid": "album-reclaim" }, bytes(album.reclaimable_bytes)),
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
)
)
)
)
);
return el(
"div",
{ "data-testid": "archive-preview" },
el("h2", {}, "Preview"),
el(
"div",
{ class: "decision-bar" },
el(
"span",
{ class: "badge", "data-testid": "destination-identity" },
`${location.name} · ${location.media_id}`
),
el("span", { class: "badge", "data-testid": "total-albums" }, `${totals.albums} album(s)`),
el("span", { class: "badge", "data-testid": "total-assets" }, `${totals.assets} photo(s)`),
el(
"span",
{ class: "badge", "data-testid": "total-reclaim" },
`${bytes(totals.bytes)} reclaimable`
),
el(
"span",
{
class: `badge ${capacity.sufficient ? "complete" : "blocked"}`,
"data-testid": "capacity",
},
`free ${bytes(capacity.free_bytes)} · reserve ${bytes(capacity.reserve_bytes)}`
)
),
blockerList(preflight.blockers, "preflight-blockers", "preflight-blocker", "This scope cannot be archived yet"),
rows.length
? el(
"table",
{ class: "grid", "data-testid": "archive-albums" },
el(
"thead",
{},
el(
"tr",
{},
...["Album", "Folder", "Destination", "Transfer", "Photos", "Reclaims", "State"].map(
(label) => el("th", { scope: "col" }, label)
)
)
),
el("tbody", {}, ...rows)
)
: el(
"p",
{ class: "muted", "data-testid": "no-albums" },
"No album has a verified upload whose bytes are still unchanged, so nothing may be archived."
)
);
}
function blockerList(blockers, containerId, itemId, title) {
if (!blockers || !blockers.length) return null;
return el(
"div",
{ class: "alert", role: "alert", "data-testid": containerId },
el("strong", {}, title),
el(
"ul",
{},
...blockers.map((blocker) =>
el("li", { "data-testid": itemId, "data-code": blocker.code }, `${blocker.code}: ${blocker.message}`)
)
)
);
}
// ── confirmation ─────────────────────────────────────────────────────────────
function confirmBlock(preflight, location) {
const ready = preflight.state === "ready";
const totals = preflight.totals;
return el(
"div",
{ class: "card", "data-testid": "confirm" },
el("h2", {}, "Confirm"),
el(
"p",
{ class: "muted", "data-testid": "confirm-token" },
`Preflight ${preflight.token.slice(0, 20)}… · ${location.name}`
),
el(
"p",
{ "data-testid": "archive-note" },
"Archiving removes each original from the library — but only after its copy on " +
"the medium has been written, hashed, and recorded in the manifest. The photos " +
"stay searchable and deduplicable while the medium is away, and can be restored " +
"from this view."
),
el(
"div",
{ class: "toolbar" },
el(
"button",
{
class: "primary",
"data-testid": "start-archive",
disabled: ready ? false : "disabled",
title: ready ? false : "resolve the blockers above first",
onclick: () =>
run(async () => {
const plan = await api.createArchivePlan({
location_id: location.id,
token: preflight.token,
});
const started = await api.applyArchivePlan(plan.id);
watch(started.job.id, plan.id, "archive");
return { archiving: plan.asset_count };
}),
},
`Archive ${totals.ready_albums} album(s) · reclaim ${bytes(totals.bytes)}`
)
)
);
}
// ── plans and progress ───────────────────────────────────────────────────────
function planList(runs, selectedId) {
if (!runs.length) {
return el("p", { class: "muted", "data-testid": "no-plans" }, "Nothing has been archived yet.");
}
return el(
"div",
{ "data-testid": "archive-plans" },
el("h2", {}, "Runs"),
el(
"table",
{ class: "grid", "data-testid": "plans" },
el(
"thead",
{},
el(
"tr",
{},
...["Created", "Direction", "State", "Photos", "Bytes"].map((l) =>
el("th", { scope: "col" }, l)
)
)
),
el(
"tbody",
{},
...runs.map((plan) =>
el(
"tr",
{
"data-testid": "plan-row",
"data-plan": plan.id,
"data-direction": plan.direction,
"aria-current": plan.id === selectedId ? "true" : false,
},
el(
"td",
{},
el("a", { class: "link", href: `#/archive?plan=${encodeURIComponent(plan.id)}` }, plan.created_at || plan.id)
),
el("td", { "data-testid": "plan-direction" }, plan.direction),
el("td", {}, el("span", { class: `badge ${plan.state}`, "data-testid": "plan-state" }, plan.state)),
el("td", {}, String(plan.asset_count)),
el("td", {}, bytes(plan.byte_size))
)
)
)
)
);
}
function planDetail(plan) {
const operations = plan.operations || [];
const counts = {};
for (const operation of operations) {
counts[operation.journal_state] = (counts[operation.journal_state] || 0) + 1;
}
return el(
"div",
{ class: "card", "data-testid": "plan-detail", "data-plan": plan.id },
el("h2", {}, `${plan.direction === "restore" ? "Restore" : "Archive"} run ${plan.created_at || plan.id}`),
// Transfer, verification, and removal are separate answers to separate
// questions: what has moved, what is proven, and what is already gone.
el(
"div",
{ class: "decision-bar", "data-testid": "plan-progress" },
el("span", { class: `badge ${plan.state}`, "data-testid": "detail-state" }, plan.state),
...PHASES.map(([key, label, title]) =>
el(
"span",
{ class: `badge ${key}`, "data-testid": `count-${label}`, title },
`${label}: ${counts[key] ?? 0}`
)
)
),
el(
"table",
{ class: "grid", "data-testid": "operations" },
el(
"thead",
{},
el(
"tr",
{},
...["Photo", "Destination", "Phase", "Attempts", "Problem"].map((l) =>
el("th", { scope: "col" }, l)
)
)
),
el(
"tbody",
{},
...operations.map((operation) =>
el(
"tr",
{ "data-testid": "operation-row", "data-asset-id": operation.asset_id },
el("td", { class: "path", "data-testid": "operation-source" }, operation.source_path),
el("td", { class: "path", "data-testid": "operation-destination" }, operation.destination_path),
el(
"td",
{},
el(
"span",
{ class: `badge ${operation.journal_state}`, "data-testid": "operation-phase" },
phaseLabel(operation.journal_state)
)
),
el("td", {}, String(operation.attempt_count)),
el(
"td",
{ class: "muted", "data-testid": "operation-error", "data-code": operation.error_code || "" },
operation.error_code ? `${operation.error_code}: ${operation.error_message || ""}` : "—"
)
)
)
)
)
);
}
function phaseLabel(state) {
const found = PHASES.find(([key]) => key === state);
return found ? found[1] : state;
}
// ── recovery ─────────────────────────────────────────────────────────────────
// What an interrupted run left behind, straight from the journal plus the files on
// disk. Only the operations the server itself classified as resolvable get an
// action; ambiguous ones are shown with their evidence and no button.
function recoverySection(recovery) {
const operations = recovery ? recovery.operations : [];
if (!operations.length) return null;
const manual = recovery.manual || [];
const resolvable = operations.length - manual.length;
return el(
"div",
{ class: "card", "data-testid": "archive-recovery" },
el("h2", {}, "Interrupted work"),
el(
"p",
{ "data-testid": "recovery-summary" },
`${operations.length} operation(s) did not finish · ${resolvable} resolvable · ` +
`${manual.length} need a decision`
),
el(
"table",
{ class: "grid", "data-testid": "recovery-operations" },
el(
"thead",
{},
el(
"tr",
{},
...["Photo", "Phase", "Verdict", "Evidence"].map((l) => el("th", { scope: "col" }, l))
)
),
el(
"tbody",
{},
...operations.map((verdict) =>
el(
"tr",
{
"data-testid": "recovery-row",
"data-classification": verdict.classification,
"data-direction": verdict.direction,
},
el("td", { class: "path", "data-testid": "recovery-source" }, verdict.source_path),
el("td", {}, el("span", { class: "badge" }, phaseLabel(verdict.journal_state))),
el(
"td",
{},
el(
"span",
{
class: `badge ${verdict.classification === "manual" ? "blocked" : "attention"}`,
"data-testid": "recovery-verdict",
},
verdict.classification
)
),
el(
"td",
{ class: "muted", "data-testid": "recovery-reason" },
`${verdict.reason} (source ${verdict.source_exists ? "present" : "absent"}, ` +
`archive copy ${verdict.destination_matches ? "verified" : verdict.destination_exists ? "different bytes" : "absent"})`
)
)
)
)
),
manual.length
? el(
"div",
{ class: "alert", role: "alert", "data-testid": "recovery-manual" },
`${manual.length} operation(s) cannot be resolved from the evidence. Nothing will be ` +
"removed or retried for them here: inspect the medium and the library, then decide."
)
: null,
el(
"div",
{ class: "toolbar" },
resolvable
? el(
"button",
{
class: "primary",
"data-testid": "resolve-recovery",
onclick: () => run(() => api.resolveArchiveRecovery()),
},
`Finish ${resolvable} recoverable operation(s)`
)
: el(
"span",
{ class: "muted", "data-testid": "no-safe-recovery" },
"No operation can be finished safely from here."
)
)
);
}
// ── archived assets ──────────────────────────────────────────────────────────
// Browsing what is already archived, including while the medium is away: the
// retained protected preview and the recorded hashes are the evidence, so the row
// stays complete and honest instead of turning into a missing file.
function archivedSection(assets, locations) {
const archived = assets.filter((asset) => asset.availability_state !== "active");
if (!archived.length) return null;
const names = Object.fromEntries(locations.map((location) => [location.id, location.name]));
return el(
"div",
{ "data-testid": "archived-assets" },
el("h2", {}, "Archived photos"),
el(
"table",
{ class: "grid", "data-testid": "archived" },
el(
"thead",
{},
el(
"tr",
{},
...["Preview", "Archived as", "Medium", "Availability", "Size"].map((l) =>
el("th", { scope: "col" }, l)
)
)
),
el(
"tbody",
{},
...archived.map((asset) =>
el(
"tr",
{ "data-testid": "archived-row", "data-asset-id": asset.id },
el(
"td",
{},
el("img", {
"data-testid": "archived-preview",
width: 96,
loading: "lazy",
src: api.thumbnailUrl(asset.id, 256),
alt: `Preview of ${asset.archive_path || asset.id}`,
})
),
el("td", { class: "path", "data-testid": "archived-path" }, asset.archive_path || "—"),
el(
"td",
{ "data-testid": "archived-medium" },
names[asset.archive_location_id] || asset.archive_location_id || "—"
),
el(
"td",
{},
el(
"span",
{
class: `badge ${asset.availability_state === "archived_online" ? "complete" : "attention"}`,
"data-testid": "archived-availability",
},
AVAILABILITY_LABEL[asset.availability_state] || asset.availability_state
)
),
el("td", {}, bytes(asset.byte_size))
)
)
)
)
);
}
// ── restore ──────────────────────────────────────────────────────────────────
function restoreSection(restore, location) {
const ready = restore.state === "ready";
const items = restore.items || [];
if (!items.length && !restore.blockers.length) return null;
return el(
"div",
{ class: "card", "data-testid": "restore" },
el("h2", {}, "Restore"),
el(
"p",
{ "data-testid": "restore-note" },
"Restoring copies the archived bytes back into the library and leaves the archive " +
"copy where it is. A name that is already taken is never overwritten: the photo " +
"comes back beside it under a visibly different name."
),
blockerList(restore.blockers, "restore-blockers", "restore-blocker", "This restore cannot run yet"),
items.length
? el(
"table",
{ class: "grid", "data-testid": "restore-items" },
el(
"thead",
{},
el(
"tr",
{},
...["Archived as", "Comes back as", "Size", "State"].map((l) =>
el("th", { scope: "col" }, l)
)
)
),
el(
"tbody",
{},
...items.map((item) =>
el(
"tr",
{ "data-testid": "restore-row", "data-asset-id": item.asset_id },
el("td", { class: "path", "data-testid": "restore-source" }, item.archive_path),
el("td", { class: "path", "data-testid": "restore-destination" }, item.destination_path || "—"),
el("td", {}, bytes(item.byte_size)),
el(
"td",
{},
item.blockers.length
? el(
"span",
{
class: "badge blocked",
"data-testid": "restore-item-blocker",
"data-code": item.blockers[0].code,
},
item.blockers[0].code
)
: el("span", { class: "badge ready", "data-testid": "restore-item-state" }, "ready")
)
)
)
)
)
: null,
el(
"div",
{ class: "toolbar" },
el(
"button",
{
class: "primary",
"data-testid": "start-restore",
disabled: ready ? false : "disabled",
title: ready ? false : "the medium and every archived copy must check out first",
onclick: () =>
run(async () => {
const plan = await api.createRestorePlan({
location_id: location.id,
token: restore.token,
});
const started = await api.applyRestorePlan(plan.id);
watch(started.job.id, plan.id, "restore");
return { restoring: plan.asset_count };
}),
},
`Restore ${items.length} photo(s) from ${location.name}`
)
)
);
}
// ── 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. The plan panel is refreshed on its own tick because the
// journal advances per file, not per job event; a full re-render would re-run
// preflight (which re-hashes the library), so that happens once when the job ends.
const REFRESH_MS = 500;
function watch(jobId, planId, kind) {
activity = [`Started ${kind} job ${jobId}`];
const tick = setInterval(() => refreshPlan(planId, kind), REFRESH_MS);
subscribeJob(jobId, {
onEvent: (event) => {
activity.push(`${event.type}${event.message ? ": " + event.message : ""}`);
const log = document.querySelector('[data-testid="archive-activity"]');
if (log) log.textContent = activity.join("\n");
},
onDone: () => {
clearInterval(tick);
activity.push("done");
render();
},
});
}
async function refreshPlan(planId, kind) {
const node = document.querySelector(`[data-testid="plan-detail"][data-plan="${planId}"]`);
if (!node) return; // the user navigated away from the running plan
try {
const plan = await (kind === "restore" ? api.getRestorePlan(planId) : api.getArchivePlan(planId));
node.replaceWith(planDetail(plan));
} 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": "archive-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 moved or removed; ` +
"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": "archive-error" },
`Failed: ${outcome.error.message}`
);
}
const result = outcome.result || {};
const message =
result.archiving !== undefined
? `Archiving ${result.archiving} photo(s). Originals are removed only after their copies verify.`
: result.restoring !== undefined
? `Restoring ${result.restoring} photo(s) into the library.`
: "Done — the state below is the server's.";
return el("div", { class: "alert", role: "status", "data-testid": "archive-result" }, message);
}
// ── formatting ───────────────────────────────────────────────────────────────
function bytes(value) {
if (value == null) return "unknown";
const units = ["B", "kB", "MB", "GB", "TB"];
let size = value;
let unit = 0;
while (size >= 1000 && unit < units.length - 1) {
size /= 1000;
unit += 1;
}
return `${unit === 0 ? size : size.toFixed(1)} ${units[unit]}`;
}

View File

@@ -0,0 +1,37 @@
"""Restore plans and archive divergence (US06-04).
Revision ID: 0014_restore_plans
Revises: 0013_protected_thumbnails
Create Date: 2026-08-16
Restore reuses the archive plan and journal tables: the crash-safe question is the
same one in the opposite direction (copy, verify, publish, register), so the rows
gain a ``direction`` instead of a parallel pair of tables. ``archive_divergent_at``
records the moment an archived copy was proven to hold bytes that are not the ones
the database recorded — a restore must never silently accept a different file.
"""
import sqlalchemy as sa
from alembic import op
revision = "0014_restore_plans"
down_revision = "0013_protected_thumbnails"
branch_labels = None
depends_on = None
def upgrade() -> None:
for table in ("archive_plans", "archive_operations"):
op.add_column(
table,
sa.Column("direction", sa.String(), nullable=False, server_default="archive"),
)
op.add_column(
"assets", sa.Column("archive_divergent_at", sa.DateTime(timezone=True), nullable=True)
)
def downgrade() -> None:
op.drop_column("assets", "archive_divergent_at")
for table in ("archive_plans", "archive_operations"):
op.drop_column(table, "direction")

View File

@@ -1,4 +1,4 @@
"""Archive location, preflight, and plan API (US06-01, US06-02).
"""Archive location, preflight, plan, and restore API (US06-01, US06-02, US06-04).
Registering a location writes a marker onto the medium; preflight is a command
rather than a read, because it probes the destination, hashes the scope, and issues
@@ -13,10 +13,11 @@ from fastapi import APIRouter, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, ARCHIVE_PLAN
from photo_pipeline.jobs.domain_handlers import ARCHIVE_LOCK, ARCHIVE_PLAN, RESTORE_PLAN
from photo_pipeline.services.archives import ArchiveError, ArchiveService
from photo_pipeline.services.archive_transfer import ArchiveTransferService
from photo_pipeline.services.jobs import JobBlocked, JobService
from photo_pipeline.services.restores import RestoreService
router = APIRouter(tags=["archives"])
@@ -42,10 +43,24 @@ class CreatePlanRequest(PreflightRequest):
token: str
class RestoreRequest(BaseModel):
location_id: str
# ``None`` means every asset archived at this location.
asset_ids: list[str] | None = None
class CreateRestoreRequest(RestoreRequest):
token: str
def _service(request: Request) -> ArchiveService:
return ArchiveService(request.app.state.session_factory, config=request.app.state.config)
def _restores(request: Request) -> RestoreService:
return RestoreService(request.app.state.session_factory, config=request.app.state.config)
def _transfers(request: Request) -> ArchiveTransferService:
return ArchiveTransferService(
request.app.state.session_factory, config=request.app.state.config
@@ -138,6 +153,76 @@ def apply_plan(plan_id: str, request: Request):
return {"plan_id": plan_id, "job": job}
@router.post("/restore-preflight")
def restore_preflight(body: RestoreRequest, request: Request):
"""Validate restoring archived assets back into the library. Nothing moves."""
try:
return _restores(request).preflight(body.location_id, body.asset_ids)
except ArchiveError as error:
return _error(error)
@router.post("/restore-plans", status_code=201)
def create_restore_plan(body: CreateRestoreRequest, request: Request):
try:
return _restores(request).create(body.location_id, body.asset_ids, token=body.token)
except ArchiveError as error:
return _error(error)
@router.get("/restore-plans")
def list_restore_plans(request: Request) -> dict:
return {"plans": _restores(request).list()}
@router.get("/restore-plans/{plan_id}")
def get_restore_plan(plan_id: str, request: Request):
plan = _restores(request).get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown restore plan {plan_id}"))
return plan
@router.post("/restore-plans/{plan_id}/apply")
def apply_restore_plan(plan_id: str, request: Request):
"""Queue the restore on the archiver lane — the same single lane as archiving,
because both move the same originals."""
service = _restores(request)
plan = service.get(plan_id)
if plan is None:
return _error(ArchiveError("unknown_plan", f"unknown restore plan {plan_id}"))
unresolved = [row for row in service.journal.incomplete() if row["plan_id"] != plan_id]
if unresolved:
return _error(
ArchiveError(
"archive_pending",
f"an unresolved archive operation ({unresolved[0]['id']}) must be recovered",
)
)
try:
job = JobService(request.app.state.session_factory).enqueue(
RESTORE_PLAN,
lock=ARCHIVE_LOCK,
idempotency_key=f"restore:{plan_id}:{plan['version']}",
items=[plan_id],
)
except JobBlocked as error:
return JSONResponse(
status_code=409, content={"error": {"code": error.code, "message": str(error)}}
)
return {"plan_id": plan_id, "job": job}
@router.get("/restore-recovery")
def restore_recovery_status(request: Request) -> dict:
return _restores(request).recovery_status()
@router.post("/restore-recovery/resolve")
def resolve_restore_recovery(request: Request) -> dict:
return _restores(request).recover()
@router.get("/archive-recovery")
def recovery_status(request: Request) -> dict:
"""What an interrupted transfer left behind, straight from journal + disk."""

View File

@@ -1,9 +1,9 @@
"""Domain job handlers: safety scoring, content analysis, uploads, archive
transfers (US02-06, US05-02, US06-02).
transfers, restores (US02-06, US05-02, US06-02, US06-04).
Importing this module registers the ``safety_score``, ``analysis``,
``upload_batch``, and ``archive_plan`` job types so the generic worker can run them
per item. Each handler delegates to its service, which owns the real work and the
``upload_batch``, ``archive_plan``, and ``restore_plan`` job types so the generic
worker can run them per item. Each handler delegates to its service, which owns the real work and the
privacy gate. Handlers are idempotent: re-scoring or re-analyzing one asset is safe
after an interrupted attempt, an upload batch refuses to re-run an attempt whose
outcome is unknown, and an archive plan skips items it already completed.
@@ -20,6 +20,7 @@ SAFETY_SCORE = "safety_score"
ANALYSIS = "analysis"
UPLOAD_BATCH = "upload_batch"
ARCHIVE_PLAN = "archive_plan"
RESTORE_PLAN = "restore_plan"
# Both mutate the library's metadata/derived state; one at a time (concept §one job).
LIBRARY_WRITE_LOCK = "library_write"
# The uploader lane: one album batch at a time (concept §16).
@@ -70,7 +71,22 @@ def _archive_plan_item(plan_id: str, ctx: JobContext) -> None:
raise RuntimeError(f"archive plan {plan_id}: {result['failed']} item(s) failed")
def _restore_plan_item(plan_id: str, ctx: JobContext) -> None:
"""One item = one restore plan. A restore removes nothing, so an item failure
simply leaves that asset archived (US06-04)."""
from photo_pipeline.config import Config
from photo_pipeline.services.restores import RestoreService
config = ctx.config if ctx.config is not None else Config.from_env()
result = RestoreService(ctx.session_factory, config=config).apply(
plan_id, worker_id=ctx.worker_id
)
if result["failed"]:
raise RuntimeError(f"restore plan {plan_id}: {result['failed']} item(s) failed")
register(SAFETY_SCORE, _safety_score_item)
register(ANALYSIS, _analysis_item)
register(UPLOAD_BATCH, _upload_batch_item)
register(ARCHIVE_PLAN, _archive_plan_item)
register(RESTORE_PLAN, _restore_plan_item)

View File

@@ -66,6 +66,8 @@ class ArchivePlan(Base):
# The preflight token this plan was approved against; re-verified before apply.
token: Mapped[str] = mapped_column(String, nullable=False)
albums: Mapped[str | None] = mapped_column(String) # JSON array
# archive | restore — the same journal read in the opposite direction (US06-04).
direction: Mapped[str] = mapped_column(String, nullable=False, default="archive")
# planned | applying | complete | failed
state: Mapped[str] = mapped_column(String, nullable=False, default="planned")
@@ -100,6 +102,9 @@ class ArchiveOperation(Base):
album: Mapped[str] = mapped_column(String, nullable=False)
asset_id: Mapped[str] = mapped_column(ForeignKey("assets.id"), nullable=False, index=True)
# archive: library → medium. restore: medium → library (US06-04). ``source_path``
# and ``destination_path`` always mean "from" and "to" for this direction.
direction: Mapped[str] = mapped_column(String, nullable=False, default="archive")
source_path: Mapped[str] = mapped_column(String, nullable=False)
destination_path: Mapped[str] = mapped_column(String, nullable=False)
# Relative to the location root, because the medium can be mounted elsewhere.

View File

@@ -43,6 +43,10 @@ class Asset(Base):
# link is written and read by the archive service (US06-02).
archive_location_id: Mapped[str | None] = mapped_column(String)
archive_path: Mapped[str | None] = mapped_column(String)
# Set when the archived copy was proven to hold bytes other than the recorded
# ones (US06-04). Restore refuses such an asset instead of accepting a different
# file; cleared as soon as a verification matches again.
archive_divergent_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
# Duplicate canonical link: NULL when the asset is itself canonical or undecided.
canonical_asset_id: Mapped[str | None] = mapped_column(ForeignKey("assets.id"))

View File

@@ -15,6 +15,11 @@ planned → transferring → verified → removing → complete
↘ ↘ ↘ failed
```
A restore (US06-04) uses the same rows with ``direction='restore'``: it copies from
the medium back into the library and removes nothing, so it goes ``verified →
complete`` directly. ``source_path``/``destination_path`` always mean "from"/"to",
which is why the evidence table below needs no direction of its own.
- ``transferring`` — intent recorded; a temporary copy may exist, the destination
may or may not have been published. Nothing has been removed.
- ``verified`` — the archived bytes exist at their final path, hash exactly as
@@ -73,10 +78,22 @@ ALLOWED_TRANSITIONS = {
ArchiveState.FAILED: {ArchiveState.PLANNED, ArchiveState.TRANSFERRING},
}
# A restore removes nothing, so it has no ``removing`` step: a verified published
# copy is the whole job (US06-04). Keeping this as a separate table means the
# archive direction still cannot reach ``complete`` without going through removal.
RESTORE_TRANSITIONS = {
**ALLOWED_TRANSITIONS,
ArchiveState.VERIFIED: {ArchiveState.COMPLETE, ArchiveState.FAILED},
}
TERMINAL_STATES = frozenset({ArchiveState.COMPLETE})
# States where this item may already have touched the filesystem.
UNSAFE_STATES = frozenset({ArchiveState.TRANSFERRING, ArchiveState.VERIFIED, ArchiveState.REMOVING})
# Which way the bytes move. Same rows, same evidence table, opposite direction.
ARCHIVE = "archive"
RESTORE = "restore"
RESUMABLE = "resumable"
FORWARD = "forward"
MANUAL = "manual"
@@ -94,8 +111,9 @@ class JournalConflict(JournalError):
"""Fencing check failed; a newer owner has taken over this operation."""
def can_transition(current: str, target: str) -> bool:
return target in ALLOWED_TRANSITIONS.get(current, set())
def can_transition(current: str, target: str, direction: str = ARCHIVE) -> bool:
table = RESTORE_TRANSITIONS if direction == RESTORE else ALLOWED_TRANSITIONS
return target in table.get(current, set())
def _now() -> datetime:
@@ -119,7 +137,7 @@ class ArchiveJournal:
if row.journal_state in TERMINAL_STATES:
raise InvalidTransition(f"{row.journal_state} is terminal")
if row.journal_state != ArchiveState.TRANSFERRING and not can_transition(
row.journal_state, ArchiveState.TRANSFERRING
row.journal_state, ArchiveState.TRANSFERRING, row.direction
):
raise InvalidTransition(f"{row.journal_state} -> {ArchiveState.TRANSFERRING}")
if row.journal_state != ArchiveState.TRANSFERRING:
@@ -160,7 +178,7 @@ class ArchiveJournal:
if row.journal_state == target:
session.commit()
return _operation_dict(row) # idempotent
if not can_transition(row.journal_state, target):
if not can_transition(row.journal_state, target, row.direction):
raise InvalidTransition(f"{row.journal_state} -> {target}")
row.journal_state = target
@@ -192,16 +210,18 @@ class ArchiveJournal:
)
return [_operation_dict(row) for row in rows]
def incomplete(self) -> list[dict]:
def incomplete(self, *, direction: str | None = None) -> list[dict]:
"""Every operation left in a non-terminal, non-planned state — the work a
restart has to reason about."""
restart has to reason about. Without ``direction`` this spans archives and
restores, because either one half-done blocks the other."""
with self._session_factory() as session:
stmt = select(ArchiveOperation).where(
ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED])
)
if direction is not None:
stmt = stmt.where(ArchiveOperation.direction == direction)
rows = session.scalars(
select(ArchiveOperation)
.where(
ArchiveOperation.journal_state.not_in([*TERMINAL_STATES, ArchiveState.PLANNED])
)
.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence)
stmt.order_by(ArchiveOperation.plan_id, ArchiveOperation.sequence)
)
return [_operation_dict(row) for row in rows]
@@ -231,6 +251,7 @@ class ArchiveJournal:
return {
"operation_id": operation_id,
"plan_id": row["plan_id"],
"direction": row["direction"],
"album": row["album"],
"asset_id": row["asset_id"],
"source_path": row["source_path"],
@@ -243,8 +264,8 @@ class ArchiveJournal:
"destination_matches": destination_matches,
}
def classify_all(self) -> list[dict]:
return [self.classify(row["id"]) for row in self.incomplete()]
def classify_all(self, *, direction: str | None = None) -> list[dict]:
return [self.classify(row["id"]) for row in self.incomplete(direction=direction)]
def blocks_mutation(self) -> bool:
"""True when any item may have the library half-archived."""
@@ -317,6 +338,7 @@ def _operation_dict(row: ArchiveOperation) -> dict:
return {
"id": row.id,
"plan_id": row.plan_id,
"direction": row.direction,
"sequence": row.sequence,
"album": row.album,
"asset_id": row.asset_id,

View File

@@ -54,6 +54,7 @@ from sqlalchemy.orm import sessionmaker
from photo_pipeline.config import Config
from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath
from photo_pipeline.services.archive_journal import (
ARCHIVE,
MANUAL,
RESUMABLE,
ArchiveJournal,
@@ -115,6 +116,7 @@ class ArchiveTransferService:
location_id=location_id,
token=token,
albums=json.dumps(albums) if albums is not None else None,
direction=ARCHIVE,
state="planned",
schema_version=MANIFEST_VERSION,
asset_count=preflight["totals"]["assets"],
@@ -132,6 +134,7 @@ class ArchiveTransferService:
ArchiveOperation(
id=str(uuid.uuid4()),
plan_id=plan_id,
direction=ARCHIVE,
sequence=sequence,
album=album["album"],
asset_id=asset["asset_id"],
@@ -161,7 +164,11 @@ class ArchiveTransferService:
def list(self) -> list[dict]:
with self._session_factory() as session:
rows = session.scalars(select(ArchivePlan).order_by(ArchivePlan.created_at))
rows = session.scalars(
select(ArchivePlan)
.where(ArchivePlan.direction == ARCHIVE)
.order_by(ArchivePlan.created_at)
)
return [_plan_dict(row) for row in rows]
# ── apply ─────────────────────────────────────────────────────────────────
@@ -262,7 +269,7 @@ class ArchiveTransferService:
if same_filesystem:
os.rename(source, destination)
else:
self._copy_and_publish(operation, source, destination)
copy_verify_publish(source, destination, operation["expected_sha256"])
_fsync_dir(destination.parent)
# 4. The published file is the archive only once it hashes as recorded.
@@ -282,27 +289,6 @@ class ArchiveTransferService:
# 5. Only now may the active source go.
self._finish(self.journal.get(operation["id"]), location, token=token, worker_id=worker_id)
def _copy_and_publish(self, operation: dict, source: Path, destination: Path) -> None:
"""Cross-filesystem: copy to a temporary file beside the destination, prove
its bytes, then publish it atomically. The source is still untouched."""
temp = destination.with_name(f"{TEMP_PREFIX}{uuid.uuid4().hex}{TEMP_SUFFIX}")
try:
with open(source, "rb") as src, open(temp, "wb") as out:
shutil.copyfileobj(src, out, 1024 * 1024)
out.flush()
os.fsync(out.fileno())
if sha256_file(temp) != operation["expected_sha256"]:
raise PreconditionFailed("copy_mismatch", f"{source} copied with wrong bytes")
if destination.exists():
raise PreconditionFailed(
"destination_exists", f"{destination} appeared during the transfer"
)
# ponytail: rename after an exists() check. The archiver lane is single
# and local; use O_EXCL/link-based publish if a second writer ever exists.
os.rename(temp, destination)
finally:
temp.unlink(missing_ok=True)
def _finish(self, operation: dict, location: dict, *, token: int, worker_id: str) -> None:
"""Drive an item whose archive copy is durable through removal and
bookkeeping. Every step is idempotent, so recovery may replay it."""
@@ -457,7 +443,7 @@ class ArchiveTransferService:
"""
results = {"resumed": 0, "completed": 0, "manual": 0}
touched: set[str] = set()
for verdict in self.journal.classify_all():
for verdict in self.journal.classify_all(direction=ARCHIVE):
operation = self.journal.get(verdict["operation_id"])
touched.add(operation["plan_id"])
token = (operation["fencing_token"] or 0) + 1
@@ -484,7 +470,7 @@ class ArchiveTransferService:
return results
def recovery_status(self) -> dict:
verdicts = self.journal.classify_all()
verdicts = self.journal.classify_all(direction=ARCHIVE)
return {
"operations": verdicts,
"manual": [v for v in verdicts if v["classification"] == MANUAL],
@@ -528,6 +514,33 @@ class ArchiveTransferService:
# ── module helpers ───────────────────────────────────────────────────────────
def copy_verify_publish(source: Path, destination: Path, expected_sha256: str) -> None:
"""Copy to a temporary file beside the destination, prove its bytes, then publish
it atomically. The source is never touched, so a failure costs nothing.
Shared by archiving (library → medium) and restoring (medium → library, US06-04):
both need the same promise that a published file is either complete and correct
or not there at all.
"""
temp = destination.with_name(f"{TEMP_PREFIX}{uuid.uuid4().hex}{TEMP_SUFFIX}")
try:
with open(source, "rb") as src, open(temp, "wb") as out:
shutil.copyfileobj(src, out, 1024 * 1024)
out.flush()
os.fsync(out.fileno())
if sha256_file(temp) != expected_sha256:
raise PreconditionFailed("copy_mismatch", f"{source} copied with wrong bytes")
if destination.exists():
raise PreconditionFailed(
"destination_exists", f"{destination} appeared during the transfer"
)
# ponytail: rename after an exists() check. The archiver lane is single and
# local; use O_EXCL/link-based publish if a second writer ever exists.
os.rename(temp, destination)
finally:
temp.unlink(missing_ok=True)
def _same_filesystem(source: Path, destination_dir: Path) -> bool:
"""Proven at run time from the actual devices, never from the plan's preview."""
try:
@@ -618,6 +631,7 @@ def _plan_dict(plan: ArchivePlan) -> dict:
"id": plan.id,
"location_id": plan.location_id,
"token": plan.token,
"direction": plan.direction,
"albums": json.loads(plan.albums) if plan.albums else None,
"state": plan.state,
"schema_version": plan.schema_version,

View File

@@ -0,0 +1,648 @@
"""RestoreService — plan and execute safe restores (US06-04).
Restore is archiving read backwards, with one decisive difference: it removes
nothing. The archived copy stays on its medium, so every failure mode here costs
at most a discarded temporary file. What restore must never do is *lose identity*
— the asset that comes back is the same asset, with its duplicate decision, safety
review, analysis, and upload history intact — or *overwrite* something in the
active library.
Preflight proves, per concept §9 "Restore":
- the recorded medium is mounted and is the right one (marker ``media_id``);
- every selected asset is archived, its archive copy exists, and it hashes to
exactly the bytes the database recorded — a mismatch is ``divergent`` and is
refused, never silently accepted as "the file";
- the destination lies inside the library, outside ``_IGNORE/``, and is free; a
taken path is answered with a collision-free name, never an overwrite;
- the library filesystem has room for the scope plus the configured reserve;
- no rename, archive, or restore lease is holding the lane.
Blocker codes: ``no_library_root``, ``location_offline``, ``wrong_volume``,
``unsafe_destination``, ``library_not_writable``, ``insufficient_capacity``,
``lock_conflict``, ``rename_pending``, ``archive_pending``, ``empty_scope``,
``not_archived``, ``archive_missing``, ``bytes_changed``.
Per item the sequence is:
```
journal.begin (transferring) ← intent persisted BEFORE any disk change
recheck: medium, hash, free destination, asset still archived
copy to a temporary file beside the destination, fsync, hash it back
atomically publish into the library
journal → verified
current_path = destination, availability = active, path occurrence opened
journal → complete
```
Like archiving, the confirmation token is derived from the report, so a changed
scope, a swapped medium, or a destination that filled up invalidates it.
"""
from __future__ import annotations
import hashlib
import json
import os
import shutil
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.jobs.domain_handlers import ARCHIVE_LOCK, LIBRARY_WRITE_LOCK, UPLOAD_LOCK
from photo_pipeline.models import ArchiveLocation, ArchiveOperation, ArchivePlan, Asset, AssetPath
from photo_pipeline.path_policy import PathPolicyError, is_excluded, normalize_root, resolve_within
from photo_pipeline.services import availability
from photo_pipeline.services.archive_journal import (
MANUAL,
RESTORE,
RESUMABLE,
ArchiveJournal,
ArchiveState,
)
from photo_pipeline.services.archive_transfer import (
_clean_temp_files,
_fsync_dir,
_plan_dict,
copy_verify_publish,
)
from photo_pipeline.services.archives import ArchiveError
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.jobs import JobService
from photo_pipeline.services.rename_apply import PreconditionFailed, maybe_fault
from photo_pipeline.services.rename_journal import RenameJournal
PREFLIGHT_VERSION = 1
TOKEN_PREFIX = f"r{PREFLIGHT_VERSION}"
# What a restored file is called when its original name is taken. The suffix is
# visible on purpose: a restore that quietly reuses a name is indistinguishable
# from an overwrite.
RESTORED_SUFFIX = "restored"
LOCKS = (LIBRARY_WRITE_LOCK, UPLOAD_LOCK, ARCHIVE_LOCK)
APPLYABLE_PLAN_STATES = frozenset({"planned", "applying", "failed", "complete"})
def _now() -> datetime:
return datetime.now(timezone.utc)
def _issue(code: str, message: str) -> dict:
return {"code": code, "message": message}
class RestoreService:
def __init__(self, session_factory: sessionmaker, *, config: Config) -> None:
self._session_factory = session_factory
self._config = config
self._roots = tuple(normalize_root(root) for root in config.library_roots)
self.journal = ArchiveJournal(session_factory)
# ── preflight ─────────────────────────────────────────────────────────────
def preflight(self, location_id: str, asset_ids: list[str] | None = None) -> dict:
"""Validate a restore scope and issue its token. Nothing is written."""
with self._session_factory() as session:
location = session.get(ArchiveLocation, location_id)
if location is None:
raise ArchiveError("unknown_location", f"unknown archive location {location_id!r}")
root = Path(location.root)
online = availability.location_online(location)
marker = availability.read_marker(root)
report = {
"schema_version": PREFLIGHT_VERSION,
"location": {
"id": location.id,
"name": location.name,
"root": str(root),
"media_id": location.media_id,
"state": _location_state(root, marker, location.media_id),
},
"blockers": [],
}
items = self._items(session, location, asset_ids, reachable=online)
report["blockers"] += self._destination_blockers(report["location"]["state"], root)
report["blockers"] += self._lock_blockers()
report["items"] = items
report["totals"] = {
"assets": len(items),
"blocked": sum(1 for item in items if item["blockers"]),
"bytes": sum(item["byte_size"] or 0 for item in items),
}
report["capacity"] = self._capacity(report["totals"]["bytes"])
if not report["capacity"]["sufficient"]:
report["blockers"].append(
_issue(
"insufficient_capacity",
f"{report['totals']['bytes']} B plus a "
f"{self._config.archive_free_space_reserve_bytes} B reserve do not fit in "
f"{report['capacity']['free_bytes']} B of free space",
)
)
if not items:
report["blockers"].append(
_issue("empty_scope", "no archived assets are in the selected scope")
)
report["state"] = (
"ready"
if not report["blockers"] and not report["totals"]["blocked"]
else "blocked"
)
report["token"] = _token(report)
report["generated_at"] = _now().isoformat()
return report
def verify_token(self, token: str, location_id: str, asset_ids: list[str] | None = None) -> bool:
return bool(token) and token == self.preflight(location_id, asset_ids)["token"]
def _items(
self, session, location: ArchiveLocation, asset_ids: list[str] | None, *, reachable: bool
) -> list[dict]:
stmt = select(Asset).where(Asset.archive_location_id == location.id)
if asset_ids is None:
# A restored asset keeps its archive link; the default scope is only what
# is still archived, so restoring twice is an empty scope, not a blocker.
stmt = stmt.where(Asset.availability_state.in_(availability.ARCHIVED))
else:
stmt = stmt.where(Asset.id.in_(asset_ids))
assets = list(session.scalars(stmt.order_by(Asset.archive_path)))
if asset_ids is not None:
unknown = sorted(set(asset_ids) - {asset.id for asset in assets})
if unknown:
raise ArchiveError(
"unknown_asset", f"not archived at this location: {', '.join(unknown)}"
)
taken: set[str] = set()
return [self._item(asset, location, reachable=reachable, taken=taken) for asset in assets]
def _item(self, asset: Asset, location: ArchiveLocation, *, reachable: bool, taken: set) -> dict:
source = Path(location.root) / (asset.archive_path or "")
blockers: list[dict] = []
archive_sha256 = None
if asset.availability_state not in availability.ARCHIVED:
blockers.append(
_issue("not_archived", f"asset {asset.id} is {asset.availability_state}")
)
if reachable:
if not source.exists():
blockers.append(_issue("archive_missing", f"{source} is not on the medium"))
else:
archive_sha256 = sha256_file(source)
if asset.current_sha256 and archive_sha256 != asset.current_sha256:
blockers.append(
_issue(
"bytes_changed",
f"{source} holds bytes that are not the recorded ones; "
"the archived copy is divergent",
)
)
destination, destination_blockers = self._destination(asset, taken)
blockers += destination_blockers
if destination is not None:
taken.add(str(destination))
return {
"asset_id": asset.id,
"archive_path": asset.archive_path,
"source_path": str(source),
"destination_path": str(destination) if destination else None,
"expected_sha256": asset.current_sha256,
"archive_sha256": archive_sha256,
"byte_size": asset.byte_size,
"availability_state": asset.availability_state,
"blockers": blockers,
}
def _destination(self, asset: Asset, taken: set) -> tuple[Path | None, list[dict]]:
"""A free path inside the library that mirrors the archived layout.
Restoring onto an existing file is never an option, so a taken name is
answered with ``name (restored).ext`` — visible, ordinary, and impossible to
confuse with an overwrite.
"""
if not self._roots:
return None, [_issue("no_library_root", "no library root is configured")]
root = self._roots[0]
try:
candidate = resolve_within(root, root / (asset.archive_path or ""))
except PathPolicyError as error:
return None, [_issue("unsafe_destination", str(error))]
if is_excluded(candidate):
return None, [
_issue("unsafe_destination", f"{candidate} is inside an excluded (_IGNORE/) tree")
]
return _free_path(candidate, taken), []
def _destination_blockers(self, state: str, root: Path) -> list[dict]:
blockers: list[dict] = []
if not self._roots:
blockers.append(_issue("no_library_root", "no library root is configured"))
elif not os.access(self._roots[0], os.W_OK):
blockers.append(
_issue("library_not_writable", f"{self._roots[0]} is not writable")
)
if state == "offline":
blockers.append(
_issue("location_offline", f"the archive medium is not mounted at {root}")
)
elif state == "wrong_volume":
blockers.append(_issue("wrong_volume", f"{root} holds a different archive medium"))
return blockers
def _lock_blockers(self) -> list[dict]:
blockers: list[dict] = []
jobs = JobService(self._session_factory)
for lock in LOCKS:
held = jobs.blockers(lock)
if held:
blockers.append(
_issue("lock_conflict", f"the {lock} lane is busy: job {held[0]['id']}")
)
if RenameJournal(self._session_factory).blocks_mutation():
blockers.append(
_issue("rename_pending", "an unresolved rename must be recovered before restoring")
)
if self.journal.blocks_mutation():
blockers.append(
_issue(
"archive_pending",
"an unresolved archive or restore must be recovered before restoring",
)
)
return blockers
def _capacity(self, required: int) -> dict:
reserve = self._config.archive_free_space_reserve_bytes
free = shutil.disk_usage(self._roots[0]).free if self._roots else None
return {
"required_bytes": required,
"reserve_bytes": reserve,
"free_bytes": free,
"sufficient": free is not None and free >= required + reserve,
}
# ── plans ─────────────────────────────────────────────────────────────────
def create(self, location_id: str, asset_ids: list[str] | None = None, *, token: str) -> dict:
preflight = self.preflight(location_id, asset_ids)
if not token or token != preflight["token"]:
raise ArchiveError("stale_token", "the restore preflight changed since it was approved")
if preflight["state"] != "ready":
codes = ", ".join(sorted({issue["code"] for issue in preflight["blockers"]})) or "-"
blocked = sorted(
{issue["code"] for item in preflight["items"] for issue in item["blockers"]}
)
raise ArchiveError(
"blocked", f"the restore scope is blocked: {', '.join(blocked) or codes}"
)
plan_id = str(uuid.uuid4())
with self._session_factory() as session:
session.add(
ArchivePlan(
id=plan_id,
location_id=location_id,
token=token,
albums=json.dumps(asset_ids) if asset_ids is not None else None,
direction=RESTORE,
state="planned",
schema_version=PREFLIGHT_VERSION,
asset_count=preflight["totals"]["assets"],
byte_size=preflight["totals"]["bytes"],
)
)
session.flush()
for sequence, item in enumerate(preflight["items"]):
session.add(
ArchiveOperation(
id=str(uuid.uuid4()),
plan_id=plan_id,
direction=RESTORE,
sequence=sequence,
album=Path(item["archive_path"]).parent.name or "(root)",
asset_id=item["asset_id"],
source_path=item["source_path"],
destination_path=item["destination_path"],
archive_path=item["archive_path"],
expected_sha256=item["expected_sha256"],
byte_size=item["byte_size"],
journal_state=ArchiveState.PLANNED,
)
)
session.commit()
return self.get(plan_id)
def get(self, plan_id: str) -> dict | None:
with self._session_factory() as session:
plan = session.get(ArchivePlan, plan_id)
if plan is None or plan.direction != RESTORE:
return None
report = _plan_dict(plan)
report["operations"] = self.journal.operations(plan_id)
return report
def list(self) -> list[dict]:
with self._session_factory() as session:
rows = session.scalars(
select(ArchivePlan)
.where(ArchivePlan.direction == RESTORE)
.order_by(ArchivePlan.created_at)
)
return [_plan_dict(row) for row in rows]
# ── apply ─────────────────────────────────────────────────────────────────
def apply(
self, plan_id: str, *, expected_version: int | None = None, worker_id: str = "restore"
) -> dict:
plan = self._require_plan(plan_id)
if expected_version is not None and plan["version"] != expected_version:
raise ArchiveError(
"stale_plan",
f"plan {plan_id} is at version {plan['version']}, expected {expected_version}",
)
if plan["state"] not in APPLYABLE_PLAN_STATES:
raise ArchiveError("invalid_state", f"plan {plan_id} is {plan['state']}")
blocking = [row for row in self.journal.incomplete() if row["plan_id"] != plan_id]
if blocking:
raise ArchiveError(
"archive_pending",
f"another archive operation is unresolved ({blocking[0]['id']}); recover it first",
)
token = self._claim_plan(plan_id)
location = self._location(plan["location_id"])
restored = failed = skipped = 0
for operation in self.journal.operations(plan_id):
if operation["journal_state"] == ArchiveState.COMPLETE:
skipped += 1
continue
try:
if operation["journal_state"] == ArchiveState.VERIFIED:
self._finish(operation, token=token)
else:
self._restore_one(operation, location, token=token, worker_id=worker_id)
restored += 1
except PreconditionFailed as error:
self._fail(operation, token, error.code, str(error))
failed += 1
except Exception as error: # unexpected: record and stop touching disk
self._fail(operation, token, "restore_error", str(error))
failed += 1
state = self.journal.sync_plan_state(plan_id)
return {
"plan_id": plan_id,
"restored": restored,
"failed": failed,
"skipped": skipped,
"state": state,
}
def _restore_one(self, operation: dict, location: dict, *, token: int, worker_id: str) -> None:
source = Path(operation["source_path"])
destination = Path(operation["destination_path"])
# 1. Intent first; from here a crash is resolvable from journal + disk.
self.journal.begin(operation["id"], worker_id=worker_id, fencing_token=token)
maybe_fault(ArchiveState.TRANSFERRING)
# 2. Recheck against the medium and the library as they are right now.
self._recheck(operation, source, destination, location)
destination.parent.mkdir(parents=True, exist_ok=True)
# 3. Always copy: the archived original stays on its medium.
copy_verify_publish(source, destination, operation["expected_sha256"])
_fsync_dir(destination.parent)
if sha256_file(destination) != operation["expected_sha256"]:
raise PreconditionFailed(
"restore_mismatch", f"{destination} does not hold the expected bytes"
)
self.journal.transition(operation["id"], ArchiveState.VERIFIED, fencing_token=token)
maybe_fault(ArchiveState.VERIFIED)
self._finish(self.journal.get(operation["id"]), token=token)
def _finish(self, operation: dict, *, token: int) -> None:
"""Publish the restored file to the database. Idempotent, so recovery may
replay it after a crash between the copy and the bookkeeping."""
destination = Path(operation["destination_path"])
if not destination.exists() or sha256_file(destination) != operation["expected_sha256"]:
raise PreconditionFailed(
"restore_unverified", f"{destination} is not a verified restored copy"
)
self._record_restored(operation, destination)
self.journal.transition(operation["id"], ArchiveState.COMPLETE, fencing_token=token)
maybe_fault(ArchiveState.COMPLETE)
def _recheck(self, operation: dict, source: Path, destination: Path, location: dict) -> None:
root = Path(location["root"])
if not root.is_dir() or not (root / availability.MARKER_NAME).exists():
raise PreconditionFailed("location_offline", f"{root} is not the archive medium")
if not source.exists():
raise PreconditionFailed("archive_missing", f"{source} is not on the medium")
if source.is_symlink() or destination.is_symlink():
raise PreconditionFailed("symlink", "refusing to restore through a symlink")
if destination.exists():
# Never overwrite: the plan's free path was taken since it was made.
raise PreconditionFailed(
"destination_exists", f"destination {destination} is occupied"
)
if not self._inside_library(destination):
raise PreconditionFailed(
"destination_escape", f"{destination} is outside the library roots"
)
if sha256_file(source) != operation["expected_sha256"]:
self._mark_divergent(operation["asset_id"])
raise PreconditionFailed(
"bytes_changed", f"{source} changed since the plan was approved"
)
with self._session_factory() as session:
asset = session.get(Asset, operation["asset_id"])
if asset is None or asset.availability_state not in availability.ARCHIVED:
raise PreconditionFailed(
"not_archived", f"asset {operation['asset_id']} is no longer archived"
)
def _inside_library(self, destination: Path) -> bool:
for root in self._roots:
try:
resolve_within(root, destination)
return True
except PathPolicyError:
continue
return False
# ── database ──────────────────────────────────────────────────────────────
def _record_restored(self, operation: dict, destination: Path) -> None:
"""The bytes are back in the library: open the new active occurrence and set
availability. Identity, decisions, and history are untouched — that is the
entire point of restoring rather than re-importing."""
now = _now()
with self._session_factory() as session:
asset = session.get(Asset, operation["asset_id"])
if asset is None:
raise PreconditionFailed(
"asset_missing", f"asset {operation['asset_id']} no longer exists"
)
# A restored asset may be returning to a path it once held, so only an
# *open* occurrence counts as already registered — that is what keeps
# recovery idempotent without collapsing the path history.
recorded = session.scalar(
select(AssetPath).where(
AssetPath.asset_id == asset.id,
AssetPath.path == str(destination),
AssetPath.valid_until.is_(None),
)
)
if recorded is None: # idempotent: recovery may replay this
session.add(
AssetPath(
asset_id=asset.id,
path=str(destination),
valid_from=now,
reason="restore",
)
)
asset.current_path = str(destination)
asset.availability_state = availability.ACTIVE
asset.missing_at = None
# The archive copy stays where it is; keeping the link means a restored
# asset still knows which medium holds its archived bytes.
asset.archive_divergent_at = None
asset.state_version += 1
asset.updated_at = now
session.commit()
def _mark_divergent(self, asset_id: str) -> None:
"""Record that the archived copy is not the recorded file. Durable, because
the next restore attempt must not rediscover this from scratch."""
with self._session_factory() as session:
asset = session.get(Asset, asset_id)
if asset is None:
return
asset.archive_divergent_at = _now()
asset.state_version += 1
session.commit()
# ── recovery ──────────────────────────────────────────────────────────────
def recover(self, *, worker_id: str = "restore-recovery") -> dict:
"""Resolve every incomplete restore from journal + disk evidence.
A restore never removed anything, so ``resumable`` simply discards the
temporary debris and re-plans the item; ``forward`` finishes the bookkeeping
for a published file; ``manual`` is left untouched and keeps blocking.
"""
results = {"resumed": 0, "completed": 0, "manual": 0}
touched: set[str] = set()
for verdict in self.journal.classify_all(direction=RESTORE):
operation = self.journal.get(verdict["operation_id"])
touched.add(operation["plan_id"])
token = (operation["fencing_token"] or 0) + 1
if verdict["classification"] == MANUAL:
results["manual"] += 1
continue
if verdict["classification"] == RESUMABLE:
_clean_temp_files(Path(operation["destination_path"]).parent)
self.journal.transition(operation["id"], ArchiveState.PLANNED, fencing_token=token)
results["resumed"] += 1
continue
try:
self._finish(operation, token=token)
results["completed"] += 1
except PreconditionFailed as error:
self._fail(operation, token, error.code, str(error))
results["manual"] += 1
for plan_id in touched:
self.journal.sync_plan_state(plan_id)
return results
def recovery_status(self) -> dict:
verdicts = self.journal.classify_all(direction=RESTORE)
return {
"operations": verdicts,
"manual": [v for v in verdicts if v["classification"] == MANUAL],
"blocks_mutation": self.journal.blocks_mutation(),
}
# ── helpers ───────────────────────────────────────────────────────────────
def _fail(self, operation: dict, token: int, code: str, message: str) -> None:
self.journal.transition(
operation["id"], ArchiveState.FAILED, fencing_token=token, error=(code, message)
)
def _require_plan(self, plan_id: str) -> dict:
plan = self.get(plan_id)
if plan is None:
raise ArchiveError("unknown_plan", f"unknown restore plan {plan_id!r}")
return plan
def _location(self, location_id: str) -> dict:
with self._session_factory() as session:
location = session.get(ArchiveLocation, location_id)
if location is None:
raise ArchiveError("unknown_location", f"unknown archive location {location_id!r}")
return {"id": location.id, "root": location.root, "media_id": location.media_id}
def _claim_plan(self, plan_id: str) -> int:
with self._session_factory() as session:
plan = session.get(ArchivePlan, plan_id)
plan.version += 1
plan.state = "applying"
plan.updated_at = _now()
token = plan.version
session.commit()
return token
# ── module helpers ───────────────────────────────────────────────────────────
def _location_state(root: Path, marker: dict | None, media_id: str) -> str:
if not root.is_dir() or marker is None:
return "offline"
return "online" if marker.get("media_id") == media_id else "wrong_volume"
def _free_path(candidate: Path, taken: set) -> Path:
"""``a.jpg`` → ``a (restored).jpg`` → ``a (restored 2).jpg`` …
``taken`` holds the destinations already claimed by earlier items of the same
plan, so two restores in one scope cannot plan the same path.
"""
if not candidate.exists() and str(candidate) not in taken:
return candidate
stem, suffix = candidate.stem, candidate.suffix
attempt = 1
while True:
label = RESTORED_SUFFIX if attempt == 1 else f"{RESTORED_SUFFIX} {attempt}"
alternative = candidate.with_name(f"{stem} ({label}){suffix}")
if not alternative.exists() and str(alternative) not in taken:
return alternative
attempt += 1
def _token(report: dict) -> str:
"""Digest of everything the report asserts about the scope and the medium.
Free space is excluded: it drifts constantly without changing what a restore
would do, and the capacity verdict itself is part of the digest.
"""
payload = {key: value for key, value in report.items() if key not in ("generated_at", "token")}
payload["capacity"] = {
key: value for key, value in payload["capacity"].items() if key != "free_bytes"
}
digest = hashlib.sha256(
json.dumps(payload, sort_keys=True, ensure_ascii=False, default=str).encode("utf-8")
).hexdigest()
return f"{TOKEN_PREFIX}:{digest}"

View File

@@ -393,6 +393,115 @@ def mark_upload_ready(seeded: Seeded, *, unverified: tuple[str, ...] = ()) -> No
session.commit()
# ── Phase F: an archivable album and a mountable fake medium ─────────────────
def mark_uploaded(seeded: Seeded, *, album: str = "rome") -> None:
"""Give every seeded photo the verified upload evidence archiving requires.
Archiving refuses anything Immich is not proven to hold, and that proof is an
upload batch — recorded here as fixture state so the archive journeys do not
have to re-run an upload they are not testing.
"""
from sqlalchemy import select
from photo_pipeline.models import Asset, UploadBatch, UploadItem
from photo_pipeline.services.hashing import sha256_file
with session_factory(seeded) as sf:
with sf() as session:
batch_id = str(uuid.uuid4())
session.add(
UploadBatch(
id=batch_id,
album=album,
folder=str(seeded.lib / album),
album_name=album,
state="succeeded",
preflight_token="v1:e2e",
outcome_state="verified",
created_at=NOW,
)
)
for asset in session.scalars(select(Asset)):
if not asset.current_path:
continue
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset.id,
path=asset.current_path,
sha256=sha256_file(asset.current_path),
sha1="0" * 40,
state="sent",
outcome="uploaded",
)
)
session.commit()
class ArchiveStack:
"""A seeded, archivable library plus the server, the worker, and a fake medium.
The medium is an ordinary directory whose marker file makes it identifiable;
``unmount()`` takes that marker away, which is exactly what the application sees
when an external disk is unplugged.
"""
MARKER = ".photo-pipeline-archive.json"
def __init__(self, tmp_path: Path, seeded: Seeded) -> None:
self.tmp_path = tmp_path
self.seeded = seeded
self.archive = tmp_path / "archive"
self.archive.mkdir(exist_ok=True)
self.server: Server | None = None
self.worker: subprocess.Popen | None = None
self.base = ""
def start(self, *, worker: bool = True, extra_env: dict[str, str] | None = None) -> "ArchiveStack":
env = {"PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": "0", **(extra_env or {})}
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 register(self, name: str = "external") -> dict:
response = httpx.post(
f"{self.base}/api/v1/archive-locations",
json={"name": name, "root": str(self.archive)},
timeout=20,
)
response.raise_for_status()
return response.json()
def unmount(self) -> None:
(self.archive / self.MARKER).rename(self.archive / f"{self.MARKER}.away")
def remount(self) -> None:
(self.archive / f"{self.MARKER}.away").rename(self.archive / self.MARKER)
def plans(self) -> list[dict]:
return httpx.get(f"{self.base}/api/v1/archive-plans", timeout=20).json()["plans"]
def assets(self) -> list[dict]:
return httpx.get(
f"{self.base}/api/v1/inventory/assets", params={"limit": 200}, timeout=20
).json()["items"]
def restart_server(self) -> None:
self.server.stop()
self.server.start()
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()
class UploadStack:
"""A seeded, upload-ready library plus the server, worker, and fake Immich."""

View File

@@ -0,0 +1,302 @@
"""Browser journeys for the archive view (US06-05).
Archiving is the only stage that removes originals, so these journeys check the
two things a browser must never get wrong about it: that the preview names exactly
what would leave and where it would go, and that nothing offers an action the
server would refuse. The medium is a real directory whose marker makes it
identifiable; unmounting it is what an unplugged disk looks like from here.
Nothing is mocked inside the browser: the transfer runs in the real worker process
and the assertions read the filesystem afterwards.
"""
from __future__ import annotations
import pytest
from playwright.sync_api import expect
from tests.e2e._pipeline_harness import (
ArchiveStack,
mark_uploaded,
seed_album,
session_factory,
wait_until,
)
pytestmark = pytest.mark.phase_f # part of the Phase F acceptance gate (US06-06)
TIMEOUT = 10
RUN_TIMEOUT = 30_000
@pytest.fixture
def stack(tmp_path):
seeded = seed_album(tmp_path)
mark_uploaded(seeded)
running = ArchiveStack(tmp_path, seeded)
try:
yield running
finally:
running.stop()
def _open(page, stack) -> None:
page.goto(f"{stack.base}/app/#/archive")
page.get_by_test_id("archive-locations").wait_for()
def _archive(page, stack) -> None:
"""Confirm the archive and wait for the worker to finish the run."""
_open(page, stack)
page.get_by_test_id("start-archive").click()
expect(page.get_by_test_id("detail-state")).to_have_text("complete", timeout=RUN_TIMEOUT)
def _archived_paths(stack) -> list[str]:
return sorted(p.name for p in (stack.archive / "rome").glob("*.jpg"))
# ── preview and confirmation ─────────────────────────────────────────────────
def test_the_preview_names_the_scope_destination_and_reclaimable_bytes(page, stack):
errors = []
page.on("console", lambda m: errors.append(m.text) if m.type == "error" else None)
stack.start(worker=False)
location = stack.register()
_open(page, stack)
row = page.get_by_test_id("location-row").first
expect(row.get_by_test_id("location-label")).to_have_text("external")
expect(row.get_by_test_id("location-media")).to_have_text(location["media_id"])
expect(row.get_by_test_id("location-state")).to_have_text("online")
album = page.get_by_test_id("archive-album-row").first
expect(album.get_by_test_id("album-name")).to_have_text("rome")
expect(album.get_by_test_id("album-destination")).to_have_text(str(stack.archive / "rome"))
expect(album.get_by_test_id("album-assets")).to_have_text("2")
expect(album.get_by_test_id("album-state")).to_have_text("ready")
# Same filesystem here, so the transfer method is the move path — and it is
# named, because copy-verify-remove and move fail differently.
expect(album.get_by_test_id("album-method")).to_contain_text("move")
expect(page.get_by_test_id("destination-identity")).to_contain_text(location["media_id"])
expect(page.get_by_test_id("capacity")).to_contain_text("reserve")
expect(page.get_by_test_id("start-archive")).to_contain_text("Archive 1 album(s) · reclaim")
expect(page.get_by_test_id("start-archive")).to_be_enabled()
assert stack.plans() == [], "previewing may not create anything"
assert errors == [], f"console errors: {errors}"
def test_an_offline_medium_blocks_the_confirmation_and_says_what_to_mount(page, stack):
stack.start(worker=False)
location = stack.register()
stack.unmount()
_open(page, stack)
expect(page.get_by_test_id("location-state")).to_have_text("offline")
expect(page.get_by_test_id("preflight-blocker").first).to_have_attribute(
"data-code", "location_offline"
)
instruction = page.get_by_test_id("mount-instruction")
expect(instruction).to_contain_text("external")
expect(instruction).to_contain_text(location["media_id"])
expect(page.get_by_test_id("start-archive")).to_be_disabled()
def test_an_unuploaded_album_is_blocked_with_its_reason(page, stack):
# No upload evidence at all: archiving would remove the only copy.
seeded = stack.seeded
with session_factory(seeded) as sf:
from sqlalchemy import delete
from photo_pipeline.models import UploadItem
with sf() as session:
session.execute(delete(UploadItem))
session.commit()
stack.start(worker=False)
stack.register()
_open(page, stack)
expect(page.get_by_test_id("album-state")).to_have_text("blocked")
expect(page.get_by_test_id("album-blocker").first).to_have_attribute(
"data-code", "partial_scope"
)
expect(page.get_by_test_id("start-archive")).to_be_disabled()
# ── running, progress, reload ────────────────────────────────────────────────
def test_a_confirmed_archive_runs_and_separates_transfer_verify_and_removal(page, stack):
stack.start()
stack.register()
_archive(page, stack)
expect(page.get_by_test_id("count-complete")).to_have_text("complete: 2")
expect(page.get_by_test_id("count-failed")).to_have_text("failed: 0")
expect(page.get_by_test_id("count-transfer")).to_have_text("transfer: 0")
expect(page.get_by_test_id("count-verified")).to_have_text("verified: 0")
expect(page.get_by_test_id("count-removing")).to_have_text("removing: 0")
expect(page.get_by_test_id("operation-phase").first).to_have_text("complete")
# The library really lost the originals and the medium really holds them.
assert _archived_paths(stack) == ["a.jpg", "b.jpg"]
assert not (stack.seeded.lib / "rome" / "a.jpg").exists()
assert {a["availability_state"] for a in stack.assets()} == {"archived_online"}
def test_the_run_survives_a_reload_because_the_state_is_the_servers(page, stack):
stack.start()
stack.register()
_archive(page, stack)
page.reload()
page.get_by_test_id("plan-detail").wait_for()
expect(page.get_by_test_id("detail-state")).to_have_text("complete")
expect(page.get_by_test_id("count-complete")).to_have_text("complete: 2")
# And the run stays addressable by its own URL.
plan_id = stack.plans()[0]["id"]
page.goto(f"{stack.base}/app/#/archive?plan={plan_id}")
expect(page.get_by_test_id("plan-detail")).to_have_attribute("data-plan", plan_id)
# ── recovery ─────────────────────────────────────────────────────────────────
def test_an_interrupted_transfer_is_shown_with_evidence_and_a_safe_action(page, stack):
stack.start(worker=False)
stack.register()
_open(page, stack)
page.get_by_test_id("start-archive").click()
expect(page.get_by_test_id("archive-result")).to_be_visible()
# No worker ran, so every item is still planned; leave one mid-transfer as a
# crash would: intent recorded, nothing published, source intact.
_interrupt(stack)
page.reload()
page.get_by_test_id("archive-recovery").wait_for()
row = page.get_by_test_id("recovery-row").first
expect(row).to_have_attribute("data-classification", "resumable")
expect(row.get_by_test_id("recovery-reason")).to_contain_text("source present")
page.get_by_test_id("resolve-recovery").click()
expect(page.get_by_test_id("archive-recovery")).to_have_count(0)
assert (stack.seeded.lib / "rome" / "a.jpg").exists(), "recovery may not move anything"
def test_ambiguous_evidence_offers_no_action_at_all(page, stack):
stack.start(worker=False)
stack.register()
_open(page, stack)
page.get_by_test_id("start-archive").click()
expect(page.get_by_test_id("archive-result")).to_be_visible()
# Journal says the copy is verified, but the medium holds nothing: no evidence
# supports either finishing or retrying this item.
_interrupt(stack, state="verified")
page.reload()
page.get_by_test_id("archive-recovery").wait_for()
expect(page.get_by_test_id("recovery-row").first).to_have_attribute(
"data-classification", "manual"
)
expect(page.get_by_test_id("recovery-manual")).to_be_visible()
expect(page.get_by_test_id("resolve-recovery")).to_have_count(0)
expect(page.get_by_test_id("no-safe-recovery")).to_be_visible()
def _interrupt(stack, *, state: str = "transferring") -> None:
"""Leave the plan's first operation in a non-terminal journal state."""
with session_factory(stack.seeded) as sf:
from photo_pipeline.services.archive_journal import ArchiveJournal
journal = ArchiveJournal(sf)
plan_id = stack.plans()[0]["id"]
operation = journal.operations(plan_id)[0]
journal.begin(operation["id"], worker_id="crashed", fencing_token=1)
if state != "transferring":
journal.transition(operation["id"], state, fencing_token=1)
# ── offline browsing and restore ─────────────────────────────────────────────
def test_archived_photos_stay_browsable_while_the_medium_is_away(page, stack):
stack.start()
stack.register()
_archive(page, stack)
stack.unmount()
page.reload()
page.get_by_test_id("archived-assets").wait_for()
row = page.get_by_test_id("archived-row").first
expect(row.get_by_test_id("archived-availability")).to_have_text("archived · medium away")
expect(row.get_by_test_id("archived-path")).to_contain_text("rome/")
expect(row.get_by_test_id("archived-medium")).to_have_text("external")
expect(page.get_by_test_id("mount-instruction")).to_contain_text("external")
# The retained protected preview is served even though the original is gone.
preview = row.get_by_test_id("archived-preview")
assert preview.evaluate("img => img.complete && img.naturalWidth > 0")
# Restoring is impossible right now, and says so instead of failing later.
expect(page.get_by_test_id("start-restore")).to_be_disabled()
def test_restoring_brings_the_photos_back_without_losing_identity(page, stack):
stack.start()
stack.register()
_archive(page, stack)
before = {asset["id"] for asset in stack.assets()}
page.reload()
page.get_by_test_id("restore").wait_for()
expect(page.get_by_test_id("restore-row").first.get_by_test_id("restore-destination")).to_contain_text(
str(stack.seeded.lib / "rome")
)
page.get_by_test_id("start-restore").click()
wait_until(
lambda: all(a["availability_state"] == "active" for a in stack.assets()),
timeout=30,
)
assert {asset["id"] for asset in stack.assets()} == before # same identities
assert (stack.seeded.lib / "rome" / "a.jpg").exists()
assert _archived_paths(stack) == ["a.jpg", "b.jpg"] # the archive copy stays
def test_a_taken_name_is_restored_beside_it_never_over_it(page, stack):
stack.start()
stack.register()
_archive(page, stack)
squatter = stack.seeded.lib / "rome" / "a.jpg"
squatter.parent.mkdir(parents=True, exist_ok=True)
squatter.write_bytes(b"a different photo now lives here")
page.reload()
page.get_by_test_id("restore").wait_for()
destinations = page.get_by_test_id("restore-destination").all_inner_texts()
assert any("(restored)" in text for text in destinations), destinations
page.get_by_test_id("start-restore").click()
wait_until(
lambda: all(a["availability_state"] == "active" for a in stack.assets()),
timeout=30,
)
assert squatter.read_bytes() == b"a different photo now lives here"
assert (stack.seeded.lib / "rome" / "a (restored).jpg").exists()
# ── keyboard ─────────────────────────────────────────────────────────────────
def test_the_whole_archive_can_be_confirmed_from_the_keyboard(page, stack):
stack.start()
stack.register()
_open(page, stack)
button = page.get_by_test_id("start-archive")
button.focus()
expect(button).to_be_focused()
page.keyboard.press("Enter")
expect(page.get_by_test_id("detail-state")).to_have_text("complete", timeout=RUN_TIMEOUT)
assert _archived_paths(stack) == ["a.jpg", "b.jpg"]

View File

@@ -0,0 +1,417 @@
"""Planning and executing safe restores (US06-04).
Restoring is the one archive operation that can *add* a file to the library, so
every case here asks two questions: did the right bytes come back under the right
identity, and did anything already in the library get touched? The media are real
directories, the hashes are real, and the failure paths assert that the archived
copy is still exactly where it was — a restore that fails must cost nothing.
"""
import shutil
import uuid
from datetime import datetime, timezone
import numpy as np
import pytest
from fastapi.testclient import TestClient
from PIL import Image
from sqlalchemy import select
from photo_pipeline.api.app import create_app
from photo_pipeline.config import Config
from photo_pipeline.db import create_db_engine, create_session_factory, run_migrations
from photo_pipeline.models import Asset, AssetPath, SafetyReview, UploadBatch, UploadItem
from photo_pipeline.services import availability
from photo_pipeline.services.archive_journal import ArchiveState
from photo_pipeline.services.archive_transfer import ArchiveTransferService
from photo_pipeline.services.archives import MARKER_NAME, ArchiveError, ArchiveService
from photo_pipeline.services.hashing import sha256_file
from photo_pipeline.services.inventory import InventoryService
from photo_pipeline.services.restores import RestoreService
pytestmark = pytest.mark.phase_f # part of the Phase F acceptance gate (US06-06)
NOW = datetime(2026, 1, 1, tzinfo=timezone.utc)
# ── environment ──────────────────────────────────────────────────────────────
def _env(tmp_path):
(tmp_path / "data").mkdir(exist_ok=True)
lib = tmp_path / "lib"
lib.mkdir(exist_ok=True)
archive = tmp_path / "archive"
archive.mkdir(exist_ok=True)
config = Config.from_env(
{
"PHOTO_PIPELINE_DATA_DIR": str(tmp_path / "data"),
"PHOTO_PIPELINE_LIBRARY_ROOTS": str(lib),
"PHOTO_PIPELINE_ARCHIVE_FREE_SPACE_RESERVE_BYTES": "0",
}
)
run_migrations(config.database_url)
return config, create_session_factory(create_db_engine(config.database_url)), lib, archive
def structured(path, seed, size=(192, 144)):
path.parent.mkdir(parents=True, exist_ok=True)
rng = np.random.default_rng(seed)
w, h = size
base = np.zeros((h, w, 3), dtype=np.uint8)
for _ in range(5):
x0 = int(rng.integers(0, w - 40))
y0 = int(rng.integers(0, h - 40))
base[y0 : y0 + 40, x0 : x0 + 40] = rng.integers(0, 256, 3)
Image.fromarray(base).save(path, quality=95)
return path
def _archived(sf, config, lib, archive, album="rome", seeds=(1, 2)):
"""A real album taken all the way through archiving, ready to be restored."""
folder = lib / album
for index, seed in enumerate(seeds):
structured(folder / f"{index}.jpg", seed)
scan = InventoryService(sf).scan(lib)
with sf() as session:
batch_id = str(uuid.uuid4())
session.add(
UploadBatch(
id=batch_id,
album=album,
folder=str(folder),
album_name=album,
state="succeeded",
preflight_token="v1:test",
outcome_state="verified",
created_at=NOW,
)
)
for path, asset_id in scan.asset_ids.items():
session.add(
UploadItem(
batch_id=batch_id,
asset_id=asset_id,
path=path,
sha256=sha256_file(path),
sha1="0" * 40,
state="sent",
outcome="uploaded",
)
)
# A decision that must survive the whole round trip.
session.add(
SafetyReview(
id=str(uuid.uuid4()),
asset_id=asset_id,
decision="sfw",
score=0.01,
reviewer="test",
)
)
session.commit()
service = ArchiveService(sf, config=config)
location = service.register("external", str(archive))
token = service.preflight(location["id"])["token"]
transfers = ArchiveTransferService(sf, config=config)
plan = transfers.create(location["id"], None, token=token)
transfers.apply(plan["id"])
return location, scan.asset_ids
def _restore(sf, config, location_id, asset_ids=None):
service = RestoreService(sf, config=config)
token = service.preflight(location_id, asset_ids)["token"]
plan = service.create(location_id, asset_ids, token=token)
return service, plan, service.apply(plan["id"])
def _unmount(archive):
(archive / MARKER_NAME).rename(archive / f"{MARKER_NAME}.away")
def _assets(sf):
with sf() as session:
return {asset.id: asset for asset in session.scalars(select(Asset))}
def _codes(report):
return {issue["code"] for issue in report["blockers"]} | {
issue["code"] for item in report["items"] for issue in item["blockers"]
}
# ── preflight ────────────────────────────────────────────────────────────────
def test_preflight_blocks_offline_medium(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, _ = _archived(sf, config, lib, archive)
_unmount(archive)
report = RestoreService(sf, config=config).preflight(location["id"])
assert report["state"] == "blocked"
assert "location_offline" in _codes(report)
def test_preflight_blocks_wrong_volume(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, _ = _archived(sf, config, lib, archive)
(archive / MARKER_NAME).write_text('{"media_id": "someone-elses-disk"}', encoding="utf-8")
report = RestoreService(sf, config=config).preflight(location["id"])
assert report["state"] == "blocked"
assert "wrong_volume" in _codes(report)
def test_preflight_blocks_changed_archive_bytes(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
asset_id = next(iter(ids.values()))
with sf() as session:
archived_file = archive / session.get(Asset, asset_id).archive_path
archived_file.write_bytes(b"not the photo that was archived")
report = RestoreService(sf, config=config).preflight(location["id"])
assert report["state"] == "blocked"
assert "bytes_changed" in _codes(report)
with pytest.raises(ArchiveError) as error:
_restore(sf, config, location["id"])
assert error.value.code == "blocked"
def test_preflight_blocks_insufficient_capacity(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, _ = _archived(sf, config, lib, archive)
greedy = config.model_copy(
update={"archive_free_space_reserve_bytes": 1 << 62} # more than any disk has
)
report = RestoreService(sf, config=greedy).preflight(location["id"])
assert report["state"] == "blocked"
assert "insufficient_capacity" in _codes(report)
def test_token_changes_with_the_scope(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive)
service = RestoreService(sf, config=config)
whole = service.preflight(location["id"])["token"]
partial = service.preflight(location["id"], [sorted(ids.values())[0]])["token"]
assert whole != partial
assert service.verify_token(whole, location["id"])
assert not service.verify_token(partial, location["id"])
# ── restore ──────────────────────────────────────────────────────────────────
def test_restore_returns_bytes_identity_and_decisions(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive)
archived_hashes = {
asset_id: asset.current_sha256 for asset_id, asset in _assets(sf).items()
}
service, plan, result = _restore(sf, config, location["id"])
assert (result["restored"], result["failed"], result["state"]) == (2, 0, "complete")
for asset_id, asset in _assets(sf).items():
assert asset.availability_state == availability.ACTIVE
assert asset.current_path == str(lib / asset.archive_path)
assert sha256_file(asset.current_path) == archived_hashes[asset_id]
# The archived copy is a copy: restoring never empties the medium.
assert (archive / asset.archive_path).exists()
assert asset.archive_location_id == location["id"]
with sf() as session:
# Identity and decisions survived: same ids, same reviews, new occurrence.
assert set(ids.values()) == {a.id for a in session.scalars(select(Asset))}
assert {r.decision for r in session.scalars(select(SafetyReview))} == {"sfw"}
occurrences = [
row.reason
for row in session.scalars(
select(AssetPath).where(AssetPath.asset_id == sorted(ids.values())[0])
)
]
assert "restore" in occurrences
def test_restore_never_overwrites_a_collision(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
asset_id = next(iter(ids.values()))
with sf() as session:
archive_path = session.get(Asset, asset_id).archive_path
occupied = lib / archive_path
occupied.parent.mkdir(parents=True, exist_ok=True)
occupied.write_bytes(b"a different photo already lives here")
before = occupied.read_bytes()
service, plan, result = _restore(sf, config, location["id"])
assert result["failed"] == 0
assert occupied.read_bytes() == before # untouched
restored = _assets(sf)[asset_id].current_path
assert restored != str(occupied)
assert "(restored)" in restored
assert sha256_file(restored) == sha256_file(archive / archive_path)
def test_apply_refuses_a_destination_taken_after_planning(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
service = RestoreService(sf, config=config)
token = service.preflight(location["id"])["token"]
plan = service.create(location["id"], None, token=token)
# Someone drops a file exactly where the plan intends to publish.
destination = plan["operations"][0]["destination_path"]
from pathlib import Path
Path(destination).parent.mkdir(parents=True, exist_ok=True)
Path(destination).write_bytes(b"squatter")
result = service.apply(plan["id"])
assert result["failed"] == 1
operation = service.journal.operations(plan["id"])[0]
assert operation["journal_state"] == ArchiveState.FAILED
assert operation["error_code"] == "destination_exists"
assert Path(destination).read_bytes() == b"squatter"
assert _assets(sf)[next(iter(ids.values()))].availability_state == availability.ARCHIVED_ONLINE
def test_changed_archive_bytes_mark_the_asset_divergent(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
asset_id = next(iter(ids.values()))
service = RestoreService(sf, config=config)
token = service.preflight(location["id"])["token"]
plan = service.create(location["id"], None, token=token)
# The medium's copy is edited after the plan was approved.
with sf() as session:
archived_file = archive / session.get(Asset, asset_id).archive_path
archived_file.write_bytes(b"edited on the shelf")
result = service.apply(plan["id"])
assert result["failed"] == 1
operation = service.journal.operations(plan["id"])[0]
assert operation["error_code"] == "bytes_changed"
asset = _assets(sf)[asset_id]
assert asset.archive_divergent_at is not None # durable divergence
assert asset.availability_state == availability.ARCHIVED_ONLINE
assert asset.current_path is None # nothing was published
# ── interruption and idempotency ─────────────────────────────────────────────
def test_interrupted_before_publishing_is_resumable(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
service = RestoreService(sf, config=config)
token = service.preflight(location["id"])["token"]
plan = service.create(location["id"], None, token=token)
operation = service.journal.operations(plan["id"])[0]
# Model a kill right after the intent was written: nothing published yet.
service.journal.begin(operation["id"], worker_id="killed", fencing_token=1)
status = service.recovery_status()
assert status["operations"][0]["classification"] == "resumable"
assert service.recover() == {"resumed": 1, "completed": 0, "manual": 0}
assert service.journal.operations(plan["id"])[0]["journal_state"] == ArchiveState.PLANNED
result = service.apply(plan["id"])
assert result["failed"] == 0
assert _assets(sf)[next(iter(ids.values()))].availability_state == availability.ACTIVE
def test_interrupted_after_publishing_is_finished_by_recovery(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
asset_id = next(iter(ids.values()))
service = RestoreService(sf, config=config)
token = service.preflight(location["id"])["token"]
plan = service.create(location["id"], None, token=token)
operation = service.journal.operations(plan["id"])[0]
# Model a kill between the published copy and the database update.
from pathlib import Path
destination = Path(operation["destination_path"])
destination.parent.mkdir(parents=True, exist_ok=True)
shutil.copy2(operation["source_path"], destination)
service.journal.begin(operation["id"], worker_id="killed", fencing_token=1)
service.journal.transition(operation["id"], ArchiveState.VERIFIED, fencing_token=1)
assert service.recovery_status()["operations"][0]["classification"] == "forward"
assert service.recover()["completed"] == 1
asset = _assets(sf)[asset_id]
assert asset.availability_state == availability.ACTIVE
assert asset.current_path == str(destination)
# Repeated recovery and a repeated apply converge on the same state.
assert service.recover() == {"resumed": 0, "completed": 0, "manual": 0}
again = service.apply(plan["id"])
assert (again["skipped"], again["failed"]) == (1, 0)
assert _assets(sf)[asset_id].current_path == str(destination)
def test_restored_state_survives_restart_and_rescan(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive)
_restore(sf, config, location["id"])
restarted = create_session_factory(create_db_engine(config.database_url))
InventoryService(restarted).scan(lib)
assets = _assets(restarted)
assert set(assets) == set(ids.values()) # no new identities from the rescan
for asset in assets.values():
assert asset.availability_state == availability.ACTIVE
assert asset.missing_at is None
# Nothing is archived at that location any more, so there is nothing to restore.
again = RestoreService(restarted, config=config).preflight(location["id"])
assert _codes(again) == {"empty_scope"}
def test_restore_api_round_trip(tmp_path):
config, sf, lib, archive = _env(tmp_path)
location, ids = _archived(sf, config, lib, archive, seeds=(1,))
with TestClient(create_app(config)) as client:
report = client.post(
"/api/v1/restore-preflight", json={"location_id": location["id"]}
).json()
assert report["state"] == "ready"
stale = client.post(
"/api/v1/restore-plans",
json={"location_id": location["id"], "token": "r1:not-the-token"},
)
assert stale.status_code == 409
created = client.post(
"/api/v1/restore-plans",
json={"location_id": location["id"], "token": report["token"]},
)
assert created.status_code == 201
plan_id = created.json()["id"]
assert created.json()["direction"] == "restore"
# The plan is visible and applying it queues work on the archiver lane.
assert client.get(f"/api/v1/restore-plans/{plan_id}").status_code == 200
queued = client.post(f"/api/v1/restore-plans/{plan_id}/apply")
assert queued.status_code == 200
assert queued.json()["job"]["job_type"] == "restore_plan"
assert queued.json()["job"]["lock_key"] == "archive"
assert client.get("/api/v1/restore-recovery").json()["manual"] == []
assert _assets(sf)[next(iter(ids.values()))].availability_state == (
availability.ARCHIVED_ONLINE # the worker, not the request, does the work
)

View File

@@ -134,6 +134,12 @@
],
"US06-03": [
"tests/integration/test_offline_assets.py"
],
"US06-04": [
"tests/integration/test_restore.py"
],
"US06-05": [
"tests/e2e/test_archive_ui.py"
]
}
}