From c4873840d4057b919a00619df2663d99ac644cc1 Mon Sep 17 00:00:00 2001 From: Nezumi-2711 Date: Fri, 21 Aug 2026 12:24:59 +0700 Subject: [PATCH] feat: add api for manage bucket --- docs/limitations.md | 4 +- docs/openapi.yaml | 377 +++++++++++++++++++++++++++++ src/cors.ts | 2 +- src/google-drive.ts | 96 ++++++-- src/multipart-core.ts | 188 +++++++++++++++ src/object-api.ts | 546 ++++++++++++++++++++++++++++++++++++++++++ src/router.ts | 165 +++---------- src/status-api.ts | 21 +- test/cors.test.ts | 2 +- test/fake-drive.ts | 216 +++++++++++++++++ test/objects.test.ts | 346 ++++++++++++++++++++++++++ test/s3.test.ts | 190 +-------------- 12 files changed, 1808 insertions(+), 345 deletions(-) create mode 100644 src/multipart-core.ts create mode 100644 src/object-api.ts create mode 100644 test/fake-drive.ts create mode 100644 test/objects.test.ts diff --git a/docs/limitations.md b/docs/limitations.md index 3ef8106..46fc4c9 100644 --- a/docs/limitations.md +++ b/docs/limitations.md @@ -24,11 +24,11 @@ Bucket `?acl`, `?versioning`, and `?location` requests are currently treated as ## Directories and objects -Directories are physical Google Drive folders created from object-key path segments. The Worker does not create zero-byte directory marker objects. An empty prefix can therefore exist as a Drive folder without an S3 marker object. +Directories are physical Google Drive folders created from object-key path segments. The Worker does not create zero-byte directory marker objects. An empty prefix can therefore exist as a Drive folder without an S3 marker object. Folder deletion via `/api/objects/folder` refuses non-empty folders by default, and supports atomic recursive deletion into Google Drive trash when `recursive=1` is provided (avoiding partial deletion via enumeration). ## Multipart behavior -Parts must be uploaded strictly in consecutive order beginning at part 1. Parts are immutable after they have committed. An incomplete upload expires after 24 hours. `ETAG_STYLE=md5` returns Google Drive's file MD5 for completed multipart objects; `ETAG_STYLE=multipart` is available only for clients that require an S3-style composite ETag. +Parts must be uploaded strictly in consecutive order beginning at part 1. Parts are immutable after they have committed. Default part size recommended and returned by the API is 8 MiB (aligned to Google Drive 256 KiB chunks). An incomplete upload expires after 24 hours. `ETAG_STYLE=md5` returns Google Drive's file MD5 for completed multipart objects; `ETAG_STYLE=multipart` is available only for clients that require an S3-style composite ETag. ## Quotas and plan limits diff --git a/docs/openapi.yaml b/docs/openapi.yaml index dcf3e5b..15e8bea 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -301,6 +301,383 @@ paths: content: application/json: schema: { $ref: '#/components/schemas/AuthError' } + /api/objects: + get: + tags: [Objects API] + operationId: listObjectsApi + summary: List objects and folders in a bucket (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: prefix + schema: { type: string, default: "" } + - in: query + name: delimiter + schema: { type: string, default: "/" } + responses: + '200': + description: Object and folder listing. + content: + application/json: + schema: + type: object + properties: + bucket: { type: string } + prefix: { type: string } + delimiter: { type: string, nullable: true } + folders: + type: array + items: + type: object + properties: + prefix: { type: string } + name: { type: string } + objects: + type: array + items: + type: object + properties: + key: { type: string } + name: { type: string } + size: { type: integer } + contentType: { type: string } + lastModified: { type: string, nullable: true } + etag: { type: string } + truncated: { type: boolean } + '400': + description: Missing query parameters. + '401': + description: Unauthorized. + '404': + description: Bucket not found. + delete: + tags: [Objects API] + operationId: deleteObjectApi + summary: Delete an object (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + responses: + '204': + description: Object deleted. + '401': + description: Unauthorized. + '404': + description: Bucket or object not found. + /api/objects/metadata: + get: + tags: [Objects API] + operationId: getObjectMetadataApi + summary: Get object metadata (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + responses: + '200': + description: Object metadata. + content: + application/json: + schema: + type: object + properties: + key: { type: string } + name: { type: string } + size: { type: integer } + contentType: { type: string } + lastModified: { type: string, nullable: true } + etag: { type: string } + '404': + description: Object not found. + /api/objects/content: + get: + tags: [Objects API] + operationId: getObjectContentApi + summary: Download object content (Dashboard Session Auth or Ticket) + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + - in: query + name: ticket + schema: { type: string } + responses: + '200': + description: Binary object data. + '206': + description: Partial content. + '404': + description: Object not found. + put: + tags: [Objects API] + operationId: putObjectContentApi + summary: Direct single-part upload (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + responses: + '200': + description: Object uploaded. + content: + application/json: + schema: + type: object + properties: + key: { type: string } + etag: { type: string } + size: { type: integer } + /api/objects/folder: + post: + tags: [Objects API] + operationId: createFolderApi + summary: Create a folder (Dashboard Session Auth) + security: + - sessionBearer: [] + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [bucket, prefix] + properties: + bucket: { type: string } + prefix: { type: string } + responses: + '201': + description: Folder created. + '409': + description: Folder already exists. + delete: + tags: [Objects API] + operationId: deleteFolderApi + summary: Delete a folder (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: prefix + required: true + schema: { type: string } + - in: query + name: recursive + schema: { type: string, enum: ["0", "1"] } + responses: + '204': + description: Folder deleted. + '409': + description: Folder is not empty. + /api/objects/download-ticket: + post: + tags: [Objects API] + operationId: createDownloadTicket + summary: Create a temporary single-use download ticket (Dashboard Session Auth) + security: + - sessionBearer: [] + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [bucket, key] + properties: + bucket: { type: string } + key: { type: string } + responses: + '201': + description: Download ticket created. + content: + application/json: + schema: + type: object + properties: + ticket: { type: string } + downloadUrl: { type: string } + expiresIn: { type: integer } + /api/objects/uploads: + post: + tags: [Objects API] + operationId: initiateMultipartApi + summary: Initiate a multipart upload (Dashboard Session Auth) + security: + - sessionBearer: [] + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [bucket, key] + properties: + bucket: { type: string } + key: { type: string } + contentType: { type: string } + responses: + '201': + description: Upload initiated. + content: + application/json: + schema: + type: object + properties: + uploadId: { type: string } + bucket: { type: string } + key: { type: string } + partSize: { type: integer } + delete: + tags: [Objects API] + operationId: abortMultipartApi + summary: Abort a multipart upload (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + - in: query + name: uploadId + required: true + schema: { type: string } + responses: + '204': + description: Upload aborted. + /api/objects/uploads/part: + put: + tags: [Objects API] + operationId: uploadPartApi + summary: Upload a part (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + - in: query + name: uploadId + required: true + schema: { type: string } + - in: query + name: partNumber + required: true + schema: { type: integer } + responses: + '200': + description: Part uploaded. + content: + application/json: + schema: + type: object + properties: + partNumber: { type: integer } + etag: { type: string } + /api/objects/uploads/complete: + post: + tags: [Objects API] + operationId: completeMultipartApi + summary: Complete a multipart upload (Dashboard Session Auth) + security: + - sessionBearer: [] + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [bucket, key, uploadId, parts] + properties: + bucket: { type: string } + key: { type: string } + uploadId: { type: string } + parts: + type: array + items: + type: object + required: [partNumber, etag] + properties: + partNumber: { type: integer } + etag: { type: string } + totalSize: { type: integer } + responses: + '200': + description: Upload completed. + content: + application/json: + schema: + type: object + properties: + key: { type: string } + etag: { type: string } + /api/objects/uploads/parts: + get: + tags: [Objects API] + operationId: listPartsApi + summary: List uploaded parts (Dashboard Session Auth) + security: + - sessionBearer: [] + parameters: + - in: query + name: bucket + required: true + schema: { type: string } + - in: query + name: key + required: true + schema: { type: string } + - in: query + name: uploadId + required: true + schema: { type: string } + responses: + '200': + description: Parts list. /{bucket}: parameters: - $ref: '#/components/parameters/Bucket' diff --git a/src/cors.ts b/src/cors.ts index 1c4e357..833ee81 100644 --- a/src/cors.ts +++ b/src/cors.ts @@ -2,7 +2,7 @@ import type { Env } from "./types"; const ALLOWED_METHODS = "GET, HEAD, PUT, POST, PATCH, DELETE, OPTIONS"; const DEFAULT_ALLOWED_HEADERS = "Authorization, Content-Type, Content-Length, Content-MD5, Content-Encoding, Range, x-amz-content-sha256, x-amz-copy-source, x-amz-date, x-amz-decoded-content-length, x-amz-mp-object-size, x-amz-security-token"; -const EXPOSED_HEADERS = "ETag, Content-Range, Content-Length, Last-Modified, Accept-Ranges, x-amz-request-id"; +const EXPOSED_HEADERS = "ETag, Content-Range, Content-Length, Last-Modified, Accept-Ranges, Content-Disposition, x-amz-request-id"; function configuredOrigins(env: Env): string[] { return (env.CORS_ALLOWED_ORIGINS ?? "") diff --git a/src/google-drive.ts b/src/google-drive.ts index 022b9b3..ff98ee9 100644 --- a/src/google-drive.ts +++ b/src/google-drive.ts @@ -20,10 +20,15 @@ function driveLiteral(value: string): string { return value.replace(/\\/g, "\\\\").replace(/'/g, "\\'"); } -function driveFilesUrl(q: string, fields: string): string { +function driveFilesUrl(q: string, fields: string, pageToken?: string): string { const url = new URL("https://www.googleapis.com/drive/v3/files"); url.searchParams.set("q", q); - url.searchParams.set("fields", fields); + url.searchParams.set("pageSize", "1000"); + const fieldsWithPageToken = fields.includes("nextPageToken") ? fields : `nextPageToken,${fields}`; + url.searchParams.set("fields", fieldsWithPageToken); + if (pageToken) { + url.searchParams.set("pageToken", pageToken); + } return url.toString(); } @@ -136,13 +141,23 @@ export async function getOrCreateFolder(accessToken: string, folderName: string, } export async function listFolderChildren(accessToken: string, folderId: string, fields = "files(id,name,mimeType,size,modifiedTime,createdTime,appProperties)"): Promise { - const listRes = await fetch(driveFilesUrl(`'${driveLiteral(folderId)}' in parents and trashed=false`, fields), { - headers: { Authorization: `Bearer ${accessToken}` }, - }); - if (!listRes.ok) throw new Error(`Drive list failed: ${await listRes.text()}`); + const allFiles: GoogleDriveFile[] = []; + let pageToken: string | undefined; - const data: GoogleDriveSearchResponse = await listRes.json(); - return data.files || []; + do { + const listRes = await fetch(driveFilesUrl(`'${driveLiteral(folderId)}' in parents and trashed=false`, fields, pageToken), { + headers: { Authorization: `Bearer ${accessToken}` }, + }); + if (!listRes.ok) throw new Error(`Drive list failed: ${await listRes.text()}`); + + const data: GoogleDriveSearchResponse = await listRes.json(); + if (data.files) { + allFiles.push(...data.files); + } + pageToken = data.nextPageToken; + } while (pageToken); + + return allFiles; } export async function folderHasChildren(accessToken: string, folderId: string): Promise { @@ -176,6 +191,31 @@ export async function updateDriveFile(accessToken: string, fileId: string, body: return (await res.json()) as GoogleDriveFile; } +/** Resolves an S3 object key to its existing parent folder ID without creating folders. Returns null if parent hierarchy doesn't exist. */ +export async function resolvePathToExistingFolderAndFile(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<{ parentFolderId: string; fileName: string } | null> { + let currentFolderId = await findBucketFolderId(accessToken, bucket, env); + if (!currentFolderId) return null; + + const parts = objectKey.split("/").filter((p) => p); + if (parts.length === 0) { + throw new Error("Invalid object key"); + } + + const fileName = parts[parts.length - 1]; + const directories = parts.slice(0, -1); + + for (const dir of directories) { + const nextFolderId = await findFolderId(accessToken, dir, currentFolderId); + if (!nextFolderId) return null; + currentFolderId = nextFolderId; + } + + return { + parentFolderId: currentFolderId, + fileName: fileName, + }; +} + /** Resolves an S3 object key to its parent folder ID, creating the directory hierarchy as needed. */ export async function resolvePathToFolderAndFile(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<{ parentFolderId: string; fileName: string }> { let currentFolderId = await getBucketFolderId(accessToken, bucket, env); @@ -265,7 +305,11 @@ export async function findFileInFolder(accessToken: string, folderId: string, fi } export async function streamDownloadFromDrive(accessToken: string, bucket: string, objectKey: string, env: Env, range?: string): Promise { - const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env); + const resolved = await resolvePathToExistingFolderAndFile(accessToken, bucket, objectKey, env); + if (!resolved) { + throw new Error("File not found"); + } + const { parentFolderId, fileName } = resolved; const file = await findFileInFolder(accessToken, parentFolderId, fileName); if (!file) { @@ -305,7 +349,11 @@ export async function streamDownloadFromDrive(accessToken: string, bucket: strin } export async function deleteFromDrive(accessToken: string, bucket: string, objectKey: string, env: Env): Promise { - const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env); + const resolved = await resolvePathToExistingFolderAndFile(accessToken, bucket, objectKey, env); + if (!resolved) { + throw new Error("File not found"); + } + const { parentFolderId, fileName } = resolved; const file = await findFileInFolder(accessToken, parentFolderId, fileName); if (!file) { @@ -323,7 +371,11 @@ export async function deleteFromDrive(accessToken: string, bucket: string, objec } export async function getFileMetadata(accessToken: string, bucket: string, objectKey: string, env: Env): Promise { - const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env); + const resolved = await resolvePathToExistingFolderAndFile(accessToken, bucket, objectKey, env); + if (!resolved) { + throw new Error("File not found"); + } + const { parentFolderId, fileName } = resolved; const file = await findFileInFolder(accessToken, parentFolderId, fileName); if (!file) { @@ -340,13 +392,23 @@ export async function getFileMetadata(accessToken: string, bucket: string, objec } async function listChildren(accessToken: string, folderId: string): Promise { - const listRes = await fetch(driveFilesUrl(`'${driveLiteral(folderId)}' in parents and trashed=false`, "files(id,name,mimeType,size,modifiedTime,md5Checksum)"), { - headers: { Authorization: `Bearer ${accessToken}` }, - }); - if (!listRes.ok) throw new Error(`Drive list failed: ${await listRes.text()}`); + const allFiles: GoogleDriveFile[] = []; + let pageToken: string | undefined; - const data: GoogleDriveSearchResponse = await listRes.json(); - return data.files || []; + do { + const listRes = await fetch(driveFilesUrl(`'${driveLiteral(folderId)}' in parents and trashed=false`, "files(id,name,mimeType,size,modifiedTime,md5Checksum)", pageToken), { + headers: { Authorization: `Bearer ${accessToken}` }, + }); + if (!listRes.ok) throw new Error(`Drive list failed: ${await listRes.text()}`); + + const data: GoogleDriveSearchResponse = await listRes.json(); + if (data.files) { + allFiles.push(...data.files); + } + pageToken = data.nextPageToken; + } while (pageToken); + + return allFiles; } /** Splits an S3 prefix into the directory portion (real Drive folder path) and the partial name filter for the final segment. */ diff --git a/src/multipart-core.ts b/src/multipart-core.ts new file mode 100644 index 0000000..eaeef3c --- /dev/null +++ b/src/multipart-core.ts @@ -0,0 +1,188 @@ +import { decodedContentLength, isAwsChunked, pumpBody } from "./aws-chunked"; +import { createSession, nextDriveOffset } from "./drive-resumable"; +import { findFileInFolder, resolvePathToFolderAndFile } from "./google-drive"; +import type { DriveUploadResult, Env } from "./types"; + +export const MAX_COMPLETE_XML = 4 * 1024 * 1024; +export const DEFAULT_PART_SIZE = 8 * 1024 * 1024; // 8 MiB aligned to 256 KiB + +export type CoreError = + | { kind: "error"; code: "NoSuchUpload"; status: 404; message?: string } + | { kind: "error"; code: "InvalidArgument"; status: 400; message?: string } + | { kind: "error"; code: "NotImplemented"; status: 501; message?: string } + | { kind: "error"; code: "EntityTooLarge"; status: 400; message?: string } + | { kind: "error"; code: "MalformedXML"; status: 400; message?: string } + | { kind: "error"; code: "InvalidPart" | "InvalidPartOrder"; status: 400; message: string } + | { kind: "error"; code: "SlowDown"; status: 503; message?: string; retryAfter?: string } + | { kind: "error"; code: "InternalError"; status: 500; message: string }; + +export function etag(file: { id: string; md5Checksum?: string }): string { + return file.md5Checksum || file.id; +} + +export function multipartStub(env: Env, uploadId: string) { + return env.MPU.getByName(uploadId); +} + +export function encodeUploadId(bucket: string, key: string): string { + const bytes = crypto.getRandomValues(new Uint8Array(16)); + const random = btoa(String.fromCharCode(...bytes)) + .replace(/\+/g, "-") + .replace(/\//g, "_") + .replace(/=+$/, ""); + return `${random}.${encodeURIComponent(bucket)}.${encodeURIComponent(key)}`; +} + +export function uploadIdMatches(uploadId: string, bucket: string, key: string): boolean { + const parts = uploadId.split("."); + if (parts.length < 3) return false; + try { + return decodeURIComponent(parts[1]) === bucket && decodeURIComponent(parts.slice(2).join(".")) === key; + } catch { + return false; + } +} + +export function parsePositiveInt(value: string | null, fallback?: number): number | undefined { + if (value === null) return fallback; + if (!/^\d+$/.test(value)) return undefined; + const number = Number(value); + return Number.isSafeInteger(number) ? number : undefined; +} + +export async function drainBody(body: ReadableStream | null): Promise { + if (body) await body.pipeTo(new WritableStream()); +} + +export async function partEtag(uploadId: string, partNumber: number, partLen: number, fileOffset: number): Promise { + const value = new TextEncoder().encode(`${uploadId}:${partNumber}:${partLen}:${fileOffset}`); + const digest = new Uint8Array(await crypto.subtle.digest("SHA-256", value)).subarray(0, 16); + return Array.from(digest, (byte) => byte.toString(16).padStart(2, "0")).join(""); +} + +export async function multipartObjectEtag(metadata: DriveUploadResult, partEtags: string[], style: Env["ETAG_STYLE"]): Promise { + if (style !== "multipart" || partEtags.length === 0) return etag(metadata); + const bytes = new Uint8Array(partEtags.length * 16); + for (let index = 0; index < partEtags.length; index++) { + const value = partEtags[index]; + for (let byte = 0; byte < 16; byte++) bytes[index * 16 + byte] = Number.parseInt(value.slice(byte * 2, byte * 2 + 2), 16); + } + const digest = await crypto.subtle.digest("MD5", bytes); + return `${Array.from(new Uint8Array(digest), (byte) => byte.toString(16).padStart(2, "0")).join("")}-${partEtags.length}`; +} + +export async function createMultipartUpload( + env: Env, + accessToken: string, + bucket: string, + key: string, + mimeType: string, +): Promise<{ uploadId: string; partSize: number } | CoreError> { + if (env.ALLOW_MULTIPART !== "true") return { kind: "error", code: "NotImplemented", status: 501 }; + const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, key, env); + const existing = await findFileInFolder(accessToken, parentFolderId, fileName); + const uploadUrl = await createSession(accessToken, { name: fileName, parents: [parentFolderId], mimeType, existingFileId: existing?.id }); + const uploadId = encodeUploadId(bucket, key); + const initialized = await multipartStub(env, uploadId).init({ uploadUrl, bucket, key, mimeType, parentFolderId, fileName, existingFileId: existing?.id }); + if (!initialized) throw new Error("Failed to initialize multipart state"); + return { uploadId, partSize: DEFAULT_PART_SIZE }; +} + +export async function uploadPartCore( + request: Request, + env: Env, + accessToken: string, + bucket: string, + key: string, + uploadId: string, + partNumber: number | undefined, +): Promise<{ etag: string } | CoreError> { + if (!uploadIdMatches(uploadId, bucket, key)) return { kind: "error", code: "NoSuchUpload", status: 404 }; + if (!partNumber || partNumber > 10_000) return { kind: "error", code: "InvalidArgument", status: 400, message: "partNumber must be between 1 and 10000" }; + if (request.headers.has("x-amz-copy-source")) return { kind: "error", code: "NotImplemented", status: 501 }; + const length = decodedContentLength(request); + if (length === undefined) return { kind: "error", code: "InvalidArgument", status: 400, message: "Content-Length or x-amz-decoded-content-length is required" }; + + const stub = multipartStub(env, uploadId); + const requestId = crypto.randomUUID(); + const onAbort = () => void stub.cancelWaiter(requestId); + request.signal.addEventListener("abort", onAbort, { once: true }); + const lease = await stub.beginPart(requestId, partNumber, length); + request.signal.removeEventListener("abort", onAbort); + if (lease.kind === "slowdown") return { kind: "error", code: "SlowDown", status: 503, retryAfter: "1" }; + if (lease.kind === "error") return { kind: "error", code: lease.code, status: lease.code === "NoSuchUpload" ? 404 : 500, message: lease.message }; + if (lease.kind === "committed") { + await drainBody(request.body); + return { etag: lease.etag }; + } + + try { + const fixed = lease.sendLen === 0 ? null : new FixedLengthStream(lease.sendLen, { highWaterMark: 1 << 20 }); + const drivePromise = fixed + ? fetch(lease.uploadUrl, { + method: "PUT", + headers: { + Authorization: `Bearer ${accessToken}`, + "Content-Length": lease.sendLen.toString(), + "Content-Range": `bytes ${lease.driveOffset}-${lease.driveOffset + lease.sendLen - 1}/*`, + }, + body: fixed.readable, + duplex: "half", + } as RequestInit) + : null; + const sink = fixed?.writable ?? new WritableStream(); + const pumped = await pumpBody(request.body, sink.getWriter(), { + awsChunked: isAwsChunked(request), + expectedLength: length, + skipBytes: lease.skipBytes, + maxBytes: lease.sendLen, + prefix: lease.carry, + }); + const driveResponse = drivePromise ? await drivePromise : null; + if (driveResponse && driveResponse.status !== 308) throw new Error(`Drive part upload returned ${driveResponse.status}: ${await driveResponse.text()}`); + const newOffset = driveResponse ? nextDriveOffset(driveResponse) : lease.driveOffset; + const value = await partEtag(uploadId, partNumber, length, lease.driveOffset + lease.carry.byteLength); + if (!(await stub.endPart(partNumber, newOffset, pumped.tail, value, length))) throw new Error("Multipart lease was lost before commit"); + return { etag: value }; + } catch (error) { + await stub.failPart(partNumber); + throw error; + } +} + +export async function completeMultipartUpload( + env: Env, + bucket: string, + key: string, + uploadId: string, + parts: Array<{ partNumber: number; etag: string }>, + expectedTotal?: number, +): Promise<{ etag: string; key: string } | CoreError> { + if (!uploadIdMatches(uploadId, bucket, key)) return { kind: "error", code: "NoSuchUpload", status: 404 }; + const result = await multipartStub(env, uploadId).complete(parts, expectedTotal); + if (result.kind === "error") return { kind: "error", code: result.code, status: result.code === "NoSuchUpload" ? 404 : 400, message: result.message }; + const value = await multipartObjectEtag(result.metadata, result.partEtags, env.ETAG_STYLE); + return { etag: value, key }; +} + +export async function abortMultipartUpload(env: Env, bucket: string, key: string, uploadId: string): Promise { + if (!uploadIdMatches(uploadId, bucket, key)) return { kind: "error", code: "NoSuchUpload", status: 404 }; + const success = await multipartStub(env, uploadId).abort(); + if (!success) return { kind: "error", code: "NoSuchUpload", status: 404 }; + return true; +} + +export async function listMultipartParts( + env: Env, + bucket: string, + key: string, + uploadId: string, + marker = 0, + maxParts = 1000, +): Promise { + if (!uploadIdMatches(uploadId, bucket, key)) return { kind: "error", code: "NoSuchUpload", status: 404 }; + if (marker < 0 || maxParts < 0 || maxParts > 1000) return { kind: "error", code: "InvalidArgument", status: 400 }; + const result = await multipartStub(env, uploadId).listParts(marker, maxParts); + if (!result) return { kind: "error", code: "NoSuchUpload", status: 404 }; + return result; +} diff --git a/src/object-api.ts b/src/object-api.ts new file mode 100644 index 0000000..b8e21df --- /dev/null +++ b/src/object-api.ts @@ -0,0 +1,546 @@ +import { jsonResponse } from "./auth-api"; +import { findBucketRecord } from "./bucket-registry"; +import { + deleteFromDrive, + findFileInFolder, + findFolderId, + getAccessToken, + getFileMetadata, + getOrCreateFolder, + listObjects, + resolvePathToExistingFolderAndFile, + resolvePathToFolderAndFile, + streamDownloadFromDrive, + streamUploadToDrive, + updateDriveFile, +} from "./google-drive"; +import { + abortMultipartUpload, + completeMultipartUpload, + createMultipartUpload, + etag, + listMultipartParts, + parsePositiveInt, + uploadPartCore, +} from "./multipart-core"; +import type { Env } from "./types"; + +export async function invalidateObjectCaches(env: Env, bucket: string, parentId?: string, name?: string): Promise { + await env.FOLDER_CACHE.delete(`bucket-stats:${bucket}`); + if (parentId && name) { + await env.FOLDER_CACHE.delete(`${parentId}/${name}`); + } +} + +async function sha256Hex(text: string): Promise { + const data = new TextEncoder().encode(text); + const hash = await crypto.subtle.digest("SHA-256", data); + return Array.from(new Uint8Array(hash), (b) => b.toString(16).padStart(2, "0")).join(""); +} + +export async function handleTicketDownload(request: Request, env: Env): Promise { + const url = new URL(request.url); + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + const ticket = url.searchParams.get("ticket") ?? ""; + + if (!bucket || !key || !ticket) { + return jsonResponse({ message: "bucket, key, and ticket are required" }, 400); + } + + const ticketHash = await sha256Hex(ticket); + const ticketKey = `dl:${ticketHash}`; + const stored = await env.AUTH_KV.get(ticketKey); + + if (!stored) { + return jsonResponse({ message: "Invalid or expired download ticket" }, 403); + } + + try { + const payload = JSON.parse(stored) as { bucket: string; key: string }; + if (payload.bucket !== bucket || payload.key !== key) { + return jsonResponse({ message: "Ticket does not match bucket and key" }, 403); + } + } catch { + return jsonResponse({ message: "Invalid download ticket payload" }, 403); + } + + // Single use ticket + await env.AUTH_KV.delete(ticketKey); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) { + return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + } + + try { + const accessToken = await getAccessToken(env); + const range = request.headers.get("Range") ?? undefined; + const result = await streamDownloadFromDrive(accessToken, bucket, key, env, range); + + const filename = key.split("/").filter(Boolean).pop() || "download"; + const headers = new Headers({ + "Content-Type": result.contentType, + "Content-Length": result.contentLength || result.size.toString(), + ETag: `"${etag(result)}"`, + "Accept-Ranges": "bytes", + "Content-Disposition": `attachment; filename="${encodeURIComponent(filename)}"`, + }); + + if (result.contentRange) headers.set("Content-Range", result.contentRange); + if (result.modifiedTime) headers.set("Last-Modified", new Date(result.modifiedTime).toUTCString()); + + return new Response(result.body, { + status: result.status, + headers, + }); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + if (message.includes("not found") || message.includes("File not found")) { + return jsonResponse({ message: "Object not found" }, 404); + } + return jsonResponse({ message: "Failed to download object" }, 500); + } +} + +export async function handleObjectRoutes(request: Request, env: Env, subSegments: string[]): Promise { + const url = new URL(request.url); + const method = request.method; + const path = subSegments.join("/"); + + // 1. GET /api/objects (List objects) OR DELETE /api/objects (Delete object) + if (path === "") { + const bucket = url.searchParams.get("bucket") ?? ""; + if (!bucket) return jsonResponse({ message: "bucket query parameter is required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + if (method === "GET") { + const prefix = url.searchParams.get("prefix") ?? ""; + const delimiter = url.searchParams.get("delimiter") ?? undefined; + try { + const accessToken = await getAccessToken(env); + const { contents, commonPrefixes, truncated } = await listObjects(accessToken, bucket, prefix, env, delimiter); + return jsonResponse( + { + bucket, + prefix, + delimiter: delimiter ?? null, + folders: commonPrefixes.map((p) => { + const trimmed = p.endsWith("/") ? p.slice(0, -1) : p; + const name = trimmed.split("/").pop() || trimmed; + return { prefix: p, name }; + }), + objects: contents.map((obj) => { + const name = obj.key.split("/").pop() || obj.key; + return { + key: obj.key, + name, + size: parseInt(obj.size || "0", 10), + contentType: obj.mimeType || "application/octet-stream", + lastModified: obj.modifiedTime || null, + etag: etag(obj), + }; + }), + truncated, + }, + 200, + ); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + if (method === "DELETE") { + const key = url.searchParams.get("key") ?? ""; + if (!key) return jsonResponse({ message: "key query parameter is required" }, 400); + + try { + const accessToken = await getAccessToken(env); + await deleteFromDrive(accessToken, bucket, key, env); + await invalidateObjectCaches(env, bucket); + return new Response(null, { status: 204 }); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + if (message.includes("not found") || message.includes("File not found")) { + return jsonResponse({ message: "Object not found" }, 404); + } + return jsonResponse({ message }, 500); + } + } + + return jsonResponse({ message: "Method Not Allowed" }, 405); + } + + // 2. GET /api/objects/metadata?bucket=&key= + if (path === "metadata") { + if (method !== "GET") return jsonResponse({ message: "Method Not Allowed" }, 405); + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + if (!bucket || !key) return jsonResponse({ message: "bucket and key query parameters are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const accessToken = await getAccessToken(env); + const meta = await getFileMetadata(accessToken, bucket, key, env); + const name = key.split("/").pop() || key; + return jsonResponse( + { + key, + name, + size: meta.size, + contentType: meta.mimeType, + lastModified: meta.modifiedTime, + etag: etag(meta), + }, + 200, + ); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + if (message.includes("not found") || message.includes("File not found")) { + return jsonResponse({ message: "Object not found" }, 404); + } + return jsonResponse({ message }, 500); + } + } + + // 3. GET /api/objects/content?bucket=&key= (Stream download) OR PUT /api/objects/content?bucket=&key= (Direct upload) + if (path === "content") { + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + if (!bucket || !key) return jsonResponse({ message: "bucket and key query parameters are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + if (method === "GET") { + try { + const accessToken = await getAccessToken(env); + const range = request.headers.get("Range") ?? undefined; + const result = await streamDownloadFromDrive(accessToken, bucket, key, env, range); + + const filename = key.split("/").filter(Boolean).pop() || "download"; + const headers = new Headers({ + "Content-Type": result.contentType, + "Content-Length": result.contentLength || result.size.toString(), + ETag: `"${etag(result)}"`, + "Accept-Ranges": "bytes", + "Content-Disposition": `attachment; filename="${encodeURIComponent(filename)}"`, + }); + + if (result.contentRange) headers.set("Content-Range", result.contentRange); + if (result.modifiedTime) headers.set("Last-Modified", new Date(result.modifiedTime).toUTCString()); + + return new Response(result.body, { + status: result.status, + headers, + }); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + if (message.includes("not found") || message.includes("File not found")) { + return jsonResponse({ message: "Object not found" }, 404); + } + return jsonResponse({ message }, 500); + } + } + + if (method === "PUT") { + try { + const accessToken = await getAccessToken(env); + const contentType = request.headers.get("Content-Type") || "application/octet-stream"; + const result = await streamUploadToDrive(accessToken, request, bucket, key, contentType, env); + await invalidateObjectCaches(env, bucket); + return jsonResponse( + { + key, + etag: etag(result), + size: parseInt(result.size || "0", 10), + }, + 200, + ); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + return jsonResponse({ message: "Method Not Allowed" }, 405); + } + + // 4. POST /api/objects/folder OR DELETE /api/objects/folder + if (path === "folder") { + if (method === "POST") { + let body: { bucket?: string; prefix?: string }; + try { + body = (await request.json()) as { bucket?: string; prefix?: string }; + } catch { + return jsonResponse({ message: "Invalid JSON body" }, 400); + } + const bucket = body.bucket ?? ""; + let prefix = body.prefix ?? ""; + if (!bucket || !prefix) return jsonResponse({ message: "bucket and prefix are required" }, 400); + + if (!prefix.endsWith("/")) prefix = `${prefix}/`; + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const accessToken = await getAccessToken(env); + // Check if folder already exists + const existingResolved = await resolvePathToExistingFolderAndFile(accessToken, bucket, `${prefix}dummy`, env); + if (existingResolved) { + return jsonResponse({ message: "Folder already exists" }, 409); + } + + // Create folder hierarchy + const resolved = await resolvePathToFolderAndFile(accessToken, bucket, `${prefix}dummy`, env); + await invalidateObjectCaches(env, bucket); + return jsonResponse({ prefix }, 201); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + if (method === "DELETE") { + const bucket = url.searchParams.get("bucket") ?? ""; + let prefix = url.searchParams.get("prefix") ?? ""; + const recursive = url.searchParams.get("recursive") === "1"; + if (!bucket || !prefix) return jsonResponse({ message: "bucket and prefix query parameters are required" }, 400); + + if (!prefix.endsWith("/")) prefix = `${prefix}/`; + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const accessToken = await getAccessToken(env); + // Walk to the target folder + const parts = prefix.split("/").filter(Boolean); + if (parts.length === 0) { + return jsonResponse({ message: "Cannot delete bucket root via folder delete" }, 400); + } + + let currentFolderId: string | null = bucketRecord.folderId; + let parentOfTarget: string | null = null; + const targetFolderName = parts[parts.length - 1]; + + for (let i = 0; i < parts.length; i++) { + const segment = parts[i]; + if (!currentFolderId) break; + if (i === parts.length - 1) { + parentOfTarget = currentFolderId; + } + currentFolderId = await findFolderId(accessToken, segment, currentFolderId); + } + + if (!currentFolderId) { + return jsonResponse({ message: "Folder not found" }, 404); + } + + if (!recursive) { + const { contents, commonPrefixes } = await listObjects(accessToken, bucket, prefix, env, "/"); + if (contents.length > 0 || commonPrefixes.length > 0) { + return jsonResponse({ message: "Folder is not empty" }, 409); + } + } + + // Update folder to trashed=true + await updateDriveFile(accessToken, currentFolderId, { trashed: true }); + await invalidateObjectCaches(env, bucket, parentOfTarget ?? undefined, targetFolderName); + return new Response(null, { status: 204 }); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + return jsonResponse({ message: "Method Not Allowed" }, 405); + } + + // 5. Download ticket: POST /api/objects/download-ticket { bucket, key } + if (path === "download-ticket") { + if (method !== "POST") return jsonResponse({ message: "Method Not Allowed" }, 405); + let body: { bucket?: string; key?: string }; + try { + body = (await request.json()) as { bucket?: string; key?: string }; + } catch { + return jsonResponse({ message: "Invalid JSON body" }, 400); + } + + const bucket = body.bucket ?? ""; + const key = body.key ?? ""; + if (!bucket || !key) return jsonResponse({ message: "bucket and key are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + const tokenBytes = crypto.getRandomValues(new Uint8Array(24)); + const ticket = btoa(String.fromCharCode(...tokenBytes)) + .replace(/\+/g, "-") + .replace(/\//g, "_") + .replace(/=+$/, ""); + + const ticketHash = await sha256Hex(ticket); + await env.AUTH_KV.put(`dl:${ticketHash}`, JSON.stringify({ bucket, key }), { expirationTtl: 120 }); + + const downloadUrl = `/api/objects/content?bucket=${encodeURIComponent(bucket)}&key=${encodeURIComponent(key)}&ticket=${ticket}`; + return jsonResponse({ ticket, downloadUrl, expiresIn: 120 }, 201); + } + + // 6. Multipart upload routes + // POST /api/objects/uploads -> initiate + // DELETE /api/objects/uploads?bucket=&key=&uploadId= -> abort + if (path === "uploads") { + if (method === "POST") { + let body: { bucket?: string; key?: string; contentType?: string }; + try { + body = (await request.json()) as { bucket?: string; key?: string; contentType?: string }; + } catch { + return jsonResponse({ message: "Invalid JSON body" }, 400); + } + const bucket = body.bucket ?? ""; + const key = body.key ?? ""; + const contentType = body.contentType || "application/octet-stream"; + if (!bucket || !key) return jsonResponse({ message: "bucket and key are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const accessToken = await getAccessToken(env); + const result = await createMultipartUpload(env, accessToken, bucket, key, contentType); + if ("kind" in result && result.kind === "error") { + return jsonResponse({ message: result.message || result.code }, result.status); + } + return jsonResponse( + { + uploadId: (result as { uploadId: string }).uploadId, + bucket, + key, + partSize: (result as { partSize: number }).partSize, + }, + 201, + ); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + if (method === "DELETE") { + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + const uploadId = url.searchParams.get("uploadId") ?? ""; + if (!bucket || !key || !uploadId) return jsonResponse({ message: "bucket, key, and uploadId query parameters are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + const result = await abortMultipartUpload(env, bucket, key, uploadId); + if (result !== true) { + return jsonResponse({ message: "Upload not found" }, 404); + } + return new Response(null, { status: 204 }); + } + + return jsonResponse({ message: "Method Not Allowed" }, 405); + } + + // PUT /api/objects/uploads/part?bucket=&key=&uploadId=&partNumber= + if (path === "uploads/part") { + if (method !== "PUT") return jsonResponse({ message: "Method Not Allowed" }, 405); + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + const uploadId = url.searchParams.get("uploadId") ?? ""; + const partNumber = parsePositiveInt(url.searchParams.get("partNumber")); + if (!bucket || !key || !uploadId || !partNumber) return jsonResponse({ message: "bucket, key, uploadId, and partNumber query parameters are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const accessToken = await getAccessToken(env); + const result = await uploadPartCore(request, env, accessToken, bucket, key, uploadId, partNumber); + if ("kind" in result && result.kind === "error") { + if (result.code === "SlowDown") { + return new Response(JSON.stringify({ message: "SlowDown" }), { + status: 503, + headers: { "Content-Type": "application/json", "Retry-After": result.retryAfter ?? "1" }, + }); + } + return jsonResponse({ message: result.message || result.code }, result.status); + } + return jsonResponse({ partNumber, etag: (result as { etag: string }).etag }, 200); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + // POST /api/objects/uploads/complete + if (path === "uploads/complete") { + if (method !== "POST") return jsonResponse({ message: "Method Not Allowed" }, 405); + let body: { bucket?: string; key?: string; uploadId?: string; parts?: Array<{ partNumber: number; etag: string }>; totalSize?: number }; + try { + body = (await request.json()) as typeof body; + } catch { + return jsonResponse({ message: "Invalid JSON body" }, 400); + } + + const bucket = body.bucket ?? ""; + const key = body.key ?? ""; + const uploadId = body.uploadId ?? ""; + const parts = body.parts ?? []; + if (!bucket || !key || !uploadId || !Array.isArray(parts)) return jsonResponse({ message: "bucket, key, uploadId, and parts array are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + try { + const result = await completeMultipartUpload(env, bucket, key, uploadId, parts, body.totalSize); + if ("kind" in result && result.kind === "error") { + return jsonResponse({ message: result.message || result.code }, result.status); + } + await invalidateObjectCaches(env, bucket); + return jsonResponse({ key: (result as { key: string }).key, etag: (result as { etag: string }).etag }, 200); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + // GET /api/objects/uploads/parts?bucket=&key=&uploadId= + if (path === "uploads/parts") { + if (method !== "GET") return jsonResponse({ message: "Method Not Allowed" }, 405); + const bucket = url.searchParams.get("bucket") ?? ""; + const key = url.searchParams.get("key") ?? ""; + const uploadId = url.searchParams.get("uploadId") ?? ""; + if (!bucket || !key || !uploadId) return jsonResponse({ message: "bucket, key, and uploadId query parameters are required" }, 400); + + const bucketRecord = await findBucketRecord(env, bucket); + if (!bucketRecord) return jsonResponse({ message: `Bucket '${bucket}' not found` }, 404); + + const marker = parsePositiveInt(url.searchParams.get("marker"), 0); + const maxParts = parsePositiveInt(url.searchParams.get("maxParts"), 1000); + if (marker === undefined || maxParts === undefined || maxParts > 1000) return jsonResponse({ message: "Invalid pagination parameters" }, 400); + + try { + const result = await listMultipartParts(env, bucket, key, uploadId, marker, maxParts); + if ("kind" in result && result.kind === "error") { + return jsonResponse({ message: result.message || result.code }, result.status); + } + return jsonResponse(result, 200); + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + return jsonResponse({ message }, 500); + } + } + + return jsonResponse({ message: "Not Found" }, 404); +} diff --git a/src/router.ts b/src/router.ts index bc0e0ca..04b7469 100644 --- a/src/router.ts +++ b/src/router.ts @@ -1,142 +1,46 @@ -import { decodedContentLength, isAwsChunked, pumpBody } from "./aws-chunked"; -import { createSession, nextDriveOffset } from "./drive-resumable"; -import { deleteFromDrive, findFileInFolder, getFileMetadata, listObjects, resolvePathToFolderAndFile, streamDownloadFromDrive, streamUploadToDrive } from "./google-drive"; +import { + MAX_COMPLETE_XML, + abortMultipartUpload, + completeMultipartUpload, + createMultipartUpload, + etag, + listMultipartParts, + parsePositiveInt, + uploadPartCore, +} from "./multipart-core"; +import { deleteFromDrive, getFileMetadata, listObjects, streamDownloadFromDrive, streamUploadToDrive } from "./google-drive"; import { S3Exception, s3Error } from "./s3-errors"; import { completeMultipartUploadResult, generateListBucketResult, initiateMultipartUploadResult, listMultipartUploadsResult, listPartsResult, parseCompleteMultipartUpload } from "./s3-xml"; -import type { DriveUploadResult, Env } from "./types"; - -const MAX_COMPLETE_XML = 4 * 1024 * 1024; - -function etag(file: { id: string; md5Checksum?: string }): string { - return file.md5Checksum || file.id; -} +import type { Env } from "./types"; function xmlResponse(body: string, status = 200): Response { return new Response(body, { status, headers: { "Content-Type": "application/xml", "Cache-Control": "no-transform" } }); } -function multipartStub(env: Env, uploadId: string) { - return env.MPU.getByName(uploadId); -} - -function encodeUploadId(bucket: string, key: string): string { - const bytes = crypto.getRandomValues(new Uint8Array(16)); - const random = btoa(String.fromCharCode(...bytes)) - .replace(/\+/g, "-") - .replace(/\//g, "_") - .replace(/=+$/, ""); - return `${random}.${encodeURIComponent(bucket)}.${encodeURIComponent(key)}`; -} - -function uploadIdMatches(uploadId: string, bucket: string, key: string): boolean { - const parts = uploadId.split("."); - if (parts.length < 3) return false; - try { - return decodeURIComponent(parts[1]) === bucket && decodeURIComponent(parts.slice(2).join(".")) === key; - } catch { - return false; - } -} - -function parsePositiveInt(value: string | null, fallback?: number): number | undefined { - if (value === null) return fallback; - if (!/^\d+$/.test(value)) return undefined; - const number = Number(value); - return Number.isSafeInteger(number) ? number : undefined; -} - -async function drainBody(body: ReadableStream | null): Promise { - if (body) await body.pipeTo(new WritableStream()); -} - -async function partEtag(uploadId: string, partNumber: number, partLen: number, fileOffset: number): Promise { - const value = new TextEncoder().encode(`${uploadId}:${partNumber}:${partLen}:${fileOffset}`); - const digest = new Uint8Array(await crypto.subtle.digest("SHA-256", value)).subarray(0, 16); - return Array.from(digest, (byte) => byte.toString(16).padStart(2, "0")).join(""); -} - -async function multipartObjectEtag(metadata: DriveUploadResult, partEtags: string[], style: Env["ETAG_STYLE"]): Promise { - if (style !== "multipart" || partEtags.length === 0) return etag(metadata); - const bytes = new Uint8Array(partEtags.length * 16); - for (let index = 0; index < partEtags.length; index++) { - const value = partEtags[index]; - for (let byte = 0; byte < 16; byte++) bytes[index * 16 + byte] = Number.parseInt(value.slice(byte * 2, byte * 2 + 2), 16); - } - const digest = await crypto.subtle.digest("MD5", bytes); - return `${Array.from(new Uint8Array(digest), (byte) => byte.toString(16).padStart(2, "0")).join("")}-${partEtags.length}`; -} - async function createMultipart(request: Request, env: Env, accessToken: string, bucket: string, key: string): Promise { - if (env.ALLOW_MULTIPART !== "true") return s3Error("NotImplemented", 501, undefined, `/${bucket}/${key}`); const mimeType = request.headers.get("Content-Type") || "application/octet-stream"; - const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, key, env); - const existing = await findFileInFolder(accessToken, parentFolderId, fileName); - const uploadUrl = await createSession(accessToken, { name: fileName, parents: [parentFolderId], mimeType, existingFileId: existing?.id }); - const uploadId = encodeUploadId(bucket, key); - const initialized = await multipartStub(env, uploadId).init({ uploadUrl, bucket, key, mimeType, parentFolderId, fileName, existingFileId: existing?.id }); - if (!initialized) throw new Error("Failed to initialize multipart state"); - return xmlResponse(initiateMultipartUploadResult(bucket, key, uploadId)); + const result = await createMultipartUpload(env, accessToken, bucket, key, mimeType); + if ("kind" in result && result.kind === "error") { + return s3Error(result.code, result.status, result.message, `/${bucket}/${key}`); + } + return xmlResponse(initiateMultipartUploadResult(bucket, key, (result as { uploadId: string }).uploadId)); } async function uploadPart(request: Request, env: Env, accessToken: string, bucket: string, key: string, url: URL): Promise { const uploadId = url.searchParams.get("uploadId") ?? ""; const partNumber = parsePositiveInt(url.searchParams.get("partNumber")); - if (!uploadIdMatches(uploadId, bucket, key)) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); - if (!partNumber || partNumber > 10_000) return s3Error("InvalidArgument", 400, "partNumber must be between 1 and 10000", `/${bucket}/${key}`); - if (request.headers.has("x-amz-copy-source")) return s3Error("NotImplemented", 501, undefined, `/${bucket}/${key}`); - const length = decodedContentLength(request); - if (length === undefined) return s3Error("InvalidArgument", 400, "Content-Length or x-amz-decoded-content-length is required", `/${bucket}/${key}`); - - const stub = multipartStub(env, uploadId); - const requestId = crypto.randomUUID(); - const onAbort = () => void stub.cancelWaiter(requestId); - request.signal.addEventListener("abort", onAbort, { once: true }); - const lease = await stub.beginPart(requestId, partNumber, length); - request.signal.removeEventListener("abort", onAbort); - if (lease.kind === "slowdown") return s3Error("SlowDown", 503, undefined, `/${bucket}/${key}`, false, { "Retry-After": "1" }); - if (lease.kind === "error") return s3Error(lease.code, lease.code === "NoSuchUpload" ? 404 : 500, lease.message, `/${bucket}/${key}`); - if (lease.kind === "committed") { - await drainBody(request.body); - return new Response(null, { status: 200, headers: { ETag: `"${lease.etag}"` } }); - } - - try { - const fixed = lease.sendLen === 0 ? null : new FixedLengthStream(lease.sendLen, { highWaterMark: 1 << 20 }); - const drivePromise = fixed - ? fetch(lease.uploadUrl, { - method: "PUT", - headers: { - Authorization: `Bearer ${accessToken}`, - "Content-Length": lease.sendLen.toString(), - "Content-Range": `bytes ${lease.driveOffset}-${lease.driveOffset + lease.sendLen - 1}/*`, - }, - body: fixed.readable, - duplex: "half", - } as RequestInit) - : null; - const sink = fixed?.writable ?? new WritableStream(); - const pumped = await pumpBody(request.body, sink.getWriter(), { - awsChunked: isAwsChunked(request), - expectedLength: length, - skipBytes: lease.skipBytes, - maxBytes: lease.sendLen, - prefix: lease.carry, - }); - const driveResponse = drivePromise ? await drivePromise : null; - if (driveResponse && driveResponse.status !== 308) throw new Error(`Drive part upload returned ${driveResponse.status}: ${await driveResponse.text()}`); - const newOffset = driveResponse ? nextDriveOffset(driveResponse) : lease.driveOffset; - const value = await partEtag(uploadId, partNumber, length, lease.driveOffset + lease.carry.byteLength); - if (!(await stub.endPart(partNumber, newOffset, pumped.tail, value, length))) throw new Error("Multipart lease was lost before commit"); - return new Response(null, { status: 200, headers: { ETag: `"${value}"` } }); - } catch (error) { - await stub.failPart(partNumber); - throw error; + const result = await uploadPartCore(request, env, accessToken, bucket, key, uploadId, partNumber); + if ("kind" in result && result.kind === "error") { + if (result.code === "SlowDown") { + return s3Error("SlowDown", 503, undefined, `/${bucket}/${key}`, false, { "Retry-After": result.retryAfter ?? "1" }); + } + return s3Error(result.code, result.status, result.message, `/${bucket}/${key}`); } + return new Response(null, { status: 200, headers: { ETag: `"${(result as { etag: string }).etag}"` } }); } async function completeMultipart(request: Request, env: Env, bucket: string, key: string, url: URL): Promise { const uploadId = url.searchParams.get("uploadId") ?? ""; - if (!uploadIdMatches(uploadId, bucket, key)) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); const contentLength = parsePositiveInt(request.headers.get("Content-Length")); if (contentLength !== undefined && contentLength > MAX_COMPLETE_XML) return s3Error("EntityTooLarge", 400, "Completion XML exceeds 4 MiB", `/${bucket}/${key}`); const text = await request.text(); @@ -148,27 +52,30 @@ async function completeMultipart(request: Request, env: Env, bucket: string, key return s3Error("MalformedXML", 400, undefined, `/${bucket}/${key}`); } const expectedTotal = parsePositiveInt(request.headers.get("x-amz-mp-object-size")); - const result = await multipartStub(env, uploadId).complete(parts, expectedTotal); - if (result.kind === "error") return s3Error(result.code, result.code === "NoSuchUpload" ? 404 : 400, result.message, `/${bucket}/${key}`); - const value = await multipartObjectEtag(result.metadata, result.partEtags, env.ETAG_STYLE); - return xmlResponse(completeMultipartUploadResult(bucket, key, value)); + const result = await completeMultipartUpload(env, bucket, key, uploadId, parts, expectedTotal); + if ("kind" in result && result.kind === "error") { + return s3Error(result.code, result.status, result.message, `/${bucket}/${key}`); + } + return xmlResponse(completeMultipartUploadResult(bucket, key, (result as { etag: string }).etag)); } async function abortMultipart(env: Env, bucket: string, key: string, url: URL): Promise { const uploadId = url.searchParams.get("uploadId") ?? ""; - if (!uploadIdMatches(uploadId, bucket, key) || !(await multipartStub(env, uploadId).abort())) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); + const result = await abortMultipartUpload(env, bucket, key, uploadId); + if (result !== true) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); return new Response(null, { status: 204 }); } async function listParts(env: Env, bucket: string, key: string, url: URL): Promise { const uploadId = url.searchParams.get("uploadId") ?? ""; - if (!uploadIdMatches(uploadId, bucket, key)) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); const marker = parsePositiveInt(url.searchParams.get("part-number-marker"), 0); const maxParts = parsePositiveInt(url.searchParams.get("max-parts"), 1000); if (marker === undefined || maxParts === undefined || maxParts > 1000) return s3Error("InvalidArgument", 400, undefined, `/${bucket}/${key}`); - const result = await multipartStub(env, uploadId).listParts(marker, maxParts); - if (!result) return s3Error("NoSuchUpload", 404, undefined, `/${bucket}/${key}`); - return xmlResponse(listPartsResult(bucket, key, uploadId, result)); + const result = await listMultipartParts(env, bucket, key, uploadId, marker, maxParts); + if ("kind" in result && result.kind === "error") { + return s3Error(result.code, result.status, result.message, `/${bucket}/${key}`); + } + return xmlResponse(listPartsResult(bucket, key, uploadId, result as import("./types").MultipartPartsList)); } async function putObject(request: Request, env: Env, accessToken: string, bucket: string, key: string): Promise { diff --git a/src/status-api.ts b/src/status-api.ts index b3a30df..1f94214 100644 --- a/src/status-api.ts +++ b/src/status-api.ts @@ -2,6 +2,7 @@ import { jsonResponse, verifySessionToken } from "./auth-api"; import { handleBucketRoutes, handleImportCandidatesRoute, handleImportRoute } from "./bucket-api"; import { getBucketRegistry } from "./bucket-registry"; import { getAccessToken, getDriveAbout, getRootFolderId, listObjects } from "./google-drive"; +import { handleObjectRoutes, handleTicketDownload } from "./object-api"; import type { DriveAbout, Env } from "./types"; export const API_PATH_PREFIX = "api"; @@ -75,11 +76,6 @@ export async function handleApi(request: Request, env: Env, subPath: string): Pr return jsonResponse({ message: "Dashboard authentication is not configured" }, 503); } - const isValid = await verifySessionToken(request, env); - if (!isValid) { - return jsonResponse({ message: "Session expired or invalid" }, 401); - } - if (!env.DRIVE_ROOT_FOLDER || env.DRIVE_ROOT_FOLDER.trim() === "") { return jsonResponse({ message: "Storage root folder is not configured" }, 503); } @@ -87,9 +83,18 @@ export async function handleApi(request: Request, env: Env, subPath: string): Pr const segments = subPath.split("/").filter(Boolean); const firstSegment = segments[0] || ""; const remainingSegments = segments.slice(1); - const url = new URL(request.url); + // Ticket download does not require session token (uses 120s random token) + if (firstSegment === "objects" && remainingSegments[0] === "content" && url.searchParams.has("ticket")) { + return await handleTicketDownload(request, env); + } + + const isValid = await verifySessionToken(request, env); + if (!isValid) { + return jsonResponse({ message: "Session expired or invalid" }, 401); + } + if (firstSegment === "status") { if (request.method !== "GET") return jsonResponse({ message: "Method Not Allowed" }, 405); return await handleStatus(env); @@ -103,6 +108,10 @@ export async function handleApi(request: Request, env: Env, subPath: string): Pr return await handleBucketRoutes(request, env, remainingSegments); } + if (firstSegment === "objects") { + return await handleObjectRoutes(request, env, remainingSegments); + } + if (firstSegment === "import-candidates") { return await handleImportCandidatesRoute(request, env); } diff --git a/test/cors.test.ts b/test/cors.test.ts index 8569e28..c11b682 100644 --- a/test/cors.test.ts +++ b/test/cors.test.ts @@ -58,7 +58,7 @@ describe("CORS", () => { expect(response.status).toBe(204); expect(response.headers.get("Access-Control-Allow-Origin")).toBe(ORIGIN); - expect(response.headers.get("Access-Control-Allow-Methods")).toBe("GET, HEAD, PUT, POST, DELETE, OPTIONS"); + expect(response.headers.get("Access-Control-Allow-Methods")).toBe("GET, HEAD, PUT, POST, PATCH, DELETE, OPTIONS"); expect(response.headers.get("Access-Control-Allow-Headers")).toBe("content-type,x-amz-date"); expect(response.headers.get("Access-Control-Max-Age")).toBe("86400"); expect(response.headers.get("Vary")).toBe("Origin"); diff --git a/test/fake-drive.ts b/test/fake-drive.ts new file mode 100644 index 0000000..3da22d7 --- /dev/null +++ b/test/fake-drive.ts @@ -0,0 +1,216 @@ +import { expect } from "vitest"; + +export interface StoredFile { + id: string; + name: string; + parent: string; + mimeType: string; + data: Uint8Array; + md5Checksum: string; + modifiedTime: string; + trashed?: boolean; +} + +export const FAKE_MODIFIED_TIME = "2026-08-16T08:00:00.000Z"; + +export interface StoredFolder { + id: string; + name: string; + parent: string; + trashed?: boolean; +} + +export const FOLDER_MIME = "application/vnd.google-apps.folder"; + +export interface UploadSession { + id: string; + fileId?: string; + name: string; + parent: string; + mimeType: string; + committed: Uint8Array; +} + +export class FakeDrive { + readonly files = new Map(); + readonly folders = new Map(); + readonly sessions = new Map(); + private nextId = 1; + + async handle(input: string | URL | Request, init?: RequestInit): Promise { + const request = input instanceof Request ? input : new Request(input, init); + const url = new URL(request.url); + if (url.hostname === "oauth2.googleapis.com") return Response.json({ access_token: "token", expires_in: 3600 }); + if (url.hostname !== "www.googleapis.com") return new Response("Not Found", { status: 404 }); + + if (url.pathname === "/drive/v3/files" && request.method === "GET") return this.search(url); + if (url.pathname === "/drive/v3/files" && request.method === "POST") return this.createMetadata(await request.json>()); + if (url.pathname.startsWith("/drive/v3/files/") && url.searchParams.get("alt") === "media") return this.download(url, request); + if (url.pathname.startsWith("/drive/v3/files/") && request.method === "PATCH") { + const fileId = url.pathname.split("/").at(-1); + if (!fileId) return new Response("Not Found", { status: 404 }); + const body = await request.json>(); + return this.patchFile(fileId, body, url); + } + if (url.pathname.startsWith("/drive/v3/files/") && request.method === "DELETE") { + const fileId = url.pathname.split("/").at(-1); + if (!fileId) return new Response("Not Found", { status: 404 }); + this.files.delete(fileId); + this.folders.delete(fileId); + return new Response(null, { status: 204 }); + } + if (url.pathname.startsWith("/upload/drive/v3/files") && url.searchParams.get("uploadType") === "resumable") return this.initialize(url, request); + if (url.pathname.startsWith("/upload/session/")) return this.upload(url, request); + return new Response("Not Found", { status: 404 }); + } + + private search(url: URL): Response { + const q = url.searchParams.get("q") ?? ""; + const hasNameFilter = /name='/.test(q); + const name = /name='((?:\\.|[^'])*)'/.exec(q)?.[1]?.replace(/\\'/g, "'").replace(/\\\\/g, "\\"); + const parent = /'([^']+)' in parents/.exec(q)?.[1] ?? "root"; + + if (q.includes(FOLDER_MIME)) { + const match = [...this.folders.values()].find((folder) => (!hasNameFilter || folder.name === name) && folder.parent === parent); + return Response.json({ files: match ? [{ id: match.id, name: match.name, mimeType: FOLDER_MIME }] : [] }); + } + + const files = [...this.files.values()].filter((file) => (!hasNameFilter || file.name === name) && file.parent === parent); + const folders = [...this.folders.values()].filter((folder) => (!hasNameFilter || folder.name === name) && folder.parent === parent); + return Response.json({ + files: [...files.map((file) => ({ ...file, size: String(file.data.byteLength), data: undefined })), ...folders.map((folder) => ({ id: folder.id, name: folder.name, mimeType: FOLDER_MIME }))], + }); + } + + private createMetadata(metadata: Record): Response { + const mimeType = String(metadata.mimeType ?? "application/octet-stream"); + const parent = String((metadata.parents as string[] | undefined)?.[0] ?? "root"); + if (mimeType === FOLDER_MIME) { + const id = `folder-${this.nextId++}`; + this.folders.set(id, { id, name: String(metadata.name), parent, trashed: false }); + return Response.json({ id, name: metadata.name, mimeType }); + } + const id = `file-${this.nextId++}`; + const file: StoredFile = { + id, + name: String(metadata.name), + parent, + mimeType, + data: new Uint8Array(), + md5Checksum: "d41d8cd98f00b204e9800998ecf8427e", + modifiedTime: FAKE_MODIFIED_TIME, + trashed: false, + }; + this.files.set(id, file); + return Response.json({ ...file, size: "0", data: undefined }); + } + + private patchFile(fileId: string, body: Record, url: URL): Response { + if (this.folders.has(fileId)) { + const folder = this.folders.get(fileId)!; + if (body.trashed !== undefined) folder.trashed = Boolean(body.trashed); + if (body.name !== undefined) folder.name = String(body.name); + return Response.json({ id: folder.id, name: folder.name, mimeType: FOLDER_MIME, trashed: folder.trashed }); + } + if (this.files.has(fileId)) { + const file = this.files.get(fileId)!; + if (body.trashed !== undefined) file.trashed = Boolean(body.trashed); + if (body.name !== undefined) file.name = String(body.name); + return Response.json({ id: file.id, name: file.name, mimeType: file.mimeType, size: String(file.data.byteLength), trashed: file.trashed }); + } + return new Response("Not Found", { status: 404 }); + } + + private async initialize(url: URL, request: Request): Promise { + const metadata = await request.json<{ name: string; parents?: string[] }>(); + const pathId = url.pathname.split("/").at(-1); + const fileId = pathId === "files" ? undefined : pathId; + const id = `session-${this.nextId++}`; + this.sessions.set(id, { + id, + fileId, + name: metadata.name, + parent: metadata.parents?.[0] ?? (fileId ? this.files.get(fileId)?.parent : "") ?? "", + mimeType: request.headers.get("X-Upload-Content-Type") ?? "application/octet-stream", + committed: new Uint8Array(), + }); + return new Response(null, { status: 200, headers: { Location: `https://www.googleapis.com/upload/session/${id}` } }); + } + + private async upload(url: URL, request: Request): Promise { + const id = url.pathname.split("/").at(-1); + if (!id) return new Response(null, { status: 404 }); + const session = this.sessions.get(id); + if (!session) return new Response(null, { status: 404 }); + if (request.method === "DELETE") { + this.sessions.delete(id); + return new Response(null, { status: 499 }); + } + const range = request.headers.get("Content-Range"); + if (range === "bytes */*") return new Response(null, { status: 308, headers: this.rangeHeaders(session.committed.byteLength) }); + const body = new Uint8Array(await request.arrayBuffer()); + if (!range) return this.finalize(session, body); + const match = /^bytes (\d+)-(\d+)\/(\*|\d+)$/.exec(range); + if (!match) return new Response("Bad range", { status: 400 }); + const start = Number(match[1]); + const end = Number(match[2]); + const total = match[3] === "*" ? null : Number(match[3]); + expect(start).toBe(session.committed.byteLength); + expect(end - start + 1).toBe(body.byteLength); + if (total === null) expect(body.byteLength % (256 * 1024)).toBe(0); + session.committed = concat(session.committed, body); + if (total === null) return new Response(null, { status: 308, headers: this.rangeHeaders(session.committed.byteLength) }); + expect(session.committed.byteLength).toBe(total); + return this.finalize(session, session.committed); + } + + private rangeHeaders(length: number): HeadersInit { + return length === 0 ? {} : { Range: `bytes=0-${length - 1}` }; + } + + private finalize(session: UploadSession, data: Uint8Array): Response { + const id = session.fileId ?? `file-${this.nextId++}`; + const file: StoredFile = { id, name: session.name, parent: session.parent, mimeType: session.mimeType, data, md5Checksum: fakeMd5(data), modifiedTime: FAKE_MODIFIED_TIME, trashed: false }; + this.files.set(id, file); + this.sessions.delete(session.id); + return Response.json({ ...file, size: String(data.byteLength), data: undefined }); + } + + private download(url: URL, request: Request): Response { + const fileId = url.pathname.split("/").at(-1); + if (!fileId) return new Response(null, { status: 404 }); + const file = this.files.get(fileId); + if (!file) return new Response(null, { status: 404 }); + const range = request.headers.get("Range"); + if (!range) return new Response(file.data, { headers: { "Content-Length": String(file.data.byteLength) } }); + const match = /^bytes=(\d+)-(\d+)?$/.exec(range); + if (!match) return new Response(null, { status: 416 }); + const start = Number(match[1]); + const end = Math.min(file.data.byteLength - 1, match[2] ? Number(match[2]) : file.data.byteLength - 1); + const data = file.data.slice(start, end + 1); + return new Response(data, { status: 206, headers: { "Content-Length": String(data.byteLength), "Content-Range": `bytes ${start}-${end}/${file.data.byteLength}` } }); + } +} + +export function concat(a: Uint8Array, b: Uint8Array): Uint8Array { + const output = new Uint8Array(a.byteLength + b.byteLength); + output.set(a); + output.set(b, a.byteLength); + return output; +} + +export function fakeMd5(data: Uint8Array): string { + let state = 0x811c9dc5; + for (const byte of data) state = Math.imul(state ^ byte, 0x01000193); + return (state >>> 0).toString(16).padStart(8, "0").repeat(4); +} + +export function bytes(length: number, seed = 17): Uint8Array { + const output = new Uint8Array(length); + let state = seed; + for (let index = 0; index < length; index++) { + state = (Math.imul(state, 1664525) + 1013904223) | 0; + output[index] = state >>> 24; + } + return output; +} diff --git a/test/objects.test.ts b/test/objects.test.ts new file mode 100644 index 0000000..9c759b5 --- /dev/null +++ b/test/objects.test.ts @@ -0,0 +1,346 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +import worker from "../src/index"; +import type { Env } from "../src/types"; +import { bytes, FakeDrive, fakeMd5 } from "./fake-drive"; + +import { env } from "cloudflare:test"; + +const ENV = env as unknown as Env; +const ENDPOINT = "https://s3-api.example.com"; +const CTX = { waitUntil: vi.fn(), passThroughOnException: vi.fn() } as unknown as ExecutionContext; + +let drive: FakeDrive; +let validToken: string; + +beforeEach(async () => { + drive = new FakeDrive(); + const rootFolderId = "folder-root"; + drive.folders.set(rootFolderId, { id: rootFolderId, name: "s3-storage", parent: "root" }); + const testBucketId = "folder-test-bucket"; + drive.folders.set(testBucketId, { id: testBucketId, name: "test-bucket", parent: rootFolderId }); + const emptyBucketId = "folder-empty-bucket"; + drive.folders.set(emptyBucketId, { id: emptyBucketId, name: "empty-bucket", parent: rootFolderId }); + + vi.stubGlobal( + "fetch", + vi.fn((input, init) => drive.handle(input, init)), + ); + await ENV.AUTH_KV.delete("google_access_token"); + for (const { name } of (await ENV.FOLDER_CACHE.list()).keys) await ENV.FOLDER_CACHE.delete(name); + + // Create session token + const token = "test-session-token-12345"; + const data = new TextEncoder().encode(token); + const hash = await crypto.subtle.digest("SHA-256", data); + const hashHex = Array.from(new Uint8Array(hash), (b) => b.toString(16).padStart(2, "0")).join(""); + await ENV.AUTH_KV.put(`session:${hashHex}`, JSON.stringify({ createdAt: Date.now() }), { expirationTtl: 3600 }); + validToken = token; +}); + +describe("Object API (/api/objects)", () => { + it("rejects requests without valid auth", async () => { + const res = await worker.fetch(new Request(`${ENDPOINT}/api/objects?bucket=test-bucket`), ENV, CTX); + expect(res.status).toBe(401); + }); + + it("returns 404 for unknown bucket", async () => { + const res = await worker.fetch( + new Request(`${ENDPOINT}/api/objects?bucket=nonexistent`, { + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(res.status).toBe(404); + const data = await res.json<{ message: string }>(); + expect(data.message).toContain("not found"); + }); + + it("lists objects and folders with delimiter", async () => { + // Create files in test-bucket + drive.files.set("file-1", { + id: "file-1", + name: "root-file.txt", + parent: "folder-test-bucket", + mimeType: "text/plain", + data: new TextEncoder().encode("hello"), + md5Checksum: fakeMd5(new TextEncoder().encode("hello")), + modifiedTime: "2026-08-16T08:00:00.000Z", + trashed: false, + }); + drive.folders.set("folder-sub", { + id: "folder-sub", + name: "photos", + parent: "folder-test-bucket", + trashed: false, + }); + drive.files.set("file-2", { + id: "file-2", + name: "pic.jpg", + parent: "folder-sub", + mimeType: "image/jpeg", + data: new Uint8Array([1, 2, 3]), + md5Checksum: fakeMd5(new Uint8Array([1, 2, 3])), + modifiedTime: "2026-08-16T08:00:00.000Z", + trashed: false, + }); + + const res = await worker.fetch( + new Request(`${ENDPOINT}/api/objects?bucket=test-bucket&prefix=&delimiter=/`, { + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(res.status).toBe(200); + const data = await res.json<{ + bucket: string; + prefix: string; + delimiter: string; + folders: Array<{ prefix: string; name: string }>; + objects: Array<{ key: string; name: string; size: number; contentType: string }>; + truncated: boolean; + }>(); + + expect(data.bucket).toBe("test-bucket"); + expect(data.folders).toEqual([{ prefix: "photos/", name: "photos" }]); + expect(data.objects).toHaveLength(1); + expect(data.objects[0].name).toBe("root-file.txt"); + expect(data.objects[0].key).toBe("root-file.txt"); + expect(data.truncated).toBe(false); + + // List nested folder + const subRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects?bucket=test-bucket&prefix=photos/&delimiter=/`, { + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(subRes.status).toBe(200); + const subData = await subRes.json<{ + folders: Array<{ prefix: string; name: string }>; + objects: Array<{ key: string; name: string; size: number }>; + }>(); + expect(subData.folders).toEqual([]); + expect(subData.objects).toHaveLength(1); + expect(subData.objects[0].name).toBe("pic.jpg"); + expect(subData.objects[0].key).toBe("photos/pic.jpg"); + }); + + it("gets object metadata", async () => { + drive.files.set("file-1", { + id: "file-1", + name: "test.txt", + parent: "folder-test-bucket", + mimeType: "text/plain", + data: new TextEncoder().encode("content"), + md5Checksum: fakeMd5(new TextEncoder().encode("content")), + modifiedTime: "2026-08-16T08:00:00.000Z", + trashed: false, + }); + + const res = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/metadata?bucket=test-bucket&key=test.txt`, { + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(res.status).toBe(200); + const data = await res.json<{ key: string; name: string; size: number; contentType: string }>(); + expect(data.key).toBe("test.txt"); + expect(data.name).toBe("test.txt"); + expect(data.size).toBe(7); + expect(data.contentType).toBe("text/plain"); + }); + + it("direct PUT and GET content", async () => { + const payload = new Uint8Array([10, 20, 30, 40]); + const putRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/content?bucket=test-bucket&key=binary.bin`, { + method: "PUT", + headers: { + Authorization: `Bearer ${validToken}`, + "Content-Type": "application/octet-stream", + }, + body: payload, + }), + ENV, + CTX, + ); + expect(putRes.status).toBe(200); + const putData = await putRes.json<{ key: string; etag: string; size: number }>(); + expect(putData.key).toBe("binary.bin"); + + const getRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/content?bucket=test-bucket&key=binary.bin`, { + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(getRes.status).toBe(200); + expect(getRes.headers.get("Content-Disposition")).toContain("binary.bin"); + const body = new Uint8Array(await getRes.arrayBuffer()); + expect(body).toEqual(payload); + }); + + it("creates and deletes folders with empty / non-empty checks", async () => { + // Create folder + const createRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/folder`, { + method: "POST", + headers: { Authorization: `Bearer ${validToken}`, "Content-Type": "application/json" }, + body: JSON.stringify({ bucket: "test-bucket", prefix: "my-folder/" }), + }), + ENV, + CTX, + ); + expect(createRes.status).toBe(201); + + // Put a file inside folder + await worker.fetch( + new Request(`${ENDPOINT}/api/objects/content?bucket=test-bucket&key=my-folder/item.txt`, { + method: "PUT", + headers: { Authorization: `Bearer ${validToken}` }, + body: "hello", + }), + ENV, + CTX, + ); + + // Try deleting non-empty folder without recursive=1 -> 409 + const delRes409 = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/folder?bucket=test-bucket&prefix=my-folder/`, { + method: "DELETE", + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(delRes409.status).toBe(409); + + // Delete with recursive=1 -> 204 + const delRes204 = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/folder?bucket=test-bucket&prefix=my-folder/&recursive=1`, { + method: "DELETE", + headers: { Authorization: `Bearer ${validToken}` }, + }), + ENV, + CTX, + ); + expect(delRes204.status).toBe(204); + }); + + it("supports download tickets without session token", async () => { + await worker.fetch( + new Request(`${ENDPOINT}/api/objects/content?bucket=test-bucket&key=ticket-doc.pdf`, { + method: "PUT", + headers: { Authorization: `Bearer ${validToken}` }, + body: "pdf contents here", + }), + ENV, + CTX, + ); + + // Create download ticket + const ticketRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/download-ticket`, { + method: "POST", + headers: { Authorization: `Bearer ${validToken}`, "Content-Type": "application/json" }, + body: JSON.stringify({ bucket: "test-bucket", key: "ticket-doc.pdf" }), + }), + ENV, + CTX, + ); + expect(ticketRes.status).toBe(201); + const { ticket, downloadUrl } = await ticketRes.json<{ ticket: string; downloadUrl: string }>(); + expect(ticket).toBeDefined(); + + // Download without Authorization header using ticket + const dlRes = await worker.fetch(new Request(`${ENDPOINT}${downloadUrl}`), ENV, CTX); + expect(dlRes.status).toBe(200); + expect(await dlRes.text()).toBe("pdf contents here"); + + // Second attempt with same ticket should fail (single use) + const dlRes2 = await worker.fetch(new Request(`${ENDPOINT}${downloadUrl}`), ENV, CTX); + expect(dlRes2.status).toBe(403); + }); + + it("performs full multipart upload flow via JSON API with non-aligned parts", async () => { + // 1. Initiate multipart + const initRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/uploads`, { + method: "POST", + headers: { Authorization: `Bearer ${validToken}`, "Content-Type": "application/json" }, + body: JSON.stringify({ bucket: "test-bucket", key: "large-video.mp4", contentType: "video/mp4" }), + }), + ENV, + CTX, + ); + expect(initRes.status).toBe(201); + const { uploadId, partSize } = await initRes.json<{ uploadId: string; partSize: number }>(); + expect(uploadId).toBeDefined(); + expect(partSize).toBe(8 * 1024 * 1024); + + // 2. Upload part 1: 300 KiB (non-aligned to 256 KiB) + const part1Data = bytes(300 * 1024, 11); + const part1Res = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/uploads/part?bucket=test-bucket&key=large-video.mp4&uploadId=${encodeURIComponent(uploadId)}&partNumber=1`, { + method: "PUT", + headers: { Authorization: `Bearer ${validToken}`, "Content-Length": String(part1Data.byteLength) }, + body: part1Data, + }), + ENV, + CTX, + ); + expect(part1Res.status).toBe(200); + const part1Json = await part1Res.json<{ partNumber: number; etag: string }>(); + + // 3. Upload part 2: 200 KiB (non-aligned) + const part2Data = bytes(200 * 1024, 22); + const part2Res = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/uploads/part?bucket=test-bucket&key=large-video.mp4&uploadId=${encodeURIComponent(uploadId)}&partNumber=2`, { + method: "PUT", + headers: { Authorization: `Bearer ${validToken}`, "Content-Length": String(part2Data.byteLength) }, + body: part2Data, + }), + ENV, + CTX, + ); + expect(part2Res.status).toBe(200); + const part2Json = await part2Res.json<{ partNumber: number; etag: string }>(); + + // 4. Complete multipart + const completeRes = await worker.fetch( + new Request(`${ENDPOINT}/api/objects/uploads/complete`, { + method: "POST", + headers: { Authorization: `Bearer ${validToken}`, "Content-Type": "application/json" }, + body: JSON.stringify({ + bucket: "test-bucket", + key: "large-video.mp4", + uploadId, + parts: [ + { partNumber: 1, etag: part1Json.etag }, + { partNumber: 2, etag: part2Json.etag }, + ], + }), + }), + ENV, + CTX, + ); + expect(completeRes.status).toBe(200); + const completeJson = await completeRes.json<{ key: string; etag: string }>(); + expect(completeJson.key).toBe("large-video.mp4"); + + // Verify stored data + const storedFile = [...drive.files.values()].find((f) => f.name === "large-video.mp4"); + expect(storedFile).toBeDefined(); + const expectedBytes = new Uint8Array(part1Data.byteLength + part2Data.byteLength); + expectedBytes.set(part1Data, 0); + expectedBytes.set(part2Data, part1Data.byteLength); + expect(storedFile!.data).toEqual(expectedBytes); + }); +}); diff --git a/test/s3.test.ts b/test/s3.test.ts index 0dd2d21..a64989a 100644 --- a/test/s3.test.ts +++ b/test/s3.test.ts @@ -4,6 +4,7 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; import { decodedBodyChunks } from "../src/aws-chunked"; import worker from "../src/index"; import type { Env } from "../src/types"; +import { bytes, FAKE_MODIFIED_TIME, fakeMd5, FakeDrive } from "./fake-drive"; import { env } from "cloudflare:test"; @@ -11,195 +12,6 @@ const ENV = env as unknown as Env; const ENDPOINT = "https://s3-api.example.com"; const CTX = { waitUntil: vi.fn(), passThroughOnException: vi.fn() } as unknown as ExecutionContext; -interface StoredFile { - id: string; - name: string; - parent: string; - mimeType: string; - data: Uint8Array; - md5Checksum: string; - modifiedTime: string; -} - -const FAKE_MODIFIED_TIME = "2026-08-16T08:00:00.000Z"; - -interface StoredFolder { - id: string; - name: string; - parent: string; -} - -const FOLDER_MIME = "application/vnd.google-apps.folder"; - -interface UploadSession { - id: string; - fileId?: string; - name: string; - parent: string; - mimeType: string; - committed: Uint8Array; -} - -class FakeDrive { - readonly files = new Map(); - readonly folders = new Map(); - readonly sessions = new Map(); - private nextId = 1; - - async handle(input: string | URL | Request, init?: RequestInit): Promise { - const request = input instanceof Request ? input : new Request(input, init); - const url = new URL(request.url); - if (url.hostname === "oauth2.googleapis.com") return Response.json({ access_token: "token", expires_in: 3600 }); - if (url.hostname !== "www.googleapis.com") return new Response("Not Found", { status: 404 }); - - if (url.pathname === "/drive/v3/files" && request.method === "GET") return this.search(url); - if (url.pathname === "/drive/v3/files" && request.method === "POST") return this.createMetadata(await request.json>()); - if (url.pathname.startsWith("/drive/v3/files/") && url.searchParams.get("alt") === "media") return this.download(url, request); - if (url.pathname.startsWith("/drive/v3/files/") && request.method === "DELETE") { - const fileId = url.pathname.split("/").at(-1); - if (!fileId) return new Response("Not Found", { status: 404 }); - this.files.delete(fileId); - return new Response(null, { status: 204 }); - } - if (url.pathname.startsWith("/upload/drive/v3/files") && url.searchParams.get("uploadType") === "resumable") return this.initialize(url, request); - if (url.pathname.startsWith("/upload/session/")) return this.upload(url, request); - return new Response("Not Found", { status: 404 }); - } - - private search(url: URL): Response { - const q = url.searchParams.get("q") ?? ""; - const hasNameFilter = /name='/.test(q); - const name = /name='((?:\\.|[^'])*)'/.exec(q)?.[1]?.replace(/\\'/g, "'").replace(/\\\\/g, "\\"); - const parent = /'([^']+)' in parents/.exec(q)?.[1] ?? "root"; - - if (q.includes(FOLDER_MIME)) { - const match = [...this.folders.values()].find((folder) => folder.name === name && folder.parent === parent); - return Response.json({ files: match ? [{ id: match.id, name: match.name, mimeType: FOLDER_MIME }] : [] }); - } - - const files = [...this.files.values()].filter((file) => (!hasNameFilter || file.name === name) && file.parent === parent); - const folders = [...this.folders.values()].filter((folder) => (!hasNameFilter || folder.name === name) && folder.parent === parent); - return Response.json({ - files: [...files.map((file) => ({ ...file, size: String(file.data.byteLength), data: undefined })), ...folders.map((folder) => ({ id: folder.id, name: folder.name, mimeType: FOLDER_MIME }))], - }); - } - - private createMetadata(metadata: Record): Response { - const mimeType = String(metadata.mimeType ?? "application/octet-stream"); - const parent = String((metadata.parents as string[] | undefined)?.[0] ?? "root"); - if (mimeType === FOLDER_MIME) { - const id = `folder-${this.nextId++}`; - this.folders.set(id, { id, name: String(metadata.name), parent }); - return Response.json({ id, name: metadata.name, mimeType }); - } - const id = `file-${this.nextId++}`; - const file: StoredFile = { - id, - name: String(metadata.name), - parent, - mimeType, - data: new Uint8Array(), - md5Checksum: "d41d8cd98f00b204e9800998ecf8427e", - modifiedTime: FAKE_MODIFIED_TIME, - }; - this.files.set(id, file); - return Response.json({ ...file, size: "0", data: undefined }); - } - - private async initialize(url: URL, request: Request): Promise { - const metadata = await request.json<{ name: string; parents?: string[] }>(); - const pathId = url.pathname.split("/").at(-1); - const fileId = pathId === "files" ? undefined : pathId; - const id = `session-${this.nextId++}`; - this.sessions.set(id, { - id, - fileId, - name: metadata.name, - parent: metadata.parents?.[0] ?? (fileId ? this.files.get(fileId)?.parent : "") ?? "", - mimeType: request.headers.get("X-Upload-Content-Type") ?? "application/octet-stream", - committed: new Uint8Array(), - }); - return new Response(null, { status: 200, headers: { Location: `https://www.googleapis.com/upload/session/${id}` } }); - } - - private async upload(url: URL, request: Request): Promise { - const id = url.pathname.split("/").at(-1); - if (!id) return new Response(null, { status: 404 }); - const session = this.sessions.get(id); - if (!session) return new Response(null, { status: 404 }); - if (request.method === "DELETE") { - this.sessions.delete(id); - return new Response(null, { status: 499 }); - } - const range = request.headers.get("Content-Range"); - if (range === "bytes */*") return new Response(null, { status: 308, headers: this.rangeHeaders(session.committed.byteLength) }); - const body = new Uint8Array(await request.arrayBuffer()); - if (!range) return this.finalize(session, body); - const match = /^bytes (\d+)-(\d+)\/(\*|\d+)$/.exec(range); - if (!match) return new Response("Bad range", { status: 400 }); - const start = Number(match[1]); - const end = Number(match[2]); - const total = match[3] === "*" ? null : Number(match[3]); - expect(start).toBe(session.committed.byteLength); - expect(end - start + 1).toBe(body.byteLength); - if (total === null) expect(body.byteLength % (256 * 1024)).toBe(0); - session.committed = concat(session.committed, body); - if (total === null) return new Response(null, { status: 308, headers: this.rangeHeaders(session.committed.byteLength) }); - expect(session.committed.byteLength).toBe(total); - return this.finalize(session, session.committed); - } - - private rangeHeaders(length: number): HeadersInit { - return length === 0 ? {} : { Range: `bytes=0-${length - 1}` }; - } - - private finalize(session: UploadSession, data: Uint8Array): Response { - const id = session.fileId ?? `file-${this.nextId++}`; - const file: StoredFile = { id, name: session.name, parent: session.parent, mimeType: session.mimeType, data, md5Checksum: fakeMd5(data), modifiedTime: FAKE_MODIFIED_TIME }; - this.files.set(id, file); - this.sessions.delete(session.id); - return Response.json({ ...file, size: String(data.byteLength), data: undefined }); - } - - private download(url: URL, request: Request): Response { - const fileId = url.pathname.split("/").at(-1); - if (!fileId) return new Response(null, { status: 404 }); - const file = this.files.get(fileId); - if (!file) return new Response(null, { status: 404 }); - const range = request.headers.get("Range"); - if (!range) return new Response(file.data, { headers: { "Content-Length": String(file.data.byteLength) } }); - const match = /^bytes=(\d+)-(\d+)?$/.exec(range); - if (!match) return new Response(null, { status: 416 }); - const start = Number(match[1]); - const end = Math.min(file.data.byteLength - 1, match[2] ? Number(match[2]) : file.data.byteLength - 1); - const data = file.data.slice(start, end + 1); - return new Response(data, { status: 206, headers: { "Content-Length": String(data.byteLength), "Content-Range": `bytes ${start}-${end}/${file.data.byteLength}` } }); - } -} - -function concat(a: Uint8Array, b: Uint8Array): Uint8Array { - const output = new Uint8Array(a.byteLength + b.byteLength); - output.set(a); - output.set(b, a.byteLength); - return output; -} - -function fakeMd5(data: Uint8Array): string { - let state = 0x811c9dc5; - for (const byte of data) state = Math.imul(state ^ byte, 0x01000193); - return (state >>> 0).toString(16).padStart(8, "0").repeat(4); -} - -function bytes(length: number, seed = 17): Uint8Array { - const output = new Uint8Array(length); - let state = seed; - for (let index = 0; index < length; index++) { - state = (Math.imul(state, 1664525) + 1013904223) | 0; - output[index] = state >>> 24; - } - return output; -} - async function signed(path: string, init: RequestInit): Promise { const aws = new AwsClient({ accessKeyId: ENV.ACCESS_KEY, secretAccessKey: ENV.SECRET_KEY, region: ENV.REGION, service: "s3" }); const bodyLength = typeof init.body === "string" ? new TextEncoder().encode(init.body).byteLength : init.body instanceof Uint8Array ? init.body.byteLength : undefined;