Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
f4722a0
feat: add Loom import job tables
cursoragent Oct 7, 2026
591137b
feat: run Loom CSV imports as durable server-side jobs
cursoragent Oct 7, 2026
bb06457
fix: only transcribe imported Loom videos once someone watches them
cursoragent Oct 7, 2026
9aca43c
feat: guided Loom bulk import with live progress and upgrade flow
cursoragent Oct 7, 2026
8c3c031
improve: keep the Loom import page fast at 2,000 videos
cursoragent Oct 7, 2026
72f07db
test: tie the Loom import delta test to the poll overlap window
cursoragent Oct 7, 2026
99c9b80
improve: count only importable videos in the Loom import heading
cursoragent Oct 7, 2026
a9aee2a
docs: describe the 2,000 video Loom import and AI on first view
cursoragent Oct 7, 2026
04d0ebf
improve: tidy Loom import actions, dark mode tints and hydration
cursoragent Oct 7, 2026
fb71e35
fix: keep the Loom import workflows inside the lean workflow runtime
cursoragent Oct 7, 2026
e0b4aae
chore: format AI entitlement helper
cursoragent Oct 7, 2026
7d04f8c
improve: hide Loom thumbnails that fail to load
cursoragent Oct 7, 2026
a35e566
chore: format generated Loom import migration metadata
cursoragent Oct 7, 2026
605a487
fix: limit how many Loom imports one person can start per hour
cursoragent Oct 7, 2026
95caa40
fix: harden Loom import starts, dispatch locking and client errors
cursoragent Oct 7, 2026
ed48f7f
fix: claim stuck Loom import recoveries and settle rows finished outs…
cursoragent Oct 7, 2026
cfecf47
fix: keep the Loom import job page in order and in sync with the job
cursoragent Oct 7, 2026
5ff19b2
fix: keep the first video in CSVs of bare Loom ids
cursoragent Oct 7, 2026
8a0bdc9
fix: give each desktop audio timing test its own temp folder
cursoragent Oct 8, 2026
d3c4e1b
feat: add a dispatch lock, fairness column and queue indexes for Loom…
cursoragent Oct 8, 2026
a93561f
feat: run Loom CSV imports through one fair, capacity-aware queue
cursoragent Oct 8, 2026
36e3bec
test: load test Loom imports with 50 people and 100,000 rows
cursoragent Oct 8, 2026
4a06aad
fix: free up to 500 silent Loom import slots per recovery run in one …
cursoragent Oct 8, 2026
a5415a0
fix: show a Loom video retried from its page as copying in its bulk i…
cursoragent Oct 8, 2026
08bb6be
fix: restart a stalled Loom upload once so recovery can't keep it hol…
cursoragent Oct 8, 2026
1d37a5c
fix: give long videos on the media server time in proportion to their…
cursoragent Oct 8, 2026
894d65c
fix: wait for long Loom imports as long as the media server allows
cursoragent Oct 8, 2026
556b7f8
improve: wait out longer Loom rate limits before failing link checks
cursoragent Oct 8, 2026
c8c29cc
test: load test 29 people importing at once with long recordings
cursoragent Oct 8, 2026
2b50fbe
fix: stop media jobs that stop making progress instead of capping lon…
cursoragent Oct 8, 2026
efd8c27
fix: wait on video processing for as long as it makes progress
cursoragent Oct 8, 2026
d6532e7
improve: poll Loom import progress less often and not in background tabs
cursoragent Oct 8, 2026
50a9d39
fix: keep regular upload links valid for as long as processing makes …
cursoragent Oct 8, 2026
1bc9e91
fix: count decoded frames as progress while checking long WebM record…
cursoragent Oct 8, 2026
8dedf3a
fix: ask the media server about a quiet job before starting another copy
cursoragent Oct 8, 2026
8fa4b21
fix: restart a quiet video only after its job stays unseen for ten mi…
cursoragent Oct 8, 2026
baedf7b
fix: keep regular upload progress live while throttling bulk import u…
cursoragent Oct 8, 2026
cccb387
fix: keep the import page refreshing as often as before while it is v…
cursoragent Oct 8, 2026
a986bf2
chore: rebuild the preview deployment
cursoragent Oct 8, 2026
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
7 changes: 5 additions & 2 deletions apps/desktop-gpui/src/editor_preparing/audio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -174,13 +174,16 @@ mod tests {

impl TestDirectory {
fn new() -> Self {
static NEXT_DIRECTORY: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let path = std::env::temp_dir().join(format!(
"cap-gpui-preparing-audio-{}-{nonce}",
std::process::id()
"cap-gpui-preparing-audio-{}-{nonce}-{}",
std::process::id(),
NEXT_DIRECTORY.fetch_add(1, std::sync::atomic::Ordering::Relaxed),
));
std::fs::create_dir(&path).unwrap();
Self(path)
Expand Down
37 changes: 37 additions & 0 deletions apps/media-server/src/__tests__/lib/idle-timeout.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
import { describe, expect, test } from "bun:test";
import { withIdleTimeout } from "../../lib/media-common";

const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));

describe("withIdleTimeout", () => {
test("lets slow work finish for as long as it keeps reporting progress", async () => {
const result = await withIdleTimeout(async (touch) => {
for (let step = 0; step < 8; step++) {
await wait(30);
touch();
}
return "done";
}, 80);
expect(result).toBe("done");
});

test("stops work that goes quiet and runs its cleanup", async () => {
let cleaned = false;
const started = performance.now();
await expect(
withIdleTimeout(
async (touch) => {
touch();
await wait(1_000);
return "late";
},
60,
() => {
cleaned = true;
},
),
).rejects.toThrow("Stopped making progress");
expect(performance.now() - started).toBeLessThan(500);
expect(cleaned).toBe(true);
});
});
58 changes: 58 additions & 0 deletions apps/media-server/src/__tests__/lib/job-manager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,14 @@ import {
createJob,
deleteJob,
getJob,
JOB_PROGRESS_STALL_MS,
type JobProgress,
markJobProgress,
type RecordingWorkerAcknowledgement,
sendWebhook,
touchJob,
updateJob,
watchJobProgress,
} from "../../lib/job-manager";

const createdJobs: string[] = [];
Expand Down Expand Up @@ -521,4 +524,59 @@ describe("job cleanup", () => {
expect(currentJob?.phase).toBe("error");
expect(currentJob?.error).toContain("maximum lifetime of 60 minutes");
});

test("keeps a job that is still making progress running for hours", () => {
const job = createTrackedJob("job-long-but-moving");
const now = Date.now();
job.phase = "processing";
watchJobProgress(job.jobId);
job.createdAt = now - 3 * 60 * 60 * 1000;
job.updatedAt = now;
job.progressAt = now - 60_000;

expect(cleanupExpiredJobs()).toBe(0);
expect(getJob(job.jobId)?.phase).toBe("processing");
});

test("fails a job that stopped making progress, however short it is", () => {
const job = createTrackedJob("job-short-but-stuck");
const now = Date.now();
job.phase = "processing";
watchJobProgress(job.jobId);
job.createdAt = now - 20 * 60 * 1000;
job.updatedAt = now;
job.progressAt = now - JOB_PROGRESS_STALL_MS - 60_000;

expect(cleanupExpiredJobs()).toBe(1);
expect(getJob(job.jobId)).toMatchObject({
phase: "error",
message: "Processing failed (stalled)",
});
expect(getJob(job.jobId)?.error).toContain("stopped making progress");
});

test("counts rising progress and new phases as progress, not heartbeats", () => {
const job = createTrackedJob("job-progress-signals");
job.phase = "processing";
job.progress = 20;
watchJobProgress(job.jobId);
const old = Date.now() - 10 * 60 * 1000;

job.progressAt = old;
touchJob(job.jobId);
updateJob(job.jobId, { message: "Still here" });
updateJob(job.jobId, { progress: 20 });
expect(job.progressAt).toBe(old);

updateJob(job.jobId, { progress: 21 });
expect(job.progressAt).toBeGreaterThan(old);

job.progressAt = old;
updateJob(job.jobId, { phase: "uploading" });
expect(job.progressAt).toBeGreaterThan(old);

job.progressAt = old;
markJobProgress(job.jobId);
expect(job.progressAt).toBeGreaterThan(old);
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import {
repairContainer,
uploadFileToS3,
uploadFileToStorage,
uploadTimeoutMs,
uploadToS3,
} from "../../lib/media-video";

Expand Down Expand Up @@ -1008,6 +1009,80 @@ describe("processVideo integration tests", () => {
expect(existsSync(tempFile.path)).toBe(false);
}, 60000);

test("lets an encode run past its idle limit while frames keep advancing", async () => {
const dir = mkdtempSync(join(tmpdir(), "advancing-encode-"));
const longInput = join(dir, "input.mp4");
execFileSync("ffmpeg", [
"-hide_banner",
"-v",
"error",
"-f",
"lavfi",
"-i",
"testsrc2=size=1280x720:rate=30",
"-t",
"60",
"-c:v",
"libx264",
"-preset",
"ultrafast",
"-pix_fmt",
"yuv420p",
"-y",
longInput,
]);
const idleTimeoutMs = 1_500;
try {
const metadata = await probeVideo(`file://${longInput}`);
const progressUpdates: number[] = [];
const started = performance.now();

const tempFile = await processVideo(
longInput,
metadata,
{ maxWidth: 960, maxHeight: 540, preset: "medium", idleTimeoutMs },
(progress) => progressUpdates.push(progress),
);
tempFiles.push(tempFile.path);

expect(performance.now() - started).toBeGreaterThan(idleTimeoutMs * 2);
expect(progressUpdates.length).toBeGreaterThan(2);
for (let index = 1; index < progressUpdates.length; index++) {
expect(progressUpdates[index]).toBeGreaterThan(
progressUpdates[index - 1] ?? -1,
);
}
await tempFile.cleanup();
} finally {
rmSync(dir, { recursive: true, force: true });
}
}, 120000);

test("stops an encode whose input stops producing frames", async () => {
const metadata = await probeVideo(`file://${TEST_VIDEO_WITH_AUDIO}`);
const dir = mkdtempSync(join(tmpdir(), "stalled-encode-"));
const stalledInput = join(dir, "input.mp4");
execFileSync("mkfifo", [stalledInput]);
try {
const started = performance.now();
await expect(
processVideo(stalledInput, metadata, {
remuxOnly: true,
idleTimeoutMs: 1_000,
}),
).rejects.toThrow("Stopped making progress");
expect(performance.now() - started).toBeLessThan(15_000);
} finally {
rmSync(dir, { recursive: true, force: true });
}
}, 30000);

test("gives large uploads time in proportion to their size", () => {
expect(uploadTimeoutMs(0)).toBe(10 * 60 * 1000);
expect(uploadTimeoutMs(200 * 1024 * 1024)).toBe(10 * 60 * 1000);
expect(uploadTimeoutMs(3 * 1024 ** 3)).toBe(3 * 1024 * 1000);
});

test("respects CRF setting", async () => {
const metadata = await probeVideo(`file://${TEST_VIDEO_WITH_AUDIO}`);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -208,4 +208,37 @@ describe("raw recording input validation", () => {
"timed out",
);
});

test("reports decoding progress while it checks a clean recording", async () => {
let reports = 0;
await expect(
validateVideoInput(clean, undefined, undefined, {
idleTimeoutMs: 5_000,
onProgress: () => reports++,
}),
).resolves.toBeUndefined();
expect(reports).toBeGreaterThan(0);
});

test("still rejects a damaged recording while tracking progress", async () => {
await expect(
validateVideoInput(corrupt, undefined, undefined, {
idleTimeoutMs: 5_000,
onProgress: () => {},
}),
).rejects.toThrow("original upload has been preserved");
});

test("stops a check whose input stops producing frames", async () => {
const stalled = join(directory, "stalled.webm");
execFileSync("mkfifo", [stalled]);
let reports = 0;
await expect(
validateVideoInput(stalled, undefined, undefined, {
idleTimeoutMs: 1_000,
onProgress: () => reports++,
}),
).rejects.toThrow("Stopped making progress");
expect(reports).toBe(0);
}, 30_000);
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
import { afterAll, afterEach, describe, expect, spyOn, test } from "bun:test";
import { copyFileSync, mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import app from "../../app";
import * as jobManager from "../../lib/job-manager";
import { deleteJob, getJob } from "../../lib/job-manager";
import * as mediaVideo from "../../lib/media-video";

const fixture = join(import.meta.dir, "..", "fixtures", "test-with-audio.mp4");
const directory = mkdtempSync(join(tmpdir(), "cap-progress-webhooks-"));
const received: { videoId: string; phase: string; message?: string }[] = [];
const receiver = Bun.serve({
port: 0,
async fetch(request) {
received.push(await request.json());
return Response.json({ success: true });
},
});

afterAll(() => {
receiver.stop(true);
rmSync(directory, { recursive: true, force: true });
});

const spies: { mockRestore: () => void }[] = [];
afterEach(() => {
for (const spy of spies.splice(0)) spy.mockRestore();
});

async function encodeWithRapidProgress(priority: "normal" | "bulk") {
const videoId = `progress-${priority}`;
const input = join(directory, `${priority}-input.mp4`);
const output = join(directory, `${priority}-output.mp4`);
copyFileSync(fixture, input);
copyFileSync(fixture, output);
spies.push(
spyOn(jobManager, "canAcceptNewVideoProcess").mockReturnValue(true),
spyOn(mediaVideo, "downloadVideoToTemp").mockResolvedValue({
path: input,
cleanup: async () => {},
}),
spyOn(mediaVideo, "processVideo").mockImplementation(
async (_input, _metadata, _options, onProgress) => {
for (let step = 1; step <= 40; step++) {
onProgress?.(step * 2.5, `Encoding: ${step * 2.5}%`);
await Bun.sleep(60);
}
return { path: output, cleanup: async () => {} };
},
),
spyOn(mediaVideo, "uploadFileToS3").mockResolvedValue({}),
);
const secret = process.env.MEDIA_SERVER_WEBHOOK_SECRET ?? "test-secret";
process.env.MEDIA_SERVER_WEBHOOK_SECRET = secret;
const response = await app.fetch(
new Request("http://localhost/video/process", {
method: "POST",
headers: {
"Content-Type": "application/json",
"x-media-server-secret": secret,
},
body: JSON.stringify({
videoId,
userId: "progress-test",
videoUrl: "https://example.com/raw.mp4",
outputPresignedUrl: "https://example.com/result.mp4",
webhookUrl: `http://127.0.0.1:${receiver.port}/progress`,
webhookSecret: secret,
priority,
}),
}),
);
expect(response.status).toBe(200);
const { jobId } = (await response.json()) as { jobId: string };
const deadline = Date.now() + 15_000;
while (Date.now() < deadline && getJob(jobId)?.phase !== "complete") {
await Bun.sleep(20);
}
expect(getJob(jobId)?.phase).toBe("complete");
deleteJob(jobId);
return received.filter(
(update) =>
update.videoId === videoId && update.message?.startsWith("Encoding"),
).length;
}

describe("encoding progress webhooks", () => {
test("keeps regular uploads updating every second for the share page", async () => {
const sent = await encodeWithRapidProgress("normal");
expect(sent).toBeGreaterThanOrEqual(2);
expect(sent).toBeLessThanOrEqual(4);
}, 30_000);

test("sends bulk imports one update per five seconds", async () => {
expect(await encodeWithRapidProgress("bulk")).toBe(1);
}, 30_000);
});
Loading
Loading