diff --git a/src/containerConfig.ts b/src/containerConfig.ts index 785fb45..8a528f7 100644 --- a/src/containerConfig.ts +++ b/src/containerConfig.ts @@ -1,5 +1,3 @@ -import { Worker } from 'worker_threads'; -import path from 'path'; import config from 'config'; import pino from 'pino'; import client from 'prom-client'; @@ -7,7 +5,6 @@ import { instanceCachingFactory, container, Lifecycle } from 'tsyringe'; import { DependencyContainer } from 'tsyringe/dist/typings/types'; import { trace } from '@opentelemetry/api'; import jsLogger, { LoggerOptions } from '@map-colonies/js-logger'; -import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; import { getOtelMixin } from '@map-colonies/telemetry'; import { SERVICES, SERVICE_NAME } from './common/constants'; import { InjectionObject, registerDependencies } from './common/dependencyRegistration'; @@ -15,10 +12,8 @@ import { IConfig } from './common/interfaces'; import { tracing } from './common/tracing'; import DEMTerrainCacheManager from './heights/models/DEMTerrainCacheManager'; import { heightsRouterFactory, HEIGHTS_ROUTER_SYMBOL } from './heights/routes/heightsRouter'; - import { CatalogRecords } from './heights/models/catalogRecords'; -import { isSame } from './heights/utilities'; -import { WorkerEvent } from './workerCatalogRecords'; +import { CatalogSyncManager } from './heights/models/catalogSyncManager'; export interface RegisterOptions { override?: InjectionObject[]; @@ -28,66 +23,13 @@ export interface RegisterOptions { export const CATALOG_RECORDS_MAP = Symbol('CATALOG_RECORDS_MAP'); export const PRODUCT_METADATA_FIELDS = Symbol('PRODUCT_METADATA_FIELDS'); export const DEM_TERRAIN_CACHE_MANAGER = Symbol('DEM_TERRAIN_CACHE_MANAGER'); +export const CATALOG_SYNC_MANAGER = Symbol('CATALOG_SYNC_MANAGER'); export const registerExternalValues = async (options?: RegisterOptions): Promise => { const loggerConfig = config.get('telemetry.logger'); // @ts-expect-error the signature is wrong const logger = jsLogger({ ...loggerConfig, mixin: getOtelMixin(), timestamp: pino.stdTimeFunctions.isoTime }); - const initCSWWorker = (): void => { - const worker = new Worker(path.resolve(__dirname, './workerCatalogRecords.js')); - - // Listen for updates from the worker - // eslint-disable-next-line - worker.on('message', async (event: WorkerEvent) => { - const data = event; - const dataValue = data.value as PycswDemCatalogRecord[]; - let catalogRecordsServiceInstance, demTerrainCacheManager; - - switch (data.action) { - case 'updateValue': - catalogRecordsServiceInstance = container.resolve(CATALOG_RECORDS_MAP); - demTerrainCacheManager = container.resolve(DEM_TERRAIN_CACHE_MANAGER); - - if (!isSame(dataValue, Object.values(catalogRecordsServiceInstance.getValue()))) { - catalogRecordsServiceInstance.setValue(Object.fromEntries(dataValue.map((record) => [record.id as string, record]))); - await demTerrainCacheManager.initProviders(dataValue); - - logger.info({ - msg: `CatalogRecords UPDATED - ${dataValue.length} records fetched`, - location: '[registerExternalValues]', - }); - } - break; - case 'error': - logger.error({ - msg: `FETCH CatalogRecords ERROR`, - ...dataValue, - location: '[registerExternalValues]', - }); - break; - } - }); - - worker.on('error', (event: WorkerEvent) => { - logger.error({ - msg: `CatalogRecords ERROR`, - ...event, - location: '[registerExternalValues]', - }); - }); - - worker.on('exit', (event: WorkerEvent) => { - logger.error({ - msg: `CatalogRecords EXIT`, - ...event, - location: '[registerExternalValues]', - }); - }); - }; - - initCSWWorker(); - const productMetadataFields = config.get('productMetadataFields').split(','); tracing.start(); @@ -114,12 +56,14 @@ export const registerExternalValues = async (options?: RegisterOptions): Promise { token: CATALOG_RECORDS_MAP, provider: { useClass: CatalogRecords }, options: { lifecycle: Lifecycle.Singleton } }, { token: PRODUCT_METADATA_FIELDS, provider: { useValue: productMetadataFields } }, { token: DEM_TERRAIN_CACHE_MANAGER, provider: { useClass: DEMTerrainCacheManager }, options: { lifecycle: Lifecycle.Singleton } }, + { token: CATALOG_SYNC_MANAGER, provider: { useClass: CatalogSyncManager }, options: { lifecycle: Lifecycle.Singleton } }, { token: HEIGHTS_ROUTER_SYMBOL, provider: { useFactory: heightsRouterFactory } }, { token: 'onSignal', provider: { useValue: { useValue: async (): Promise => { + container.resolve(CATALOG_SYNC_MANAGER).stop(); await Promise.all([tracing.stop()]); }, }, @@ -127,5 +71,14 @@ export const registerExternalValues = async (options?: RegisterOptions): Promise }, ]; - return Promise.resolve(registerDependencies(dependencies, options?.override, options?.useChild)); + const registeredContainer = registerDependencies(dependencies, options?.override, options?.useChild); + + registeredContainer + .resolve(CATALOG_SYNC_MANAGER) + .start( + registeredContainer.resolve(CATALOG_RECORDS_MAP), + registeredContainer.resolve(DEM_TERRAIN_CACHE_MANAGER) + ); + + return Promise.resolve(registeredContainer); }; diff --git a/src/heights/models/DEMTerrainCacheManager.ts b/src/heights/models/DEMTerrainCacheManager.ts index dfb204f..31f47f3 100644 --- a/src/heights/models/DEMTerrainCacheManager.ts +++ b/src/heights/models/DEMTerrainCacheManager.ts @@ -2,6 +2,7 @@ import { inject, injectable } from 'tsyringe'; import { IConfig } from 'config'; import { Logger } from '@map-colonies/js-logger'; import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; +import PromisePool from '@supercharge/promise-pool/dist'; import { HeightProviders } from '../interfaces'; import { SERVICES } from '../../common/constants'; import GeotiffHeightProvider from './geotiffHeightProvider'; @@ -21,25 +22,27 @@ export default class DEMTerrainCacheManager { const geotiffRecords = demCatalogRecords.filter((record) => record.links?.some((link) => link.protocol === GEOTIFF_PROTOCOL)); const samplingConcurrency = Number(this.config.get('samplingConcurrency')); - for (const record of geotiffRecords) { - const link = record.links?.find((currentLink) => currentLink.protocol === GEOTIFF_PROTOCOL); - if (!link) { - continue; - } - - try { - const objectUrl = this.transformRouteToObjectUrl(link.url as string); - const { url, headers } = this.buildAuthenticatedUrl(objectUrl); - heightProviders[record.id as string] = await GeotiffHeightProvider.fromUrl(url, headers, samplingConcurrency); - } catch (err) { + // Open providers concurrently but bounded — each fromUrl is an independent gateway round-trip. + // A single record's failure is isolated (logged, skipped) so the rest still register. + await PromisePool.for(geotiffRecords) + .withConcurrency(samplingConcurrency) + .handleError((err, record) => { this.logger.error({ msg: 'Failed to open geotiff provider; skipping record', recordId: record.id, err, location: '[DEMTerrainCacheManager] [initProviders]', }); - } - } + }) + .process(async (record) => { + const link = record.links?.find((currentLink) => currentLink.protocol === GEOTIFF_PROTOCOL); + if (!link) { + return; + } + const objectUrl = this.transformRouteToObjectUrl(link.url as string); + const { url, headers } = this.buildAuthenticatedUrl(objectUrl); + heightProviders[record.id as string] = await GeotiffHeightProvider.fromUrl(url, headers, samplingConcurrency); + }); this.heightProviders = heightProviders; } diff --git a/src/heights/models/catalogSyncManager.ts b/src/heights/models/catalogSyncManager.ts new file mode 100644 index 0000000..3dc5f4c --- /dev/null +++ b/src/heights/models/catalogSyncManager.ts @@ -0,0 +1,95 @@ +import { inject, injectable } from 'tsyringe'; +import { Logger } from '@map-colonies/js-logger'; +import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; +import { IConfig } from '../../common/interfaces'; +import { SERVICES } from '../../common/constants'; +import { CswClientWrapper } from '../../common/csw/cswClientWrapper'; +import { IService } from '../../common/csw/utils'; +import { isSame } from '../utilities'; +import { CatalogRecords } from './catalogRecords'; +import DEMTerrainCacheManager from './DEMTerrainCacheManager'; + +const START_RECORD = 1; +const END_RECORD = 1000; + +// isSame checksums a JSON serialization, so array order matters. CSW may return the same +// records in a different order between polls — normalize by id so reordering alone doesn't +// trigger a needless provider rebuild. +const sortById = (records: PycswDemCatalogRecord[]): PycswDemCatalogRecord[] => + [...records].sort((first, second) => (first.id as string).localeCompare(second.id as string)); + +@injectable() +export class CatalogSyncManager { + private readonly cswClient: CswClientWrapper; + private readonly intervalMs: number; + private timer: NodeJS.Timeout | undefined; + private stopped = false; + private catalogRecords: CatalogRecords | undefined; + private cacheManager: DEMTerrainCacheManager | undefined; + + public constructor(@inject(SERVICES.CONFIG) private readonly config: IConfig, @inject(SERVICES.LOGGER) private readonly logger: Logger) { + this.cswClient = new CswClientWrapper( + 'mc:MCDEMRecord', + PycswDemCatalogRecord.getPyCSWMappings(), + 'http://schema.mapcolonies.com/dem', + this.config.get('csw') + ); + this.intervalMs = this.config.get('synchRecordsInterval'); + } + + public start(catalogRecords: CatalogRecords, cacheManager: DEMTerrainCacheManager): void { + this.catalogRecords = catalogRecords; + this.cacheManager = cacheManager; + this.stopped = false; + void this.syncOnce(); + } + + public stop(): void { + this.stopped = true; + if (this.timer) { + clearTimeout(this.timer); + this.timer = undefined; + } + } + + private async syncOnce(): Promise { + const catalogRecords = this.catalogRecords; + const cacheManager = this.cacheManager; + + try { + const records = await this.fetchCatalogRecords(); + if (catalogRecords && cacheManager && !isSame(sortById(records), sortById(Object.values(catalogRecords.getValue())))) { + // Rebuild providers before publishing the new catalog. Otherwise, during initProviders' + // await window a reader sees the new catalog paired with stale providers, and a removed + // record resolves to an undefined catalog entry (throws on .footprint access). + await cacheManager.initProviders(records); + catalogRecords.setValue(Object.fromEntries(records.map((record) => [record.id as string, record]))); + this.logger.info({ msg: `CatalogRecords UPDATED - ${records.length} records fetched`, location: '[CatalogSyncManager]' }); + } + } catch (err) { + this.logger.error({ msg: 'FETCH CatalogRecords ERROR', err, location: '[CatalogSyncManager]' }); + } finally { + if (!this.stopped) { + this.timer = setTimeout(() => void this.syncOnce(), this.intervalMs); + } + } + } + + private async fetchCatalogRecords(): Promise { + return this.cswClient.getRecords(START_RECORD, END_RECORD, { + filter: [ + // DEM profile links carry the object protocol. We match records exposing a GEOTIFF link + // (COG served from the S3 gateway) via a LIKE filter on the LINKS field. + { + field: 'mc:links', + like: 'GEOTIFF', + }, + { + field: 'mc:productStatus', + eq: 'PUBLISHED', + }, + ], + sort: undefined, + }); + } +} diff --git a/src/workerCatalogRecords.ts b/src/workerCatalogRecords.ts deleted file mode 100644 index a33c2e2..0000000 --- a/src/workerCatalogRecords.ts +++ /dev/null @@ -1,60 +0,0 @@ -import { parentPort } from 'worker_threads'; -import config from 'config'; -import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; -import { CswClientWrapper } from './common/csw/cswClientWrapper'; -import { IService } from './common/csw/utils'; - -const SYNCH_RECORDS_INTERVAL = config.get('synchRecordsInterval'); - -const cswClient = new CswClientWrapper( - 'mc:MCDEMRecord', - PycswDemCatalogRecord.getPyCSWMappings(), - 'http://schema.mapcolonies.com/dem', - config.get('csw') -); - -const START_RECORD = 1; -const END_RECORD = 1000; - -const getCatalogRecords = async (): Promise => { - const res = await cswClient.getRecords(START_RECORD, END_RECORD, { - filter: [ - // ******* DEM profile links carry the object protocol. We match records exposing a GEOTIFF link - // ******* (COG served from the S3 gateway) via a LIKE filter on the LINKS field. - // { - // field: 'mc:hasTerrain', - // eq: 'True', - // }, - { - field: 'mc:links', - like: 'GEOTIFF', - }, - { - field: 'mc:productStatus', - eq: 'PUBLISHED', - }, - ], - sort: undefined, - }); - - return res; -}; - -// eslint-disable-next-line -(async function updateValuePeriodically(): Promise { - try { - const newValue = await getCatalogRecords(); - parentPort?.postMessage({ action: 'updateValue', value: newValue }); - } catch (err) { - parentPort?.postMessage({ action: 'error', value: err }); - } - - // eslint-disable-next-line - setTimeout(updateValuePeriodically, SYNCH_RECORDS_INTERVAL); -})(); - -export interface WorkerEvent { - action: 'updateValue' | 'error'; - // eslint-disable-next-line - value: any; -} diff --git a/tests/integration/heights/heights.spec.ts b/tests/integration/heights/heights.spec.ts index af335a8..30dca19 100644 --- a/tests/integration/heights/heights.spec.ts +++ b/tests/integration/heights/heights.spec.ts @@ -7,7 +7,7 @@ import { container } from 'tsyringe'; import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; import { getApp } from '../../../src/app'; import { SERVICES } from '../../../src/common/constants'; -import { CATALOG_RECORDS_MAP, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig'; +import { CATALOG_RECORDS_MAP, CATALOG_SYNC_MANAGER, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig'; import { GetHeightsPointsRequest, GetHeightsPointsResponse } from '../../../src/heights/controllers/heightsController'; import { PosWithHeight, TerrainTypes } from '../../../src/heights/interfaces'; import { CatalogRecords } from '../../../src/heights/models/catalogRecords'; @@ -40,7 +40,10 @@ describe('heights', function () { productMetadataFields = config.get('productMetadataFields').split(','); const app = await getApp({ - override: [{ token: SERVICES.LOGGER, provider: { useValue: jsLogger({ enabled: false }) } }], + override: [ + { token: SERVICES.LOGGER, provider: { useValue: jsLogger({ enabled: false }) } }, + { token: CATALOG_SYNC_MANAGER, provider: { useValue: { start: (): void => undefined, stop: (): void => undefined } } }, + ], }); requestSender = new HeightsRequestSender(app as Application); diff --git a/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts b/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts index 4a2aa6a..2944335 100644 --- a/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts +++ b/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts @@ -28,4 +28,19 @@ describe('DEMTerrainCacheManager', () => { expect(Object.keys(manager.heightProviders)).toEqual(['good']); expect(errorSpy).toHaveBeenCalled(); }); + + it('registers a provider for every geotiff record', async () => { + const records = [ + { id: 'a', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/a.tif' }] }, + { id: 'b', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/b.tif' }] }, + { id: 'c', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/c.tif' }] }, + ] as unknown as PycswDemCatalogRecord[]; + + jest.spyOn(GeotiffHeightProvider, 'fromUrl').mockResolvedValue({} as GeotiffHeightProvider); + + const manager = new DEMTerrainCacheManager(config, jsLogger({ enabled: false })); + await manager.initProviders(records); + + expect(Object.keys(manager.heightProviders).sort((first, second) => first.localeCompare(second))).toEqual(['a', 'b', 'c']); + }); }); diff --git a/tests/unit/heights/models/catalogSyncManager.spec.ts b/tests/unit/heights/models/catalogSyncManager.spec.ts new file mode 100644 index 0000000..d8c444b --- /dev/null +++ b/tests/unit/heights/models/catalogSyncManager.spec.ts @@ -0,0 +1,138 @@ +import config from 'config'; +import jsLogger from '@map-colonies/js-logger'; +import { PycswDemCatalogRecord } from '@map-colonies/mc-model-types'; +import { CswClientWrapper } from '../../../../src/common/csw/cswClientWrapper'; +import { CatalogSyncManager } from '../../../../src/heights/models/catalogSyncManager'; +import { CatalogRecords } from '../../../../src/heights/models/catalogRecords'; +import DEMTerrainCacheManager from '../../../../src/heights/models/DEMTerrainCacheManager'; + +const RECORDS = [{ id: 'r1', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/r1.tif' }] }] as unknown as PycswDemCatalogRecord[]; + +function makeCacheManager(): { cacheManager: DEMTerrainCacheManager; initProviders: jest.Mock } { + const initProviders = jest.fn().mockResolvedValue(undefined); + return { cacheManager: { initProviders } as unknown as DEMTerrainCacheManager, initProviders }; +} + +describe('CatalogSyncManager', () => { + let manager: CatalogSyncManager | undefined; + + afterEach(() => { + manager?.stop(); + manager = undefined; + jest.restoreAllMocks(); + jest.useRealTimers(); + }); + + it('fetches records, updates the cache, and rebuilds providers on first run', async () => { + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue(RECORDS); + const catalogRecords = new CatalogRecords(); + const { cacheManager, initProviders } = makeCacheManager(); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(catalogRecords, cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(Object.keys(catalogRecords.getValue())).toEqual(['r1']); + expect(initProviders).toHaveBeenCalledWith(RECORDS); + }); + + it('does not rebuild providers when the fetched records are unchanged', async () => { + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue(RECORDS); + const catalogRecords = new CatalogRecords(); + catalogRecords.setValue(Object.fromEntries(RECORDS.map((r) => [r.id as string, r]))); + const { cacheManager, initProviders } = makeCacheManager(); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(catalogRecords, cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(initProviders).not.toHaveBeenCalled(); + }); + + it('logs and survives a fetch error without throwing', async () => { + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockRejectedValue(new Error('csw down')); + const logger = jsLogger({ enabled: false }); + const errorSpy = jest.spyOn(logger, 'error'); + const { cacheManager, initProviders } = makeCacheManager(); + + manager = new CatalogSyncManager(config, logger); + manager.start(new CatalogRecords(), cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(errorSpy).toHaveBeenCalled(); + expect(initProviders).not.toHaveBeenCalled(); + }); + + it('schedules the next run with the configured interval after a cycle completes', async () => { + const setTimeoutSpy = jest.spyOn(global, 'setTimeout'); + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue(RECORDS); + const interval = config.get('synchRecordsInterval'); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(new CatalogRecords(), makeCacheManager().cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + const scheduledWithInterval = setTimeoutSpy.mock.calls.filter((call) => call[1] === interval); + expect(scheduledWithInterval.length).toBeGreaterThanOrEqual(1); + }); + + it('stop() clears the pending timer so no further run is scheduled', async () => { + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue(RECORDS); + const clearTimeoutSpy = jest.spyOn(global, 'clearTimeout'); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(new CatalogRecords(), makeCacheManager().cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + manager.stop(); + + expect(clearTimeoutSpy).toHaveBeenCalled(); + }); + + it('reschedules the next run even after a fetch error', async () => { + const setTimeoutSpy = jest.spyOn(global, 'setTimeout'); + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockRejectedValue(new Error('csw down')); + const interval = config.get('synchRecordsInterval'); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(new CatalogRecords(), makeCacheManager().cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + const scheduledWithInterval = setTimeoutSpy.mock.calls.filter((call) => call[1] === interval); + expect(scheduledWithInterval.length).toBeGreaterThanOrEqual(1); + }); + + it('does not rebuild providers when the same records return in a different order', async () => { + const twoRecords = [ + { id: 'a', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/a.tif' }] }, + { id: 'b', links: [{ protocol: 'GEOTIFF', url: 'https://gw/cogs/b.tif' }] }, + ] as unknown as PycswDemCatalogRecord[]; + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue([twoRecords[1], twoRecords[0]]); + const catalogRecords = new CatalogRecords(); + catalogRecords.setValue(Object.fromEntries(twoRecords.map((record) => [record.id as string, record]))); + const { cacheManager, initProviders } = makeCacheManager(); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(catalogRecords, cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(initProviders).not.toHaveBeenCalled(); + }); + + it('rebuilds providers before publishing the new catalog', async () => { + jest.spyOn(CswClientWrapper.prototype, 'getRecords').mockResolvedValue(RECORDS); + const catalogRecords = new CatalogRecords(); + const initProviders = jest.fn().mockImplementation(async () => { + // the new catalog must not be visible while providers are still being (re)built + expect(Object.keys(catalogRecords.getValue())).toEqual([]); + await Promise.resolve(); + }); + const cacheManager = { initProviders } as unknown as DEMTerrainCacheManager; + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(catalogRecords, cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(initProviders).toHaveBeenCalledWith(RECORDS); + expect(Object.keys(catalogRecords.getValue())).toEqual(['r1']); + }); +});