diff --git a/src/__tests__/collections.test.ts b/src/__tests__/collections.test.ts index 90d3c978..6639b9cb 100644 --- a/src/__tests__/collections.test.ts +++ b/src/__tests__/collections.test.ts @@ -1,4 +1,9 @@ -import { registerCollections } from '../collections'; +import { + registerCollections, + getServerRecordsRegistry, + onUpdateServerRecordsRegistry, + setServerRecordsRegistry +} from '../collections'; jest.mock('../logger', () => ({ log: { @@ -10,12 +15,171 @@ jest.mock('../logger', () => ({ formatUnknownError: jest.fn(err => String(err)) })); +describe('getServerRecordsRegistry', () => { + beforeEach(() => { + const registry = getServerRecordsRegistry() as Map; + + registry.clear(); + + jest.clearAllMocks(); + }); + + it('returns the full registry Map when called without params', () => { + const registry = getServerRecordsRegistry(); + + expect(registry).toBeInstanceOf(Map); + expect((registry as Map).size).toBe(0); + }); +}); + +describe('onUpdateServerRecordsRegistry', () => { + beforeEach(() => { + jest.useFakeTimers(); + const registry = getServerRecordsRegistry() as Map; + + registry.clear(); + + jest.clearAllMocks(); + }); + + afterEach(() => jest.useRealTimers()); + + it('should return a no-op unsubscribe when callback is not a function', () => { + const unsubscribe = onUpdateServerRecordsRegistry(null as any); + + expect(unsubscribe()).toBe(false); + }); + + it('should not replay existing registry entries by default', async () => { + const response = { records: [{ id: '1' }] } as any; + + await setServerRecordsRegistry({ name: 'cached', response }); + + const handler = jest.fn(); + + onUpdateServerRecordsRegistry(handler); + + await jest.runAllTimersAsync(); + + expect(handler).not.toHaveBeenCalled(); + }); + + it('replays existing registry entries when replay is enabled', async () => { + const docs = { records: [{ id: 'docs' }] } as any; + const schemas = { records: [{ id: 'schemas' }] } as any; + + await setServerRecordsRegistry({ name: 'patternfly-docs', response: docs }); + await setServerRecordsRegistry({ name: 'patternfly-component-schemas', response: schemas }); + + const handler = jest.fn(); + + onUpdateServerRecordsRegistry(handler, { replay: true }); + + await jest.runAllTimersAsync(); + + expect(handler).toHaveBeenCalledTimes(2); + expect(handler).toHaveBeenCalledWith({ + name: 'patternfly-docs', + response: docs, + error: undefined + }); + expect(handler).toHaveBeenCalledWith({ + name: 'patternfly-component-schemas', + response: schemas, + error: undefined + }); + }); + + it('should attempt to fire the callback again after replay on a subsequent update', async () => { + const response = { records: [{ id: '1' }] } as any; + + await setServerRecordsRegistry({ name: 'repeatable', response }); + + const handler = jest.fn(); + + onUpdateServerRecordsRegistry(handler, { replay: true }); + + await jest.runAllTimersAsync(); + + expect(handler).toHaveBeenCalledTimes(1); + + await setServerRecordsRegistry({ name: 'repeatable', response }); + + expect(handler).toHaveBeenCalledTimes(2); + }); +}); + +describe('get, set, update the server records registry', () => { + beforeEach(() => { + const registry = getServerRecordsRegistry() as Map; + + registry.clear(); + + jest.clearAllMocks(); + }); + + it('should return a specific collection by name when available', async () => { + const response = { records: [{ id: '1', sourceId: 's', sourceType: 'local' }] } as any; + + await setServerRecordsRegistry({ name: 'hello', response }); + + expect(getServerRecordsRegistry({ collectionName: 'hello' })).toEqual(response); + expect(getServerRecordsRegistry({ collectionName: 'world' })).toBeUndefined(); + }); + + it('should register and unregister listeners correctly', async () => { + const handler = jest.fn(); + const unsubscribe = onUpdateServerRecordsRegistry(handler); + + await setServerRecordsRegistry({ name: 'ipsum', response: { records: [] } as any }); + + expect(handler).toHaveBeenCalledWith({ name: 'ipsum', response: { records: [] }, error: undefined }); + + expect(unsubscribe()).toBe(true); + expect(unsubscribe()).toBe(false); + + await setServerRecordsRegistry({ name: 'ipsum', response: { records: [] } as any }); + expect(handler).toHaveBeenCalledTimes(1); + }); + + it('should continue processing when a listener throws', async () => { + const faulty = jest.fn().mockRejectedValue(new Error('lorem ipsum')); + const good = jest.fn(); + + onUpdateServerRecordsRegistry(faulty); + onUpdateServerRecordsRegistry(good); + + await setServerRecordsRegistry({ name: 'sit', response: { records: [] } as any }); + expect(good).toHaveBeenCalled(); + }); + + it('should store records when name and response are provided', async () => { + await setServerRecordsRegistry({ name: 'lorem-ipsum', response: { records: [{ id: 'x' }] } as any }); + + const stored = getServerRecordsRegistry({ collectionName: 'lorem-ipsum' }); + + expect(stored).toEqual({ records: [{ id: 'x' }] }); + }); + + it('should not store or notify when response is missing', async () => { + const listener = jest.fn(); + + onUpdateServerRecordsRegistry(listener); + + await setServerRecordsRegistry({ name: 'dolor' }); + + expect(getServerRecordsRegistry({ collectionName: 'dolor' })).toBeUndefined(); + expect(listener).not.toHaveBeenCalled(); + }); +}); + describe('registerCollections', () => { beforeEach(() => { jest.clearAllMocks(); }); it('should register valid collections and call onUpdate', async () => { + jest.useFakeTimers(); const onUpdate = jest.fn(); const handler = jest.fn().mockResolvedValue({ records: [] }); const collections: any[] = [ @@ -23,12 +187,15 @@ describe('registerCollections', () => { ]; await registerCollections(collections, { onUpdate }); + await jest.runAllTimersAsync(); expect(handler).toHaveBeenCalled(); expect(onUpdate).toHaveBeenCalledWith(expect.objectContaining({ name: 'test-collection', response: { records: [] } })); + + jest.useRealTimers(); }); it('should handle isRequired and throw if it fails', async () => { diff --git a/src/collections.ts b/src/collections.ts index ce2f1636..78d3132e 100644 --- a/src/collections.ts +++ b/src/collections.ts @@ -99,7 +99,18 @@ type RegisterCollectionItem = { * @param {McpCollectionResult|undefined} [item.response] - Optional response associated with the item. * @param [item.error] - Optional error object if an error occurred during the collection process. */ -type RegisterOnUpdate = ({ name, response, error }: RegisterCollectionItem) => void; +type RegisterOnUpdate = ({ name, response, error }: RegisterCollectionItem) => void | Promise; + +/** + * Options for {@link onUpdateServerRecordsRegistry}. + * + * @property replay - When `true`, invokes the callback once for each collection already in the registry. + * Live updates after subscribe **MAY INVOKE THE CALLBACK AGAIN** for the same collection. + * Deduplication is the consumer's responsibility. + */ +type OnUpdateServerRecordsRegistryOptions = { + replay?: boolean; +}; /** * Callback invoked when required collections are loaded/updated. @@ -149,6 +160,122 @@ type RegisterCollectionsResult = { rejected: { name: string | null, reason: unknown }[]; }; +/** + * Central in-memory registry for all PatternFly collection records + */ +const serverRecordsRegistry = new Map(); + +/** + * Listeners for server records registry updates + */ +const serverRecordsRegistryListeners = new Set(); + +/** + * Invokes a server records registry listener and logs errors without rethrowing. + * + * @param callback - Listener to invoke/fire. + * @param item - Collection item passed to the listener. + */ +const invokeServerRecordsRegistryListener = async ( + callback: RegisterOnUpdate, + item: RegisterCollectionItem +) => { + try { + await callback(item); + } catch (error) { + log.error(`Error in server records registry listener:`, error); + } +}; + +/** + * Retrieves the server collections/records registry, all or for a given collection name. + * + * @param params - Optional parameters. + * @param params.collectionName - Name of the collection to retrieve. + * @returns The entire server collections/records registry, or the registry for the specified collection name + * if provided and available, otherwise returns `undefined`. + */ +const getServerRecordsRegistry = ({ collectionName }: { collectionName?: string } = {}) => { + if (collectionName) { + return serverRecordsRegistry.get(collectionName); + } + + return serverRecordsRegistry; +}; + +/** + * Executes a collection callback, invalidates any cache, and then any next-call to the functions + * blends the returned records and "re-memos" the results. + * + * @param {McpCollectionResult} collection - Collection. + */ +const setServerRecordsRegistry = async (collection: RegisterCollectionItem) => { + const { name, response } = collection || {}; + + try { + if (name && response) { + serverRecordsRegistry.set(name, response); + + for (const listener of serverRecordsRegistryListeners) { + await invokeServerRecordsRegistryListener(listener, collection); + } + + log.debug(`Storing server collection ${name} records. (${response?.records?.length})`); + } + } catch (error) { + log.error(`Failed to store server collection ${name}:`, error); + } +}; + +/** + * Register a listener callback to be fired whenever a server record in the registry is updated. + * + * @note Using the `replay` {@link OnUpdateServerRecordsRegistryOptions.replay} option means the + * callback can be fired multiple times for the same collection. Deduplication is the consumer's + * responsibility. This isn't needed if your collections are `required`. + * + * @param callback - The callback to execute on update. + * @param [options] - Subscribe options. + * @param [options.replay] - When `true`, fire the registry-level callback for each collection + * already stored in the registry. Useful for callbacks registered after the registry-level callback + * has already fired. Defaults to `false`. See {@link OnUpdateServerRecordsRegistryOptions.replay} + * @returns A function to unregister/unsubscribe the listener. + */ +const onUpdateServerRecordsRegistry = ( + callback: RegisterOnUpdate, + { replay = false }: OnUpdateServerRecordsRegistryOptions = {} +) => { + if (typeof callback !== 'function') { + log.warn('onUpdateServerRecordsRegistry: callback must be a function'); + + return () => false; + } + + serverRecordsRegistryListeners.add(callback); + + if (replay) { + void (async () => { + for (const [name, response] of serverRecordsRegistry) { + if (!serverRecordsRegistryListeners.has(callback)) { + break; + } + + await invokeServerRecordsRegistryListener(callback, { name, response, error: undefined }); + } + })(); + } + + return () => { + if (serverRecordsRegistryListeners.has(callback)) { + serverRecordsRegistryListeners.delete(callback); + + return true; + } + + return false; + }; +}; + /** * Registers a set of collections asynchronously. * @@ -160,12 +287,12 @@ type RegisterCollectionsResult = { * * @param {McpCollection[]} collections - An array of collection sources to be registered. Each source is represented as a tuple. * @param [options] - Options callback functions to handle registration events. - * @param [options.onSettle] - Callback function executed after all collection - * registrations are settled. Receives the results as an object containing settled, fulfilled, and rejected collections. - * @param [options.onUpdate] - Callback function executed for each collection - * registration update. Receives details about the collection being processed including name, response, and any error encountered. - * @param [options.onRequired] - Callback function executed when required - * collections are processed. Receives an array of results containing collection name, response, and error details. + * @param [options.onSettle] - A non-blocking consumer-facing callback executed after all collection registrations are + * settled. Receives the results as an object containing settled, fulfilled, and rejected collections. + * @param [options.onUpdate] - A non-blocking consumer-facing callback executed for each collection registration update. + * Receives details about the collection being processed, including name, response, and any error encountered. + * @param [options.onRequired] - A non-blocking consumer-facing callback executed when required collections are processed. + * Receives an array of results containing collection name, response, and error details. * @returns Resolves when all "isRequired" collections are registered and settled. * @throws {Error} If any required collection fails to register successfully. */ @@ -192,17 +319,23 @@ const registerCollections = async ( } try { - onUpdate?.({ name, response, error }); + if (response) { + await setServerRecordsRegistry({ name, response, error }); + } } catch (err) { - log.error(`Error "onUpdate" for collection ${name}: ${formatUnknownError(err)}`); + log.error(`Error "setServerRecordsRegistry" for collection ${name}: ${formatUnknownError(err)}`); } + // Fire-and-forget if it exists. Review using `Promise.try` in the future. + Promise.resolve() + .then(() => onUpdate?.({ name, response, error })) + .catch(err => log.debug(`Error calling "onUpdate": ${formatUnknownError(err)}`)); + return { name, response, isSuccess, error }; }); // Determine which collections are required and optional const required = registrationPromises.filter((_, index) => collections[index]?.[2]?.isRequired); - // const optional = registrationPromises.filter((_, index) => !collections[index]?.[2]?.isRequired); // Gatekeep on any required collections const results = await Promise.all(required); @@ -216,11 +349,10 @@ const registerCollections = async ( } } - try { - onRequired?.(results.map(({ name, response, error }) => ({ name, response, error }))); - } catch (err) { - log.error(`Error calling "onRequired": ${formatUnknownError(err)}`); - } + // Fire-and-forget if it exists. Review using `Promise.try` in the future. + Promise.resolve() + .then(() => onRequired?.(results.map(({ name, response, error }) => ({ name, response, error })))) + .catch(err => log.debug(`Error calling "onRequired": ${formatUnknownError(err)}`)); // Wait for all loaders to settle Promise.all(registrationPromises).then(allResults => { @@ -253,19 +385,21 @@ const registerCollections = async ( const returnValues = { settled, fulfilled, rejected }; - // Fire onSettle if it exists - try { - onSettle?.(returnValues); - } catch (err) { - throw new Error(`Error calling "onSettle" ${formatUnknownError(err)}`); - } + // Fire-and-forget if it exists. Review using `Promise.try` in the future. + Promise.resolve() + .then(() => onSettle?.(returnValues)) + .catch(err => log.debug(`Error calling "onSettle": ${formatUnknownError(err)}`)); }).catch(err => { - log.error(`Failed to settle collections: ${err}`); + log.debug(`Failed to settle collections: ${err}`); }); }; export { + getServerRecordsRegistry, + onUpdateServerRecordsRegistry, registerCollections, + setServerRecordsRegistry, + type OnUpdateServerRecordsRegistryOptions, type McpCollection, type McpCollectionCreator, type McpCollectionRecord, diff --git a/src/patternFly.getResources.ts b/src/patternFly.getResources.ts index be05c5b4..0444ad78 100644 --- a/src/patternFly.getResources.ts +++ b/src/patternFly.getResources.ts @@ -16,7 +16,11 @@ import { type PatternFlyMcpDocsCatalogEntry, type PatternFlyMcpDocsCatalogDoc } from './docs.embedded'; -import { type McpCollectionResult } from './collections'; +import { + onUpdateServerRecordsRegistry, + type McpCollectionResult, + type RegisterCollectionItem +} from './collections'; /** * Derive the component schema type from @patternfly/patternfly-component-schemas @@ -705,6 +709,23 @@ const setPatternFlyCollection = async ( } }; +/** + * Add listener for PatternFly collection updates, see {@link setPatternFlyCollection} + * + * @note We don't need to use the `replay` option here, all of PF collections we need are `required` + * currently, any future updates to this logic may consider adding the `replay` option. + */ +onUpdateServerRecordsRegistry(({ name, response, error }: RegisterCollectionItem) => { + if (name && response) { + setPatternFlyCollection(name, response); + log.info(`Update collection: ${name}`); + } + + if (error) { + log.error(`Update collection error "${name}": ${error}`); + } +}); + export { getPatternFlyComponentSchema, getPatternFlyMcpResources, diff --git a/src/server.ts b/src/server.ts index 7edbb26b..fa4ac8df 100644 --- a/src/server.ts +++ b/src/server.ts @@ -35,11 +35,9 @@ import { import { registerCollections, type McpCollectionCreator, - type McpCollection, - type RegisterCollectionItem + type McpCollection } from './collections'; import { composeCollections } from './server.collections'; -import { setPatternFlyCollection } from './patternFly.getResources'; /** * Server options. Equivalent to GlobalOptions. @@ -151,19 +149,7 @@ const registerServerCollections = async (collections: McpCollectionCreator[], op ] as McpCollection; }); - // Update PatternFly collections, see {@link setPatternFlyCollection} - const onUpdate = ({ name, response, error }: RegisterCollectionItem) => { - if (response) { - setPatternFlyCollection(name, response); - log.info(`Update collection: ${name}`); - } - - if (error) { - log.error(`Update collection error "${name}": ${error}`); - } - }; - - return registerCollections(updatedCollections, { onUpdate }); + return registerCollections(updatedCollections); }; /**