(
session: Composition,
adapter: PersistAdapter,
opts: PersistQueueOptions = {},
)
| 25 | } |
| 26 | |
| 27 | export 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 | } |
no test coverage detected