import type { IDBFactory } from "fake-indexeddb"; import type { IndexedDbHeartbeatScheduler } from "./indexeddb-binary-asset-lifecycle"; import { createIndexedDbOwnerCollisionError, type IndexedDbOwnerLivenessProvider, } from "./indexeddb-binary-asset-liveness"; export function createManualHeartbeatScheduler() { let callback: (() => Promise) | undefined; return { active: () => callback !== undefined, run: async () => callback?.(), scheduler: { cancel(handle) { if (callback === handle) callback = undefined; }, schedule(nextCallback) { callback = nextCallback; return nextCallback; }, } satisfies IndexedDbHeartbeatScheduler, }; } function createPause() { let markStarted!: () => void; let resume!: () => void; const started = new Promise((resolve) => { markStarted = resolve; }); const wait = new Promise((resolve) => { resume = resolve; }); return { begin: markStarted, resume, started, wait }; } export function createManualOwnerLivenessProvider() { const coordinationTails = new Map>(); const heldByScope = new Map>(); let activeCoordinationCount = 0; let acquireAttemptCount = 0; let coordinationRequestCount = 0; let queryMode: "available" | "throw" | "unknown" = "available"; let acquirePause: ReturnType | undefined; let queryPause: ReturnType | undefined; const provider: IndexedDbOwnerLivenessProvider = { async acquire(scope, ownerId) { acquireAttemptCount += 1; const current = heldByScope.get(scope); if (current?.has(ownerId)) { throw createIndexedDbOwnerCollisionError(scope, ownerId); } const pause = acquirePause; acquirePause = undefined; pause?.begin(); await pause?.wait; const held = heldByScope.get(scope) ?? new Set(); if (held.has(ownerId)) { throw createIndexedDbOwnerCollisionError(scope, ownerId); } heldByScope.set(scope, held); held.add(ownerId); let released = false; return { release: async () => { if (released) return; released = true; held.delete(ownerId); if (held.size === 0) heldByScope.delete(scope); }, }; }, async query(scope) { const snapshot = new Set(heldByScope.get(scope) ?? []); const pause = queryPause; queryPause = undefined; pause?.begin(); await pause?.wait; if (queryMode === "throw") throw new Error("lock query failed"); if (queryMode === "unknown") return null; return snapshot; }, async runExclusive(scope, operation) { coordinationRequestCount += 1; const previous = coordinationTails.get(scope) ?? Promise.resolve(); let release!: () => void; const tail = new Promise((resolve) => { release = resolve; }); coordinationTails.set(scope, tail); await previous; activeCoordinationCount += 1; try { return await operation(); } finally { activeCoordinationCount -= 1; release(); if (coordinationTails.get(scope) === tail) { coordinationTails.delete(scope); } } }, }; return { activeCoordinationCount: () => activeCoordinationCount, abandon(ownerId: string) { for (const [scope, held] of heldByScope) { held.delete(ownerId); if (held.size === 0) heldByScope.delete(scope); } }, acquireAttemptCount: () => acquireAttemptCount, coordinationRequestCount: () => coordinationRequestCount, heldCount: () => [...heldByScope.values()].reduce((count, held) => count + held.size, 0), pauseNextAcquire() { acquirePause = createPause(); return acquirePause; }, pauseNextQuery() { queryPause = createPause(); return queryPause; }, provider, setQueryMode(mode: "available" | "throw" | "unknown") { queryMode = mode; }, }; } export function requestResult( request: IDBRequest, ): Promise { return new Promise((resolve, reject) => { request.addEventListener("success", () => resolve(request.result)); request.addEventListener("error", () => reject(request.error)); }); } export function transactionDone(transaction: IDBTransaction): Promise { return new Promise((resolve, reject) => { transaction.addEventListener("complete", () => resolve()); transaction.addEventListener("abort", () => reject(transaction.error)); transaction.addEventListener("error", () => reject(transaction.error)); }); } export async function openVersionOneDatabase( indexedDB: IDBFactory, databaseName: string, ): Promise { const request = indexedDB.open(databaseName, 1); request.addEventListener("upgradeneeded", () => { request.result.createObjectStore("entries", { keyPath: "ref" }); const leases = request.result.createObjectStore("leases", { keyPath: ["leaseId", "ref"], }); leases.createIndex("leaseId", "leaseId", { unique: false }); }); return requestResult(request); }