diff --git a/graphile/graphile-cache/src/__tests__/build-readiness.test.ts b/graphile/graphile-cache/src/__tests__/build-readiness.test.ts new file mode 100644 index 0000000000..b876be741c --- /dev/null +++ b/graphile/graphile-cache/src/__tests__/build-readiness.test.ts @@ -0,0 +1,126 @@ +import { awaitGraphileBuildReadiness } from '../build-readiness'; + +interface Deferred { + promise: Promise; + resolve(value: T): void; + reject(error: Error): void; +} + +const deferred = (): Deferred => { + let resolve!: (value: T) => void; + let reject!: (error: Error) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, resolve, reject }; +}; + +const flushPromises = (): Promise => + new Promise((resolve) => setImmediate(resolve)); + +describe('awaitGraphileBuildReadiness', () => { + it('does not resolve before schema gathering and Grafserv are ready', async () => { + const schemaResult = deferred(); + const ready = deferred(); + const release = jest.fn().mockResolvedValue(undefined); + let resolved = false; + const buildPromise = awaitGraphileBuildReadiness({ + schemaResult: schemaResult.promise, + addTo: jest.fn().mockResolvedValue(undefined), + ready: () => ready.promise, + release, + }).then(() => { + resolved = true; + }); + + schemaResult.resolve({}); + await flushPromises(); + expect(resolved).toBe(false); + + ready.resolve(undefined); + await buildPromise; + expect(release).not.toHaveBeenCalled(); + }); + + it('does not start readiness checks before the adapter is attached', async () => { + const addTo = deferred(); + const ready = jest.fn().mockResolvedValue(undefined); + const buildPromise = awaitGraphileBuildReadiness({ + schemaResult: Promise.resolve({}), + addTo: () => addTo.promise, + ready, + release: jest.fn().mockResolvedValue(undefined), + }); + + await flushPromises(); + expect(ready).not.toHaveBeenCalled(); + + addTo.resolve(undefined); + await buildPromise; + expect(ready).toHaveBeenCalledTimes(1); + }); + + it('observes schema failure while adapter attachment is pending', async () => { + const schemaResult = deferred(); + const addTo = deferred(); + const release = jest.fn().mockResolvedValue(undefined); + const failure = new Error('schema build failed early'); + const buildPromise = awaitGraphileBuildReadiness({ + schemaResult: schemaResult.promise, + addTo: () => addTo.promise, + ready: jest.fn().mockResolvedValue(undefined), + release, + }); + + schemaResult.reject(failure); + await flushPromises(); + expect(release).not.toHaveBeenCalled(); + + addTo.resolve(undefined); + await expect(buildPromise).rejects.toBe(failure); + expect(release).toHaveBeenCalledTimes(1); + }); + + it('awaits failed-generation release before rejecting', async () => { + const schemaResult = deferred(); + const release = deferred(); + const releaseFn = jest.fn(() => release.promise); + const failure = new Error('schema build failed'); + let rejected = false; + const buildPromise = awaitGraphileBuildReadiness({ + schemaResult: schemaResult.promise, + addTo: jest.fn().mockResolvedValue(undefined), + ready: jest.fn().mockResolvedValue(undefined), + release: releaseFn, + }).catch((error) => { + rejected = true; + throw error; + }); + + schemaResult.reject(failure); + await flushPromises(); + expect(releaseFn).toHaveBeenCalledTimes(1); + expect(rejected).toBe(false); + + release.resolve(undefined); + await expect(buildPromise).rejects.toBe(failure); + }); + + it('preserves the build failure when cleanup also fails', async () => { + const failure = new Error('schema build failed'); + const cleanupFailure = new Error('release failed'); + const onReleaseError = jest.fn(); + + await expect( + awaitGraphileBuildReadiness({ + schemaResult: Promise.reject(failure), + addTo: jest.fn().mockResolvedValue(undefined), + ready: jest.fn().mockResolvedValue(undefined), + release: jest.fn().mockRejectedValue(cleanupFailure), + onReleaseError, + }) + ).rejects.toBe(failure); + expect(onReleaseError).toHaveBeenCalledWith(cleanupFailure); + }); +}); diff --git a/graphile/graphile-cache/src/__tests__/disposal-lifecycle.test.ts b/graphile/graphile-cache/src/__tests__/disposal-lifecycle.test.ts new file mode 100644 index 0000000000..679488ace9 --- /dev/null +++ b/graphile/graphile-cache/src/__tests__/disposal-lifecycle.test.ts @@ -0,0 +1,137 @@ +jest.mock('@pgpmjs/logger', () => ({ + Logger: jest.fn(() => ({ + debug: jest.fn(), + error: jest.fn(), + })), +})); + +import type { GraphileCacheEntry } from '../graphile-cache'; +import { + clearGraphileCache, + disposeUncachedEntry, + graphileCache, + waitForEntryDisposal, +} from '../graphile-cache'; + +interface Deferred { + promise: Promise; + resolve(value: T): void; +} + +const deferred = (): Deferred => { + let resolve!: (value: T) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +}; + +const flushPromises = (): Promise => + new Promise((resolve) => setImmediate(resolve)); + +const makeEntry = ( + cacheKey: string, + release = jest.fn().mockResolvedValue(undefined), + releasePresetServices = jest.fn().mockResolvedValue(undefined) +): GraphileCacheEntry => + ({ + pgl: { release }, + serv: {}, + handler: {}, + httpServer: { listening: false }, + cacheKey, + createdAt: Date.now(), + releasePresetServices, + }) as unknown as GraphileCacheEntry; + +describe('Graphile cache disposal lifecycle', () => { + afterEach(async () => { + await clearGraphileCache(); + }); + + it('coalesces concurrent disposal of one exact entry', async () => { + const release = jest.fn().mockResolvedValue(undefined); + const releasePresetServices = jest.fn().mockResolvedValue(undefined); + const entry = makeEntry('same-entry', release, releasePresetServices); + + const first = disposeUncachedEntry(entry); + const second = disposeUncachedEntry(entry); + + expect(second).toBe(first); + await Promise.all([first, second]); + expect(release).toHaveBeenCalledTimes(1); + expect(releasePresetServices).toHaveBeenCalledTimes(1); + }); + + it('disposes distinct generations that reuse the same cache key', async () => { + const firstRelease = jest.fn().mockResolvedValue(undefined); + const secondRelease = jest.fn().mockResolvedValue(undefined); + const first = makeEntry('shared-key', firstRelease); + const second = makeEntry('shared-key', secondRelease); + + await Promise.all([ + disposeUncachedEntry(first), + disposeUncachedEntry(second), + ]); + + expect(firstRelease).toHaveBeenCalledTimes(1); + expect(secondRelease).toHaveBeenCalledTimes(1); + }); + + it('continues cleanup and exposes the first disposal failure', async () => { + const failure = new Error('realtime stop failed'); + const release = jest.fn().mockResolvedValue(undefined); + const releasePresetServices = jest.fn().mockResolvedValue(undefined); + const entry = makeEntry('failed-cleanup', release, releasePresetServices); + entry.realtimeManager = { stop: jest.fn().mockRejectedValue(failure) }; + + await expect(disposeUncachedEntry(entry)).rejects.toBe(failure); + expect(release).toHaveBeenCalledTimes(1); + expect(releasePresetServices).toHaveBeenCalledTimes(1); + }); + + it('lets callers await an eviction through the exact entry', async () => { + const release = deferred(); + const releasePresetServices = jest.fn().mockResolvedValue(undefined); + const entry = makeEntry( + 'evicted-entry', + jest.fn(() => release.promise), + releasePresetServices + ); + graphileCache.set(entry.cacheKey, entry); + graphileCache.delete(entry.cacheKey); + + let disposed = false; + const waiting = waitForEntryDisposal(entry).then(() => { + disposed = true; + }); + await flushPromises(); + expect(disposed).toBe(false); + + release.resolve(undefined); + await waiting; + expect(disposed).toBe(true); + expect(releasePresetServices).toHaveBeenCalledTimes(1); + }); + + it('does not resolve a cache clear before resident disposal completes', async () => { + const release = deferred(); + const entry = makeEntry( + 'clear-entry', + jest.fn(() => release.promise) + ); + graphileCache.set(entry.cacheKey, entry); + + let cleared = false; + const clearing = clearGraphileCache().then(() => { + cleared = true; + }); + await flushPromises(); + expect(graphileCache.size).toBe(0); + expect(cleared).toBe(false); + + release.resolve(undefined); + await clearing; + expect(cleared).toBe(true); + }); +}); diff --git a/graphile/graphile-cache/src/__tests__/preset-services.test.ts b/graphile/graphile-cache/src/__tests__/preset-services.test.ts new file mode 100644 index 0000000000..f4b38a855b --- /dev/null +++ b/graphile/graphile-cache/src/__tests__/preset-services.test.ts @@ -0,0 +1,50 @@ +import { createPresetServicesReleaser } from '../preset-services'; + +describe('preset service ownership', () => { + it('releases unique services in reverse order exactly once', async () => { + const events: string[] = []; + const first = { + release: jest.fn(async () => { + events.push('first'); + }), + }; + const second = { + release: jest.fn(async () => { + events.push('second'); + }), + }; + const release = createPresetServicesReleaser({ + pgServices: [first, second, first], + }); + + const releases = [release(), release(), release()]; + expect(releases[0]).toBe(releases[1]); + expect(releases[1]).toBe(releases[2]); + await Promise.all(releases); + + expect(events).toEqual(['second', 'first']); + expect(first.release).toHaveBeenCalledTimes(1); + expect(second.release).toHaveBeenCalledTimes(1); + }); + + it('continues releasing services and preserves the first error', async () => { + const firstFailure = new Error('second failed'); + const first = { release: jest.fn().mockResolvedValue(undefined) }; + const second = { release: jest.fn().mockRejectedValue(firstFailure) }; + const release = createPresetServicesReleaser({ + pgServices: [first, second], + }); + + await expect(release()).rejects.toBe(firstFailure); + expect(first.release).toHaveBeenCalledTimes(1); + expect(second.release).toHaveBeenCalledTimes(1); + }); + + it('is a safe idempotent no-op when the preset has no services', async () => { + const release = createPresetServicesReleaser({}); + + const first = release(); + expect(release()).toBe(first); + await first; + }); +}); diff --git a/graphile/graphile-cache/src/build-readiness.ts b/graphile/graphile-cache/src/build-readiness.ts new file mode 100644 index 0000000000..2ab5e921b2 --- /dev/null +++ b/graphile/graphile-cache/src/build-readiness.ts @@ -0,0 +1,32 @@ +export interface GraphileBuildReadiness { + schemaResult: PromiseLike | unknown; + addTo(): PromiseLike | unknown; + ready(): PromiseLike | unknown; + release(): PromiseLike | unknown; + onReleaseError?(error: unknown): void; +} + +/** + * Resolve only after schema gathering and the HTTP adapter are ready. A failed + * generation reaches its release terminal state before the failure escapes. + */ +export const awaitGraphileBuildReadiness = async ( + build: GraphileBuildReadiness +): Promise => { + const schemaOutcome = Promise.resolve(build.schemaResult).then( + () => ({ ready: true as const }), + (error: unknown) => ({ ready: false as const, error }) + ); + try { + await build.addTo(); + const [schema] = await Promise.all([schemaOutcome, build.ready()]); + if ('error' in schema) throw schema.error; + } catch (error) { + try { + await build.release(); + } catch (releaseError) { + build.onReleaseError?.(releaseError); + } + throw error; + } +}; diff --git a/graphile/graphile-cache/src/create-instance.ts b/graphile/graphile-cache/src/create-instance.ts index 575b767589..1fe84ce156 100644 --- a/graphile/graphile-cache/src/create-instance.ts +++ b/graphile/graphile-cache/src/create-instance.ts @@ -5,7 +5,9 @@ import express from 'express'; import { grafserv } from 'grafserv/express/v4'; import { postgraphile } from 'postgraphile'; +import { awaitGraphileBuildReadiness } from './build-readiness'; import type { GraphileCacheEntry } from './graphile-cache'; +import { createPresetServicesReleaser } from './preset-services'; const log = new Logger('graphile-cache:create'); @@ -42,12 +44,47 @@ export const createGraphileInstance = async ( const { preset, cacheKey, enableRealtime = false } = opts; const pgl = postgraphile(preset); + const resolvedPreset = pgl.getResolvedPreset(); + const releasePresetServices = createPresetServicesReleaser(resolvedPreset); const serv = pgl.createServ(grafserv); const handler = express(); const httpServer = createServer(handler); - await serv.addTo(handler, httpServer); - await serv.ready(); + let failedBuildReleasePromise: Promise | null = null; + const releaseFailedBuild = (): Promise => { + if (failedBuildReleasePromise) return failedBuildReleasePromise; + failedBuildReleasePromise = (async () => { + let firstError: unknown; + let failed = false; + try { + await pgl.release(); + } catch (error) { + firstError = error; + failed = true; + } + try { + await releasePresetServices(); + } catch (error) { + if (!failed) firstError = error; + failed = true; + } + if (failed) throw firstError; + })(); + return failedBuildReleasePromise; + }; + + await awaitGraphileBuildReadiness({ + schemaResult: pgl.getSchemaResult(), + addTo: () => serv.addTo(handler, httpServer), + ready: () => serv.ready(), + release: releaseFailedBuild, + onReleaseError: (releaseError) => { + log.error( + `Failed to release PostGraphile[${cacheKey}] after build failure:`, + releaseError + ); + } + }); const entry: GraphileCacheEntry = { pgl, @@ -56,6 +93,7 @@ export const createGraphileInstance = async ( httpServer, cacheKey, createdAt: Date.now(), + releasePresetServices }; if (enableRealtime) { @@ -65,7 +103,6 @@ export const createGraphileInstance = async ( // Extract PgSubscriber and pool from the resolved preset's pgServices. // The pool is the same instance managed by pg-cache (via getPgPool) // and threaded into the preset by makePgService({ pool, schemas }). - const resolvedPreset = pgl.getResolvedPreset(); const pgService = (resolvedPreset as any).pgServices?.[0]; const pgSubscriber = pgService?.pgSubscriber ?? null; const pool = pgService?.adaptorSettings?.pool ?? null; diff --git a/graphile/graphile-cache/src/graphile-cache.ts b/graphile/graphile-cache/src/graphile-cache.ts index 83782c6a21..67dd35a59c 100644 --- a/graphile/graphile-cache/src/graphile-cache.ts +++ b/graphile/graphile-cache/src/graphile-cache.ts @@ -86,12 +86,14 @@ export interface GraphileCacheEntry { httpServer: HttpServer; cacheKey: string; createdAt: number; + /** Idempotent release for pgServices owned by this exact preset generation. */ + releasePresetServices?: () => Promise; /** Optional RealtimeManager for cursor-tracked subscription delivery */ realtimeManager?: { stop(): Promise } | null; } -// Track disposed entries to prevent double-disposal -const disposedKeys = new Set(); +const disposalPromises = new WeakMap>(); +const activeDisposals = new Set>(); // Track keys that are being manually evicted for accurate eviction reason const manualEvictionKeys = new Set(); @@ -101,42 +103,86 @@ const manualEvictionKeys = new Set(); * * Properly releases resources by: * 1. Closing the HTTP server if listening - * 2. Releasing the PostGraphile instance (which internally releases grafserv) - * - * Uses disposedKeys set to prevent double-disposal when closeAllCaches() - * explicitly disposes entries and then clear() triggers the dispose callback. + * 2. Stopping the realtime manager + * 3. Releasing PostGraphile/Grafserv and preset services */ -const disposeEntry = async (entry: GraphileCacheEntry, key: string): Promise => { - // Prevent double-disposal - if (disposedKeys.has(key)) { - return; - } - disposedKeys.add(key); - +const releaseEntry = async ( + entry: GraphileCacheEntry, + key: string +): Promise => { log.debug(`Disposing PostGraphile[${key}]`); + let firstError: unknown; + let failed = false; try { - // Close HTTP server if it's listening if (entry.httpServer?.listening) { await new Promise((resolve) => { entry.httpServer.close(() => resolve()); }); } - // Stop RealtimeManager if present (before releasing PostGraphile) + } catch (error) { + firstError = error; + failed = true; + } + try { if (entry.realtimeManager) { - try { - await entry.realtimeManager.stop(); - } catch (err) { - log.error(`Error stopping RealtimeManager for PostGraphile[${key}]:`, err); - } - } - // Release PostGraphile instance (this also releases grafserv internally) - if (entry.pgl) { - await entry.pgl.release(); + await entry.realtimeManager.stop(); } - } catch (err) { - log.error(`Error disposing PostGraphile[${key}]:`, err); - } finally { - disposedKeys.delete(key); + } catch (error) { + if (!failed) firstError = error; + failed = true; + } + try { + await entry.pgl.release(); + } catch (error) { + if (!failed) firstError = error; + failed = true; + } + try { + await entry.releasePresetServices?.(); + } catch (error) { + if (!failed) firstError = error; + failed = true; + } + if (failed) throw firstError; +}; + +/** + * Coalesce teardown by exact entry identity, not by its reusable cache key. + */ +const scheduleDisposal = ( + entry: GraphileCacheEntry, + key: string +): Promise => { + const existing = disposalPromises.get(entry); + if (existing) return existing; + + const pending = releaseEntry(entry, key); + + disposalPromises.set(entry, pending); + activeDisposals.add(pending); + void pending + .catch((error) => { + log.error(`Failed to dispose PostGraphile[${key}]:`, error); + }) + .finally(() => activeDisposals.delete(pending)); + return pending; +}; + +/** Dispose a generation that was built but never published in the cache. */ +export const disposeUncachedEntry = ( + entry: GraphileCacheEntry, + key = entry.cacheKey +): Promise => scheduleDisposal(entry, key); + +/** Await the terminal result for an entry whose disposal has been scheduled. */ +export const waitForEntryDisposal = ( + entry: GraphileCacheEntry +): Promise => disposalPromises.get(entry) ?? Promise.resolve(); + +/** Await every disposal that is active at or begins during this drain. */ +export const waitForActiveDisposals = async (): Promise => { + while (activeDisposals.size > 0) { + await Promise.allSettled([...activeDisposals]); } }; @@ -176,11 +222,7 @@ export const graphileCache = new LRUCache({ log.debug(`Evicting PostGraphile[${key}] (reason: ${reason})`); - // LRU dispose is synchronous, but v5 disposal is async - // Fire and forget the async cleanup - disposeEntry(entry, key).catch((err) => { - log.error(`Failed to dispose PostGraphile[${key}]:`, err); - }); + scheduleDisposal(entry, key); } }); @@ -245,6 +287,16 @@ const unregister = pgCache.registerCleanupCallback((pgPoolKey: string) => { // Enhanced close function that handles all caches const closePromise: { promise: Promise | null } = { promise: null }; +/** Clear all resident entries and await every exact-generation disposal. */ +export const clearGraphileCache = async (): Promise => { + for (const key of graphileCache.keys()) { + manualEvictionKeys.add(key); + } + graphileCache.clear(); + await waitForActiveDisposals(); + manualEvictionKeys.clear(); +}; + /** * Close all caches and release resources * @@ -263,27 +315,7 @@ export const closeAllCaches = async (verbose = false): Promise => { try { if (verbose) log.info('Closing all server caches...'); - // Collect all entries and dispose them properly - const entries = [...graphileCache.entries()]; - - // Mark all as manual evictions - for (const [key] of entries) { - manualEvictionKeys.add(key); - } - - const disposePromises = entries.map(([key, entry]) => - disposeEntry(entry, key) - ); - - // Wait for all disposals to complete - await Promise.allSettled(disposePromises); - - // Clear the cache after disposal (dispose callback will no-op due to disposedKeys) - graphileCache.clear(); - - // Clear disposed keys tracking after full cleanup - disposedKeys.clear(); - manualEvictionKeys.clear(); + await clearGraphileCache(); // Close pg pools await pgCache.close(); diff --git a/graphile/graphile-cache/src/index.ts b/graphile/graphile-cache/src/index.ts index 9a845fafe3..7cbed436e3 100644 --- a/graphile/graphile-cache/src/index.ts +++ b/graphile/graphile-cache/src/index.ts @@ -8,9 +8,11 @@ export { CacheEvictionEvent, // Cache stats CacheStats, + clearGraphileCache, // Clear matching entries clearMatchingEntries, closeAllCaches, + disposeUncachedEntry, // Eviction tracking EvictionReason, FIVE_MINUTES_MS, @@ -20,7 +22,9 @@ export { graphileCache, GraphileCacheEntry, // Time constants - ONE_HOUR_MS} from './graphile-cache'; + ONE_HOUR_MS, + waitForActiveDisposals, + waitForEntryDisposal} from './graphile-cache'; // Factory for creating PostGraphile v5 instances export { createGraphileInstance } from './create-instance'; diff --git a/graphile/graphile-cache/src/preset-services.ts b/graphile/graphile-cache/src/preset-services.ts new file mode 100644 index 0000000000..eaf3111c84 --- /dev/null +++ b/graphile/graphile-cache/src/preset-services.ts @@ -0,0 +1,29 @@ +interface ReleasablePresetService { + release?: () => void | Promise; +} + +/** Own and release the unique pgServices for one resolved preset generation. */ +export const createPresetServicesReleaser = (resolvedPreset: { + pgServices?: readonly ReleasablePresetService[]; +}): (() => Promise) => { + const services = [...new Set(resolvedPreset.pgServices ?? [])]; + let releasePromise: Promise | null = null; + + return (): Promise => { + if (releasePromise) return releasePromise; + releasePromise = (async () => { + let firstError: unknown; + let failed = false; + for (const service of [...services].reverse()) { + try { + await service.release?.(); + } catch (error) { + if (!failed) firstError = error; + failed = true; + } + } + if (failed) throw firstError; + })(); + return releasePromise; + }; +};