Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
350 changes: 248 additions & 102 deletions apps/server/src/application/services/backup-service.ts

Large diffs are not rendered by default.

238 changes: 230 additions & 8 deletions apps/server/src/application/services/job-transfer-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,20 @@ import { createWriteStream } from "node:fs";
import fs from "node:fs/promises";
import path from "node:path";
import { pipeline } from "node:stream/promises";
import type { Job } from "@solid-imager/core/domain/repositories/job-repository";
import { webReadableToNodeStream } from "~/infrastructure/utils/stream-utils";

const configuredTransferDirectory = process.env.SOLID_IMAGER_JOB_TRANSFER_DIR;
const isolatedRuntimeDirectory =
process.env.E2E_RUNTIME_DIR ?? process.env.DEV_STARTUP_RUNTIME_DIR;
const JobTransferDirectory = path.resolve(
process.cwd(),
".cache",
"job-transfers",
configuredTransferDirectory ??
(isolatedRuntimeDirectory
? path.join(isolatedRuntimeDirectory, ".cache", "job-transfers")
: path.join(process.cwd(), ".cache", "job-transfers")),
);
const JobArtifactTtlMs = 24 * 60 * 60 * 1000;
const JobTransferStaleFileTtlMs = 60 * 60 * 1000;

export type JobTransferMode = "json" | "zip";

Expand Down Expand Up @@ -38,6 +44,13 @@ export function getInputPath(jobId: string, mode: JobTransferMode): string {
);
}

export function getInputPartialPath(
jobId: string,
mode: JobTransferMode,
): string {
return `${getInputPath(jobId, mode)}.partial`;
}

