MCPcopy Create free account
hub / github.com/heygen-com/hyperframes / createPersistQueue

Function createPersistQueue

packages/sdk/src/persist-queue.ts:27–83  ·  view source on GitHub ↗
(
  session: Composition,
  adapter: PersistAdapter,
  opts: PersistQueueOptions = {},
)

Source from the content-addressed store, hash-verified

25}
26
27export function createPersistQueue(
28 session: Composition,
29 adapter: PersistAdapter,
30 opts: PersistQueueOptions = {},
31): PersistQueueModule {
32 const path = opts.path ?? "composition.html";
33 let pendingWrite: ReturnType<typeof setTimeout> | null = null;
34 // Promise-chain mutex: each write chains onto the prior, preventing concurrent writes.
35 let writeChain: Promise<void> = Promise.resolve();
36 let disposed = false;
37
38 function scheduleWrite(): void {
39 if (pendingWrite !== null) clearTimeout(pendingWrite);
40 pendingWrite = setTimeout(() => {
41 pendingWrite = null;
42 void doWrite();
43 }, 0);
44 }
45
46 function doWrite(): Promise<void> {
47 if (disposed) return Promise.resolve();
48 const content = session.serialize();
49 writeChain = writeChain.then(async () => {
50 if (disposed) return;
51 try {
52 await adapter.write(path, content);
53 } catch (err) {
54 const message = err instanceof Error ? err.message : String(err);
55 opts.onError?.({ error: { message, cause: err } });
56 }
57 });
58 return writeChain;
59 }
60
61 const unsubscribe = session.on("change", () => {
62 scheduleWrite();
63 });
64
65 return {
66 async flush(): Promise<void> {
67 if (pendingWrite !== null) {
68 clearTimeout(pendingWrite);
69 pendingWrite = null;
70 }
71 await doWrite();
72 },
73
74 dispose(): void {
75 disposed = true;
76 if (pendingWrite !== null) {
77 clearTimeout(pendingWrite);
78 pendingWrite = null;
79 }
80 unsubscribe();
81 },
82 };
83}

Callers 1

openCompositionFunction · 0.85

Calls 2

scheduleWriteFunction · 0.85
onMethod · 0.65

Tested by

no test coverage detected