Skip to content

Commit 1995a61

Browse files
committed
fix: batch recording source checkpoints without losing retry safety
1 parent babd5e6 commit 1995a61

3 files changed

Lines changed: 359 additions & 2 deletions

File tree

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
[
2+
{
3+
"name": "long-s3",
4+
"terminalResult": "ready",
5+
"baselineWallMs": 919500,
6+
"expectedBatchedSteps": 31,
7+
"pageDurationsMs": [
8+
464, 462, 512, 404, 577, 373, 517, 384, 672, 368, 458, 340, 592, 372, 516,
9+
422, 841, 477, 559, 358, 600, 423, 508, 369, 621, 353, 449, 384, 546, 694,
10+
442, 508, 837, 338, 439, 334, 596, 410, 436, 346, 682, 377, 442, 388, 557,
11+
375, 484, 381, 757, 433, 483, 333, 564, 368, 510, 362, 662, 340, 418, 351,
12+
538, 490, 512, 394, 928, 364, 429, 379, 547, 362, 449, 382, 668, 373, 436,
13+
387, 550, 363, 442, 377, 731, 372, 488, 398, 578, 429, 529, 344, 625, 532,
14+
469, 387, 576, 365, 501, 415, 901, 375, 523, 411, 668, 894, 496, 438, 695,
15+
414, 480, 359, 687, 393, 493, 401, 781, 439, 531, 422, 590, 416, 461, 378,
16+
702, 369, 448, 370, 549, 356, 477, 408, 1058, 383, 513, 417, 524, 368,
17+
470, 378, 653, 363, 589, 369, 545, 360, 495, 392, 787, 387, 483, 358, 621,
18+
448, 482, 377, 638, 345, 701, 417, 634, 381, 471, 504, 850, 415, 717, 489,
19+
682, 1794, 1821, 1800, 1905, 1693, 1710, 1671, 2152, 1734, 1919, 1649,
20+
1948, 1642, 1776, 1731, 2183, 1672, 1965, 1607, 1773, 1521, 1634, 1626,
21+
1744, 1566, 1607, 1470, 1664, 1640, 1744, 1500, 2105, 1534, 1659, 1452,
22+
1837, 1451, 1637, 1592, 1787, 1608, 1671, 1607, 1673, 1737, 1614, 1550,
23+
1930, 1772, 1613, 1602, 1758, 1737, 1652, 1674, 1912, 1743, 1722, 1880,
24+
1895, 1713, 1642, 1600, 2281, 1694, 1926, 1797, 2012, 1641, 1813, 1555,
25+
1969, 1804, 1870, 1570, 2104, 1655, 1985, 1704, 1935, 1824, 2122, 1694,
26+
1829, 1577, 1654, 1588, 1850, 1661, 1899, 1723, 1767, 1604, 1615, 1625,
27+
2481, 1714, 1571, 1655, 1710, 1596, 1627, 1557, 1889, 1338, 1522, 1584,
28+
1698, 1703, 1662, 1835, 1748, 1770, 1595, 1393, 1699, 1676, 1533, 1423,
29+
1954, 1445, 1595, 1551, 1611, 1457, 1709, 1453, 2335, 1532, 1519, 1394,
30+
1512, 1335, 1492, 1405, 1826, 1765, 1479, 1385, 1626, 1342, 1513, 1438,
31+
1688, 1434, 1761, 1545, 1533, 1890, 1659, 1571, 1927, 1325, 1745, 1550,
32+
1633, 1423, 1592, 1556, 1703, 1459, 1362, 1337, 1517, 1178, 786, 719, 611,
33+
702, 739, 693, 630, 966, 800, 713, 676, 920, 872, 695, 678, 684, 680, 644,
34+
924, 775, 684, 791, 634, 815, 724, 775, 688, 690, 631, 745, 973, 643, 636,
35+
691, 686, 751, 671, 762, 680, 838, 629, 608, 640, 929, 805, 708, 648, 704,
36+
677, 643, 729, 705, 675, 762, 940, 869, 689, 1050, 799, 701, 696, 632,
37+
768, 711, 723, 723, 715, 695, 660, 642, 671, 642, 735, 1002, 705, 694,
38+
668, 700, 830, 686, 640, 709, 678, 656, 716, 649, 681, 733, 750, 644, 828,
39+
701, 733, 759, 651, 873, 657, 664, 675, 667, 938, 646, 764, 756, 609,
40+
1117, 820, 627, 619, 660, 585, 622, 744, 763, 741, 769, 641, 765, 794,
41+
736, 779, 687, 942, 665, 635, 752, 618, 686, 809, 780, 623, 914, 572, 620,
42+
627, 708, 579, 729, 713, 604, 613, 592, 629, 555, 672, 644, 648, 646, 588,
43+
587, 799, 681, 644, 666, 564, 590, 795, 772, 686, 444, 485, 464, 493, 542
44+
]
45+
},
46+
{
47+
"name": "short-drive-first",
48+
"terminalResult": "ready",
49+
"baselineWallMs": 29925,
50+
"expectedBatchedSteps": 2,
51+
"pageDurationsMs": [4693, 2796, 10156, 1521, 8332]
52+
},
53+
{
54+
"name": "short-drive-retry",
55+
"terminalResult": "ready",
56+
"baselineWallMs": 26610,
57+
"expectedBatchedSteps": 2,
58+
"pageDurationsMs": [4867, 2914, 8932, 2810, 4612]
59+
},
60+
{
61+
"name": "incomplete-source",
62+
"terminalResult": "source-blocked",
63+
"baselineWallMs": 185,
64+
"expectedBatchedSteps": 1,
65+
"pageDurationsMs": [185]
66+
}
67+
]