export function getArtifactPath(jobId: string, mode: JobTransferMode): string {
return path.join(
JobTransferDirectory,
Expand All @@ -46,6 +59,13 @@ export function getArtifactPath(jobId: string, mode: JobTransferMode): string {
);
}

export function getArtifactPartialPath(
jobId: string,
mode: JobTransferMode,
): string {
return `${getArtifactPath(jobId, mode)}.partial`;
}

export function isJobTransferPath(jobId: string, targetPath: string): boolean {
const resolvedTarget = path.resolve(targetPath);
const resolvedRoot = path.resolve(JobTransferDirectory);
Expand Down Expand Up @@ -83,12 +103,20 @@ export async function persistJobInput(
file: File,
): Promise<string> {
const inputPath = getInputPath(jobId, mode);
const partialPath = getInputPartialPath(jobId, mode);
await fs.mkdir(path.dirname(inputPath), { recursive: true });
await pipeline(
webReadableToNodeStream(file.stream()),
createWriteStream(inputPath),
);
return inputPath;
await removeJobTransferFile(inputPath);
await removeJobTransferFile(partialPath);
try {
await pipeline(
webReadableToNodeStream(file.stream()),
createWriteStream(partialPath),
);
await fs.rename(partialPath, inputPath);
return inputPath;
} finally {
await removeJobTransferFile(partialPath);
}
}

export async function removeJobTransferFile(targetPath: string): Promise<void> {
Expand Down Expand Up @@ -135,3 +163,197 @@ export async function cleanupExpiredJobTransferFiles(
}
}
}

type TransferFileKind = "inputs" | "artifacts";

type JobLookup = (jobId: string) => Promise<Job | null>;

type JobTransferCleanupResult = {
removedFiles: number;
removedBytes: number;
};

const TransferFileNamePattern =
/^([0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12})\.(?:ndjson|tar)$/i;
const TarStagingDirectoryNamePattern =
/^([0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12})-export-/i;

function readTransferJobId(fileName: string): string | null {
return TransferFileNamePattern.exec(fileName)?.[1] ?? null;
}

function readTarStagingJobId(directoryName: string): string | null {
return TarStagingDirectoryNamePattern.exec(directoryName)?.[1] ?? null;
}

function readRestoreInputPath(payload: unknown): string | null {
if (
typeof payload !== "object" ||
payload === null ||
Array.isArray(payload)
) {
return null;
}
const inputPath = (payload as { inputPath?: unknown }).inputPath;
return typeof inputPath === "string" ? inputPath : null;
}

function isExpectedTransferFile(
kind: TransferFileKind,
targetPath: string,
job: Job | null,
): boolean {
if (!job) {
return false;
}

if (kind === "artifacts") {
return job.status === "completed" && job.artifactPath === targetPath;
}

return (
(job.status === "pending" || job.status === "in_progress") &&
readRestoreInputPath(job.payload) === targetPath
);
}

async function latestModificationTime(targetPath: string): Promise<number> {
let stat: Awaited<ReturnType<typeof fs.stat>>;
try {
stat = await fs.stat(targetPath);
} catch {
return 0;
}

if (!stat.isDirectory()) {
return stat.mtimeMs;
}

let latest = stat.mtimeMs;
let entries: Dirent[];
try {
entries = await fs.readdir(targetPath, { withFileTypes: true });
} catch {
return latest;
}

for (const entry of entries) {
latest = Math.max(
latest,
await latestModificationTime(path.join(targetPath, entry.name)),
);
}
return latest;
}

async function cleanupOrphanedTarStaging(
expirationTime: number,
findJob: JobLookup,
): Promise<JobTransferCleanupResult> {
const stagingDirectory = path.join(JobTransferDirectory, "..", "tar-staging");
let entries: Dirent[];
try {
entries = await fs.readdir(stagingDirectory, { withFileTypes: true });
} catch (error) {
if (isNodeErrorCode(error, "ENOENT")) {
return { removedFiles: 0, removedBytes: 0 };
}
throw error;
}

let removedFiles = 0;
let removedBytes = 0;
for (const entry of entries) {
if (!entry.isDirectory()) {
continue;
}
const targetPath = path.join(stagingDirectory, entry.name);
const jobId = readTarStagingJobId(entry.name);
if (jobId) {
const job = await findJob(jobId);
if (job?.status === "pending" || job?.status === "in_progress") {
continue;
}
}
if ((await latestModificationTime(targetPath)) > expirationTime) {
continue;
}
const stat = await fs.stat(targetPath).catch(() => null);
await fs.rm(targetPath, { recursive: true, force: true });
removedFiles++;
removedBytes += stat?.isDirectory() ? 0 : (stat?.size ?? 0);
}
Comment thread
hmjn023 marked this conversation as resolved.

return { removedFiles, removedBytes };
}

/**
* Removes transfer files left by failed, cancelled, stale, or deleted jobs.
* Completed artifacts and restore inputs still referenced by pending jobs are
* retained; age-based expiry remains handled by cleanupExpiredJobTransferFiles.
*/
export async function cleanupOrphanedJobTransferFiles(
findJob: JobLookup,
now = Date.now(),
): Promise<JobTransferCleanupResult> {
const expirationTime = now - JobTransferStaleFileTtlMs;
const jobCache = new Map<string, Job | null>();
let removedFiles = 0;
let removedBytes = 0;

for (const kind of ["inputs", "artifacts"] as const) {
const directoryPath = path.join(JobTransferDirectory, kind);
let entries: Dirent[];
try {
entries = await fs.readdir(directoryPath, { withFileTypes: true });
} catch (error) {
if (isNodeErrorCode(error, "ENOENT")) {
continue;
}
throw error;
}

for (const entry of entries) {
if (!entry.isFile()) {
continue;
}
const targetPath = path.join(directoryPath, entry.name);
const stat = await fs.stat(targetPath).catch(() => null);
if (!stat || stat.mtimeMs > expirationTime) {
continue;
}

const isPartial = entry.name.endsWith(".partial");
const jobId = readTransferJobId(
isPartial ? entry.name.slice(0, -".partial".length) : entry.name,
);
let shouldRemove = isPartial;
if (!isPartial && jobId) {
if (!jobCache.has(jobId)) {
jobCache.set(jobId, await findJob(jobId));
}
const job = jobCache.get(jobId) ?? null;
// An unknown job may belong to another isolated database/runtime.
// Keep it for age-based expiry instead of deleting user data.
shouldRemove =
job !== null && !isExpectedTransferFile(kind, targetPath, job);
}

if (!shouldRemove) {
continue;
}
await removeJobTransferFile(targetPath);
removedFiles++;
removedBytes += stat.size;
}
}

const stagingCleanup = await cleanupOrphanedTarStaging(
expirationTime,
findJob,
);
return {
removedFiles: removedFiles + stagingCleanup.removedFiles,
removedBytes: removedBytes + stagingCleanup.removedBytes,
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import { services } from "~/application/registry";
import { BackupService } from "~/application/services/backup-service";
import {
getArtifactMetadata,
getArtifactPartialPath,
isJobTransferPath,
removeJobTransferFile,
} from "~/application/services/job-transfer-storage";
Expand All @@ -27,19 +28,36 @@ export async function processSourceExportJob(job: Job): Promise<void> {
const payload = sourceExportJobPayloadSchema.parse(job.payload);
const dump = await BackupService.createDump(job.mediaSourceId, payload.mode, {
includeImages: payload.includeImages,
jobId: job.id,
});
const artifact = getArtifactMetadata(job.id, job.mediaSourceId, payload.mode);
const partialPath = getArtifactPartialPath(job.id, payload.mode);
let artifactCommitted = false;

await fs.mkdir(path.dirname(artifact.path), { recursive: true });
await pipeline(
webReadableToNodeStream(asDumpStream(dump)),
createWriteStream(artifact.path),
);
const stat = await fs.stat(artifact.path);
await services.getJobRepository().setArtifact(job.id, {
...artifact,
size: stat.size,
});
await removeJobTransferFile(artifact.path);
await removeJobTransferFile(partialPath);
try {
await pipeline(
webReadableToNodeStream(asDumpStream(dump)),
createWriteStream(partialPath),
);
await fs.rename(partialPath, artifact.path);
artifactCommitted = true;

const stat = await fs.stat(artifact.path);
await services.getJobRepository().setArtifact(job.id, {
...artifact,
size: stat.size,
});
} catch (error) {
if (artifactCommitted) {
await removeJobTransferFile(artifact.path);
}
throw error;
} finally {
await removeJobTransferFile(partialPath);
}
}

export async function processSourceRestoreJob(job: Job): Promise<unknown> {
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/infrastructure/api-clients/sources-api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ export async function fetchSourceDump(
mode: "json" | "zip" = "json",
opts?: { includeImages?: boolean },
): Promise<Blob> {
const includeImages = opts?.includeImages ?? false;
const includeImages = opts?.includeImages ?? mode === "zip";
const job = await orpc.sources.enqueueExport({
id,
mode,
Expand Down
Loading