shoc-frontend-new/src/api/media-upload-engine.ts

310 lines
9.8 KiB
TypeScript
Raw Normal View History

import {
type MediaChunkAckDto,
type MediaUploadCreateDto,
type MediaUploadSessionDto,
type UploadCategory,
type UploadSurface,
MEDIA_UPLOAD_CHUNK_SIZE_BYTES,
MediaUploadCanceledError,
MediaUploadExpiredError,
MediaUploadRejectedError,
MediaUploadScanTimeoutError,
isTerminalMediaUploadState,
} from "@/api/media-upload-contracts";
import {
type MediaUploadEndpointSet,
computeSha256Hex,
isAbortError,
isNetworkUploadError,
} from "@/api/media-upload-client";
const MAX_CREATE_ATTEMPTS = 2;
const MAX_CHUNK_RECOVERY_ATTEMPTS = 3;
const CHUNK_RECOVERY_BACKOFF_MS = 200;
const MAX_COMMIT_ATTEMPTS = 2;
const MAX_POLL_ATTEMPTS = 120;
const DEFAULT_POLL_DELAY_MS = 1_500;
export interface ResumableUploadArgs {
file: File;
endpoints: MediaUploadEndpointSet;
surface: UploadSurface;
category?: UploadCategory;
signal?: AbortSignal;
onProgress?: (percent: number) => void;
/** Defaults to a fresh random UUID per attempt. */
idempotencyKey?: string;
digest?: (buffer: ArrayBuffer) => Promise<string>;
pollDelayMs?: number;
}
export function chunkCountFor(sizeBytes: number, chunkSizeBytes: number): number {
if (sizeBytes <= 0) return 1;
return Math.ceil(sizeBytes / chunkSizeBytes);
}
export function missingChunkIndices(received: readonly number[], total: number): number[] {
const acked = new Set(received);
const missing: number[] = [];
for (let index = 0; index < total; index += 1) {
if (!acked.has(index)) {
missing.push(index);
}
}
return missing;
}
function acknowledgedBytes(
received: readonly number[],
totalBytes: number,
chunkSize: number,
): number {
let bytes = 0;
for (const index of received) {
if (index >= 0 && index < chunkCountFor(totalBytes, chunkSize)) {
bytes += Math.min(chunkSize, totalBytes - index * chunkSize);
}
}
return bytes;
}
function reportProgress(
args: ResumableUploadArgs,
received: readonly number[],
totalBytes: number,
chunkSize: number,
): void {
if (!args.onProgress || totalBytes <= 0) return;
const percent = Math.min(
100,
Math.round((acknowledgedBytes(received, totalBytes, chunkSize) / totalBytes) * 100),
);
args.onProgress(percent);
}
function assertNotAborted(signal: AbortSignal | undefined): void {
if (signal?.aborted) {
throw new DOMException("The upload was aborted.", "AbortError");
}
}
function throwForTerminalState(state: MediaUploadSessionDto["state"]): void {
if (state === "Rejected") throw new MediaUploadRejectedError();
if (state === "Expired") throw new MediaUploadExpiredError();
if (state === "Canceled") throw new MediaUploadCanceledError();
}
/** Idempotent create: a network failure retries the same idempotency key, never a new session. */
async function createSessionIdempotent(
endpoints: MediaUploadEndpointSet,
dto: MediaUploadCreateDto,
signal: AbortSignal | undefined,
): Promise<MediaUploadSessionDto> {
let lastError: unknown;
for (let attempt = 0; attempt < MAX_CREATE_ATTEMPTS; attempt += 1) {
try {
return await endpoints.create(dto, signal);
} catch (error) {
if (!isNetworkUploadError(error)) throw error;
lastError = error;
}
}
throw lastError;
}
/**
* Network failure recovery: re-read the session; terminal states settle here,
* otherwise the returned receivedChunks define what still needs sending.
*/
async function recoverSession(
endpoints: MediaUploadEndpointSet,
uploadId: string,
signal: AbortSignal | undefined,
): Promise<MediaUploadSessionDto> {
let session: MediaUploadSessionDto;
try {
session = await endpoints.getStatus(uploadId, signal);
} catch (error) {
if (isAbortError(error)) throw error;
if (isNetworkUploadError(error)) {
// Still offline — bubble so the attempt surfaces as retryable failure.
throw error;
}
throw error;
}
if (isTerminalMediaUploadState(session.state)) {
throwForTerminalState(session.state);
return session;
}
return session;
}
async function waitForScanTerminal(
args: ResumableUploadArgs,
uploadId: string,
): Promise<MediaUploadSessionDto> {
const delayMs = args.pollDelayMs ?? DEFAULT_POLL_DELAY_MS;
for (let attempt = 0; attempt < MAX_POLL_ATTEMPTS; attempt += 1) {
assertNotAborted(args.signal);
await new Promise((resolve) => setTimeout(resolve, delayMs));
assertNotAborted(args.signal);
let session: MediaUploadSessionDto;
try {
session = await args.endpoints.getStatus(uploadId, args.signal);
} catch (error) {
if (isAbortError(error)) throw error;
continue; // Transient poll failure — keep polling until the deadline.
}
if (isTerminalMediaUploadState(session.state)) {
throwForTerminalState(session.state);
return session;
}
}
throw new MediaUploadScanTimeoutError();
}
async function commitAndObserve(
args: ResumableUploadArgs,
uploadId: string,
): Promise<MediaUploadSessionDto> {
let lastError: unknown;
for (let attempt = 0; attempt < MAX_COMMIT_ATTEMPTS; attempt += 1) {
try {
const committed = await args.endpoints.commit(uploadId, args.signal);
return committed.state === "Scanning" ? waitForScanTerminal(args, uploadId) : committed;
} catch (error) {
if (!isNetworkUploadError(error)) throw error;
lastError = error;
const observed = await args.endpoints.getStatus(uploadId, args.signal);
if (observed.state === "Completed") return observed;
if (isTerminalMediaUploadState(observed.state)) {
throwForTerminalState(observed.state);
}
if (observed.state === "Scanning") return waitForScanTerminal(args, uploadId);
if (attempt + 1 < MAX_COMMIT_ATTEMPTS) {
continue; // Commit is idempotent for this uploadId.
}
}
}
throw lastError;
}
async function sendChunk(
args: ResumableUploadArgs,
uploadId: string,
index: number,
chunkSize: number,
): Promise<MediaChunkAckDto> {
const { file } = args;
const start = index * chunkSize;
const end = Math.min(start + chunkSize, file.size);
const blob = file.slice(start, end);
const buffer = await blob.arrayBuffer();
const digest = args.digest ?? computeSha256Hex;
const sha256 = await digest(buffer);
return args.endpoints.putChunk({
uploadId,
index,
blob,
start,
end,
totalBytes: file.size,
sha256,
signal: args.signal,
});
}
/**
* Upload one file through the resumable session API: create → chunks → commit → poll.
* One in-flight chunk request per file; durable progress counts only acknowledged parts.
*/
export async function uploadFileResumable(
args: ResumableUploadArgs,
): Promise<MediaUploadSessionDto> {
assertNotAborted(args.signal);
const createDto: MediaUploadCreateDto = {
idempotencyKey: args.idempotencyKey ?? crypto.randomUUID(),
fileName: args.file.name,
contentType: args.file.type,
sizeBytes: args.file.size,
surface: args.surface,
...(args.category ? { category: args.category } : {}),
};
let session = await createSessionIdempotent(args.endpoints, createDto, args.signal);
const uploadId = session.uploadId;
const chunkSize =
Number.isFinite(session.chunkSizeBytes) && session.chunkSizeBytes > 0
? session.chunkSizeBytes
: MEDIA_UPLOAD_CHUNK_SIZE_BYTES;
const totalChunks = chunkCountFor(args.file.size, chunkSize);
reportProgress(args, session.receivedChunks, args.file.size, chunkSize);
try {
const recoveryAttempts = new Map<number, number>();
for (;;) {
assertNotAborted(args.signal);
const missing = missingChunkIndices(session.receivedChunks, totalChunks);
if (missing.length === 0) break;
for (const index of missing) {
assertNotAborted(args.signal);
try {
const ack = await sendChunk(args, session.uploadId, index, chunkSize);
const acked = Array.isArray(ack?.receivedChunks) ? ack.receivedChunks : [];
if (!acked.includes(index)) {
// Treat a malformed ack as not-yet-acked so the loop cannot spin forever.
throw new Error("Chunk acknowledgement missing from server response.");
}
session = { ...session, receivedChunks: acked };
} catch (error) {
if (!isNetworkUploadError(error)) throw error;
session = await recoverSession(args.endpoints, session.uploadId, args.signal);
if (session.receivedChunks.includes(index)) {
recoveryAttempts.delete(index);
} else {
const attempts = (recoveryAttempts.get(index) ?? 0) + 1;
recoveryAttempts.set(index, attempts);
if (attempts >= MAX_CHUNK_RECOVERY_ATTEMPTS) throw error;
await new Promise((resolve) =>
setTimeout(resolve, CHUNK_RECOVERY_BACKOFF_MS * attempts),
);
assertNotAborted(args.signal);
}
break; // Re-derive missing indices from the recovered server state.
}
reportProgress(args, session.receivedChunks, args.file.size, chunkSize);
}
}
assertNotAborted(args.signal);
session = await commitAndObserve(args, uploadId);
} catch (error) {
// Cancellation aborts the in-flight request and DELETEs the known session.
if (isAbortError(error)) {
await args.endpoints.cancel(uploadId);
}
throw error;
}
throwForTerminalState(session.state);
reportProgress(args, session.receivedChunks, args.file.size, chunkSize);
return session;
}
/** Upload a queue of files with at most `maxConcurrent` files in flight. */
export async function runMediaUploadQueue<T>(
items: readonly T[],
runner: (item: T) => Promise<void>,
maxConcurrent = 2,
): Promise<void> {
let cursor = 0;
const workers = Array.from({ length: Math.min(maxConcurrent, items.length) }, async () => {
while (cursor < items.length) {
const item = items[cursor];
cursor += 1;
await runner(item);
}
});
await Promise.all(workers);
}