diff --git a/agent.md b/agent.md index 879a6ce..b88b0af 100644 --- a/agent.md +++ b/agent.md @@ -164,10 +164,12 @@ has its own limit. When the org has a Google Drive connection whose account can open the folder, clips are downloaded through it instead, as that account, which -Drive's limit on shared links doesn't touch. Those downloads run in the -background steps only. A file too big for the connection's temporary storage -(seen at 9 GB; 3.4 GB passed) comes by the original's shared link, as above: -whole, or a piece at a time when Drive refuses it whole. It is never copied inside Google first: Drive answers a header +Drive's limit on shared links doesn't touch: Google's own Drive API, signed +by the connection. Those downloads run in the background steps only. A file +over 250 MB, more than one answer through the connection carries, comes by +the original's shared link, as above: whole, or a piece at a time when Drive +refuses it whole. A piece Drive refuses on the shared link (it can, after a +few GB of one file) comes through the connection instead. It is never copied inside Google first: Drive answers a header check on a fresh copy of a big file with an empty page, and the video host refuses such a download. Copies an earlier version made are deleted once their import is over (or with the project). diff --git a/package.json b/package.json index 55fe577..65b0f1b 100644 --- a/package.json +++ b/package.json @@ -13,7 +13,7 @@ }, "dependencies": { "@clawnify/app": "^0.2.1", - "@clawnify/connections": "^0.5.0", + "@clawnify/connections": "^0.7.0", "@clawnify/db": "^0.4.2", "@clawnify/queue": "^0.1.1", "@fontsource-variable/inter": "^5.3.0", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 8f3b27d..87d2f3e 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -12,8 +12,8 @@ importers: specifier: ^0.2.1 version: 0.2.1(@cloudflare/workers-types@4.20260613.1)(@hono/zod-openapi@0.19.10(hono@4.12.25)(zod@3.25.76))(@phosphor-icons/react@2.1.10(react-dom@19.2.7(react@19.2.7))(react@19.2.7))(hono@4.12.25)(react@19.2.7) '@clawnify/connections': - specifier: ^0.5.0 - version: 0.5.0 + specifier: ^0.7.0 + version: 0.7.0 '@clawnify/db': specifier: ^0.4.2 version: 0.4.2(@cloudflare/workers-types@4.20260613.1) @@ -200,14 +200,14 @@ packages: react: optional: true - '@clawnify/connections@0.5.0': - resolution: {integrity: sha512-WYZtqObjok2mRnG8+QywBwSoswUHuq82hwTFSKghHUESlq3lLd2e6zMUWa4J+cQmLjGPitLpqj4gXUvKZYzn7Q==} + '@clawnify/connections@0.7.0': + resolution: {integrity: sha512-e5N3ItrwdHuW80tbx08bmi2Ya1CAxsCX9pnMq/LaRK4j0EMKqToztYx4a+17TGAnJ44JHpHBYwZI/tv4PSz7Qg==} '@clawnify/db@0.4.2': resolution: {integrity: sha512-s+G/5bDVFx7Na8WVETfl8EhyPcbmgs5LIE1TsM0vGUOAtPihQGlVX1xpfmg5lsMv8+OChKaHsZIMnzpXY8pyAw==} - '@clawnify/integrations@0.4.0': - resolution: {integrity: sha512-epX7r9VHczmjHJuegsn1ogRM9SG8CsYLsyF6hfxN3uV4nDYjdWazNingZnAgOI+ynVj0ZLV8x2rUnNoUUnSuQw==} + '@clawnify/integrations@0.5.0': + resolution: {integrity: sha512-asGXIsAbSe/W7RLAxW4lLJCuILhw/I+H4KGCRVCRyvxPqm87IuGytNkb3XdgbvFTqj+H73wm05fZJgFMdAO1sQ==} '@clawnify/queue@0.1.1': resolution: {integrity: sha512-6gfaz4B3rVVFWMB0IjJ0mt3ausETXpiWodRqShMBTRwVtNyCXJ4FO4/f3qc/TYtnA9ZZBPVeBLjbzg5nTqvULw==} @@ -2744,9 +2744,9 @@ snapshots: - sql.js - sqlite3 - '@clawnify/connections@0.5.0': + '@clawnify/connections@0.7.0': dependencies: - '@clawnify/integrations': 0.4.0 + '@clawnify/integrations': 0.5.0 '@clawnify/db@0.4.2(@cloudflare/workers-types@4.20260613.1)': dependencies: @@ -2782,7 +2782,7 @@ snapshots: - sql.js - sqlite3 - '@clawnify/integrations@0.4.0': {} + '@clawnify/integrations@0.5.0': {} '@clawnify/queue@0.1.1': {} diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 661850f..f5cbf37 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -2,3 +2,5 @@ allowBuilds: esbuild: true sharp: true workerd: true +minimumReleaseAgeExclude: + - '@clawnify/connections@0.7.0' diff --git a/src/server/drive.ts b/src/server/drive.ts index f7f4b65..9c1e7aa 100644 --- a/src/server/drive.ts +++ b/src/server/drive.ts @@ -4,16 +4,47 @@ // exactly like an upload, so the editor and the export never read from Drive // themselves. // -// The connection's broker never hands this app a raw Google token, so the -// bytes come the one way it allows: the download action parks the file behind -// a short-lived signed link, and the import streams that link into storage. +// Every call is Google's own Drive API, sent through the connection's broker, +// which signs it with the org's credential: the app never holds a Google +// token, and the code is the same whoever holds the credential. A file comes +// back as a short-lived link to its bytes, at most 250 MB in one answer. import { connect, describe, type ConnectionsEnv } from "@clawnify/connections"; -// Both toolkits carry the same Drive actions, prefixed with their own name. // Google Drive comes first: it is the connection a studio makes for its files. const SERVICES = ["googledrive", "googlesuper"] as const; -type Service = (typeof SERVICES)[number]; +export type DriveService = (typeof SERVICES)[number]; +type Service = DriveService; + +/** + * Where the Drive API sits under each connection's base URL: Google Drive's + * already ends in `/drive/v3`, Google Super's is the googleapis.com root. + */ +const DRIVE_API: Record = { googledrive: "", googlesuper: "/drive/v3" }; + +/** One request to the Drive API through the connection. Query values are sent as given. */ +function driveCall( + env: ConnectionsEnv, + service: Service, + method: "GET" | "DELETE", + path: string, + query: Record = {}, +) { + return connect(service, env).rawRequest({ + method, + endpoint: `${DRIVE_API[service]}${path}`, + parameters: Object.entries(query).map(([name, value]) => ({ name, value, in: "query" as const })), + }); +} + +/** The broker's words for a provider 404: " 404: …". */ +const notFound = (e: unknown) => /^\w+ 404:/.test(e instanceof Error ? e.message : String(e)); + +/** + * The most one answer through the connection can carry: the broker refuses a + * bigger file ("maximum file size limit of 250MB"). + */ +export const PROXY_FILE_BYTES = 250_000_000; /** What the picker lists. `duration` is seconds, for videos Drive has probed. */ export interface DriveFile { @@ -85,18 +116,25 @@ export async function listDriveFiles( const search = opts.search?.replace(/['\\]/g, " ").trim(); if (search) q += ` and name contains '${search}'`; const shared = opts.folderId === SHARED_WITH_ME; + // Shared items are found by the query, not by a parent folder. if (shared) q += " and sharedWithMe = true"; + else { + const parent = opts.folderId || "root"; + if (parent !== "root" && !DRIVE_FILE_ID.test(parent)) throw new Error("not a Drive folder id"); + q += ` and '${parent}' in parents`; + } const service = await requireDriveService(env); - const data = (await connect(service, env).run(`${service.toUpperCase()}_FIND_FILE`, { + const data = (await driveCall(env, service, "GET", "/files", { q, // `folder` sorts folders ahead of files, so the picker needs no re-sort. orderBy: "folder,modifiedTime desc", - pageSize: 50, + pageSize: "50", fields: "nextPageToken,files(id,name,mimeType,size,modifiedTime,thumbnailLink,videoMediaMetadata(durationMillis))", - // Shared items are found by the query, not by a parent folder. - ...(shared ? {} : { folder_id: opts.folderId || "root" }), + // Folders on a shared drive, as well as the account's own. + supportsAllDrives: "true", + includeItemsFromAllDrives: "true", ...(opts.pageToken ? { pageToken: opts.pageToken } : {}), })) as { files?: { @@ -138,11 +176,16 @@ async function driveItem( id: string, ): Promise<{ name: string; parents: string[] } | null> { const service = await requireDriveService(env); - const data = (await connect(service, env).run(`${service.toUpperCase()}_GET_FILE_METADATA`, { - fileId: id, - fields: "name,parents", - })) as { name?: string; parents?: string[] }; - return data.name ? { name: data.name, parents: data.parents ?? [] } : null; + try { + const data = (await driveCall(env, service, "GET", `/files/${encodeURIComponent(id)}`, { + fields: "name,parents", + supportsAllDrives: "true", + })) as { name?: string; parents?: string[] }; + return data.name ? { name: data.name, parents: data.parents ?? [] } : null; + } catch (e) { + if (notFound(e)) return null; + throw e; + } } export async function driveFolderName(env: ConnectionsEnv, id: string): Promise { @@ -165,8 +208,103 @@ export async function withinFolder(env: ConnectionsEnv, itemId: string, folderId return false; } -/** A signed, short-lived link to one Drive file's original bytes. */ -export async function driveDownloadLink( +export interface DriveFileInfo { + name: string; + mimeType: string; + /** Bytes. Null for a file Drive keeps no size for (a Google Doc, say). */ + size: number | null; +} + +/** + * A short-lived link to one Drive file's original bytes, through the + * connection, or `tooBig` (with what the file is) when it is over what one + * answer carries. Then it comes by its shared link instead (footage, a piece + * at a time when need be), or by {@link driveActionLink}. + */ +export async function driveFileLink( + env: ConnectionsEnv, + fileId: string, +): Promise<(DriveFileInfo & { url: string }) | (DriveFileInfo & { tooBig: true })> { + const service = await requireDriveService(env); + const path = `/files/${encodeURIComponent(fileId)}`; + const meta = (await driveCall(env, service, "GET", path, { fields: "name,mimeType,size", supportsAllDrives: "true" })) as { + name?: string; + mimeType?: string; + size?: string; + }; + const info: DriveFileInfo = { + name: meta.name || fileId, + mimeType: meta.mimeType || "application/octet-stream", + size: meta.size ? Number(meta.size) : null, + }; + // Asked for whole, a bigger file is fetched into the broker's storage before + // it is refused: ask only for what can come back. + if (info.size === null || info.size > PROXY_FILE_BYTES) return { ...info, tooBig: true }; + const file = await connect(service, env).rawFile({ + method: "GET", + endpoint: `${DRIVE_API[service]}${path}`, + parameters: [ + { name: "alt", value: "media", in: "query" }, + { name: "supportsAllDrives", value: "true", in: "query" }, + ], + }); + return { ...info, url: file.url }; +} + +/** + * One piece of a Drive file through the connection, as a Response carrying + * Drive's own status and headers and the bytes behind the broker's link. A + * refusal comes back as its status, so the caller can tell it from a hiccup; + * null when the broker or its link couldn't be reached. + * + * Drive serves these as the connected account even when its limit refuses + * the shared link's pieces: on 2026-10-10 a 15 GB file's shared link refused + * ranges past about 4 GB read, while this served the same range. + */ +export async function driveFilePiece( + env: ConnectionsEnv, + service: DriveService, + fileId: string, + start: number, + end: number, + signal: AbortSignal, +): Promise { + let file; + try { + file = await connect(service, env).rawFile({ + method: "GET", + endpoint: `${DRIVE_API[service]}/files/${encodeURIComponent(fileId)}`, + parameters: [ + { name: "alt", value: "media", in: "query" }, + { name: "supportsAllDrives", value: "true", in: "query" }, + { name: "Range", value: `bytes=${start}-${end}`, in: "header" }, + ], + }); + } catch (e) { + const status = /^\w+ (\d{3}):/.exec(e instanceof Error ? e.message : String(e))?.[1]; + return status ? new Response(null, { status: Number(status) }) : null; + } + const bytes = await fetch(file.url, { signal }).catch(() => null); + if (!bytes?.ok || !bytes.body) { + await bytes?.body?.cancel().catch(() => {}); + return null; + } + const range = file.headers["content-range"] ?? file.headers["Content-Range"]; + return new Response(bytes.body, { + status: file.status ?? 206, + headers: { "content-type": file.contentType, ...(range ? { "content-range": range } : {}) }, + }); +} + +/** + * A whole file of any size up to the broker's own storage (3.4 GB passed, 9 GB + * did not) through its download action, which parks the file behind a + * short-lived link. The one Drive call still made as an action: for a file + * over what one proxy answer carries, where nothing reads it in pieces (an + * import into the media library). See the platform's + * docs/internal/connections-architecture.md §19.5 and §23. + */ +export async function driveActionLink( env: ConnectionsEnv, fileId: string, ): Promise<{ url: string; name: string; mimeType: string }> { @@ -182,5 +320,5 @@ export async function driveDownloadLink( /** Delete a file the app made in the connected account (a copy an earlier version made for an import). */ export async function driveRemove(env: ConnectionsEnv, fileId: string): Promise { const service = await requireDriveService(env); - await connect(service, env).run(`${service.toUpperCase()}_GOOGLE_DRIVE_DELETE_FOLDER_OR_FILE_ACTION`, { fileId, supportsAllDrives: true }); + await driveCall(env, service, "DELETE", `/files/${encodeURIComponent(fileId)}`, { supportsAllDrives: "true" }); } diff --git a/src/server/footage.ts b/src/server/footage.ts index c643f2f..577ed3d 100644 --- a/src/server/footage.ts +++ b/src/server/footage.ts @@ -21,7 +21,16 @@ import { query, get, run } from "./db"; import { deleteMedia, importMedia, mediaState, openMediaUpload, prepareMedia, startTranscode, transcodeState, type MediaConfig } from "./media"; import { refusalDetail } from "./refusal"; import { directDownloadUrl, folderListingUrl, judgeLinkResponse, listFolderVideos, type FolderVideo } from "./drive-link"; -import { MAX_RELAY_BYTES, RELAY_AT_ONCE, RELAY_BUDGET_MS, RELAY_EXPIRY_MARGIN_MS, rangedSize, relayPieces } from "./relay"; +import { + MAX_RELAY_BYTES, + RELAY_AT_ONCE, + RELAY_BUDGET_MS, + RELAY_EXPIRY_MARGIN_MS, + rangedSize, + relayPieces, + sharedLinkReader, + type PieceReader, +} from "./relay"; const DEFAULT_SERVICES_URL = "https://services.clawnify.com"; @@ -342,6 +351,11 @@ export interface DriveSource { download(fileId: string): Promise<{ url: string; mimeType: string } | { error: string }>; /** Delete a copy an earlier version made in the connected account for an import. */ remove(fileId: string): Promise; + /** + * One piece of a file through the connection: Drive's answer, with its own + * status and headers. It serves pieces the shared link refuses. + */ + piece?(fileId: string, start: number, end: number, signal: AbortSignal): Promise; } /** @@ -505,7 +519,7 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste // counts from here, so it doesn't matter what ran before. if (opts.relay && relaying.length) { const until = Date.now() + RELAY_BUDGET_MS; - await Promise.all(relaying.slice(0, RELAY_AT_ONCE).map((r) => relayClip(cfg, r, until))); + await Promise.all(relaying.slice(0, RELAY_AT_ONCE).map((r) => relayClip(cfg, r, until, opts.drive))); } // 2. Start imports while there is room, together. A refusal that is about @@ -719,13 +733,19 @@ async function openRelay(cfg: MediaConfig, r: WorkRow): Promise { +/** + * Move a clip coming in a piece at a time on, and settle it once it is all in. + * Pieces come from the shared link, or, where Drive refuses those, through the + * org's Drive connection when it has one. + */ +async function relayClip(cfg: MediaConfig, r: WorkRow, until: number, drive?: DriveSource): Promise { const uid = r.upload_uid!; const expiring = r.upload_expires !== null && Date.parse(r.upload_expires) - Date.now() < RELAY_EXPIRY_MARGIN_MS; + const readers: PieceReader[] = [sharedLinkReader(r.drive_file_id)]; + if (drive?.piece) readers.push((start, end, signal) => drive.piece!(r.drive_file_id, start, end, signal)); const result = expiring ? ({ state: "gone" } as const) - : await relayPieces({ url: r.upload_url!, fileId: r.drive_file_id, size: r.upload_size! }, until, async (received) => { + : await relayPieces({ url: r.upload_url!, size: r.upload_size! }, readers, until, async (received) => { await setRow(r.id, { upload_done: received }); }); if (result.state === "moving") return; @@ -787,8 +807,8 @@ async function startImport( // Already coming in a piece at a time, back from a wait on Drive: it // carries on from where its upload got to, on the next delivery. if (r.upload_uid) return true; - // Through the org's connection first. A file it can't hand over (too big for - // the connector's temporary storage) comes by the original's shared link, + // Through the org's connection first. A file it can't hand over (over the + // 250 MB one answer through it carries) comes by the original's shared link, // with its waits. Not by a copy in the connected account: Drive answers a // header check on a fresh copy of a big file with an empty page while the // file itself downloads, and the video host, which checks first, refuses diff --git a/src/server/index.ts b/src/server/index.ts index 75be045..39a5332 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -23,7 +23,9 @@ import { import { DRIVE_FILE_ID, SHARED_WITH_ME, - driveDownloadLink, + driveActionLink, + driveFileLink, + driveFilePiece, driveRemove, driveFolderName, driveStatus, @@ -314,7 +316,11 @@ app.post("/api/drive/import", async (c) => { return c.json({ error: `that file is outside ${limit.name}` }, 403); } - const file = await driveDownloadLink(c.env, b.fileId); + // Through the connection's Drive API. A file over what one answer carries + // (250 MB) takes the broker's download action instead: nothing here reads + // a library import in pieces. + const link = await driveFileLink(c.env, b.fileId); + const file = "tooBig" in link ? await driveActionLink(c.env, b.fileId) : link; // A video goes to the media service: it fetches the link itself, so nothing // passes through this app and no size ceiling applies. Stills and sound are @@ -940,13 +946,18 @@ async function addFolder( /** The org's Google Drive connection as a footage source, when it has one. */ async function driveSource(env: Bindings): Promise { if (!env.CREDENTIALS) return undefined; - const status = await driveStatus(env).catch(() => ({ connected: false })); - if (!status.connected) return undefined; + const status = await driveStatus(env).catch(() => null); + if (!status?.connected || !status.service) return undefined; + const service = status.service; const message = (e: unknown) => (e instanceof Error ? e.message : String(e)); return { + piece: (fileId, start, end, signal) => driveFilePiece(env, service, fileId, start, end, signal), async download(fileId) { try { - const file = await driveDownloadLink(env, fileId); + const file = await driveFileLink(env, fileId); + // Footage always has its shared link: a bigger file comes by that, + // a piece at a time when Drive won't hand it over whole. + if ("tooBig" in file) return { error: "over what one answer through the connection carries (250 MB)" }; return { url: file.url, mimeType: file.mimeType }; } catch (e) { return { error: message(e) }; diff --git a/src/server/relay.ts b/src/server/relay.ts index b234180..4e53e90 100644 --- a/src/server/relay.ts +++ b/src/server/relay.ts @@ -145,6 +145,18 @@ async function sendPiece( return res.ok && Number.isInteger(next) && next > offset ? next : null; } +/** + * Reads one piece of the file: Drive's answer, to be checked and streamed on. + * Null when it couldn't be asked (asked again on the next step). + */ +export type PieceReader = (start: number, end: number, signal: AbortSignal) => Promise; + +/** The file's shared link: fast, and free of any broker. */ +export function sharedLinkReader(fileId: string): PieceReader { + return (start, end, signal) => + fetch(directDownloadUrl(fileId), { headers: { Range: `bytes=${start}-${end}` }, signal }).catch(() => null); +} + export type RelayResult = /** Every byte is in. `contentType` is what Drive said the file is. */ | { state: "done"; contentType: string } @@ -161,9 +173,16 @@ export type RelayResult = * Append pieces of a Drive file to its open upload until the file is all in * or `until` (a Date.now() time) passes; a piece already on its way finishes. * `progress` hears each new offset. + * + * Each piece is read from the first of `readers` that serves it: the shared + * link, then the org's Drive connection. Drive's limit can come to refuse a + * shared link's pieces too, past some point in the file, while its API, asked + * as the connected account, still serves them. Once a later reader has served + * a piece, the rest of this call starts from it. */ export async function relayPieces( - upload: { url: string; fileId: string; size: number }, + upload: { url: string; size: number }, + readers: PieceReader[], until: number, progress: (received: number) => Promise, ): Promise { @@ -171,6 +190,7 @@ export async function relayPieces( if (at === "gone") return { state: "gone" }; if (at === null) return { state: "moving" }; let contentType = "video/mp4"; + let first = 0; while (at < upload.size && Date.now() < until) { const piece = nextPiece(at, upload.size)!; // One signal cuts both ends of the piece: the read from Drive and the send. @@ -178,13 +198,23 @@ export async function relayPieces( const timer = setTimeout(() => cut.abort(), PIECE_TIMEOUT_MS); let sent: Awaited>; try { - const res = await fetch(directDownloadUrl(upload.fileId), { - headers: { Range: `bytes=${piece.start}-${piece.end}` }, - signal: cut.signal, - }).catch(() => null); - if (!res) break; - const verdict = pieceVerdict(res.status, res.headers.get("content-type"), res.headers.get("content-range"), piece, upload.size); - if (verdict !== "ok" || !res.body) { + let res: Response | null = null; + let verdict: ReturnType = "refused"; + for (let i = first; i < readers.length; i++) { + res = await readers[i](piece.start, piece.end, cut.signal); + if (!res) { + verdict = "retry"; + break; + } + verdict = pieceVerdict(res.status, res.headers.get("content-type"), res.headers.get("content-range"), piece, upload.size); + if (verdict === "ok") { + first = i; + break; + } + await res.body?.cancel().catch(() => {}); + if (verdict === "retry") break; + } + if (verdict !== "ok" || !res?.body) { cut.abort(); if (verdict === "refused") return { state: "drive" }; break; diff --git a/test/drive.test.ts b/test/drive.test.ts new file mode 100644 index 0000000..5b72c1a --- /dev/null +++ b/test/drive.test.ts @@ -0,0 +1,179 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { + PROXY_FILE_BYTES, + SHARED_WITH_ME, + driveActionLink, + driveFileLink, + driveFilePiece, + driveFolderName, + driveRemove, + listDriveFiles, + withinFolder, +} from "../src/server/drive"; + +/** + * OpenVideo reaches Drive through the org's connection with Google's own Drive + * API, signed by the broker, rather than the broker's action catalogue, so the + * code is the same whoever holds the credential. These check what each call + * asks Drive for and how it reads the answer, through the real SDK. + */ + +const ORG = "org-1"; +type Req = { method: string; endpoint: string; parameters?: { name: string; value: string; in: string }[] }; +type Answer = { data: unknown; error: string | null; successful: boolean; status?: number; headers?: Record; file?: unknown }; + +function connection(service: "googledrive" | "googlesuper", answer: (req: Req) => Answer) { + const sent: Req[] = []; + const actions: unknown[][] = []; + const env = { + CLAWNIFY_ORG_ID: ORG, + CREDENTIALS: { + getToken: async () => null, + listConnected: async () => [service], + searchActions: async () => [], + executeTool: async (_s: string, action: string, args: unknown) => { + actions.push([action, args]); + return { + data: { downloaded_file_content: { s3url: "https://temp.test/whole", name: "big.mp4", mimetype: "video/mp4" } }, + error: null, + successful: true, + }; + }, + proxyRequest: async (_s: string, _o: string, req: Req) => { + sent.push(req); + return answer(req); + }, + }, + }; + return { env, sent, actions }; +} + +const param = (req: Req, name: string) => req.parameters?.find((p) => p.name === name)?.value; +const ok = (data: unknown): Answer => ({ data, error: null, successful: true, status: 200 }); + +afterEach(() => vi.unstubAllGlobals()); + +describe("browsing", () => { + it("lists a folder's children with the Drive API, folders first, across shared drives", async () => { + const c = connection("googledrive", () => + ok({ + files: [ + { id: "f1", name: "Day 1", mimeType: "application/vnd.google-apps.folder" }, + { id: "v1", name: "A.mp4", mimeType: "video/mp4", size: "1000", videoMediaMetadata: { durationMillis: "2500" } }, + ], + nextPageToken: "p2", + }), + ); + const page = await listDriveFiles(c.env, { kind: "media" }); + const req = c.sent[0]; + expect([req.method, req.endpoint]).toEqual(["GET", "/files"]); + expect(param(req, "q")).toContain("'root' in parents"); + expect(param(req, "q")).toContain("trashed = false"); + expect([param(req, "pageSize"), param(req, "supportsAllDrives"), param(req, "includeItemsFromAllDrives")]).toEqual(["50", "true", "true"]); + expect(page.folders).toEqual([{ id: SHARED_WITH_ME, name: "Shared with me" }, { id: "f1", name: "Day 1" }]); + expect(page.files[0]).toMatchObject({ id: "v1", size: 1000, duration: 2.5 }); + expect(page.nextPageToken).toBe("p2"); + }); + + it("under a Google Super connection, the same call carries the Drive API's own path", async () => { + const c = connection("googlesuper", () => ok({ files: [] })); + await listDriveFiles(c.env, { kind: "audio", folderId: "folder_id_12345" }); + expect(c.sent[0].endpoint).toBe("/drive/v3/files"); + expect(param(c.sent[0], "q")).toContain("'folder_id_12345' in parents"); + }); + + it("finds shared-with-me items by the query, not by a folder", async () => { + const c = connection("googledrive", () => ok({ files: [] })); + await listDriveFiles(c.env, { kind: "media", folderId: SHARED_WITH_ME }); + expect(param(c.sent[0], "q")).toContain("sharedWithMe = true"); + expect(param(c.sent[0], "q")).not.toContain("in parents"); + }); + + it("refuses a folder id that could end the query early", async () => { + const c = connection("googledrive", () => ok({ files: [] })); + await expect(listDriveFiles(c.env, { kind: "media", folderId: "x' or name contains '" })).rejects.toThrow("not a Drive folder id"); + expect(c.sent).toEqual([]); + }); + + it("an item Drive won't show is no name, not an error, and the folder limit walks parents", async () => { + const parents: Record = { clip_id_0000: ["mid_id_00000"], mid_id_00000: ["limit_id_000"] }; + const c = connection("googledrive", (req) => { + const id = req.endpoint.split("/").pop()!; + if (id === "gone_id_0000") return { data: { error: { code: 404 } }, error: "googledrive 404: File not found", successful: false, status: 404 }; + return ok({ name: id, parents: parents[id] ?? [] }); + }); + expect(await driveFolderName(c.env, "gone_id_0000")).toBeNull(); + expect(await withinFolder(c.env, "clip_id_0000", "limit_id_000")).toBe(true); + expect(await withinFolder(c.env, "mid_id_00000", "other_id_000")).toBe(false); + expect(param(c.sent[0], "fields")).toBe("name,parents"); + }); +}); + +describe("downloading", () => { + it("a file of up to 250 MB comes back as a link from one answer", async () => { + const c = connection("googledrive", (req) => + param(req, "alt") === "media" + ? { data: null, error: null, successful: true, status: 200, headers: {}, file: { url: "https://temp.test/clip", contentType: "video/mp4", size: 1000, expiresAt: null } } + : ok({ name: "clip.mp4", mimeType: "video/mp4", size: "1000" }), + ); + expect(await driveFileLink(c.env, "clip_id_0000")).toEqual({ name: "clip.mp4", mimeType: "video/mp4", size: 1000, url: "https://temp.test/clip" }); + expect(c.sent.map((r) => [r.endpoint, param(r, "alt") ?? param(r, "fields")])).toEqual([ + ["/files/clip_id_0000", "name,mimeType,size"], + ["/files/clip_id_0000", "media"], + ]); + }); + + it("a bigger file, or one with no size, is never fetched whole through the connection", async () => { + for (const size of [String(PROXY_FILE_BYTES + 1), undefined]) { + const c = connection("googledrive", () => ok({ name: "big.mp4", mimeType: "video/mp4", ...(size ? { size } : {}) })); + const link = await driveFileLink(c.env, "big_id_00000"); + expect(link).toMatchObject({ tooBig: true, name: "big.mp4" }); + expect(c.sent).toHaveLength(1); + } + }); + + it("the media library's fallback for a big file is the broker's download action", async () => { + const c = connection("googledrive", () => ok({})); + expect(await driveActionLink(c.env, "big_id_00000")).toEqual({ url: "https://temp.test/whole", name: "big.mp4", mimeType: "video/mp4" }); + expect(c.actions).toEqual([["GOOGLEDRIVE_DOWNLOAD_FILE", { fileId: "big_id_00000" }]]); + }); + + it("deleting a file is a DELETE on it", async () => { + const c = connection("googledrive", () => ({ data: null, error: null, successful: true, status: 204 })); + await driveRemove(c.env, "copy_id_0000"); + expect([c.sent[0].method, c.sent[0].endpoint, param(c.sent[0], "supportsAllDrives")]).toEqual(["DELETE", "/files/copy_id_0000", "true"]); + }); +}); + +describe("a piece through the connection", () => { + const signal = new AbortController().signal; + + it("carries Drive's own status and range, and the bytes behind the broker's link", async () => { + vi.stubGlobal("fetch", async (url: string) => + url === "https://temp.test/piece" ? new Response("0123456789abcdef") : new Response("no", { status: 404 }), + ); + const c = connection("googledrive", () => ({ + data: null, + error: null, + successful: true, + status: 206, + headers: { "content-range": "bytes 0-15/100" }, + file: { url: "https://temp.test/piece", contentType: "video/mp4", size: 16, expiresAt: null }, + })); + const res = await driveFilePiece(c.env, "googledrive", "clip_id_0000", 0, 15, signal); + expect(res!.status).toBe(206); + expect(res!.headers.get("content-range")).toBe("bytes 0-15/100"); + expect(res!.headers.get("content-type")).toBe("video/mp4"); + expect(await res!.text()).toBe("0123456789abcdef"); + const req = c.sent[0]; + expect([req.endpoint, param(req, "alt"), param(req, "Range")]).toEqual(["/files/clip_id_0000", "media", "bytes=0-15"]); + expect(req.parameters?.find((p) => p.name === "Range")?.in).toBe("header"); + }); + + it("a refusal comes back as its status; a broker that can't be reached as nothing", async () => { + let c = connection("googledrive", () => ({ data: null, error: "googledrive 403: limit", successful: false, status: 403 })); + expect((await driveFilePiece(c.env, "googledrive", "clip_id_0000", 0, 15, signal))!.status).toBe(403); + c = connection("googledrive", () => ({ data: null, error: "proxy failed", successful: false })); + expect(await driveFilePiece(c.env, "googledrive", "clip_id_0000", 0, 15, signal)).toBeNull(); + }); +}); diff --git a/test/e2e/run.mjs b/test/e2e/run.mjs index ae2c281..601ee80 100644 --- a/test/e2e/run.mjs +++ b/test/e2e/run.mjs @@ -125,6 +125,11 @@ globalThis.fetch = async (input, init = {}) => { ? new Response("xx", { status: 206, headers: { "content-type": "video/mp4", "content-range": "bytes 0-1/123456789" } }) : new Response("x".repeat(64), { status: 200, headers: { "content-type": "video/mp4", "content-length": "123456789" } }); } + // The broker's link to one piece it fetched (see proxyRequest below). + if (u.hostname === "temp.r2.test" && u.pathname.startsWith("/piece/")) { + const [a, b] = u.pathname.split("/").pop().split("-").map(Number); + return new Response(zeros(b - a + 1), { status: 200, headers: { "content-type": "video/mp4" } }); + } // The media service's resumable upload, with the video host's rules: a // named user agent ("error code: 1010" otherwise), a piece only at the // upload's own offset, at most 200 MiB, and every piece but the last at @@ -838,27 +843,57 @@ const deleted = [...world.media.values()].filter((v) => v.deleted).length; console.log(`10 ok: project deleted in two calls, ${deleted} media copies deleted, a refused one on the second`); // 12. With the org's Google Drive connection, clips come in through it, not -// the shared link: on deliveries only (a download through it is slow). A file -// the connection can't hand over comes by the original's shared link from -// then on, never by a copy in the connected account. +// the shared link: on deliveries only (a download through it is slow). The +// calls are Google's own Drive API through the broker (`proxyRequest`); no +// Composio action is run. A file over what one answer through it carries +// (250 MB) comes by the original's shared link from then on, never by a copy +// in the connected account. +const CONN_SIZES = { conn_small_________________: 1000, conn_big___________________: 9_000_000_000, conn_high__________________: 2000 }; env.CREDENTIALS = { async listConnected() { return ["googledrive"]; }, - async executeTool(service, action, args) { + async executeTool(service, action) { (world.connActions ??= []).push(action); - if (action === "GOOGLEDRIVE_GOOGLE_DRIVE_DELETE_FOLDER_OR_FILE_ACTION") { + return { data: null, error: "footage runs no action", successful: false }; + }, + // The broker as it answers since the proxy passes the provider's status and + // a file's link through: Google Drive's base already ends in /drive/v3. + async proxyRequest(service, orgId, req) { + (world.proxied ??= []).push(`${req.method} ${req.endpoint}`); + const q = Object.fromEntries((req.parameters ?? []).map((p) => [p.name, p.value])); + const m = /^\/files\/([^/?]+)$/.exec(req.endpoint); + if (service !== "googledrive" || !m) return { data: null, error: `googledrive 404: no ${req.endpoint}`, successful: false, status: 404 }; + const id = decodeURIComponent(m[1]); + if (req.method === "DELETE") { world.removeTries = (world.removeTries ?? 0) + 1; if (world.refuseRemoves > 0) { world.refuseRemoves--; - return { successful: false, error: "User rate limit exceeded" }; + return { data: { error: { code: 403 } }, error: "googledrive 403: User rate limit exceeded", successful: false, status: 403 }; } - (world.removed ??= []).push(args.fileId); - return { successful: true, data: {} }; + (world.removed ??= []).push(id); + return { data: null, error: null, successful: true, status: 204 }; + } + if (q.alt === "media" && q.Range) { + // One piece, as the broker answers a ranged read: a link to its bytes, + // with Drive's own status and range. + const [, a, b] = /^bytes=(\d+)-(\d+)$/.exec(q.Range); + (world.proxiedRanges ??= []).push(Number(a)); + return { + data: null, + error: null, + successful: true, + status: 206, + headers: { "content-range": `bytes ${a}-${b}/${CONN_SIZES[id]}` }, + file: { url: `https://temp.r2.test/piece/${id}/${a}-${b}`, contentType: "video/mp4", size: Number(b) - Number(a) + 1, expiresAt: null }, + }; } - world.connCalls = (world.connCalls ?? 0) + 1; - if (args.fileId === "conn_big___________________") return { successful: false, error: "Insufficient disk space to download file" }; - return { successful: true, data: { downloaded_file_content: { s3url: `https://temp.r2.test/${args.fileId}`, name: "x.MP4", mimetype: "video/mp4" } } }; + if (q.alt === "media") { + world.connCalls = (world.connCalls ?? 0) + 1; + const size = CONN_SIZES[id] ?? 1000; + return { data: null, error: null, successful: true, status: 200, file: { url: `https://temp.r2.test/${id}`, contentType: "video/mp4", size, expiresAt: null } }; + } + return { data: { name: "x.MP4", mimeType: "video/mp4", size: String(CONN_SIZES[id] ?? 1000) }, error: null, successful: true, status: 200 }; }, }; env.CLAWNIFY_ORG_ID = "org1"; @@ -881,13 +916,16 @@ const small = f2.items.find((i) => i.name === "A_0001.MP4"); const big = f2.items.find((i) => i.name === "Interview.MP4"); assert.equal(small.status, "importing", JSON.stringify(small)); assert.ok(world.importUrls.includes("https://temp.r2.test/conn_small_________________"), "imported from the connection's link"); -// Too big for the connection's download: the original's shared link, the same -// delivery, and no copy, permission or folder made in the connected account. +// Over what one answer through the connection carries: the original's shared +// link, the same delivery. It was never asked for whole through the broker +// (only its size), and no copy, permission or folder was made anywhere. assert.equal(big.status, "importing", JSON.stringify(big)); assert.ok(world.importUrls.includes("https://drive.usercontent.google.com/download?id=conn_big___________________&export=download&confirm=t"), "imported from the original's shared link"); assert.equal(db.prepare("SELECT link_only, copy_id FROM project_footage WHERE drive_file_id = 'conn_big___________________'").get().link_only, 2); -assert.ok(!world.connActions.some((a) => /COPY_FILE|CREATE_PERMISSION|CREATE_FOLDER/.test(a)), JSON.stringify(world.connActions)); -const callsBefore = world.connCalls; +const bigCalls = () => world.proxied.filter((c) => c.includes("conn_big")); +assert.deepEqual(bigCalls(), ["GET /files/conn_big___________________"], "only its size was asked for"); +assert.deepEqual(world.connActions ?? [], [], "footage runs no Composio action"); +const callsBefore = bigCalls().length; // A copy an earlier version made for this clip is deleted once the import is // over (Drive refuses the first try: it stays on the clip and goes on the // next), and the connection's download isn't tried again for the file it @@ -896,12 +934,44 @@ const leftover = "copy_left__________________"; db.prepare("UPDATE project_footage SET copy_id = ? WHERE drive_file_id = 'conn_big___________________'").run(leftover); world.refuseRemoves = 1; for (let i = 0; i < 6; i++) await deliver2(); -assert.equal(world.connCalls, callsBefore); +assert.equal(bigCalls().length, callsBefore, "the connection isn't asked again for the file it couldn't carry"); assert.equal(db.prepare("SELECT status FROM project_footage WHERE drive_file_id = 'conn_big___________________'").get().status, "ready"); assert.deepEqual(world.removed, [leftover], "the leftover copy is deleted once the import is over"); assert.equal(world.removeTries, 2, "a refused delete keeps the copy on the clip and is tried again"); assert.equal(db.prepare("SELECT copy_id FROM project_footage WHERE drive_file_id = 'conn_big___________________'").get().copy_id, null); -console.log("12 ok: imports through the Drive connection on deliveries; a file too big for it comes by its shared link, no copy made; a leftover copy is deleted, after a refused delete too"); +assert.deepEqual(world.connActions ?? [], [], "still no Composio action"); +console.log("12 ok: imports through the Drive API over the connection on deliveries; a file over 250 MB comes by its shared link, never fetched whole through the broker, no copy made; a leftover copy is deleted, after a refused delete too"); + +// 12b. A big file coming in by pieces, whose shared link Drive starts refusing +// part-way (past about 4 GB read, live): the pieces it refuses come through +// the connection, which still serves them, and the clip never waits. +{ + const P = 200 * 2 ** 20; + const size = 2 * P + 3 * 2 ** 20; + world.big[fid("big7")] = { size, refuseFrom: P }; + CONN_SIZES[fid("big7")] = size; + world.drive[fid("day7")] = page("Day 7", [fileEntry(fid("big7"), "INTERVIEW_E.MP4")]); + r = await call("POST", "/api/projects", { folder: `https://drive.google.com/drive/folders/${fid("day7")}` }); + const p7 = r.data.id; + const deliver7 = async () => { + const body = JSON.stringify({ project_id: p7 }); + const res = await app.request("https://open-video.apps.clawnify.com/api/footage/step", { method: "POST", headers: { "content-type": "application/json", ...(await signed(body)) }, body }, env, ctx); + assert.equal(res.status, 200, await res.text()); + }; + const row7 = () => db.prepare("SELECT status, upload_uid, asset_id, error FROM project_footage WHERE project_id = ?").get(p7); + await deliver7(); // over 250 MB for the connection: shared link, quota page, an upload opens + const u7 = row7().upload_uid; + assert.ok(u7, JSON.stringify(row7())); + world.proxiedRanges = []; + await deliver7(); // the pieces + assert.deepEqual(world.uploads.get(u7).patches, [0, P, 2 * P], "every piece in, in order"); + assert.deepEqual(world.proxiedRanges, [P, 2 * P], "the refused pieces came through the connection"); + assert.equal(row7().status, "importing"); + assert.ok(row7().asset_id, "and the clip is in, without waiting on Drive"); + r = await call("DELETE", `/api/projects/${p7}`); + assert.equal(r.status, 200, JSON.stringify(r.data)); + console.log("12b ok: pieces the shared link refuses part-way come through the Drive connection; the clip never waits"); +} // 13. A source over the video host's bitrate cap goes back in line marked for // re-encoding; the re-encode becomes a media id, and the clip is ready. diff --git a/test/relay.test.ts b/test/relay.test.ts index cb3f067..38eac1f 100644 --- a/test/relay.test.ts +++ b/test/relay.test.ts @@ -1,5 +1,5 @@ import { afterEach, describe, expect, it, vi } from "vitest"; -import { MAX_RELAY_BYTES, PIECE_BYTES, nextPiece, pieceVerdict, rangedSize, relayPieces } from "../src/server/relay"; +import { MAX_RELAY_BYTES, PIECE_BYTES, nextPiece, pieceVerdict, rangedSize, relayPieces, sharedLinkReader } from "../src/server/relay"; afterEach(() => vi.unstubAllGlobals()); @@ -36,6 +36,8 @@ function fakeUpload(size: number, at = 0) { at, /** Where each piece that landed started. */ patches: [] as number[], + /** Where each piece asked of the shared link started. */ + asked: [] as number[], headers: [] as Headers[], refuseRange: (_start: number): "quota" | number | null => null, patchFails: (_at: number): number | "throw" | null => null, @@ -50,6 +52,7 @@ function fakeUpload(size: number, at = 0) { const [, a, b] = /^bytes=(\d+)-(\d+)$/.exec(h.get("range") ?? "")!; const start = Number(a); const end = Number(b); + state.asked.push(start); const refusal = state.refuseRange(start); if (refusal === "quota") { return new Response("Google Drive - Quota exceeded", { @@ -159,14 +162,15 @@ describe("rangedSize", () => { describe("relayPieces", () => { const SIZE = 2 * PIECE_BYTES + 3 * MiB; - const file = { url: UPLOAD, fileId: "file-id", size: SIZE }; + const file = { url: UPLOAD, size: SIZE }; + const link = [sharedLinkReader("file-id")]; const later = () => Date.now() + 60_000; const quiet = async () => {}; it("sends the file in order, a piece at a time, and says when it is all in", async () => { const up = fakeUpload(SIZE); const seen: number[] = []; - const result = await relayPieces(file, later(), async (n) => { + const result = await relayPieces(file, link, later(), async (n) => { seen.push(n); }); expect(result).toEqual({ state: "done", contentType: "video/mp4" }); @@ -181,37 +185,37 @@ describe("relayPieces", () => { it("carries on from where the upload stands, not from a count of its own", async () => { const up = fakeUpload(SIZE, PIECE_BYTES + 7); - expect((await relayPieces(file, later(), quiet)).state).toBe("done"); + expect((await relayPieces(file, link, later(), quiet)).state).toBe("done"); expect(up.patches).toEqual([PIECE_BYTES + 7, 2 * PIECE_BYTES + 7]); expect(up.at).toBe(SIZE); }); it("starts no piece once its time is up", async () => { const up = fakeUpload(SIZE); - expect(await relayPieces(file, Date.now() - 1, quiet)).toEqual({ state: "moving" }); + expect(await relayPieces(file, link, Date.now() - 1, quiet)).toEqual({ state: "moving" }); expect(up.patches).toEqual([]); }); it("stops for Drive when it refuses a piece, keeping what landed", async () => { const up = fakeUpload(SIZE); up.refuseRange = (start) => (start >= PIECE_BYTES ? "quota" : null); - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "drive" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "drive" }); expect(up.at).toBe(PIECE_BYTES); }); it("leaves a hiccup to the next step: Drive busy, a lost connection, a piece the upload turned away", async () => { let up = fakeUpload(SIZE); up.refuseRange = () => 503; - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "moving" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "moving" }); up = fakeUpload(SIZE); up.patchFails = (at) => (at === PIECE_BYTES ? "throw" : null); - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "moving" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "moving" }); expect(up.at).toBe(PIECE_BYTES); up = fakeUpload(SIZE); up.patchFails = () => 500; - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "moving" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "moving" }); expect(up.at).toBe(0); }); @@ -224,23 +228,67 @@ describe("relayPieces", () => { up.at = PIECE_BYTES; // someone else's first piece return 409; }; - expect((await relayPieces(file, later(), quiet)).state).toBe("done"); + expect((await relayPieces(file, link, later(), quiet)).state).toBe("done"); expect(up.patches).toEqual([PIECE_BYTES, 2 * PIECE_BYTES]); }); it("gives up on a piece the upload turns away for good, with the host's reason", async () => { const up = fakeUpload(SIZE); up.patchFails = () => 400; - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "refused", detail: "400: Decoding Error" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "refused", detail: "400: Decoding Error" }); expect(up.at).toBe(0); }); + /** A second reader, as the Drive connection serves pieces: its own answer per range. */ + const connection = (answer: (start: number) => "ok" | number | null) => { + const served: number[] = []; + const reader = async (start: number, end: number) => { + const a = answer(start); + if (a === null) return null; + if (typeof a === "number") return new Response(null, { status: a }); + served.push(start); + return new Response(bytes(end - start + 1), { + status: 206, + headers: { "content-type": "video/mp4", "content-range": `bytes ${start}-${end}/${SIZE}` }, + }); + }; + return { reader, served }; + }; + + it("reads a piece the shared link refuses through the connection, and stays with it", async () => { + const up = fakeUpload(SIZE); + up.refuseRange = (start) => (start >= PIECE_BYTES ? "quota" : null); + const conn = connection(() => "ok"); + expect((await relayPieces(file, [...link, conn.reader], later(), quiet)).state).toBe("done"); + expect(up.patches).toEqual([0, PIECE_BYTES, 2 * PIECE_BYTES]); + expect(conn.served).toEqual([PIECE_BYTES, 2 * PIECE_BYTES]); + expect(up.asked, "once the connection served a piece, the refusing link isn't asked again").toEqual([0, PIECE_BYTES]); + }); + + it("waits on Drive only when every reader refuses the piece", async () => { + const up = fakeUpload(SIZE); + up.refuseRange = (start) => (start >= PIECE_BYTES ? "quota" : null); + const conn = connection(() => 403); + expect(await relayPieces(file, [...link, conn.reader], later(), quiet)).toEqual({ state: "drive" }); + expect(up.at).toBe(PIECE_BYTES); + }); + + it("leaves the connection's hiccup to the next step", async () => { + const up = fakeUpload(SIZE); + up.refuseRange = (start) => (start >= PIECE_BYTES ? "quota" : null); + let conn = connection(() => null); + expect(await relayPieces(file, [...link, conn.reader], later(), quiet)).toEqual({ state: "moving" }); + conn = connection(() => 503); + expect(await relayPieces(file, [...link, conn.reader], later(), quiet)).toEqual({ state: "moving" }); + expect(up.at).toBe(PIECE_BYTES); + }); + it("says when the upload can't take more", async () => { let up = fakeUpload(SIZE); up.headStatus = 404; - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "gone" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "gone" }); up = fakeUpload(SIZE); up.patchFails = () => 410; - expect(await relayPieces(file, later(), quiet)).toEqual({ state: "gone" }); + expect(await relayPieces(file, link, later(), quiet)).toEqual({ state: "gone" }); }); });