apps/web/__tests__/unit/finalize-desktop-recording.test.ts

Lines changed: 258 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ import type {
55
DesktopRecordingAttempt,
66
DesktopRecordingJob,
77
} from "@/lib/desktop-recording-jobs";
8+
import type { DesktopRecordingSourceCheckpoint } from "@/lib/desktop-recording-source-checkpoint";
9+
import type { RecordingVerification } from "@/lib/desktop-recording-verification";
810
import { MediaProcessingBudgetError } from "@/lib/media-processing-budget";
11+
import sourceCommitTimings from "../fixtures/source-commit-timings.json";
912

1013
const mocks = vi.hoisted(() => ({
1114
state: vi.fn(),
@@ -465,8 +468,262 @@ describe("short durable processing polls", () => {
465468
});
466469
});
467470

471+
describe("bounded source commitment steps", () => {
472+
beforeEach(() => {
473+
withCurrent({ source: null, state: "committing" });
474+
});
475+
476+
function advanceCheckpoint(
477+
_video: unknown,
478+
checkpoint: DesktopRecordingSourceCheckpoint,
479+
) {
480+
return { checkpoint: { ...checkpoint, revision: checkpoint.revision + 1 } };
481+
}
482+
483+
it.each(sourceCommitTimings)(
484+
"replays the $name production timings without repeating or skipping a page",
485+
async ({ pageDurationsMs, expectedBatchedSteps, terminalResult }) => {
486+
const revisions: number[] = [];
487+
mocks.commitSource.mockImplementation(async (_video, checkpoint) => {
488+
const duration = pageDurationsMs[checkpoint.revision];
489+
if (duration === undefined) throw new Error("Unexpected source page");
490+
vi.setSystemTime(Date.now() + duration);
491+
revisions.push(checkpoint.revision);
492+
if (checkpoint.revision + 1 < pageDurationsMs.length)
493+
return advanceCheckpoint(_video, checkpoint);
494+
if (terminalResult === "source-blocked")
495+
throw Object.assign(new Error("Manifest is incomplete"), {
496+
code: "source-incomplete",
497+
});
498+
return { source };
499+
});
500+
let batches = 0;
501+
let result: Awaited<ReturnType<typeof commitDesktopRecordingAttempt>>;
502+
do {
503+
if (++batches > pageDurationsMs.length)
504+
throw new Error("Source preparation did not terminate");
505+
result = await commitDesktopRecordingAttempt(fixture);
506+
} while (result === "progress");
507+
expect(result).toBe(terminalResult);
508+
expect(batches).toBe(expectedBatchedSteps);
509+
expect(revisions).toEqual(pageDurationsMs.map((_, index) => index));
510+
expect(mocks.saveCheckpoint).toHaveBeenCalledTimes(
511+
pageDurationsMs.length - 1,
512+
);
513+
expect(mocks.persist).toHaveBeenCalledTimes(
514+
terminalResult === "ready" ? 1 : 0,
515+
);
516+
},
517+
);
518+
519+
it("persists every page before advancing and commits only once after a lost step response", async () => {
520+
const revisions: number[] = [];
521+
mocks.commitSource.mockImplementation(async (_video, checkpoint) => {
522+
expect(current?.output).toEqual(
523+
checkpoint.revision === 0 ? null : checkpoint,
524+
);
525+
revisions.push(checkpoint.revision);
526+
return checkpoint.revision === 4
527+
? { source }
528+
: advanceCheckpoint(_video, checkpoint);
529+
});
530+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
531+
expect(revisions).toEqual([0, 1, 2, 3, 4]);
532+
expect(mocks.saveCheckpoint).toHaveBeenCalledTimes(4);
533+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
534+
expect(mocks.commitSource).toHaveBeenCalledTimes(5);
535+
expect(mocks.persist).toHaveBeenCalledOnce();
536+
expect(mocks.fetch).not.toHaveBeenCalled();
537+
});
538+
539+
it.each([
540+
[0, 32],
541+
[-1, 32],
542+
[1_000, 15],
543+
[14_999, 2],
544+
[15_000, 1],
545+
[180_000, 1],
546+
])("yields after bounded work with %i ms pages", async (duration, pages) => {
547+
mocks.commitSource.mockImplementation(async (_video, checkpoint) => {
548+
vi.setSystemTime(Date.now() + duration);
549+
return advanceCheckpoint(_video, checkpoint);
550+
});
551+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
552+
expect(mocks.commitSource).toHaveBeenCalledTimes(pages);
553+
expect(mocks.saveCheckpoint).toHaveBeenCalledTimes(pages);
554+
expect(current?.output).toMatchObject({ revision: pages });
555+
expect(mocks.persist).not.toHaveBeenCalled();
556+
});
557+
558+
it("includes checkpoint persistence latency in the batch budget", async () => {
559+
mocks.commitSource.mockImplementation(advanceCheckpoint);
560+
mocks.saveCheckpoint.mockImplementation(async (_attempt, checkpoint) => {
561+
withCurrent({ output: checkpoint });
562+
vi.setSystemTime(Date.now() + 15_000);
563+
return true;
564+
});
565+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
566+
expect(mocks.commitSource).toHaveBeenCalledOnce();
567+
expect(current?.output).toMatchObject({ revision: 1 });
568+
});
569+
570+
it("resumes from the last saved page when a later storage operation fails", async () => {
571+
mocks.commitSource
572+
.mockImplementationOnce(advanceCheckpoint)
573+
.mockImplementationOnce(advanceCheckpoint)
574+
.mockRejectedValueOnce(new Error("Storage unavailable"));
575+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
576+
expect(current?.output).toMatchObject({ revision: 2 });
577+
expect(mocks.blocked).not.toHaveBeenCalled();
578+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
579+
expect(
580+
mocks.commitSource.mock.calls.map((call) => call[1].revision),
581+
).toEqual([0, 1, 2, 2]);
582+
expect(mocks.persist).toHaveBeenCalledOnce();
583+
});
584+
585+
it("preserves a fresh durable retry budget after each saved page", async () => {
586+
const failedRevisions = new Set<number>();
587+
mocks.commitSource.mockImplementation(async (_video, checkpoint) => {
588+
if (
589+
checkpoint.revision > 0 &&
590+
!failedRevisions.has(checkpoint.revision)
591+
) {
592+
failedRevisions.add(checkpoint.revision);
593+
throw new Error("Intermittent storage failure");
594+
}
595+
return checkpoint.revision === 5
596+
? { source }
597+
: advanceCheckpoint(_video, checkpoint);
598+
});
599+
for (let revision = 1; revision <= 5; revision++) {
600+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
601+
expect(current?.output).toMatchObject({ revision });
602+
}
603+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
604+
expect(mocks.saveCheckpoint).toHaveBeenCalledTimes(5);
605+
expect(mocks.retry).not.toHaveBeenCalled();
606+
});
607+
608+
it("throws persistent failures when the new batch cannot make progress", async () => {
609+
mocks.commitSource.mockImplementationOnce(advanceCheckpoint);
610+
mocks.commitSource.mockRejectedValue(new Error("Storage unavailable"));
611+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
612+
await expect(commitDesktopRecordingAttempt(fixture)).rejects.toThrow(
613+
"Storage unavailable",
614+
);
615+
expect(current?.output).toMatchObject({ revision: 1 });
616+
expect(mocks.commitSource).toHaveBeenCalledTimes(3);
617+
});
618+
619+
it.each([false, true])(
620+
"reconciles a failed checkpoint response with database commit=%s",
621+
async (committed) => {
622+
mocks.commitSource.mockImplementationOnce(advanceCheckpoint);
623+
mocks.saveCheckpoint.mockImplementationOnce(
624+
async (_attempt, checkpoint) => {
625+
if (committed) withCurrent({ output: checkpoint });
626+
throw new Error("Database response lost");
627+
},
628+
);
629+
await expect(commitDesktopRecordingAttempt(fixture)).rejects.toThrow(
630+
"Database response lost",
631+
);
632+
expect(mocks.commitSource).toHaveBeenCalledOnce();
633+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
634+
expect(mocks.commitSource.mock.calls[1]?.[1]).toMatchObject({
635+
revision: committed ? 1 : 0,
636+
});
637+
},
638+
);
639+
640+
it("stops immediately when checkpoint persistence loses its ownership fence", async () => {
641+
mocks.commitSource.mockImplementation(advanceCheckpoint);
642+
mocks.saveCheckpoint.mockResolvedValueOnce(false);
643+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("superseded");
644+
expect(mocks.commitSource).toHaveBeenCalledOnce();
645+
expect(mocks.persist).not.toHaveBeenCalled();
646+
});
647+
648+
it.each(["deleted", "replaced", "lease-expired"])(
649+
"stops between pages when the job is %s",
650+
async (change) => {
651+
mocks.commitSource.mockImplementationOnce(advanceCheckpoint);
652+
mocks.saveCheckpoint.mockImplementationOnce(
653+
async (_attempt, checkpoint) => {
654+
withCurrent({ output: checkpoint });
655+
if (change === "deleted") current = null;
656+
else if (change === "replaced")
657+
withCurrent({ attemptId: "replacement" });
658+
else mocks.checkpoint.mockResolvedValueOnce(null);
659+
return true;
660+
},
661+
);
662+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("superseded");
663+
expect(mocks.commitSource).toHaveBeenCalledOnce();
664+
expect(mocks.persist).not.toHaveBeenCalled();
665+
},
666+
);
667+
668+
it("refreshes the video and late verification before the next page", async () => {
669+
const verification: RecordingVerification = {
670+
version: 1,
671+
artifact: { kind: "segments", manifestSha256: source.manifestSha256 },
672+
requiredAudio: true,
673+
};
674+
const updatedVideo = { id: "video", storageIntegrationId: "updated" };
675+
mocks.commitSource.mockImplementationOnce(async (_video, checkpoint) => {
676+
withCurrent({ verification });
677+
mocks.databaseRows.mockResolvedValueOnce([updatedVideo]);
678+
return advanceCheckpoint(_video, checkpoint);
679+
});
680+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("ready");
681+
expect(mocks.commitSource.mock.calls[0]?.[2]).toBeUndefined();
682+
expect(mocks.commitSource.mock.calls[1]?.slice(0, 3)).toEqual([
683+
updatedVideo,
684+
expect.objectContaining({ revision: 1 }),
685+
verification,
686+
]);
687+
});
688+
689+
it("keeps heartbeat ownership checks on later pages", async () => {
690+
mocks.commitSource
691+
.mockImplementationOnce(advanceCheckpoint)
692+
.mockImplementationOnce(
693+
async (_video, _checkpoint, _verification, pulse) => {
694+
await pulse();
695+
return { source };
696+
},
697+
);
698+
mocks.heartbeat.mockResolvedValueOnce(false);
699+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("progress");
700+
expect(current?.output).toMatchObject({ revision: 1 });
701+
expect(mocks.persist).not.toHaveBeenCalled();
702+
mocks.checkpoint.mockResolvedValueOnce(null);
703+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("superseded");
704+
expect(mocks.commitSource).toHaveBeenCalledTimes(2);
705+
});
706+
707+
it.each([
708+
"source-incomplete",
709+
"source-missing",
710+
"source-changed",
711+
"source-invalid",
712+
])("preserves %s rejection after earlier pages were saved", async (code) => {
713+
mocks.commitSource
714+
.mockImplementationOnce(advanceCheckpoint)
715+
.mockRejectedValueOnce(Object.assign(new Error(code), { code }));
716+
expect(await commitDesktopRecordingAttempt(fixture)).toBe("source-blocked");
717+
expect(mocks.blocked).toHaveBeenCalledWith(
718+
expect.objectContaining({ errorCode: code }),
719+
);
720+
expect(mocks.persist).not.toHaveBeenCalled();
721+
expect(mocks.fetch).not.toHaveBeenCalled();
722+
});
723+
});
724+
468725
describe("source commitment and media request compatibility", () => {
469-
it("checkpoints each source batch in a separate durable step before processing", async () => {
726+
it("checkpoints each source page before processing", async () => {
470727
withCurrent({ source: null, state: "committing", leaseExpiresAt: null });
471728
mocks.commitSource.mockImplementationOnce(async (_video, checkpoint) => ({
472729
checkpoint: { ...checkpoint, revision: 1, phase: "copy" },

0 commit comments

Comments
 (0)