| 13 | let queueing: Promise<Map<object, unknown>> | undefined; |
| 14 | |
| 15 | async function dequeue() { |
| 16 | if (queueing) return queueing; |
| 17 | return (queueing = new Promise((resolve) => |
| 18 | raf(async () => { |
| 19 | const db = await connection; |
| 20 | const result = new Map<object, unknown>(); |
| 21 | await db |
| 22 | .transaction() |
| 23 | .execute(async (trx: any) => { |
| 24 | const current: any = |
| 25 | trigger && (await selectVersion.bind(trx)().execute()).current; |
| 26 | for (const [id, query] of queue.entries()) { |
| 27 | const rows = await query(trx).catch((x) => ({ [error]: x })); |
| 28 | result.set(id, rows); |
| 29 | } |
| 30 | trigger?.( |
| 31 | (await changesSince.bind(trx)(current).execute()) as string, |
| 32 | ); |
| 33 | }) |
| 34 | .catch((reason) => { |
| 35 | if (String(reason).includes("driver has already been destroyed")) { |
| 36 | return; |
| 37 | } |
| 38 | throw reason; |
| 39 | }); |
| 40 | queue.clear(); |
| 41 | queueing = undefined; |
| 42 | resolve(result); |
| 43 | }), |
| 44 | )); |
| 45 | } |
| 46 | |
| 47 | return { |
| 48 | enqueue<T extends any[], R>( |