Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 13 additions & 7 deletions agent.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ otherwise), `next_offset`, and `items`:
{
"id": "…", "name": "B_0012.MP4", "folder": "Day 1/Cam B",
"status": "ready", // waiting | importing | ready | failed
"received": null, // while it comes in a piece at a time: { "bytes", "size" }
"asset": { "id": "…", "duration": 41.2 }, // the EDL uses "asset:<id>"
"log_status": "done", // preparing | running | done | failed
"log": {
Expand Down Expand Up @@ -149,19 +150,24 @@ each file's Drive link.
cut from the kept soundbites" places them with its `add_highlight` operation.

Google Drive limits how often a file shared with the link can be downloaded,
and on a heavy day it refuses the files of a whole folder for a while. A clip it
refuses is not failed: it waits (`status: "waiting"` with `retry_at`, counted in
and on a heavy day it refuses the files of a whole folder for a while. It
refuses the whole file but still serves pieces of it, so a refused file of up
to 30 GB comes in a piece at a time instead, in the background steps only:
`status: "importing"`, with `received` saying how much is in. A big file takes
a while this way. A clip Drive won't serve even a piece of, or one over 30 GB,
is not failed: it waits (`status: "waiting"` with `retry_at`, counted in
`counts.drive_waiting`) and is tried again by itself, further apart each time,
for about a day, after which it fails with the reason. `imports_paused` says so
when every waiting clip is waiting on Drive. A copy of the folder in another
Drive account has its own limit.
for about a day, after which it fails with the reason. One that was coming in
by pieces carries on from where it got to. `imports_paused` says so when every
waiting clip is waiting on Drive. A copy of the folder in another Drive account
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, with the
waits above. It is never copied inside Google first: Drive answers a header
(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
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).
Expand Down
11 changes: 11 additions & 0 deletions schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -219,3 +219,14 @@ ALTER TABLE project_footage ADD COLUMN copy_id TEXT;
-- other refusal sends it back to be re-encoded.
ALTER TABLE project_footage ADD COLUMN transcode INTEGER NOT NULL DEFAULT 0;
ALTER TABLE project_footage ADD COLUMN transcode_job TEXT;

-- A file Google Drive won't hand over whole (its download limit) still serves
-- byte ranges, so it comes in a piece at a time, into an upload on the media
-- service (src/server/relay.ts). Set while that upload is open: its id and
-- link, when the link expires, the file's size, and how much is in so far
-- (for showing; the upload itself says where it stands).
ALTER TABLE project_footage ADD COLUMN upload_uid TEXT;
ALTER TABLE project_footage ADD COLUMN upload_url TEXT;
ALTER TABLE project_footage ADD COLUMN upload_expires TEXT;
ALTER TABLE project_footage ADD COLUMN upload_size INTEGER;
ALTER TABLE project_footage ADD COLUMN upload_done INTEGER;
5 changes: 5 additions & 0 deletions src/client/edit.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,8 @@ interface FootageItem {
error: string | null;
/** Set while it waits on Google Drive: when it is tried again. */
retry_at?: string | null;
/** Set while a file Drive won't hand over whole comes in a piece at a time. */
received?: { bytes: number; size: number } | null;
asset: Asset | null;
log_status: "preparing" | "running" | "done" | "failed" | null;
log_error: string | null;
Expand Down Expand Up @@ -2055,6 +2057,9 @@ function footageLine(i: FootageItem): string {
return `Waiting for Google Drive, which is limiting downloads of this file. Trying again at ${at}`;
}
if (i.status === "waiting") return "Waiting to import";
if (i.status === "importing" && i.received) {
return `Importing: ${Math.floor((i.received.bytes / i.received.size) * 100)}% of ${fmtBytes(i.received.size)}`;
}
if (i.status === "importing") return "Importing";
if (i.status === "failed") return `Not imported: ${i.error ?? "unknown error"}`;
if (i.log_status === "done" && i.log) return i.log.summary;
Expand Down
181 changes: 161 additions & 20 deletions src/server/footage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,9 @@
// A shoot arrives as a shared folder: a folder per camera, a few hundred clips
// a day. Every video in it becomes a row of project_footage, which moves on by
// itself:
// waiting → importing (the media service pulls it from Drive) → ready
// waiting → importing (the media service pulls it from Drive, or, for a
// file Drive won't hand over whole, this app sends it a piece at
// a time: src/server/relay.ts) → ready
// and once ready, its log:
// preparing (the platform makes the copy analysis reads) → running → done
//
Expand All @@ -16,9 +18,10 @@

import { enqueueJob, type QueueEnv } from "@clawnify/queue";
import { query, get, run } from "./db";
import { deleteMedia, importMedia, mediaState, prepareMedia, startTranscode, transcodeState, type MediaConfig } from "./media";
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";

const DEFAULT_SERVICES_URL = "https://services.clawnify.com";

Expand Down Expand Up @@ -324,6 +327,10 @@ interface WorkRow {
copy_id: string | null;
transcode: number;
transcode_job: string | null;
upload_uid: string | null;
upload_url: string | null;
upload_expires: string | null;
upload_size: number | null;
}

/**
Expand Down Expand Up @@ -354,6 +361,8 @@ export interface StepOptions {
drive?: DriveSource;
/** False: start no imports this step (a read, while imports go through the connection, which is slow). */
startImports?: boolean;
/** Move files coming in a piece at a time a few pieces on: deliveries only, as it takes a while. */
relay?: boolean;
}

export interface StepOutcome {
Expand Down Expand Up @@ -409,7 +418,8 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste
}
const rows = await query<WorkRow>(
`SELECT f.id, f.status, f.drive_file_id, f.name, f.folder, f.language, f.asset_id, f.log_status, f.log_job,
f.updated_at, f.retry_at, f.drive_tries, f.link_only, f.copy_id, f.transcode, f.transcode_job, a.media_uid, a.duration
f.updated_at, f.retry_at, f.drive_tries, f.link_only, f.copy_id, f.transcode, f.transcode_job,
f.upload_uid, f.upload_url, f.upload_expires, f.upload_size, a.media_uid, a.duration
FROM project_footage f LEFT JOIN assets a ON a.id = f.asset_id
WHERE f.project_id = ?
AND (f.status IN ('waiting', 'importing')
Expand All @@ -427,6 +437,7 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste

// 1. Imports that finished, or never will.
const importing = rows.filter((r) => r.status === "importing");
const relaying: WorkRow[] = [];
await Promise.all(
importing.map(async (r) => {
// Being re-encoded: once it has a media id it imports like any other.
Expand All @@ -452,6 +463,12 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste
}
return;
}
// Coming in a piece at a time: moved on below, on deliveries. It has no
// media id until the last piece is in, and is not a stale claim.
if (r.upload_uid) {
relaying.push(r);
return;
}
if (!r.media_uid) {
if (age(r.updated_at) > STALE_CLAIM_MS) await setRow(r.id, { status: "waiting" }, " AND status = 'importing'");
return;
Expand Down Expand Up @@ -484,6 +501,13 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste
}),
);

// 1b. Files Drive won't hand over whole: a few more pieces each. The time
// 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)));
}

// 2. Start imports while there is room, together. A refusal that is about
// the org (storage full) puts its clip back in line and starts no more.
const room = opts.startImports === false ? 0 : STEP_LIMITS.importing - importing.filter((r) => r.status === "importing").length;
Expand All @@ -499,15 +523,7 @@ export async function stepFootage(cfg: MediaConfig, projectId: string, opts: Ste
const started = await startImport(cfg, r, opts.drive);
if (started === true) return;
if (started.driveLimit) {
// Drive lifts its limit within a day: the clip waits and is tried
// again by itself, further apart each time, then gives up.
const tries = r.drive_tries + 1;
if (tries > DRIVE_RETRY_MS.length) {
await setRow(r.id, { status: "failed", error: DRIVE_QUOTA, drive_tries: tries });
} else {
const at = new Date(Date.now() + DRIVE_RETRY_MS[tries - 1]).toISOString();
await setRow(r.id, { status: "waiting", error: DRIVE_WAIT, drive_tries: tries, retry_at: at });
}
await waitOnDrive(cfg, r);
return;
}
if (started.orgWide) {
Expand Down Expand Up @@ -640,6 +656,120 @@ const DRIVE_WAIT_ALL =
/** How long a clip Drive refused waits before each new try: about a day in all. */
const DRIVE_RETRY_MS = [30, 60, 120, 240, 240, 240, 240, 240].map((m) => m * 60_000);

/**
* Drive refused the clip. It lifts its limit within a day, so the clip waits
* and is tried again by itself, further apart each time, then gives up. A clip
* coming in a piece at a time keeps its upload while it waits, and carries on
* from there.
*/
async function waitOnDrive(cfg: MediaConfig, r: WorkRow): Promise<void> {
const tries = r.drive_tries + 1;
if (tries > DRIVE_RETRY_MS.length) {
// An upload that can't be deleted now stays on the row: retrying the clip
// or deleting the project deletes it.
await dropUpload(cfg, r);
await setRow(r.id, { status: "failed", error: DRIVE_QUOTA, drive_tries: tries });
r.status = "failed";
} else {
const at = new Date(Date.now() + DRIVE_RETRY_MS[tries - 1]).toISOString();
await setRow(r.id, { status: "waiting", error: DRIVE_WAIT, drive_tries: tries, retry_at: at });
Object.assign(r, { status: "waiting", drive_tries: tries, retry_at: at });
}
}

const NO_UPLOAD = { upload_uid: null, upload_url: null, upload_expires: null, upload_size: null, upload_done: null };

/**
* Delete a clip's open upload on the media service, then forget it. False
* when the service wouldn't: it stays on the row, to be deleted later, rather
* than forgotten while it still holds storage.
*/
async function dropUpload(cfg: MediaConfig, r: WorkRow): Promise<boolean> {
if (!r.upload_uid) return true;
try {
await deleteMedia(cfg, r.upload_uid);
} catch {
return false;
}
await setRow(r.id, NO_UPLOAD);
Object.assign(r, NO_UPLOAD);
return true;
}

/**
* Drive won't hand the file over whole: open an upload for it to come in a
* piece at a time instead. Null when that can't help, so the clip waits on
* Drive: it won't serve even a piece, the file is over the video host's limit,
* or the clip is marked for a re-encode, which reads the whole file.
*/
async function openRelay(cfg: MediaConfig, r: WorkRow): Promise<true | { orgWide: boolean; detail: string } | null> {
if (r.transcode) return null;
const size = await rangedSize(r.drive_file_id);
if (!size || size > MAX_RELAY_BYTES) return null;
await dropFailedCopy(cfg, r);
const opened = await openMediaUpload(cfg, size, r.name);
if ("failure" in opened) return ORG_LIMITS.has(opened.failure.error) ? { orgWide: true, detail: opened.failure.detail } : null;
await setRow(r.id, {
upload_uid: opened.id,
upload_url: opened.uploadUrl,
upload_expires: opened.expiresAt,
upload_size: size,
upload_done: 0,
});
return true;
}

/** Move a clip coming in a piece at a time on, and settle it once it is all in. */
async function relayClip(cfg: MediaConfig, r: WorkRow, until: number): Promise<void> {
const uid = r.upload_uid!;
const expiring = r.upload_expires !== null && Date.parse(r.upload_expires) - Date.now() < RELAY_EXPIRY_MARGIN_MS;
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 setRow(r.id, { upload_done: received });
});
if (result.state === "moving") return;
if (result.state === "done") return finishRelay(r, uid, result.contentType);
if (result.state === "drive") return waitOnDrive(cfg, r);
if (result.state === "refused") {
// Sending it again would be refused again: the clip fails, saying why.
await dropUpload(cfg, r);
await setRow(r.id, { status: "failed", error: `the video host refused a piece of this file (${result.detail})` });
r.status = "failed";
return;
}
// The upload can't take more. One whose last piece landed without being
// recorded (the step ended first) is the video: it is recorded now.
const s = await mediaState(cfg, uid);
if (!("failure" in s) && s.media.state !== "pendingupload") return finishRelay(r, uid, "video/mp4");
// Otherwise it starts again with a new upload, counted like a refusal so a
// file that never gets in gives up in the end.
if (!(await dropUpload(cfg, r))) return;
const tries = r.drive_tries + 1;
if (tries > DRIVE_RETRY_MS.length) {
await setRow(r.id, { status: "failed", error: "the import did not finish", drive_tries: tries });
r.status = "failed";
} else {
await setRow(r.id, { status: "waiting", drive_tries: tries });
Object.assign(r, { status: "waiting", drive_tries: tries });
}
}

/** The last piece is in: the upload is the clip's video, imported like any other. */
async function finishRelay(r: WorkRow, uid: string, contentType: string): Promise<void> {
const key = `media/${uid}`;
// Two steps that overlap can both get here: the key is unique, so one asset.
await run("INSERT INTO assets (key, name, content_type, size, media_uid) VALUES (?, ?, ?, ?, ?) ON CONFLICT(key) DO NOTHING", [
key,
r.name,
contentType,
r.upload_size ?? 0,
uid,
]);
const asset = await get<{ id: string }>("SELECT id FROM assets WHERE key = ?", [key]);
await setRow(r.id, { asset_id: asset!.id, ...NO_UPLOAD });
}

/**
* Check the Drive file serves video, then have the media service pull it.
*
Expand All @@ -654,6 +784,9 @@ async function startImport(
r: WorkRow,
drive?: DriveSource,
): Promise<true | { orgWide: boolean; detail: string; driveLimit?: true }> {
// 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,
// with its waits. Not by a copy in the connected account: Drive answers a
Expand All @@ -674,7 +807,10 @@ async function startImport(
if (!page) await probe.body?.cancel();
const verdict = judgeLinkResponse(probe.status, served, probe.headers.get("content-range"), probe.headers.get("content-length"));
if (!verdict.ok) {
if (/<title>[^<]*quota exceeded/i.test(page)) return { orgWide: false, detail: DRIVE_QUOTA, driveLimit: true };
// Drive's limit refuses the whole file but still serves pieces of it.
if (/<title>[^<]*quota exceeded/i.test(page)) {
return (await openRelay(cfg, r)) ?? { orgWide: false, detail: DRIVE_QUOTA, driveLimit: true };
}
return { orgWide: false, detail: verdict.reason ?? "that file isn't a video" };
}

Expand All @@ -684,6 +820,17 @@ async function startImport(
/** The video host's refusal of a source over its bitrate cap. */
const BITRATE_REFUSED = /bitrate exceeded/i;

/**
* A copy from an attempt that failed is of no use and still counts against
* the org's storage: it goes before the clip is imported again.
*/
async function dropFailedCopy(cfg: MediaConfig, r: WorkRow): Promise<void> {
if (!r.asset_id) return;
if (r.media_uid) await deleteMedia(cfg, r.media_uid).catch(() => {});
await run("DELETE FROM assets WHERE id = ?", [r.asset_id]);
await setRow(r.id, { asset_id: null });
}

/** Have the media service pull the clip from `url`, and record it as an asset. */
async function importFrom(
cfg: MediaConfig,
Expand All @@ -692,13 +839,7 @@ async function importFrom(
type: string,
size: number,
): Promise<true | { orgWide: boolean; detail: string }> {
// A copy from an attempt that failed is of no use and still counts against
// the org's storage: it goes before the clip is imported again.
if (r.asset_id) {
if (r.media_uid) await deleteMedia(cfg, r.media_uid).catch(() => {});
await run("DELETE FROM assets WHERE id = ?", [r.asset_id]);
await setRow(r.id, { asset_id: null });
}
await dropFailedCopy(cfg, r);
// Marked for re-encoding: the media id comes later, from the re-encode.
if (r.transcode) {
const started = await startTranscode(cfg, url, r.name);
Expand Down
Loading
Loading