refactor: update the structure of the code

This commit is contained in:
2026-08-11 18:02:24 +07:00
parent ff3ccc2078
commit 5dd168dc98
20 changed files with 7192 additions and 4421 deletions
+19 -1
View File
@@ -70,7 +70,25 @@ https://developers.cloudflare.com/workers/configuration/secrets/#via-the-dashboa
| `ALLOWED_BUCKETS` | Set the buckets allowed, separated by `,`. A directory with the bucket name will be created directly under Google Drive. | | `ALLOWED_BUCKETS` | Set the buckets allowed, separated by `,`. A directory with the bucket name will be created directly under Google Drive. |
| `PUBLIC_READ_BUCKETS` | *(Optional)* Buckets that allow unauthenticated GET/HEAD access without signature, separated by `,`. Write operations (PUT/POST/DELETE) still require authentication. Must be a subset of `ALLOWED_BUCKETS`. | | `PUBLIC_READ_BUCKETS` | *(Optional)* Buckets that allow unauthenticated GET/HEAD access without signature, separated by `,`. Write operations (PUT/POST/DELETE) still require authentication. Must be a subset of `ALLOWED_BUCKETS`. |
### 4. Enable Multipart Uploads
### 4. CORS Configuration Large uploads require S3 Multipart Upload because Cloudflare's request-size limit applies before a request reaches the Worker. Multipart support uses one SQLite-backed Durable Object per upload and streams each admitted part directly into one Google Drive resumable session.
After deploying the `MultipartUploadDO` migration, change `ALLOW_MULTIPART` to `"true"` in `wrangler.jsonc`. Recommended AWS CLI settings:
```ini
[default]
s3 =
multipart_chunksize = 16MB
max_concurrent_requests = 3
request_checksum_calculation = when_required
```
Multipart parts must use consecutive numbers starting at 1 and are immutable after commit. The completed object's ETag defaults to Google Drive's real MD5 rather than S3's composite multipart ETag; set `ETAG_STYLE` to `"multipart"` only for clients that require the composite form. Existing objects may be re-evaluated once when their old Drive-ID ETag changes to MD5.
Google Drive's free tier has 15 GB total storage, and Google applies a 750 GB daily upload limit.
### 5. CORS Configuration
If you need to configure CORS, set up your own domain for Workers and use Cloudflare's Response Header Transform Rules to add the necessary headers. If you need to configure CORS, set up your own domain for Workers and use Cloudflare's Response Header Transform Rules to add the necessary headers.
https://developers.cloudflare.com/rules/transform/response-header-modification/ https://developers.cloudflare.com/rules/transform/response-header-modification/
+2 -2
View File
@@ -1,5 +1,5 @@
{ {
"$schema": "https://biomejs.dev/schemas/2.3.11/schema.json", "$schema": "https://biomejs.dev/schemas/2.5.7/schema.json",
"vcs": { "vcs": {
"enabled": true, "enabled": true,
"clientKind": "git", "clientKind": "git",
@@ -18,7 +18,7 @@
"linter": { "linter": {
"enabled": true, "enabled": true,
"rules": { "rules": {
"recommended": true, "preset": "recommended",
"performance": { "performance": {
"noImgElement": "off" "noImgElement": "off"
}, },
+6 -6
View File
@@ -12,12 +12,12 @@
"format": "biome format --write" "format": "biome format --write"
}, },
"devDependencies": { "devDependencies": {
"@aws-sdk/client-s3": "^3.1031.0", "@aws-sdk/client-s3": "^3.1106.0",
"@aws-sdk/s3-request-presigner": "^3.1000.0", "@aws-sdk/s3-request-presigner": "^3.1106.0",
"@biomejs/biome": "^2.4.12", "@biomejs/biome": "^2.5.7",
"@cloudflare/vitest-pool-workers": "^0.12.1", "@cloudflare/vitest-pool-workers": "^0.20.3",
"aws4fetch": "^1.0.20", "aws4fetch": "^1.0.20",
"typescript": "^6.0.3", "typescript": "^7.0.2",
"vitest": "~4.1.4" "vitest": "~4.1.10"
} }
} }
+1268 -2491
View File
File diff suppressed because it is too large Load Diff
+158
View File
@@ -0,0 +1,158 @@
const MAX_CONTROL_LINE = 128;
type DecoderState = "HEADER" | "DATA" | "CRLF" | "TRAILER" | "DONE";
function asBytes(value: Uint8Array | ArrayBuffer): Uint8Array {
return value instanceof Uint8Array ? value : new Uint8Array(value);
}
/** Iterates decoded payload views without scanning or copying payload bytes. */
export async function* decodedBodyChunks(stream: ReadableStream<Uint8Array> | null, awsChunked: boolean): AsyncGenerator<Uint8Array> {
if (!stream) return;
const reader = stream.getReader();
if (!awsChunked) {
try {
while (true) {
const { done, value } = await reader.read();
if (done) return;
if (value.byteLength > 0) yield asBytes(value);
}
} finally {
reader.releaseLock();
}
}
let state: DecoderState = "HEADER";
let remaining = 0;
let control: number[] = [];
let emitted = 0;
const consumeControlByte = (byte: number): string | null => {
control.push(byte);
if (control.length > MAX_CONTROL_LINE) throw new Error("Invalid aws-chunked control line");
const length = control.length;
if (length < 2 || control[length - 2] !== 13 || control[length - 1] !== 10) return null;
const line = new TextDecoder().decode(new Uint8Array(control.slice(0, -2)));
control = [];
return line;
};
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
const chunk = asBytes(value);
let offset = 0;
while (offset < chunk.byteLength) {
if (state === "DATA") {
const length = Math.min(remaining, chunk.byteLength - offset);
if (length > 0) {
emitted += length;
remaining -= length;
yield chunk.subarray(offset, offset + length);
offset += length;
}
if (remaining === 0) state = "CRLF";
continue;
}
const line = consumeControlByte(chunk[offset++]);
if (line === null) continue;
if (state === "HEADER") {
const sizeText = line.split(";", 1)[0];
if (!/^[0-9a-fA-F]+$/.test(sizeText)) throw new Error("Invalid aws-chunked chunk size");
remaining = Number.parseInt(sizeText, 16);
if (!Number.isSafeInteger(remaining)) throw new Error("aws-chunked chunk is too large");
state = remaining === 0 ? "TRAILER" : "DATA";
} else if (state === "CRLF") {
if (line !== "") throw new Error("Invalid aws-chunked data terminator");
state = "HEADER";
} else if (state === "TRAILER" && line === "") {
state = "DONE";
if (offset !== chunk.byteLength) throw new Error("Unexpected bytes after aws-chunked trailer");
} else if (state === "DONE") {
throw new Error("Unexpected bytes after aws-chunked trailer");
}
}
}
} finally {
reader.releaseLock();
}
if (state !== "DONE") throw new Error("Truncated aws-chunked body");
return emitted;
}
export async function pumpBody(
stream: ReadableStream<Uint8Array> | null,
writer: WritableStreamDefaultWriter<ArrayBuffer | ArrayBufferView>,
options: { awsChunked: boolean; expectedLength?: number; skipBytes?: number; maxBytes?: number; prefix?: Uint8Array } = { awsChunked: false },
): Promise<{ written: number; tail: Uint8Array }> {
let decoded = 0;
let written = 0;
let skip = options.skipBytes ?? 0;
const tailParts: Uint8Array[] = [];
let tailLength = 0;
const consume = async (input: Uint8Array): Promise<void> => {
let chunk = input;
if (skip >= chunk.byteLength) {
skip -= chunk.byteLength;
return;
}
if (skip > 0) {
chunk = chunk.subarray(skip);
skip = 0;
}
const remaining = options.maxBytes === undefined ? chunk.byteLength : Math.max(0, options.maxBytes - written);
const send = chunk.subarray(0, remaining);
if (send.byteLength > 0) {
await writer.write(send);
written += send.byteLength;
}
if (send.byteLength < chunk.byteLength) {
const tail = chunk.subarray(send.byteLength);
tailParts.push(tail);
tailLength += tail.byteLength;
}
};
try {
if (options.prefix) await consume(options.prefix);
for await (const chunk of decodedBodyChunks(stream, options.awsChunked)) {
decoded += chunk.byteLength;
await consume(chunk);
}
if (skip !== 0) throw new Error("Body is shorter than the committed Drive range");
if (options.expectedLength !== undefined && decoded !== options.expectedLength) throw new Error(`Decoded body length ${decoded} does not match x-amz-decoded-content-length ${options.expectedLength}`);
await writer.close();
} catch (error) {
await writer.abort(error).catch(() => undefined);
throw error;
}
const tail = new Uint8Array(tailLength);
let offset = 0;
for (const part of tailParts) {
tail.set(part, offset);
offset += part.byteLength;
}
return { written, tail };
}
export function isAwsChunked(request: Request): boolean {
return (
request.headers
.get("content-encoding")
?.split(",")
.some((encoding) => encoding.trim().toLowerCase() === "aws-chunked") ?? false
);
}
export function decodedContentLength(request: Request): number | undefined {
const value = request.headers.get("x-amz-decoded-content-length") ?? request.headers.get("content-length");
if (value === null || !/^\d+$/.test(value)) return undefined;
const length = Number(value);
return Number.isSafeInteger(length) ? length : undefined;
}
+126
View File
@@ -0,0 +1,126 @@
import type { Env } from "./types";
function encodeRFC3986(str: string): string {
return encodeURIComponent(str).replace(/[!'()*]/g, (c) => `%${c.charCodeAt(0).toString(16).toUpperCase()}`);
}
function bufToHex(buf: ArrayBuffer): string {
return Array.from(new Uint8Array(buf))
.map((b) => b.toString(16).padStart(2, "0"))
.join("");
}
async function hmacSha256(key: string | ArrayBuffer, data: string): Promise<ArrayBuffer> {
const keyData = typeof key === "string" ? new TextEncoder().encode(key) : key;
const cryptoKey = await crypto.subtle.importKey("raw", keyData, { name: "HMAC", hash: "SHA-256" }, false, ["sign"]);
return await crypto.subtle.sign("HMAC", cryptoKey, new TextEncoder().encode(data));
}
async function sha256(data: string): Promise<string> {
const hash = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(data));
return bufToHex(hash);
}
async function getSigningKey(secret: string, date: string, region: string, service: string): Promise<ArrayBuffer> {
const kDate = await hmacSha256(`AWS4${secret}`, date);
const kRegion = await hmacSha256(kDate, region);
const kService = await hmacSha256(kRegion, service);
return await hmacSha256(kService, "aws4_request");
}
async function createCanonicalRequest(request: Request, isQueryAuth: boolean): Promise<string> {
const url = new URL(request.url);
const method = request.method;
const canonicalUri = url.pathname || "/";
const params = Array.from(url.searchParams.entries())
.filter(([key]) => key !== "X-Amz-Signature")
.sort(([a], [b]) => {
if (a < b) return -1;
if (a > b) return 1;
return 0;
})
.map(([key, val]) => `${encodeRFC3986(key)}=${encodeRFC3986(val)}`)
.join("&");
let signedHeadersList: string[];
if (isQueryAuth) {
signedHeadersList = (url.searchParams.get("X-Amz-SignedHeaders") ?? "host").split(";");
} else {
const authHeader = request.headers.get("Authorization") ?? "";
const match = authHeader.match(/SignedHeaders=([^,\s]+)/);
signedHeadersList = match ? match[1].split(";") : ["host"];
}
const canonicalHeaders = signedHeadersList
.map((h) => {
const headerName = h.toLowerCase();
let headerValue = "";
if (headerName === "host") {
headerValue = url.hostname;
const port = url.port;
if (port && !((url.protocol === "https:" && port === "443") || (url.protocol === "http:" && port === "80"))) {
headerValue += `:${port}`;
}
} else {
headerValue = request.headers.get(headerName)?.trim() ?? "";
}
return `${headerName}:${headerValue}\n`;
})
.join("");
const signedHeaders = signedHeadersList.join(";");
const payloadHash = request.headers.get("x-amz-content-sha256") ?? "UNSIGNED-PAYLOAD";
return [method, canonicalUri, params, canonicalHeaders, signedHeaders, payloadHash].join("\n");
}
/** Verifies an AWS Signature V4 signature carried in either the Authorization header or presigned query params. */
export async function verifySignature(request: Request, env: Env): Promise<boolean> {
const url = new URL(request.url);
const headers = request.headers;
const isQueryAuth = url.searchParams.has("X-Amz-Algorithm");
let algorithm: string;
if (isQueryAuth) {
algorithm = url.searchParams.get("X-Amz-Algorithm") ?? "";
} else {
const authHeader = headers.get("Authorization") ?? "";
algorithm = authHeader.split(" ")[0];
}
if (!algorithm?.includes("AWS4-HMAC-SHA256")) {
return false;
}
const datetime = (isQueryAuth ? url.searchParams.get("X-Amz-Date") : headers.get("x-amz-date")) ?? "";
if (!datetime) return false;
const date = datetime.substring(0, 8);
const canonicalRequest = await createCanonicalRequest(request, isQueryAuth);
const hashedCanonicalRequest = await sha256(canonicalRequest);
const credentialScope = `${date}/${env.REGION}/s3/aws4_request`;
const stringToSign = ["AWS4-HMAC-SHA256", datetime, credentialScope, hashedCanonicalRequest].join("\n");
const signingKey = await getSigningKey(env.SECRET_KEY, date, env.REGION, "s3");
const signature = await hmacSha256(signingKey, stringToSign);
const signatureHex = bufToHex(signature);
let expectedSignature = "";
if (isQueryAuth) {
expectedSignature = url.searchParams.get("X-Amz-Signature") ?? "";
} else {
const authHeader = headers.get("Authorization") ?? "";
const match = authHeader.match(/Signature=([a-f0-9]+)/);
expectedSignature = match ? match[1] : "";
}
return signatureHex === expectedSignature;
}
+34
View File
@@ -0,0 +1,34 @@
import type { Env } from "./types";
/**
* Checks whether the bucket is present in the ALLOWED_BUCKETS allowlist.
* Access is denied by default when the allowlist is missing or empty.
*/
export function isAllowedBucket(bucket: string, env: Env): boolean {
if (!env.ALLOWED_BUCKETS) {
return false;
}
const allowedBuckets = env.ALLOWED_BUCKETS.split(",")
.map((b) => b.trim())
.filter((b) => b);
if (allowedBuckets.length === 0) {
return false;
}
return allowedBuckets.includes(bucket);
}
/** Checks whether the bucket allows unauthenticated read access. */
export function isPublicReadBucket(bucket: string, env: Env): boolean {
if (!env.PUBLIC_READ_BUCKETS) {
return false;
}
const publicReadBuckets = env.PUBLIC_READ_BUCKETS.split(",")
.map((b) => b.trim())
.filter((b) => b);
return publicReadBuckets.includes(bucket);
}
+76
View File
@@ -0,0 +1,76 @@
import type { DriveUploadResult } from "./types";
export const DRIVE_CHUNK_SIZE = 256 * 1024;
const DRIVE_FIELDS = "id,name,size,mimeType,md5Checksum";
export function alignedSendLen(driveOffset: number, available: number): number {
if (available <= 1) return 0;
return Math.max(0, Math.floor((driveOffset + available - 1) / DRIVE_CHUNK_SIZE) * DRIVE_CHUNK_SIZE - driveOffset);
}
export async function createSession(accessToken: string, metadata: { name: string; parents: string[]; mimeType: string; existingFileId?: string }): Promise<string> {
const url = metadata.existingFileId ? `https://www.googleapis.com/upload/drive/v3/files/${metadata.existingFileId}?uploadType=resumable&fields=${encodeURIComponent(DRIVE_FIELDS)}` : `https://www.googleapis.com/upload/drive/v3/files?uploadType=resumable&fields=${encodeURIComponent(DRIVE_FIELDS)}`;
const response = await fetch(url, {
method: metadata.existingFileId ? "PATCH" : "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Type": "application/json; charset=UTF-8",
"X-Upload-Content-Type": metadata.mimeType,
},
body: JSON.stringify(metadata.existingFileId ? { name: metadata.name } : { name: metadata.name, parents: metadata.parents }),
});
const uploadUrl = response.headers.get("Location");
if (!response.ok || !uploadUrl) throw new Error(`Failed to create Drive resumable session (${response.status})`);
return uploadUrl;
}
export function nextDriveOffset(response: Response): number {
const range = response.headers.get("Range");
if (!range) return 0;
const match = /^bytes=0-(\d+)$/.exec(range);
if (!match) throw new Error("Drive returned an invalid committed range");
return Number(match[1]) + 1;
}
export async function queryStatus(uploadUrl: string, accessToken: string): Promise<number> {
const response = await fetch(uploadUrl, {
method: "PUT",
headers: { Authorization: `Bearer ${accessToken}`, "Content-Length": "0", "Content-Range": "bytes */*" },
});
if (response.status === 308) return nextDriveOffset(response);
if (response.ok) {
const metadata = await response.json<DriveUploadResult>();
return Number(metadata.size ?? 0);
}
throw new Error(`Drive resumable status query failed (${response.status})`);
}
export async function putFinalChunk(uploadUrl: string, accessToken: string, startOffset: number, total: number, body: Uint8Array): Promise<DriveUploadResult> {
const response = await fetch(uploadUrl, {
method: "PUT",
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Length": body.byteLength.toString(),
"Content-Range": `bytes ${startOffset}-${total - 1}/${total}`,
},
body,
});
if (!response.ok) throw new Error(`Drive final chunk failed (${response.status}): ${await response.text()}`);
return await response.json<DriveUploadResult>();
}
export async function createEmptyFile(accessToken: string, metadata: { name: string; parents: string[]; mimeType: string; existingFileId?: string }): Promise<DriveUploadResult> {
const url = metadata.existingFileId ? `https://www.googleapis.com/drive/v3/files/${metadata.existingFileId}?fields=${encodeURIComponent(DRIVE_FIELDS)}` : `https://www.googleapis.com/drive/v3/files?fields=${encodeURIComponent(DRIVE_FIELDS)}`;
const response = await fetch(url, {
method: metadata.existingFileId ? "PATCH" : "POST",
headers: { Authorization: `Bearer ${accessToken}`, "Content-Type": "application/json" },
body: JSON.stringify(metadata.existingFileId ? { name: metadata.name } : { name: metadata.name, parents: metadata.parents, mimeType: metadata.mimeType }),
});
if (!response.ok) throw new Error(`Drive empty-file creation failed (${response.status})`);
return await response.json<DriveUploadResult>();
}
export async function cancelSession(uploadUrl: string, accessToken: string): Promise<void> {
const response = await fetch(uploadUrl, { method: "DELETE", headers: { Authorization: `Bearer ${accessToken}` } });
if (response.status !== 499 && !response.ok && response.status !== 404) throw new Error(`Drive resumable cancellation failed (${response.status})`);
}
+270
View File
@@ -0,0 +1,270 @@
import { decodedContentLength, isAwsChunked, pumpBody } from "./aws-chunked";
import type { DriveDownloadResult, DriveFileMetadata, DriveUploadResult, Env, GoogleDriveFile, GoogleDriveSearchResponse } from "./types";
interface GoogleTokenResponse {
access_token: string;
expires_in: number;
error_description?: string;
}
interface GoogleDriveCreateResponse {
id: string;
}
const DRIVE_FIELDS = "id,name,size,mimeType,md5Checksum";
function driveLiteral(value: string): string {
return value.replace(/\\/g, "\\\\").replace(/'/g, "\\'");
}
function driveFilesUrl(q: string, fields: string): string {
const url = new URL("https://www.googleapis.com/drive/v3/files");
url.searchParams.set("q", q);
url.searchParams.set("fields", fields);
return url.toString();
}
/** Fetches an OAuth access token, using the KV-cached one when available. */
export async function getAccessToken(env: Env): Promise<string> {
const cacheKey = "google_access_token";
const cachedToken = await env.AUTH_KV.get(cacheKey);
if (cachedToken) {
return cachedToken;
}
const response = await fetch("https://oauth2.googleapis.com/token", {
method: "POST",
headers: { "Content-Type": "application/x-www-form-urlencoded" },
body: new URLSearchParams({
client_id: env.GOOGLE_CLIENT_ID,
client_secret: env.GOOGLE_CLIENT_SECRET,
refresh_token: env.GOOGLE_REFRESH_TOKEN,
grant_type: "refresh_token",
}),
});
const data: GoogleTokenResponse = await response.json();
if (!response.ok) {
throw new Error(`Token Error: ${data.error_description}`);
}
await env.AUTH_KV.put(cacheKey, data.access_token, {
expirationTtl: data.expires_in - 60,
});
return data.access_token;
}
/** Finds a folder by name under the given parent, creating it if it doesn't exist yet. */
async function getOrCreateFolder(accessToken: string, folderName: string, parentId: string | null, env: Env): Promise<string> {
// Include parentId in the cache key so folders with the same name in different parents don't collide.
const cacheKey = parentId ? `${parentId}/${folderName}` : folderName;
const cached = await env.FOLDER_CACHE.get(cacheKey);
if (cached) return cached;
const parentQuery = parentId ? ` and '${parentId}' in parents` : "";
const searchRes = await fetch(driveFilesUrl(`name='${driveLiteral(folderName)}' and mimeType='application/vnd.google-apps.folder' and trashed=false${parentQuery}`, "files(id,name)"), {
headers: { Authorization: `Bearer ${accessToken}` },
});
const searchData: GoogleDriveSearchResponse = await searchRes.json();
if (searchData.files && searchData.files.length > 0) {
const folderId = searchData.files[0].id;
await env.FOLDER_CACHE.put(cacheKey, folderId, { expirationTtl: 3600 });
return folderId;
}
const createBody: { name: string; mimeType: string; parents?: string[] } = {
name: folderName,
mimeType: "application/vnd.google-apps.folder",
};
if (parentId) {
createBody.parents = [parentId];
}
const createRes = await fetch("https://www.googleapis.com/drive/v3/files", {
method: "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Type": "application/json",
},
body: JSON.stringify(createBody),
});
const createData: GoogleDriveCreateResponse = await createRes.json();
await env.FOLDER_CACHE.put(cacheKey, createData.id, { expirationTtl: 3600 });
return createData.id;
}
/** 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 getOrCreateFolder(accessToken, bucket, null, env);
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) {
currentFolderId = await getOrCreateFolder(accessToken, dir, currentFolderId, env);
}
return {
parentFolderId: currentFolderId,
fileName: fileName,
};
}
export async function streamUploadToDrive(accessToken: string, request: Request, bucket: string, objectKey: string, mimeType: string, env: Env): Promise<DriveUploadResult> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const existing = await findFileInFolder(accessToken, parentFolderId, fileName);
const initUrl = existing ? `https://www.googleapis.com/upload/drive/v3/files/${existing.id}?uploadType=resumable&fields=${encodeURIComponent(DRIVE_FIELDS)}` : `https://www.googleapis.com/upload/drive/v3/files?uploadType=resumable&fields=${encodeURIComponent(DRIVE_FIELDS)}`;
const decodedLength = decodedContentLength(request);
// Initialize a resumable upload session.
const initRes = await fetch(initUrl, {
method: existing ? "PATCH" : "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
"X-Upload-Content-Type": mimeType,
...(decodedLength === undefined ? {} : { "X-Upload-Content-Length": decodedLength.toString() }),
"Content-Type": "application/json; charset=UTF-8",
},
body: JSON.stringify(existing ? { name: fileName } : { name: fileName, parents: [parentFolderId] }),
});
const uploadUrl = initRes.headers.get("Location");
if (!uploadUrl) {
console.error(initRes.status);
console.error(await initRes.text());
throw new Error("Failed to get upload URL");
}
let uploadRes: Response;
if (isAwsChunked(request)) {
if (decodedLength === undefined) throw new Error("x-amz-decoded-content-length is required for aws-chunked uploads");
const decoded = new FixedLengthStream(decodedLength, { highWaterMark: 1 << 20 });
const uploadPromise = fetch(uploadUrl, {
method: "PUT",
headers: { Authorization: `Bearer ${accessToken}`, "Content-Length": decodedLength.toString() },
body: decoded.readable,
duplex: "half",
} as RequestInit);
await pumpBody(request.body, decoded.writable.getWriter(), { awsChunked: true, expectedLength: decodedLength });
uploadRes = await uploadPromise;
} else {
const body = request.body ?? new Uint8Array();
uploadRes = await fetch(uploadUrl, {
method: "PUT",
headers: { Authorization: `Bearer ${accessToken}`, ...(decodedLength === undefined ? {} : { "Content-Length": decodedLength.toString() }) },
body,
duplex: "half",
} as RequestInit);
}
if (!uploadRes.ok) {
const errorText = await uploadRes.text();
throw new Error(`Upload failed: ${errorText}`);
}
return await uploadRes.json();
}
export async function findFileInFolder(accessToken: string, folderId: string, fileName: string): Promise<GoogleDriveFile | null> {
const searchRes = await fetch(driveFilesUrl(`name='${driveLiteral(fileName)}' and '${driveLiteral(folderId)}' in parents and trashed=false`, `files(${DRIVE_FIELDS})`), {
headers: { Authorization: `Bearer ${accessToken}` },
});
const data: GoogleDriveSearchResponse = await searchRes.json();
return data.files && data.files.length > 0 ? data.files[0] : null;
}
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 file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
const controller = new AbortController();
const timeout = setTimeout(() => {
controller.abort();
}, 30000);
const downloadRes = await fetch(`https://www.googleapis.com/drive/v3/files/${file.id}?alt=media`, {
headers: { Authorization: `Bearer ${accessToken}`, ...(range ? { Range: range } : {}) },
signal: controller.signal,
});
clearTimeout(timeout);
if (!downloadRes.ok) {
console.error(downloadRes.status);
console.error(await downloadRes.text());
throw new Error("Download failed");
}
return {
body: downloadRes.body!,
contentType: file.mimeType || "application/octet-stream",
size: parseInt(file.size || "0", 10),
id: file.id,
md5Checksum: file.md5Checksum,
status: downloadRes.status,
contentRange: downloadRes.headers.get("Content-Range") ?? undefined,
contentLength: downloadRes.headers.get("Content-Length") ?? undefined,
};
}
export async function deleteFromDrive(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<void> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
const deleteRes = await fetch(`https://www.googleapis.com/drive/v3/files/${file.id}`, {
method: "DELETE",
headers: { Authorization: `Bearer ${accessToken}` },
});
if (!deleteRes.ok) {
throw new Error("Delete failed");
}
}
export async function getFileMetadata(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<DriveFileMetadata> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
return {
id: file.id,
mimeType: file.mimeType || "application/octet-stream",
size: parseInt(file.size || "0", 10),
md5Checksum: file.md5Checksum,
};
}
export async function listFiles(accessToken: string, bucket: string, env: Env): Promise<GoogleDriveFile[]> {
const folderId = await getOrCreateFolder(accessToken, bucket, null, env);
const listRes = await fetch(driveFilesUrl(`'${driveLiteral(folderId)}' in parents and trashed=false`, "files(id,name,mimeType,size,modifiedTime,md5Checksum)"), {
headers: { Authorization: `Bearer ${accessToken}` },
});
const data: GoogleDriveSearchResponse = await listRes.json();
return data.files || [];
}
+28 -573
View File
@@ -1,581 +1,36 @@
/** import { verifySignature } from "./aws-signature";
* S3-Compatible API Server on Cloudflare Workers import { isAllowedBucket, isPublicReadBucket } from "./bucket-access";
* Backend: Google Drive with streaming support and nested directory structure import { getAccessToken } from "./google-drive";
*/ import { MultipartUploadDO } from "./multipart-do";
import { dispatch } from "./router";
import { S3Exception, s3Error } from "./s3-errors";
import type { Env } from "./types";
export { MultipartUploadDO };
export default { export default {
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> { async fetch(request: Request, env: Env, _ctx: ExecutionContext): Promise<Response> {
if (request.method === "OPTIONS") return new Response(null, { status: 204 });
const url = new URL(request.url);
const pathParts = url.pathname.split("/").filter(Boolean);
const bucket = pathParts[0] || "";
const objectKey = pathParts.slice(1).join("/");
const resource = url.pathname || "/";
try { try {
if (request.method === "OPTIONS") { if (!isAllowedBucket(bucket, env)) return s3Error("AccessDenied", 403, undefined, resource, request.method === "HEAD");
return new Response(null, { status: 204 });
const isPublicRead = isPublicReadBucket(bucket, env) && (request.method === "GET" || request.method === "HEAD");
if (!isPublicRead && !(await verifySignature(request, env))) {
return s3Error("SignatureDoesNotMatch", 403, undefined, resource, request.method === "HEAD");
} }
const url = new URL(request.url); return await dispatch(request, env, await getAccessToken(env), bucket, objectKey);
const method = request.method; } catch (error) {
if (error instanceof S3Exception) return s3Error(error.code, error.status, error.message, resource, request.method === "HEAD", error.headers);
// パスからバケット名とオブジェクトキーを抽出 console.error(JSON.stringify({ message: "request failed", error: error instanceof Error ? error.message : String(error), method: request.method, path: url.pathname }));
const pathParts = url.pathname.split("/").filter((p) => p); return s3Error("InternalError", 500, undefined, resource, request.method === "HEAD");
const bucket = pathParts[0] || "";
const objectKey = pathParts.slice(1).join("/");
if (!isAllowedBucket(bucket, env)) {
return new Response("Access denied to this bucket", { status: 403 });
}
const isPublicRead = isPublicReadBucket(bucket, env);
const isReadMethod = method === "GET" || method === "HEAD";
if (!(isPublicRead && isReadMethod)) {
const isValid = await verifySignature(request, env);
if (!isValid) {
return new Response("Invalid Signature", { status: 403 });
}
}
const accessToken = await getAccessToken(env);
if (method === "PUT" || method === "POST") {
// ストリーミングアップロード
if (!objectKey) {
return new Response("Object key required", { status: 400 });
}
const contentType = request.headers.get("Content-Type") || "application/octet-stream";
const result = await streamUploadToDrive(accessToken, request.body, bucket, objectKey, contentType, env);
return new Response(JSON.stringify(result), {
status: 200,
headers: {
"Content-Type": "application/json",
ETag: `"${result.id}"`, // Google DriveのファイルIDをETagとして使用
},
});
} else if (method === "GET") {
// ストリーミングダウンロード
if (!objectKey) {
// バケット(フォルダ)の一覧を返す
const files = await listFiles(accessToken, bucket, env);
return new Response(generateListBucketResult(files, bucket), {
status: 200,
headers: { "Content-Type": "application/xml" },
});
}
try {
const fileStream = await streamDownloadFromDrive(accessToken, bucket, objectKey, env);
return new Response(fileStream.body, {
status: 200,
headers: {
"Content-Type": fileStream.contentType,
"Content-Length": fileStream.size.toString(),
"Cache-Control": "s-maxage=300, no-store",
ETag: `"${fileStream.id}"`,
},
});
} catch (e) {
const error = e as Error;
if (error.message === "File not found") {
return new Response("NoSuchKey", { status: 404 });
}
throw e;
}
} else if (method === "DELETE") {
// ファイル削除
if (!objectKey) {
return new Response("Object key required", { status: 400 });
}
try {
await deleteFromDrive(accessToken, bucket, objectKey, env);
return new Response(null, { status: 204 });
} catch (e) {
const error = e as Error;
if (error.message === "File not found") {
return new Response(null, { status: 404 });
}
throw e;
}
} else if (method === "HEAD") {
// メタデータ取得
if (!objectKey) {
return new Response(null, { status: 400 });
}
try {
const metadata = await getFileMetadata(accessToken, bucket, objectKey, env);
return new Response(null, {
status: 200,
headers: {
"Content-Type": metadata.mimeType,
"Content-Length": metadata.size.toString(),
ETag: `"${metadata.id}"`,
},
});
} catch (e) {
const error = e as Error;
if (error.message === "File not found") {
return new Response(null, { status: 404 });
}
throw e;
}
}
return new Response("Method not allowed", { status: 405 });
} catch (e) {
const error = e as Error;
console.error("Error:", error);
return new Response(error.message, { status: 500 });
} }
}, },
} satisfies ExportedHandler<Env>; } satisfies ExportedHandler<Env>;
interface Env {
ACCESS_KEY: string;
SECRET_KEY: string;
REGION: string;
GOOGLE_CLIENT_ID: string;
GOOGLE_CLIENT_SECRET: string;
GOOGLE_REFRESH_TOKEN: string;
AUTH_KV: KVNamespace;
FOLDER_CACHE: KVNamespace;
ALLOWED_BUCKETS?: string;
PUBLIC_READ_BUCKETS?: string;
}
interface GoogleDriveFile {
id: string;
name: string;
mimeType: string;
size: string;
modifiedTime?: string;
}
interface GoogleDriveSearchResponse {
files?: GoogleDriveFile[];
}
function isAllowedBucket(bucket: string, env: Env): boolean {
console.log(bucket);
// 許可リストが設定されていない場合はすべて拒否
if (!env.ALLOWED_BUCKETS) {
return false;
}
const allowedBuckets = env.ALLOWED_BUCKETS.split(",")
.map((b) => b.trim())
.filter((b) => b);
// 空の許可リストの場合もすべて拒否
if (allowedBuckets.length === 0) {
return false;
}
// バケット名が許可リストに含まれているかチェック
return allowedBuckets.includes(bucket);
}
function isPublicReadBucket(bucket: string, env: Env): boolean {
if (!env.PUBLIC_READ_BUCKETS) {
return false;
}
const publicReadBuckets = env.PUBLIC_READ_BUCKETS.split(",")
.map((b) => b.trim())
.filter((b) => b);
return publicReadBuckets.includes(bucket);
}
// ========================================
// Google Drive API Functions
// ========================================
async function getAccessToken(env: Env): Promise<string> {
const cacheKey = "google_access_token";
const cachedToken = await env.AUTH_KV.get(cacheKey);
if (cachedToken) {
return cachedToken;
}
const response = await fetch("https://oauth2.googleapis.com/token", {
method: "POST",
headers: { "Content-Type": "application/x-www-form-urlencoded" },
body: new URLSearchParams({
client_id: env.GOOGLE_CLIENT_ID,
client_secret: env.GOOGLE_CLIENT_SECRET,
refresh_token: env.GOOGLE_REFRESH_TOKEN,
grant_type: "refresh_token",
}),
});
const data: any = await response.json();
if (!response.ok) {
throw new Error(`Token Error: ${data.error_description}`);
}
await env.AUTH_KV.put(cacheKey, data.access_token, {
expirationTtl: data.expires_in - 60,
});
return data.access_token;
}
async function getOrCreateFolder(accessToken: string, folderName: string, parentId: string | null, env: Env): Promise<string> {
// キャッシュキーにparentIdを含める
const cacheKey = parentId ? `${parentId}/${folderName}` : folderName;
const cached = await env.FOLDER_CACHE.get(cacheKey);
if (cached) return cached;
// フォルダを検索(親フォルダを指定)
const parentQuery = parentId ? ` and '${parentId}' in parents` : "";
const searchRes = await fetch(`https://www.googleapis.com/drive/v3/files?q=name='${encodeURIComponent(folderName)}' and mimeType='application/vnd.google-apps.folder' and trashed=false${parentQuery}`, {
headers: { Authorization: `Bearer ${accessToken}` },
});
const searchData: GoogleDriveSearchResponse = await searchRes.json();
if (searchData.files && searchData.files.length > 0) {
const folderId = searchData.files[0].id;
await env.FOLDER_CACHE.put(cacheKey, folderId, { expirationTtl: 3600 });
return folderId;
}
// フォルダが存在しない場合は作成
const createBody: any = {
name: folderName,
mimeType: "application/vnd.google-apps.folder",
};
if (parentId) {
createBody.parents = [parentId];
}
const createRes = await fetch("https://www.googleapis.com/drive/v3/files", {
method: "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
"Content-Type": "application/json",
},
body: JSON.stringify(createBody),
});
const createData: any = await createRes.json();
await env.FOLDER_CACHE.put(cacheKey, createData.id, { expirationTtl: 3600 });
return createData.id;
}
async function resolvePathToFolderAndFile(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<{ parentFolderId: string; fileName: string }> {
// バケットのルートフォルダを取得
let currentFolderId = await getOrCreateFolder(accessToken, bucket, null, env);
// objectKeyを/で分割
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) {
currentFolderId = await getOrCreateFolder(accessToken, dir, currentFolderId, env);
}
return {
parentFolderId: currentFolderId,
fileName: fileName,
};
}
async function streamUploadToDrive(accessToken: string, stream: ReadableStream | null, bucket: string, objectKey: string, mimeType: string, env: Env): Promise<any> {
if (!stream) {
throw new Error("Request body is required");
}
// パスを解析してフォルダ階層を作成
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
// Resumable uploadの初期化
const initRes = await fetch("https://www.googleapis.com/upload/drive/v3/files?uploadType=resumable", {
method: "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
"X-Upload-Content-Type": mimeType,
"Content-Type": "application/json; charset=UTF-8",
},
body: JSON.stringify({
name: fileName,
parents: [parentFolderId],
}),
});
const uploadUrl = initRes.headers.get("Location");
if (!uploadUrl) {
console.error(initRes.status);
console.error(await initRes.text());
throw new Error("Failed to get upload URL");
}
// ストリーミングアップロード
const uploadRes = await fetch(uploadUrl, {
method: "PUT",
headers: {
Authorization: `Bearer ${accessToken}`,
},
body: stream,
duplex: "half",
} as RequestInit);
if (!uploadRes.ok) {
const errorText = await uploadRes.text();
throw new Error(`Upload failed: ${errorText}`);
}
return await uploadRes.json();
}
async function findFileInFolder(accessToken: string, folderId: string, fileName: string): Promise<GoogleDriveFile | null> {
const searchRes = await fetch(`https://www.googleapis.com/drive/v3/files?q=name='${encodeURIComponent(fileName)}' and '${folderId}' in parents and trashed=false&fields=files(id,name,mimeType,size)`, {
headers: { Authorization: `Bearer ${accessToken}` },
});
const data: GoogleDriveSearchResponse = await searchRes.json();
return data.files && data.files.length > 0 ? data.files[0] : null;
}
async function streamDownloadFromDrive(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<{ body: ReadableStream; contentType: string; size: number; id: string }> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
const controller = new AbortController();
const timeout = setTimeout(() => {
controller.abort();
}, 5000);
const downloadRes = await fetch(`https://www.googleapis.com/drive/v3/files/${file.id}?alt=media`, {
headers: { Authorization: `Bearer ${accessToken}` },
signal: controller.signal,
});
clearTimeout(timeout);
if (!downloadRes.ok) {
console.error(downloadRes.status);
console.error(await downloadRes.text());
throw new Error("Download failed");
}
return {
body: downloadRes.body!,
contentType: file.mimeType || "application/octet-stream",
size: parseInt(file.size || "0"),
id: file.id,
};
}
async function deleteFromDrive(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<void> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
const deleteRes = await fetch(`https://www.googleapis.com/drive/v3/files/${file.id}`, {
method: "DELETE",
headers: { Authorization: `Bearer ${accessToken}` },
});
if (!deleteRes.ok) {
throw new Error("Delete failed");
}
}
async function getFileMetadata(accessToken: string, bucket: string, objectKey: string, env: Env): Promise<{ id: string; mimeType: string; size: number }> {
const { parentFolderId, fileName } = await resolvePathToFolderAndFile(accessToken, bucket, objectKey, env);
const file = await findFileInFolder(accessToken, parentFolderId, fileName);
if (!file) {
throw new Error("File not found");
}
return {
id: file.id,
mimeType: file.mimeType || "application/octet-stream",
size: parseInt(file.size || "0"),
};
}
async function listFiles(accessToken: string, bucket: string, env: Env): Promise<GoogleDriveFile[]> {
const folderId = await getOrCreateFolder(accessToken, bucket, null, env);
const listRes = await fetch(`https://www.googleapis.com/drive/v3/files?q='${folderId}' in parents and trashed=false&fields=files(id,name,mimeType,size,modifiedTime)`, {
headers: { Authorization: `Bearer ${accessToken}` },
});
const data: GoogleDriveSearchResponse = await listRes.json();
return data.files || [];
}
function generateListBucketResult(files: any[], bucket: string): string {
const contents = files
.map(
(f) => `
<Contents>
<Key>${escapeXml(f.name)}</Key>
<LastModified>${f.modifiedTime || new Date().toISOString()}</LastModified>
<ETag>"${f.id}"</ETag>
<Size>${f.size || 0}</Size>
<StorageClass>STANDARD</StorageClass>
</Contents>`,
)
.join("");
return `<?xml version="1.0" encoding="UTF-8"?>
<ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/">
<Name>${escapeXml(bucket)}</Name>
<Prefix></Prefix>
<MaxKeys>1000</MaxKeys>
<IsTruncated>false</IsTruncated>
${contents}
</ListBucketResult>`;
}
function escapeXml(str: string): string {
return str.replace(/&/g, "&amp;").replace(/</g, "&lt;").replace(/>/g, "&gt;").replace(/"/g, "&quot;").replace(/'/g, "&apos;");
}
// ========================================
// AWS Signature V4 Verification
// ========================================
async function verifySignature(request: Request, env: Env): Promise<boolean> {
const url = new URL(request.url);
const headers = request.headers;
const isQueryAuth = url.searchParams.has("X-Amz-Algorithm");
let algorithm: string;
if (isQueryAuth) {
algorithm = url.searchParams.get("X-Amz-Algorithm") ?? "";
} else {
const authHeader = headers.get("Authorization") ?? "";
algorithm = authHeader.split(" ")[0];
}
if (!algorithm || !algorithm.includes("AWS4-HMAC-SHA256")) {
return false;
}
const datetime = (isQueryAuth ? url.searchParams.get("X-Amz-Date") : headers.get("x-amz-date")) ?? "";
if (!datetime) return false;
const date = datetime.substring(0, 8);
const canonicalRequest = await createCanonicalRequest(request, isQueryAuth);
const hashedCanonicalRequest = await sha256(canonicalRequest);
const credentialScope = `${date}/${env.REGION}/s3/aws4_request`;
const stringToSign = ["AWS4-HMAC-SHA256", datetime, credentialScope, hashedCanonicalRequest].join("\n");
const signingKey = await getSigningKey(env.SECRET_KEY, date, env.REGION, "s3");
const signature = await hmacSha256(signingKey, stringToSign);
const signatureHex = bufToHex(signature);
let expectedSignature = "";
if (isQueryAuth) {
expectedSignature = url.searchParams.get("X-Amz-Signature") ?? "";
} else {
const authHeader = headers.get("Authorization") ?? "";
const match = authHeader.match(/Signature=([a-f0-9]+)/);
expectedSignature = match ? match[1] : "";
}
return signatureHex === expectedSignature;
}
async function createCanonicalRequest(request: Request, isQueryAuth: boolean): Promise<string> {
const url = new URL(request.url);
const method = request.method;
const canonicalUri = url.pathname || "/";
const params = Array.from(url.searchParams.entries())
.filter(([key]) => key !== "X-Amz-Signature")
.sort(([a], [b]) => {
if (a < b) return -1;
if (a > b) return 1;
return 0;
})
.map(([key, val]) => `${encodeRFC3986(key)}=${encodeRFC3986(val)}`)
.join("&");
let signedHeadersList: string[];
if (isQueryAuth) {
signedHeadersList = (url.searchParams.get("X-Amz-SignedHeaders") ?? "host").split(";");
} else {
const authHeader = request.headers.get("Authorization") ?? "";
const match = authHeader.match(/SignedHeaders=([^,\s]+)/);
signedHeadersList = match ? match[1].split(";") : ["host"];
}
const canonicalHeaders = signedHeadersList
.map((h) => {
const headerName = h.toLowerCase();
let headerValue = "";
if (headerName === "host") {
headerValue = url.hostname;
const port = url.port;
if (port && !((url.protocol === "https:" && port === "443") || (url.protocol === "http:" && port === "80"))) {
headerValue += `:${port}`;
}
} else {
headerValue = request.headers.get(headerName)?.trim() ?? "";
}
return `${headerName}:${headerValue}\n`;
})
.join("");
const signedHeaders = signedHeadersList.join(";");
const payloadHash = request.headers.get("x-amz-content-sha256") ?? (isQueryAuth ? "UNSIGNED-PAYLOAD" : "UNSIGNED-PAYLOAD");
return [method, canonicalUri, params, canonicalHeaders, signedHeaders, payloadHash].join("\n");
}
function encodeRFC3986(str: string): string {
return encodeURIComponent(str).replace(/[!'()*]/g, (c) => `%${c.charCodeAt(0).toString(16).toUpperCase()}`);
}
async function getSigningKey(secret: string, date: string, region: string, service: string): Promise<ArrayBuffer> {
const kDate = await hmacSha256("AWS4" + secret, date);
const kRegion = await hmacSha256(kDate, region);
const kService = await hmacSha256(kRegion, service);
return await hmacSha256(kService, "aws4_request");
}
async function hmacSha256(key: string | ArrayBuffer, data: string): Promise<ArrayBuffer> {
const keyData = typeof key === "string" ? new TextEncoder().encode(key) : key;
const cryptoKey = await crypto.subtle.importKey("raw", keyData, { name: "HMAC", hash: "SHA-256" }, false, ["sign"]);
return await crypto.subtle.sign("HMAC", cryptoKey, new TextEncoder().encode(data));
}
async function sha256(data: string): Promise<string> {
const hash = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(data));
return bufToHex(hash);
}
function bufToHex(buf: ArrayBuffer): string {
return Array.from(new Uint8Array(buf))
.map((b) => b.toString(16).padStart(2, "0"))
.join("");
}
+282
View File
@@ -0,0 +1,282 @@
import { alignedSendLen, cancelSession, createEmptyFile, putFinalChunk, queryStatus } from "./drive-resumable";
import { getAccessToken } from "./google-drive";
import type { DriveUploadResult, Env } from "./types";
import { DurableObject } from "cloudflare:workers";
const WAIT_TIMEOUT_MS = 20_000;
const LEASE_MS = 120_000;
const EXPIRE_AFTER_MS = 24 * 60 * 60 * 1000;
interface PartRow extends Record<string, SqlStorageValue> {
partNumber: number;
size: number;
etag: string;
endOffset: number;
ts: number;
}
interface InFlight {
requestId: string;
partNumber: number;
driveOffsetAtStart: number;
carryLen: number;
partLen: number;
sendLen: number;
leaseExpiresAt: number;
}
export type BeginPartResult = { kind: "admit"; uploadUrl: string; driveOffset: number; carry: Uint8Array; sendLen: number; skipBytes: number } | { kind: "committed"; etag: string } | { kind: "slowdown" } | { kind: "error"; code: "NoSuchUpload" | "InvalidPart" | "InternalError"; message: string };
export type CompleteResult = { kind: "complete"; metadata: DriveUploadResult; partEtags: string[] } | { kind: "error"; code: "NoSuchUpload" | "InvalidPart" | "InvalidPartOrder" | "InternalError"; message: string };
interface Waiter {
requestId: string;
partNumber: number;
partLen: number;
resolve: (result: BeginPartResult) => void;
timer: ReturnType<typeof setTimeout>;
}
interface StateValueRow extends Record<string, SqlStorageValue> {
v: ArrayBuffer | string | number | null;
}
export class MultipartUploadDO extends DurableObject<Env> {
private readonly waiters = new Map<string, Waiter>();
constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
void ctx.blockConcurrencyWhile(async () => {
this.ctx.storage.sql.exec("CREATE TABLE IF NOT EXISTS state (k TEXT PRIMARY KEY, v)");
this.ctx.storage.sql.exec("CREATE TABLE IF NOT EXISTS parts (partNumber INTEGER PRIMARY KEY, size INTEGER NOT NULL, etag TEXT NOT NULL, endOffset INTEGER NOT NULL, ts INTEGER NOT NULL)");
});
}
private getValue<T extends ArrayBuffer | string | number>(key: string): T | undefined {
return this.ctx.storage.sql.exec<StateValueRow>("SELECT v FROM state WHERE k = ?", key).toArray()[0]?.v as T | undefined;
}
private setValue(key: string, value: ArrayBuffer | string | number | null): void {
this.ctx.storage.sql.exec("INSERT INTO state (k, v) VALUES (?, ?) ON CONFLICT(k) DO UPDATE SET v = excluded.v", key, value);
}
private hasUpload(): boolean {
return this.getValue<string>("uploadUrl") !== undefined;
}
private carry(): Uint8Array {
const value = this.getValue<ArrayBuffer>("carry");
return value ? new Uint8Array(value) : new Uint8Array();
}
private inFlight(): InFlight | null {
const value = this.getValue<string>("inFlight");
return value ? (JSON.parse(value) as InFlight) : null;
}
private part(partNumber: number): PartRow | undefined {
return this.ctx.storage.sql.exec<PartRow>("SELECT partNumber, size, etag, endOffset, ts FROM parts WHERE partNumber = ?", partNumber).toArray()[0];
}
async init(input: { uploadUrl: string; bucket: string; key: string; mimeType: string; parentFolderId: string; fileName: string; existingFileId?: string }): Promise<boolean> {
if (this.hasUpload()) return false;
this.ctx.storage.sql.exec("DELETE FROM parts");
for (const [key, value] of Object.entries(input)) {
if (value !== undefined) this.setValue(key, value);
}
this.setValue("driveOffset", 0);
this.setValue("carry", new ArrayBuffer(0));
this.setValue("nextExpectedPart", 1);
this.setValue("createdAt", Date.now());
this.setValue("inFlight", null);
await this.ctx.storage.setAlarm(Date.now() + EXPIRE_AFTER_MS);
return true;
}
private async resyncExpiredLease(inFlight: InFlight, requestId: string, partLen: number): Promise<BeginPartResult | null> {
if (partLen !== inFlight.partLen) return { kind: "error", code: "InvalidPart", message: "A retried part must have the same decoded length" };
const accessToken = await getAccessToken(this.env);
const actualOffset = await queryStatus(this.getValue<string>("uploadUrl")!, accessToken);
const partStartFileOffset = inFlight.driveOffsetAtStart + inFlight.carryLen;
let skipBytes = 0;
if (actualOffset === inFlight.driveOffsetAtStart) {
// The persisted carry remains valid.
} else if (actualOffset >= partStartFileOffset) {
skipBytes = actualOffset - partStartFileOffset;
if (skipBytes > inFlight.partLen) return { kind: "error", code: "InternalError", message: "Drive committed beyond the leased part" };
this.setValue("driveOffset", actualOffset);
this.setValue("carry", new ArrayBuffer(0));
} else {
return { kind: "error", code: "InternalError", message: "Drive committed only a prefix of the multipart carry" };
}
this.setValue("inFlight", null);
return this.admit(requestId, inFlight.partNumber, inFlight.partLen, skipBytes);
}
private admit(requestId: string, partNumber: number, partLen: number, skipBytes = 0): BeginPartResult {
const carry = this.carry();
const driveOffset = this.getValue<number>("driveOffset") ?? 0;
const available = carry.byteLength + partLen - skipBytes;
const sendLen = alignedSendLen(driveOffset, available);
const inFlight: InFlight = {
requestId,
partNumber,
driveOffsetAtStart: driveOffset,
carryLen: carry.byteLength,
partLen,
sendLen,
leaseExpiresAt: Date.now() + LEASE_MS,
};
this.setValue("inFlight", JSON.stringify(inFlight));
return { kind: "admit", uploadUrl: this.getValue<string>("uploadUrl")!, driveOffset, carry, sendLen, skipBytes };
}
async beginPart(requestId: string, partNumber: number, partLen: number): Promise<BeginPartResult> {
if (!this.hasUpload()) return { kind: "error", code: "NoSuchUpload", message: "Multipart upload not found" };
if (!Number.isInteger(partNumber) || partNumber < 1 || partNumber > 10_000 || !Number.isSafeInteger(partLen) || partLen < 0) {
return { kind: "error", code: "InvalidPart", message: "Invalid part number or length" };
}
const committed = this.part(partNumber);
if (committed) return { kind: "committed", etag: committed.etag };
const nextExpected = this.getValue<number>("nextExpectedPart") ?? 1;
if (partNumber < nextExpected) return { kind: "error", code: "InvalidPart", message: "Part is no longer available" };
if (partNumber > nextExpected + 64) return { kind: "slowdown" };
const lease = this.inFlight();
if (partNumber === nextExpected && (!lease || lease.leaseExpiresAt <= Date.now())) {
if (lease) {
try {
const result = await this.resyncExpiredLease(lease, requestId, partLen);
if (result) return result;
} catch (error) {
return { kind: "error", code: "InternalError", message: error instanceof Error ? error.message : String(error) };
}
}
return this.admit(requestId, partNumber, partLen);
}
return await new Promise<BeginPartResult>((resolve) => {
const timer = setTimeout(() => {
this.waiters.delete(requestId);
resolve({ kind: "slowdown" });
}, WAIT_TIMEOUT_MS);
this.waiters.set(requestId, { requestId, partNumber, partLen, resolve, timer });
});
}
async cancelWaiter(requestId: string): Promise<void> {
const waiter = this.waiters.get(requestId);
if (!waiter) return;
clearTimeout(waiter.timer);
this.waiters.delete(requestId);
waiter.resolve({ kind: "slowdown" });
}
async failPart(partNumber: number): Promise<void> {
const lease = this.inFlight();
if (!lease || lease.partNumber !== partNumber) return;
lease.leaseExpiresAt = 0;
this.setValue("inFlight", JSON.stringify(lease));
await this.wakeWaiters();
}
async endPart(partNumber: number, newDriveOffset: number, tail: Uint8Array, etag: string, partLen: number): Promise<boolean> {
const lease = this.inFlight();
if (!lease || lease.partNumber !== partNumber) return false;
if (newDriveOffset !== lease.driveOffsetAtStart + lease.sendLen || tail.byteLength > 256 * 1024) return false;
const tailBuffer = new Uint8Array(tail).buffer;
this.ctx.storage.transactionSync(() => {
this.setValue("driveOffset", newDriveOffset);
this.setValue("carry", tailBuffer);
this.setValue("nextExpectedPart", partNumber + 1);
this.setValue("inFlight", null);
this.ctx.storage.sql.exec("INSERT INTO parts (partNumber, size, etag, endOffset, ts) VALUES (?, ?, ?, ?, ?)", partNumber, partLen, etag, newDriveOffset + tail.byteLength, Date.now());
});
await this.wakeWaiters();
return true;
}
private async wakeWaiters(): Promise<void> {
const nextExpected = this.getValue<number>("nextExpectedPart") ?? 1;
if (this.inFlight()) return;
const waiter = [...this.waiters.values()].find((candidate) => candidate.partNumber === nextExpected);
if (!waiter) return;
clearTimeout(waiter.timer);
this.waiters.delete(waiter.requestId);
waiter.resolve(this.admit(waiter.requestId, waiter.partNumber, waiter.partLen));
}
async complete(clientParts: Array<{ partNumber: number; etag: string }>, expectedTotal?: number): Promise<CompleteResult> {
if (!this.hasUpload()) return { kind: "error", code: "NoSuchUpload", message: "Multipart upload not found" };
if (this.inFlight()) return { kind: "error", code: "InvalidPart", message: "A part is still being uploaded" };
const stored = this.ctx.storage.sql.exec<PartRow>("SELECT partNumber, size, etag, endOffset, ts FROM parts ORDER BY partNumber").toArray();
if (clientParts.length !== stored.length) return { kind: "error", code: "InvalidPart", message: "The completed part list does not match uploaded parts" };
for (let index = 0; index < clientParts.length; index++) {
const client = clientParts[index];
const part = stored[index];
if (client.partNumber !== index + 1) return { kind: "error", code: "InvalidPartOrder", message: "Parts must be consecutive starting at 1" };
if (part.partNumber !== client.partNumber || part.etag !== client.etag) return { kind: "error", code: "InvalidPart", message: `Part ${client.partNumber} does not match` };
}
const driveOffset = this.getValue<number>("driveOffset") ?? 0;
const carry = this.carry();
const total = driveOffset + carry.byteLength;
if (expectedTotal !== undefined && expectedTotal !== total) return { kind: "error", code: "InvalidPart", message: "x-amz-mp-object-size does not match uploaded parts" };
try {
const accessToken = await getAccessToken(this.env);
let metadata: DriveUploadResult;
if (total === 0) {
await cancelSession(this.getValue<string>("uploadUrl")!, accessToken);
metadata = await createEmptyFile(accessToken, {
name: this.getValue<string>("fileName")!,
parents: [this.getValue<string>("parentFolderId")!],
mimeType: this.getValue<string>("mimeType")!,
existingFileId: this.getValue<string>("existingFileId"),
});
} else {
metadata = await putFinalChunk(this.getValue<string>("uploadUrl")!, accessToken, driveOffset, total, carry);
}
const partEtags = stored.map((part) => part.etag);
await this.ctx.storage.deleteAlarm();
await this.ctx.storage.deleteAll();
return { kind: "complete", metadata, partEtags };
} catch (error) {
return { kind: "error", code: "InternalError", message: error instanceof Error ? error.message : String(error) };
}
}
async listParts(marker: number, maxParts: number): Promise<{ parts: PartRow[]; nextMarker: number; truncated: boolean } | null> {
if (!this.hasUpload()) return null;
const rows = this.ctx.storage.sql.exec<PartRow>("SELECT partNumber, size, etag, endOffset, ts FROM parts WHERE partNumber > ? ORDER BY partNumber LIMIT ?", marker, maxParts + 1).toArray();
const truncated = rows.length > maxParts;
const parts = rows.slice(0, maxParts);
return { parts, nextMarker: parts.at(-1)?.partNumber ?? marker, truncated };
}
async abort(): Promise<boolean> {
if (!this.hasUpload()) return false;
const accessToken = await getAccessToken(this.env);
await cancelSession(this.getValue<string>("uploadUrl")!, accessToken);
await this.ctx.storage.deleteAlarm();
await this.ctx.storage.deleteAll();
for (const waiter of this.waiters.values()) {
clearTimeout(waiter.timer);
waiter.resolve({ kind: "error", code: "NoSuchUpload", message: "Multipart upload was aborted" });
}
this.waiters.clear();
return true;
}
async alarm(): Promise<void> {
if (!this.hasUpload()) return;
try {
await this.abort();
} catch (error) {
console.error(JSON.stringify({ message: "multipart cleanup failed", error: error instanceof Error ? error.message : String(error) }));
await this.ctx.storage.setAlarm(Date.now() + 60 * 60 * 1000);
}
}
}
+237
View File
@@ -0,0 +1,237 @@
import { decodedContentLength, isAwsChunked, pumpBody } from "./aws-chunked";
import { createSession, nextDriveOffset } from "./drive-resumable";
import { deleteFromDrive, findFileInFolder, getFileMetadata, listFiles, resolvePathToFolderAndFile, 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;
}
function xmlResponse(body: string, status = 200): Response {
return new Response(body, { status, headers: { "Content-Type": "application/xml" } });
}
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));
}
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;
}
}
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();
if (text.length > MAX_COMPLETE_XML) return s3Error("EntityTooLarge", 400, "Completion XML exceeds 4 MiB", `/${bucket}/${key}`);
let parts: Array<{ partNumber: number; etag: string }>;
try {
parts = parseCompleteMultipartUpload(text);
} catch {
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));
}
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}`);
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));
}
async function putObject(request: Request, env: Env, accessToken: string, bucket: string, key: string): Promise<Response> {
if (request.headers.has("x-amz-copy-source")) return s3Error("NotImplemented", 501, undefined, `/${bucket}/${key}`);
const result = await streamUploadToDrive(accessToken, request, bucket, key, request.headers.get("Content-Type") || "application/octet-stream", env);
return new Response(null, { status: 200, headers: { ETag: `"${etag(result)}"` } });
}
export async function dispatch(request: Request, env: Env, accessToken: string, bucket: string, key: string): Promise<Response> {
const url = new URL(request.url);
const method = request.method;
const resource = `/${bucket}${key ? `/${key}` : ""}`;
if (method === "POST" && key && url.searchParams.has("uploads")) return createMultipart(request, env, accessToken, bucket, key);
if (method === "PUT" && key && url.searchParams.has("uploadId") && url.searchParams.has("partNumber")) return uploadPart(request, env, accessToken, bucket, key, url);
if (method === "POST" && key && url.searchParams.has("uploadId")) return completeMultipart(request, env, bucket, key, url);
if (method === "DELETE" && key && url.searchParams.has("uploadId")) return abortMultipart(env, bucket, key, url);
if (method === "GET" && key && url.searchParams.has("uploadId")) return listParts(env, bucket, key, url);
if (method === "GET" && !key && url.searchParams.has("uploads")) return xmlResponse(listMultipartUploadsResult(bucket));
if (method === "PUT" && key) return putObject(request, env, accessToken, bucket, key);
if (method === "GET") {
if (!key) return xmlResponse(generateListBucketResult(await listFiles(accessToken, bucket, env), bucket));
try {
const file = await streamDownloadFromDrive(accessToken, bucket, key, env, request.headers.get("Range") ?? undefined);
const headers = new Headers({
"Content-Type": file.contentType,
"Content-Length": file.contentLength ?? file.size.toString(),
"Cache-Control": "s-maxage=300, no-store",
"Accept-Ranges": "bytes",
ETag: `"${etag(file)}"`,
});
if (file.contentRange) headers.set("Content-Range", file.contentRange);
return new Response(file.body, { status: file.status, headers });
} catch (error) {
if (error instanceof Error && error.message === "File not found") return s3Error("NoSuchKey", 404, undefined, resource);
throw error;
}
}
if (method === "HEAD") {
if (!key) return new Response(null, { status: 200 });
try {
const metadata = await getFileMetadata(accessToken, bucket, key, env);
return new Response(null, {
status: 200,
headers: { "Content-Type": metadata.mimeType, "Content-Length": metadata.size.toString(), "Accept-Ranges": "bytes", ETag: `"${etag(metadata)}"` },
});
} catch (error) {
if (error instanceof Error && error.message === "File not found") return s3Error("NoSuchKey", 404, undefined, resource, true);
throw error;
}
}
if (method === "DELETE" && key) {
try {
await deleteFromDrive(accessToken, bucket, key, env);
return new Response(null, { status: 204 });
} catch (error) {
if (error instanceof Error && error.message === "File not found") return s3Error("NoSuchKey", 404, undefined, resource);
throw error;
}
}
throw new S3Exception("MethodNotAllowed", 405);
}
+39
View File
@@ -0,0 +1,39 @@
import { escapeXml } from "./s3-xml";
export type S3ErrorCode = "AccessDenied" | "SignatureDoesNotMatch" | "NoSuchKey" | "NoSuchUpload" | "InvalidPart" | "InvalidPartOrder" | "EntityTooLarge" | "MalformedXML" | "InvalidArgument" | "MethodNotAllowed" | "NotImplemented" | "SlowDown" | "InternalError";
const DEFAULT_MESSAGES: Record<S3ErrorCode, string> = {
AccessDenied: "Access Denied",
SignatureDoesNotMatch: "The request signature we calculated does not match the signature you provided.",
NoSuchKey: "The specified key does not exist.",
NoSuchUpload: "The specified multipart upload does not exist.",
InvalidPart: "One or more of the specified parts could not be found.",
InvalidPartOrder: "The list of parts was not in ascending order.",
EntityTooLarge: "Your proposed upload exceeds the maximum allowed object size.",
MalformedXML: "The XML you provided was not well-formed or did not validate against our published schema.",
InvalidArgument: "Invalid argument.",
MethodNotAllowed: "The specified method is not allowed against this resource.",
NotImplemented: "A header you provided implies functionality that is not implemented.",
SlowDown: "Please reduce your request rate.",
InternalError: "We encountered an internal error. Please try again.",
};
export class S3Exception extends Error {
constructor(
readonly code: S3ErrorCode,
readonly status: number,
message = DEFAULT_MESSAGES[code],
readonly headers?: HeadersInit,
) {
super(message);
}
}
export function s3Error(code: S3ErrorCode, status: number, message = DEFAULT_MESSAGES[code], resource?: string, head = false, headers?: HeadersInit): Response {
const requestId = crypto.randomUUID();
const body = head ? null : `<?xml version="1.0" encoding="UTF-8"?>\n<Error><Code>${code}</Code><Message>${escapeXml(message)}</Message>${resource ? `<Resource>${escapeXml(resource)}</Resource>` : ""}<RequestId>${requestId}</RequestId></Error>`;
const responseHeaders = new Headers(headers);
responseHeaders.set("Content-Type", "application/xml");
responseHeaders.set("x-amz-request-id", requestId);
return new Response(body, { status, headers: responseHeaders });
}
+66
View File
@@ -0,0 +1,66 @@
import type { GoogleDriveFile } from "./types";
export function escapeXml(str: string): string {
return str.replace(/&/g, "&amp;").replace(/</g, "&lt;").replace(/>/g, "&gt;").replace(/"/g, "&quot;").replace(/'/g, "&apos;");
}
export function generateListBucketResult(files: GoogleDriveFile[], bucket: string): string {
const contents = files
.map(
(f) => `
<Contents>
<Key>${escapeXml(f.name)}</Key>
<LastModified>${f.modifiedTime || new Date().toISOString()}</LastModified>
<ETag>"${f.md5Checksum || f.id}"</ETag>
<Size>${f.size || 0}</Size>
<StorageClass>STANDARD</StorageClass>
</Contents>`,
)
.join("");
return `<?xml version="1.0" encoding="UTF-8"?>
<ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/">
<Name>${escapeXml(bucket)}</Name>
<Prefix></Prefix>
<MaxKeys>1000</MaxKeys>
<IsTruncated>false</IsTruncated>
${contents}
</ListBucketResult>`;
}
export function initiateMultipartUploadResult(bucket: string, key: string, uploadId: string): string {
return `<?xml version="1.0" encoding="UTF-8"?>
<InitiateMultipartUploadResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Bucket>${escapeXml(bucket)}</Bucket><Key>${escapeXml(key)}</Key><UploadId>${escapeXml(uploadId)}</UploadId></InitiateMultipartUploadResult>`;
}
export function completeMultipartUploadResult(bucket: string, key: string, etag: string): string {
return `<?xml version="1.0" encoding="UTF-8"?>
<CompleteMultipartUploadResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Location>/${escapeXml(bucket)}/${escapeXml(key)}</Location><Bucket>${escapeXml(bucket)}</Bucket><Key>${escapeXml(key)}</Key><ETag>"${escapeXml(etag)}"</ETag></CompleteMultipartUploadResult>`;
}
export function listPartsResult(bucket: string, key: string, uploadId: string, result: { parts: Array<{ partNumber: number; size: number; etag: string; ts: number }>; nextMarker: number; truncated: boolean }): string {
const parts = result.parts.map((part) => `<Part><PartNumber>${part.partNumber}</PartNumber><LastModified>${new Date(part.ts).toISOString()}</LastModified><ETag>"${escapeXml(part.etag)}"</ETag><Size>${part.size}</Size></Part>`).join("");
return `<?xml version="1.0" encoding="UTF-8"?>
<ListPartsResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Bucket>${escapeXml(bucket)}</Bucket><Key>${escapeXml(key)}</Key><UploadId>${escapeXml(uploadId)}</UploadId><NextPartNumberMarker>${result.nextMarker}</NextPartNumberMarker><IsTruncated>${result.truncated}</IsTruncated>${parts}</ListPartsResult>`;
}
export function listMultipartUploadsResult(bucket: string): string {
return `<?xml version="1.0" encoding="UTF-8"?>
<ListMultipartUploadsResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Bucket>${escapeXml(bucket)}</Bucket><KeyMarker></KeyMarker><UploadIdMarker></UploadIdMarker><NextKeyMarker></NextKeyMarker><NextUploadIdMarker></NextUploadIdMarker><MaxUploads>1000</MaxUploads><IsTruncated>false</IsTruncated></ListMultipartUploadsResult>`;
}
const PART_RE = /<Part\b[^>]*>([\s\S]*?)<\/Part>/g;
const NUM_RE = /<PartNumber>\s*(\d+)\s*<\/PartNumber>/;
const ETAG_RE = /<ETag>\s*(?:&quot;|")?([^<"&]*)(?:&quot;|")?\s*<\/ETag>/;
export function parseCompleteMultipartUpload(xml: string): Array<{ partNumber: number; etag: string }> {
const parts: Array<{ partNumber: number; etag: string }> = [];
for (const match of xml.matchAll(PART_RE)) {
const numberMatch = NUM_RE.exec(match[1]);
const etagMatch = ETAG_RE.exec(match[1]);
if (!numberMatch || !etagMatch || parts.length >= 10_000) throw new Error("MalformedXML");
parts.push({ partNumber: Number(numberMatch[1]), etag: etagMatch[1].trim() });
}
if (parts.length === 0 && !/<CompleteMultipartUpload\b[^>]*>\s*<\/CompleteMultipartUpload>/.test(xml)) throw new Error("MalformedXML");
return parts;
}
+54
View File
@@ -0,0 +1,54 @@
export interface Env {
ACCESS_KEY: string;
SECRET_KEY: string;
REGION: string;
GOOGLE_CLIENT_ID: string;
GOOGLE_CLIENT_SECRET: string;
GOOGLE_REFRESH_TOKEN: string;
AUTH_KV: KVNamespace;
FOLDER_CACHE: KVNamespace;
MPU: DurableObjectNamespace<import("./multipart-do").MultipartUploadDO>;
ALLOWED_BUCKETS?: string;
PUBLIC_READ_BUCKETS?: string;
ALLOW_MULTIPART?: string;
ETAG_STYLE?: "md5" | "multipart";
}
export interface GoogleDriveFile {
id: string;
name: string;
mimeType: string;
size: string;
modifiedTime?: string;
md5Checksum?: string;
}
export interface GoogleDriveSearchResponse {
files?: GoogleDriveFile[];
}
export interface DriveUploadResult {
id: string;
name: string;
mimeType?: string;
size?: string;
md5Checksum?: string;
}
export interface DriveDownloadResult {
body: ReadableStream;
contentType: string;
size: number;
id: string;
md5Checksum?: string;
status: number;
contentRange?: string;
contentLength?: string;
}
export interface DriveFileMetadata {
id: string;
mimeType: string;
size: number;
md5Checksum?: string;
}
+286 -1072
View File
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -1,7 +1,7 @@
{ {
"extends": "../tsconfig.json", "extends": "../tsconfig.json",
"compilerOptions": { "compilerOptions": {
"types": ["@cloudflare/vitest-pool-workers"] "types": ["@cloudflare/vitest-pool-workers/types"]
}, },
"include": ["./**/*.ts", "../worker-configuration.d.ts"], "include": ["./**/*.ts", "../worker-configuration.d.ts"],
"exclude": [] "exclude": []
+19
View File
@@ -1,6 +1,25 @@
import { cloudflareTest } from "@cloudflare/vitest-pool-workers";
import { configDefaults, defineConfig } from "vitest/config"; import { configDefaults, defineConfig } from "vitest/config";
export default defineConfig({ export default defineConfig({
plugins: [
cloudflareTest({
wrangler: { configPath: "./wrangler.jsonc" },
miniflare: {
bindings: {
ACCESS_KEY: "test-access-key",
SECRET_KEY: "test-secret-key",
REGION: "auto",
GOOGLE_CLIENT_ID: "test-client-id",
GOOGLE_CLIENT_SECRET: "test-client-secret",
GOOGLE_REFRESH_TOKEN: "test-refresh-token",
ALLOWED_BUCKETS: "test-bucket,empty-bucket,my-bucket",
ALLOW_MULTIPART: "true",
ETAG_STYLE: "md5",
},
},
}),
],
test: { test: {
exclude: [...configDefaults.exclude, ".pnpm-wrangler/**", ".pnpm-store/**", "node_modules/**"], exclude: [...configDefaults.exclude, ".pnpm-wrangler/**", ".pnpm-store/**", "node_modules/**"],
}, },
+4203 -275
View File
File diff suppressed because it is too large Load Diff
+18
View File
@@ -7,9 +7,27 @@
"name": "iris", "name": "iris",
"main": "src/index.ts", "main": "src/index.ts",
"compatibility_date": "2025-09-27", "compatibility_date": "2025-09-27",
"vars": {
"ALLOW_MULTIPART": "false",
"ETAG_STYLE": "md5"
},
"observability": { "observability": {
"enabled": true "enabled": true
}, },
"durable_objects": {
"bindings": [
{
"name": "MPU",
"class_name": "MultipartUploadDO"
}
]
},
"migrations": [
{
"tag": "v1",
"new_sqlite_classes": ["MultipartUploadDO"]
}
],
"kv_namespaces": [ "kv_namespaces": [
{ {
"binding": "AUTH_KV", "binding": "AUTH_KV",