mirror of
https://github.com/Nezumi-2711/google-drive-s3.git
synced 2026-09-22 13:38:30 +00:00
feat: add api for manage bucket
This commit is contained in:
+2
-2
@@ -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
|
||||
|
||||
|
||||
@@ -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'
|
||||
|
||||
+1
-1
@@ -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 ?? "")
|
||||
|
||||
+79
-17
@@ -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<GoogleDriveFile[]> {
|
||||
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<boolean> {
|
||||
@@ -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<DriveDownloadResult> {
|
||||
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<void> {
|
||||
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<DriveFileMetadata> {
|
||||
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<GoogleDriveFile[]> {
|
||||
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. */
|
||||
|
||||
@@ -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<void> {
|
||||
if (body) await body.pipeTo(new WritableStream());
|
||||
}
|
||||
|
||||
export async function partEtag(uploadId: string, partNumber: number, partLen: number, fileOffset: number): Promise<string> {
|
||||
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<string> {
|
||||
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<ArrayBuffer | ArrayBufferView>();
|
||||
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<boolean | CoreError> {
|
||||
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<import("./types").MultipartPartsList | CoreError> {
|
||||
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;
|
||||
}
|
||||
@@ -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<void> {
|
||||
await env.FOLDER_CACHE.delete(`bucket-stats:${bucket}`);
|
||||
if (parentId && name) {
|
||||
await env.FOLDER_CACHE.delete(`${parentId}/${name}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function sha256Hex(text: string): Promise<string> {
|
||||
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<Response> {
|
||||
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<Response> {
|
||||
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);
|
||||
}
|
||||
+36
-129
@@ -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<void> {
|
||||
if (body) await body.pipeTo(new WritableStream());
|
||||
}
|
||||
|
||||
async function partEtag(uploadId: string, partNumber: number, partLen: number, fileOffset: number): Promise<string> {
|
||||
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<string> {
|
||||
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<Response> {
|
||||
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<Response> {
|
||||
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<ArrayBuffer | ArrayBufferView>();
|
||||
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<Response> {
|
||||
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<Response> {
|
||||
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<Response> {
|
||||
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<Response> {
|
||||
|
||||
+15
-6
@@ -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);
|
||||
}
|
||||
|
||||
+1
-1
@@ -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");
|
||||
|
||||
@@ -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<string, StoredFile>();
|
||||
readonly folders = new Map<string, StoredFolder>();
|
||||
readonly sessions = new Map<string, UploadSession>();
|
||||
private nextId = 1;
|
||||
|
||||
async handle(input: string | URL | Request, init?: RequestInit): Promise<Response> {
|
||||
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<Record<string, unknown>>());
|
||||
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<Record<string, unknown>>();
|
||||
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<string, unknown>): 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<string, unknown>, 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<Response> {
|
||||
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<Response> {
|
||||
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;
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
+1
-189
@@ -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<string, StoredFile>();
|
||||
readonly folders = new Map<string, StoredFolder>();
|
||||
readonly sessions = new Map<string, UploadSession>();
|
||||
private nextId = 1;
|
||||
|
||||
async handle(input: string | URL | Request, init?: RequestInit): Promise<Response> {
|
||||
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<Record<string, unknown>>());
|
||||
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<string, unknown>): 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<Response> {
|
||||
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<Response> {
|
||||
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<Request> {
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user