From 9cbe807f06db748eac38c6386958c05990ac827e Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Tue, 22 Sep 2026 14:42:31 +0300 Subject: [PATCH 1/7] perf: open geotiff providers concurrently with bounded PromisePool --- src/heights/models/DEMTerrainCacheManager.ts | 29 ++++++++++--------- .../models/DEMTerrainCacheManager.spec.ts | 15 ++++++++++ 2 files changed, 31 insertions(+), 13 deletions(-) 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/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts b/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts index 4a2aa6a..7071742 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()).toEqual(['a', 'b', 'c']); + }); }); From 667f8576d41c05adb7a9465519aeb61b209935e3 Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Tue, 22 Sep 2026 14:53:00 +0300 Subject: [PATCH 2/7] feat: add in-process CatalogSyncManager poll loop --- src/heights/models/catalogSyncManager.ts | 86 +++++++++++++++ .../heights/models/catalogSyncManager.spec.ts | 103 ++++++++++++++++++ 2 files changed, 189 insertions(+) create mode 100644 src/heights/models/catalogSyncManager.ts create mode 100644 tests/unit/heights/models/catalogSyncManager.spec.ts diff --git a/src/heights/models/catalogSyncManager.ts b/src/heights/models/catalogSyncManager.ts new file mode 100644 index 0000000..ea2fc24 --- /dev/null +++ b/src/heights/models/catalogSyncManager.ts @@ -0,0 +1,86 @@ +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; + +@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(records, Object.values(catalogRecords.getValue()))) { + catalogRecords.setValue(Object.fromEntries(records.map((record) => [record.id as string, record]))); + await cacheManager.initProviders(records); + 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/tests/unit/heights/models/catalogSyncManager.spec.ts b/tests/unit/heights/models/catalogSyncManager.spec.ts new file mode 100644 index 0000000..bfcdc14 --- /dev/null +++ b/tests/unit/heights/models/catalogSyncManager.spec.ts @@ -0,0 +1,103 @@ +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); + }); +}); From 2c226cdf029bfe34fc907ab235dc06b6ed6dc2a1 Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Tue, 22 Sep 2026 15:01:18 +0300 Subject: [PATCH 3/7] refactor: replace CSW worker thread with in-process CatalogSyncManager --- src/containerConfig.ts | 75 +++++------------------ src/workerCatalogRecords.ts | 60 ------------------ tests/integration/heights/heights.spec.ts | 7 ++- 3 files changed, 19 insertions(+), 123 deletions(-) delete mode 100644 src/workerCatalogRecords.ts 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/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); From fce333edc449684ebe3076f987b77e26f199b09a Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Tue, 22 Sep 2026 15:11:28 +0300 Subject: [PATCH 4/7] test: use comparator in sort assertion to satisfy eslint require-array-sort-compare flags a bare .sort() on the provider-key assertion added for the concurrent initProviders test. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01ATQerkHAV3Vur82HVVFack --- tests/unit/heights/models/DEMTerrainCacheManager.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts b/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts index 7071742..2944335 100644 --- a/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts +++ b/tests/unit/heights/models/DEMTerrainCacheManager.spec.ts @@ -41,6 +41,6 @@ describe('DEMTerrainCacheManager', () => { const manager = new DEMTerrainCacheManager(config, jsLogger({ enabled: false })); await manager.initProviders(records); - expect(Object.keys(manager.heightProviders).sort()).toEqual(['a', 'b', 'c']); + expect(Object.keys(manager.heightProviders).sort((first, second) => first.localeCompare(second))).toEqual(['a', 'b', 'c']); }); }); From dd60b4e09d9c4b8953e03435507d3d15363db2d6 Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Tue, 22 Sep 2026 15:11:28 +0300 Subject: [PATCH 5/7] docs: add in-process catalog-sync refactor plan Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01ATQerkHAV3Vur82HVVFack --- ...6-09-22-catalog-sync-inprocess-refactor.md | 551 ++++++++++++++++++ 1 file changed, 551 insertions(+) create mode 100644 docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md diff --git a/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md b/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md new file mode 100644 index 0000000..eb87160 --- /dev/null +++ b/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md @@ -0,0 +1,551 @@ +# Catalog Sync In-Process Refactor Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Replace the CSW catalog-sync worker thread with an in-process polling manager, and parallelize GeoTIFF provider opening, without changing the `/points` API or behavior. + +**Architecture:** The catalog is fetched from CSW on a fixed timer and cached in-memory (singleton `CatalogRecords` + `DEMTerrainCacheManager.heightProviders`). Today a Node `worker_threads` Worker does the fetch and posts results back to the main thread, which diffs (`isSame`) and rebuilds providers. This is wrong for pure network I/O (threads are for CPU work), the fetched record set is structured-cloned across the thread boundary every cycle, and a dead worker silently freezes the cache (its `exit`/`error` handlers only log, no restart). This refactor moves the poll loop in-process into a new `CatalogSyncManager` (a `setTimeout` self-rescheduling loop that survives fetch errors), and replaces the sequential `for`-loop in `initProviders` with a bounded-concurrency `PromisePool`. + +**Tech Stack:** TypeScript, tsyringe DI, node-config, `@map-colonies/csw-client`, `@supercharge/promise-pool` (already a dependency on this branch), Jest + ts-jest. + +**Base branch:** `feat/geotiff-heights-migration` (PR #52) — NOT `master`. That PR rewrites `DEMTerrainCacheManager`, renames `initTerrainProviders`→`initProviders`, and touches `workerCatalogRecords.ts` + `containerConfig.ts`; branching off master would collide head-on and refactor soon-to-be-deleted Cesium code. + +--- + +## Prerequisite: branch off #52 + +- [ ] **Step 0: Create the working branch from the #52 head** + +```bash +git fetch origin feat/geotiff-heights-migration +git switch -c refactor/catalog-sync-inprocess origin/feat/geotiff-heights-migration +git log --oneline -1 # expect: tip of feat/geotiff-heights-migration +``` + +--- + +## File Structure + +- **Create:** `src/heights/models/catalogSyncManager.ts` — in-process CSW poll loop; owns the `CswClientWrapper`, the fetch filter, the `isSame` diff, and `setTimeout` scheduling. Injects only `CONFIG` + `LOGGER`; receives the two cache singletons via `start(...)` to avoid a circular import with `containerConfig`. +- **Create:** `tests/unit/heights/models/catalogSyncManager.spec.ts` — unit tests for the manager (update, no-op-when-unchanged, error-survives, stop()). +- **Modify:** `src/containerConfig.ts` — delete the worker (`initCSWWorker`, `Worker`, `path`, `WorkerEvent`); register `CATALOG_SYNC_MANAGER`; start it after registration; stop it in `onSignal`. +- **Modify:** `src/heights/models/DEMTerrainCacheManager.ts` — replace the sequential `for` loop in `initProviders` with a `PromisePool` bounded by `samplingConcurrency`; keep per-record error isolation. +- **Modify:** `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts` — keep the isolation test; add a "registers every record" test. +- **Modify:** `tests/integration/heights/heights.spec.ts` — override `CATALOG_SYNC_MANAGER` with a no-op stub so the real poll loop never runs during integration tests. +- **Delete:** `src/workerCatalogRecords.ts` — its `getCatalogRecords`, CSW client, `START_RECORD`/`END_RECORD`, and filter move into `CatalogSyncManager`; `WorkerEvent` is deleted with the worker. + +No config keys are added — provider-open concurrency reuses the existing `samplingConcurrency` key. No build-script or Dockerfile change (removing the worker only removes the `./workerCatalogRecords.js` runtime dependency). + +--- + +## Task 1: Parallelize `initProviders` (isolated, low-risk — do first) + +**Files:** + +- Modify: `src/heights/models/DEMTerrainCacheManager.ts` +- Test: `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts` + +- [ ] **Step 1: Add the "registers every record" failing test** + +Append this test inside the existing `describe('DEMTerrainCacheManager', ...)` block in `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts`: + +```ts +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()).toEqual(['a', 'b', 'c']); +}); +``` + +- [ ] **Step 2: Run the test to verify it PASSES against the current sequential loop** + +Run: `npm run test:unit -- --testPathPattern DEMTerrainCacheManager` +Expected: PASS. (This test is a behavior-preservation guard — it must stay green through the refactor. The isolation test is the existing safety net.) + +- [ ] **Step 3: Replace the sequential loop with a bounded `PromisePool`** + +In `src/heights/models/DEMTerrainCacheManager.ts`, add the import near the other imports: + +```ts +import PromisePool from '@supercharge/promise-pool/dist'; +``` + +Replace the entire `initProviders` method body with: + +```ts + public async initProviders(demCatalogRecords: PycswDemCatalogRecord[]): Promise { + const heightProviders: HeightProviders = {}; + + const geotiffRecords = demCatalogRecords.filter((record) => record.links?.some((link) => link.protocol === GEOTIFF_PROTOCOL)); + const samplingConcurrency = Number(this.config.get('samplingConcurrency')); + + // 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; + } +``` + +Note: `PromisePool.handleError` that returns (does not throw) skips the failed item and lets the pool continue — this preserves the existing "skip a record whose provider fails to open" contract. `transformRouteToObjectUrl`, `buildAuthenticatedUrl`, `GEOTIFF_PROTOCOL`, and `HeightProviders` are unchanged and already present in the file. + +- [ ] **Step 4: Run both DEMTerrainCacheManager tests to verify they PASS** + +Run: `npm run test:unit -- --testPathPattern DEMTerrainCacheManager` +Expected: PASS — both the isolation test (skips `bad`, keeps `good`, logs error) and the "registers every record" test. + +- [ ] **Step 5: Commit** + +```bash +git add src/heights/models/DEMTerrainCacheManager.ts tests/unit/heights/models/DEMTerrainCacheManager.spec.ts +git commit -m "perf: open geotiff providers concurrently with bounded PromisePool" +``` + +--- + +## Task 2: Create `CatalogSyncManager` with unit tests + +**Files:** + +- Create: `src/heights/models/catalogSyncManager.ts` +- Test: `tests/unit/heights/models/catalogSyncManager.spec.ts` + +- [ ] **Step 1: Write the failing unit test** + +Create `tests/unit/heights/models/catalogSyncManager.spec.ts`: + +```ts +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(): DEMTerrainCacheManager { + return { initProviders: jest.fn().mockResolvedValue(undefined) } as unknown as DEMTerrainCacheManager; +} + +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 = 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(cacheManager.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 = makeCacheManager(); + + manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); + manager.start(catalogRecords, cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(cacheManager.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 = makeCacheManager(); + + manager = new CatalogSyncManager(config, logger); + manager.start(new CatalogRecords(), cacheManager); + await new Promise((resolve) => setImmediate(resolve)); + + expect(errorSpy).toHaveBeenCalled(); + expect(cacheManager.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()); + 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()); + await new Promise((resolve) => setImmediate(resolve)); + manager.stop(); + + expect(clearTimeoutSpy).toHaveBeenCalled(); + }); +}); +``` + +> **Jest 28 note:** this repo pins `jest@^28`, which has no `jest.advanceTimersByTimeAsync` (Jest 29+). The reschedule/stop behavior is therefore verified with real timers + `setImmediate` flushing (same mechanism as the first three tests) plus `setTimeout`/`clearTimeout` spies, rather than by advancing fake timers through async cycles. `afterEach` calls `manager?.stop()`, which clears the real pending timer so no handle leaks. + +- [ ] **Step 2: Run the test to verify it FAILS** + +Run: `npm run test:unit -- --testPathPattern catalogSyncManager` +Expected: FAIL — `Cannot find module '.../catalogSyncManager'` (file not yet created). + +- [ ] **Step 3: Create the `CatalogSyncManager` implementation** + +Create `src/heights/models/catalogSyncManager.ts`: + +```ts +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; + +@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(records, Object.values(catalogRecords.getValue()))) { + catalogRecords.setValue(Object.fromEntries(records.map((record) => [record.id as string, record]))); + await cacheManager.initProviders(records); + 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, + }); + } +} +``` + +- [ ] **Step 4: Run the test to verify it PASSES** + +Run: `npm run test:unit -- --testPathPattern catalogSyncManager` +Expected: PASS — all five cases (update, no-op-unchanged, error-survives, schedules-next-run, stop-clears-timer). + +- [ ] **Step 5: Commit** + +```bash +git add src/heights/models/catalogSyncManager.ts tests/unit/heights/models/catalogSyncManager.spec.ts +git commit -m "feat: add in-process CatalogSyncManager poll loop" +``` + +--- + +## Task 3: Wire `CatalogSyncManager` into bootstrap; remove the worker + +**Files:** + +- Modify: `src/containerConfig.ts` +- Modify: `tests/integration/heights/heights.spec.ts` +- Delete: `src/workerCatalogRecords.ts` + +- [ ] **Step 1: Stub the sync manager in the integration override (write first — this is the guard)** + +In `tests/integration/heights/heights.spec.ts`, add the import alongside the existing container imports: + +```ts +import { CATALOG_RECORDS_MAP, CATALOG_SYNC_MANAGER, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig'; +``` + +(That replaces the existing `import { CATALOG_RECORDS_MAP, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig';` line — add `CATALOG_SYNC_MANAGER` to it.) + +Then extend the `getApp` override in `beforeAll` so the real poll loop never runs in tests: + +```ts +const app = await getApp({ + override: [ + { token: SERVICES.LOGGER, provider: { useValue: jsLogger({ enabled: false }) } }, + { token: CATALOG_SYNC_MANAGER, provider: { useValue: { start: (): void => undefined, stop: (): void => undefined } } }, + ], +}); +``` + +- [ ] **Step 2: Rewrite `src/containerConfig.ts` to remove the worker and start the manager** + +Replace the entire contents of `src/containerConfig.ts` with: + +```ts +import config from 'config'; +import pino from 'pino'; +import client from 'prom-client'; +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 { getOtelMixin } from '@map-colonies/telemetry'; +import { SERVICES, SERVICE_NAME } from './common/constants'; +import { InjectionObject, registerDependencies } from './common/dependencyRegistration'; +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 { CatalogSyncManager } from './heights/models/catalogSyncManager'; + +export interface RegisterOptions { + override?: InjectionObject[]; + useChild?: boolean; +} + +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 productMetadataFields = config.get('productMetadataFields').split(','); + + tracing.start(); + const tracer = trace.getTracer(SERVICE_NAME); + + const dependencies: InjectionObject[] = [ + { token: SERVICES.CONFIG, provider: { useValue: config } }, + { token: SERVICES.LOGGER, provider: { useValue: logger } }, + { token: SERVICES.TRACER, provider: { useValue: tracer } }, + { + token: SERVICES.METRICS_REGISTRY, + provider: { + useFactory: instanceCachingFactory((container) => { + const config = container.resolve(SERVICES.CONFIG); + if (config.get('telemetry.metrics.enabled')) { + client.register.setDefaultLabels({ + app: SERVICE_NAME, + }); + return client.register; + } + }), + }, + }, + { 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()]); + }, + }, + }, + }, + ]; + + 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 registeredContainer; +}; +``` + +What changed vs. the #52 version: removed `import { Worker } from 'worker_threads'`, `import path`, the `PycswDemCatalogRecord` import, `import { isSame }`, and `import { WorkerEvent }`; deleted the entire `initCSWWorker` function and its call; added the `CatalogSyncManager` import, the `CATALOG_SYNC_MANAGER` symbol, its DI registration, the `.start(...)` call after registration, and the `.stop()` call in `onSignal`. The `container` import is retained (used by `onSignal` and the metrics factory). + +- [ ] **Step 3: Delete the worker source file** + +```bash +git rm src/workerCatalogRecords.ts +``` + +- [ ] **Step 4: Verify nothing else references the worker** + +Run: `git grep -n "workerCatalogRecords\|WorkerEvent\|initCSWWorker\|worker_threads" -- src tests` +Expected: no output (empty). If anything prints, fix that reference before proceeding. + +- [ ] **Step 5: Build to confirm no dangling imports / type errors** + +Run: `npm run build` +Expected: clean exit, no TypeScript errors. + +- [ ] **Step 6: Commit** + +```bash +git add src/containerConfig.ts tests/integration/heights/heights.spec.ts +git commit -m "refactor: replace CSW worker thread with in-process CatalogSyncManager" +``` + +--- + +## Task 4: Full verification + +**Files:** none (verification only) + +- [ ] **Step 1: Run the full unit suite with coverage** + +Run: `npm run test:unit` +Expected: all suites PASS, coverage gate met (branches ≥60, functions ≥80, lines ≥80). New `catalogSyncManager.ts` is under `src/heights/models/` (in the unit coverage set) and is covered by Task 2's tests; `containerConfig.ts` is excluded from unit coverage by the `!/src/*` glob. + +- [ ] **Step 2: Run the full integration suite** + +Run: `npm run test:integration` +Expected: 8/8 PASS, coverage gate met, and — unlike before — no `CatalogRecords ERROR` / worker-exit noise in the log (the sync manager is stubbed via the override; no real CSW fetch, no leaked timers, Jest exits cleanly). + +- [ ] **Step 3: Lint** + +Run: `npm run lint` +Expected: clean. + +- [ ] **Step 4: Confirm the worker is fully gone from the build output** + +Run: `ls dist/workerCatalogRecords.js 2>/dev/null && echo "STILL PRESENT — investigate" || echo "removed OK"` +Expected: `removed OK` (a stale `dist` from a previous build is fine to ignore; a fresh `npm run build` in Task 3 Step 5 will not emit it). + +- [ ] **Step 5: Push and open the PR against the #52 branch** + +```bash +git push -u origin refactor/catalog-sync-inprocess +gh pr create --base feat/geotiff-heights-migration --repo MapColonies/dem-heights \ + --title "refactor: in-process catalog sync + concurrent provider opens (MAPCO-11560)" \ + --body "Replaces the CSW worker thread with an in-process CatalogSyncManager poll loop and parallelizes GeoTIFF provider opening. Stacked on #52." +``` + +Note: base is `feat/geotiff-heights-migration`, so this PR merges into #52, not master. If #52 merges to master first, rebase this branch onto master before merging. + +--- + +## Self-Review + +**Spec coverage:** + +- Replace worker with in-process poll → Tasks 2 (manager) + 3 (wiring/removal). ✅ +- Fix silent-death → inherent: the `setTimeout` self-reschedule in `syncOnce`'s `finally` survives fetch errors; there is no separate process to die (Task 2, error test proves survival). ✅ +- Parallelize provider opens → Task 1. ✅ +- No API/behavior change → `/points`, response shape, `openapi3.yaml` untouched; startup remains non-blocking (loop kicked off, not awaited), matching prior behavior. ✅ +- No test regressions / no leaked timers → Task 3 Step 1 stub + Task 4. ✅ + +**Placeholder scan:** No TBD/TODO/"handle errors appropriately"; every code step contains full code. ✅ + +**Type consistency:** `initProviders` (not `initTerrainProviders`); `heightProviders` property; `HeightProviders` type; `start(catalogRecords, cacheManager)` signature matches both the `containerConfig` caller and the unit-test caller; `CATALOG_SYNC_MANAGER` symbol name identical across `containerConfig`, `onSignal`, and the integration override. ✅ + +**Decisions locked:** + +- Manager injects only CONFIG + LOGGER and receives singletons via `start(...)` — deliberately avoids importing the DI symbols from `containerConfig` (which imports the manager), preventing a circular-import decoration hazard. +- Provider-open concurrency reuses `samplingConcurrency` rather than adding a new config key (records are few; avoids config sprawl across `default.json` / `custom-environment-variables.json` / helm). From f3e598be320fe57127d2ecbeb52eabb02657f227 Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Thu, 24 Sep 2026 11:52:01 +0300 Subject: [PATCH 6/7] fix: rebuild providers before publishing catalog and normalize record order Addresses PR #57 review: - Swap so initProviders runs before catalogRecords.setValue, so a reader never sees the new catalog paired with stale providers (a removed record would resolve to an undefined catalog entry mid-refresh). - Sort records by id before the isSame diff so CSW returning the same set in a different order no longer triggers a needless provider rebuild. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01ATQerkHAV3Vur82HVVFack --- src/heights/models/catalogSyncManager.ts | 13 +++++-- .../heights/models/catalogSyncManager.spec.ts | 35 +++++++++++++++++++ 2 files changed, 46 insertions(+), 2 deletions(-) diff --git a/src/heights/models/catalogSyncManager.ts b/src/heights/models/catalogSyncManager.ts index ea2fc24..3dc5f4c 100644 --- a/src/heights/models/catalogSyncManager.ts +++ b/src/heights/models/catalogSyncManager.ts @@ -12,6 +12,12 @@ 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; @@ -52,9 +58,12 @@ export class CatalogSyncManager { try { const records = await this.fetchCatalogRecords(); - if (catalogRecords && cacheManager && !isSame(records, Object.values(catalogRecords.getValue()))) { - catalogRecords.setValue(Object.fromEntries(records.map((record) => [record.id as string, record]))); + 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) { diff --git a/tests/unit/heights/models/catalogSyncManager.spec.ts b/tests/unit/heights/models/catalogSyncManager.spec.ts index bfcdc14..d8c444b 100644 --- a/tests/unit/heights/models/catalogSyncManager.spec.ts +++ b/tests/unit/heights/models/catalogSyncManager.spec.ts @@ -100,4 +100,39 @@ describe('CatalogSyncManager', () => { 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']); + }); }); From c8f70a3b7ce8d7f713d09ac1a40b5ad04e9c072d Mon Sep 17 00:00:00 2001 From: shimoncohen Date: Mon, 28 Sep 2026 12:09:21 +0300 Subject: [PATCH 7/7] docs: untrack superpowers implementation plan Co-Authored-By: Claude Opus 5.5 (1M context) --- ...6-09-22-catalog-sync-inprocess-refactor.md | 551 ------------------ 1 file changed, 551 deletions(-) delete mode 100644 docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md diff --git a/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md b/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md deleted file mode 100644 index eb87160..0000000 --- a/docs/superpowers/plans/2026-09-22-catalog-sync-inprocess-refactor.md +++ /dev/null @@ -1,551 +0,0 @@ -# Catalog Sync In-Process Refactor Implementation Plan - -> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. - -**Goal:** Replace the CSW catalog-sync worker thread with an in-process polling manager, and parallelize GeoTIFF provider opening, without changing the `/points` API or behavior. - -**Architecture:** The catalog is fetched from CSW on a fixed timer and cached in-memory (singleton `CatalogRecords` + `DEMTerrainCacheManager.heightProviders`). Today a Node `worker_threads` Worker does the fetch and posts results back to the main thread, which diffs (`isSame`) and rebuilds providers. This is wrong for pure network I/O (threads are for CPU work), the fetched record set is structured-cloned across the thread boundary every cycle, and a dead worker silently freezes the cache (its `exit`/`error` handlers only log, no restart). This refactor moves the poll loop in-process into a new `CatalogSyncManager` (a `setTimeout` self-rescheduling loop that survives fetch errors), and replaces the sequential `for`-loop in `initProviders` with a bounded-concurrency `PromisePool`. - -**Tech Stack:** TypeScript, tsyringe DI, node-config, `@map-colonies/csw-client`, `@supercharge/promise-pool` (already a dependency on this branch), Jest + ts-jest. - -**Base branch:** `feat/geotiff-heights-migration` (PR #52) — NOT `master`. That PR rewrites `DEMTerrainCacheManager`, renames `initTerrainProviders`→`initProviders`, and touches `workerCatalogRecords.ts` + `containerConfig.ts`; branching off master would collide head-on and refactor soon-to-be-deleted Cesium code. - ---- - -## Prerequisite: branch off #52 - -- [ ] **Step 0: Create the working branch from the #52 head** - -```bash -git fetch origin feat/geotiff-heights-migration -git switch -c refactor/catalog-sync-inprocess origin/feat/geotiff-heights-migration -git log --oneline -1 # expect: tip of feat/geotiff-heights-migration -``` - ---- - -## File Structure - -- **Create:** `src/heights/models/catalogSyncManager.ts` — in-process CSW poll loop; owns the `CswClientWrapper`, the fetch filter, the `isSame` diff, and `setTimeout` scheduling. Injects only `CONFIG` + `LOGGER`; receives the two cache singletons via `start(...)` to avoid a circular import with `containerConfig`. -- **Create:** `tests/unit/heights/models/catalogSyncManager.spec.ts` — unit tests for the manager (update, no-op-when-unchanged, error-survives, stop()). -- **Modify:** `src/containerConfig.ts` — delete the worker (`initCSWWorker`, `Worker`, `path`, `WorkerEvent`); register `CATALOG_SYNC_MANAGER`; start it after registration; stop it in `onSignal`. -- **Modify:** `src/heights/models/DEMTerrainCacheManager.ts` — replace the sequential `for` loop in `initProviders` with a `PromisePool` bounded by `samplingConcurrency`; keep per-record error isolation. -- **Modify:** `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts` — keep the isolation test; add a "registers every record" test. -- **Modify:** `tests/integration/heights/heights.spec.ts` — override `CATALOG_SYNC_MANAGER` with a no-op stub so the real poll loop never runs during integration tests. -- **Delete:** `src/workerCatalogRecords.ts` — its `getCatalogRecords`, CSW client, `START_RECORD`/`END_RECORD`, and filter move into `CatalogSyncManager`; `WorkerEvent` is deleted with the worker. - -No config keys are added — provider-open concurrency reuses the existing `samplingConcurrency` key. No build-script or Dockerfile change (removing the worker only removes the `./workerCatalogRecords.js` runtime dependency). - ---- - -## Task 1: Parallelize `initProviders` (isolated, low-risk — do first) - -**Files:** - -- Modify: `src/heights/models/DEMTerrainCacheManager.ts` -- Test: `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts` - -- [ ] **Step 1: Add the "registers every record" failing test** - -Append this test inside the existing `describe('DEMTerrainCacheManager', ...)` block in `tests/unit/heights/models/DEMTerrainCacheManager.spec.ts`: - -```ts -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()).toEqual(['a', 'b', 'c']); -}); -``` - -- [ ] **Step 2: Run the test to verify it PASSES against the current sequential loop** - -Run: `npm run test:unit -- --testPathPattern DEMTerrainCacheManager` -Expected: PASS. (This test is a behavior-preservation guard — it must stay green through the refactor. The isolation test is the existing safety net.) - -- [ ] **Step 3: Replace the sequential loop with a bounded `PromisePool`** - -In `src/heights/models/DEMTerrainCacheManager.ts`, add the import near the other imports: - -```ts -import PromisePool from '@supercharge/promise-pool/dist'; -``` - -Replace the entire `initProviders` method body with: - -```ts - public async initProviders(demCatalogRecords: PycswDemCatalogRecord[]): Promise { - const heightProviders: HeightProviders = {}; - - const geotiffRecords = demCatalogRecords.filter((record) => record.links?.some((link) => link.protocol === GEOTIFF_PROTOCOL)); - const samplingConcurrency = Number(this.config.get('samplingConcurrency')); - - // 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; - } -``` - -Note: `PromisePool.handleError` that returns (does not throw) skips the failed item and lets the pool continue — this preserves the existing "skip a record whose provider fails to open" contract. `transformRouteToObjectUrl`, `buildAuthenticatedUrl`, `GEOTIFF_PROTOCOL`, and `HeightProviders` are unchanged and already present in the file. - -- [ ] **Step 4: Run both DEMTerrainCacheManager tests to verify they PASS** - -Run: `npm run test:unit -- --testPathPattern DEMTerrainCacheManager` -Expected: PASS — both the isolation test (skips `bad`, keeps `good`, logs error) and the "registers every record" test. - -- [ ] **Step 5: Commit** - -```bash -git add src/heights/models/DEMTerrainCacheManager.ts tests/unit/heights/models/DEMTerrainCacheManager.spec.ts -git commit -m "perf: open geotiff providers concurrently with bounded PromisePool" -``` - ---- - -## Task 2: Create `CatalogSyncManager` with unit tests - -**Files:** - -- Create: `src/heights/models/catalogSyncManager.ts` -- Test: `tests/unit/heights/models/catalogSyncManager.spec.ts` - -- [ ] **Step 1: Write the failing unit test** - -Create `tests/unit/heights/models/catalogSyncManager.spec.ts`: - -```ts -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(): DEMTerrainCacheManager { - return { initProviders: jest.fn().mockResolvedValue(undefined) } as unknown as DEMTerrainCacheManager; -} - -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 = 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(cacheManager.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 = makeCacheManager(); - - manager = new CatalogSyncManager(config, jsLogger({ enabled: false })); - manager.start(catalogRecords, cacheManager); - await new Promise((resolve) => setImmediate(resolve)); - - expect(cacheManager.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 = makeCacheManager(); - - manager = new CatalogSyncManager(config, logger); - manager.start(new CatalogRecords(), cacheManager); - await new Promise((resolve) => setImmediate(resolve)); - - expect(errorSpy).toHaveBeenCalled(); - expect(cacheManager.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()); - 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()); - await new Promise((resolve) => setImmediate(resolve)); - manager.stop(); - - expect(clearTimeoutSpy).toHaveBeenCalled(); - }); -}); -``` - -> **Jest 28 note:** this repo pins `jest@^28`, which has no `jest.advanceTimersByTimeAsync` (Jest 29+). The reschedule/stop behavior is therefore verified with real timers + `setImmediate` flushing (same mechanism as the first three tests) plus `setTimeout`/`clearTimeout` spies, rather than by advancing fake timers through async cycles. `afterEach` calls `manager?.stop()`, which clears the real pending timer so no handle leaks. - -- [ ] **Step 2: Run the test to verify it FAILS** - -Run: `npm run test:unit -- --testPathPattern catalogSyncManager` -Expected: FAIL — `Cannot find module '.../catalogSyncManager'` (file not yet created). - -- [ ] **Step 3: Create the `CatalogSyncManager` implementation** - -Create `src/heights/models/catalogSyncManager.ts`: - -```ts -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; - -@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(records, Object.values(catalogRecords.getValue()))) { - catalogRecords.setValue(Object.fromEntries(records.map((record) => [record.id as string, record]))); - await cacheManager.initProviders(records); - 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, - }); - } -} -``` - -- [ ] **Step 4: Run the test to verify it PASSES** - -Run: `npm run test:unit -- --testPathPattern catalogSyncManager` -Expected: PASS — all five cases (update, no-op-unchanged, error-survives, schedules-next-run, stop-clears-timer). - -- [ ] **Step 5: Commit** - -```bash -git add src/heights/models/catalogSyncManager.ts tests/unit/heights/models/catalogSyncManager.spec.ts -git commit -m "feat: add in-process CatalogSyncManager poll loop" -``` - ---- - -## Task 3: Wire `CatalogSyncManager` into bootstrap; remove the worker - -**Files:** - -- Modify: `src/containerConfig.ts` -- Modify: `tests/integration/heights/heights.spec.ts` -- Delete: `src/workerCatalogRecords.ts` - -- [ ] **Step 1: Stub the sync manager in the integration override (write first — this is the guard)** - -In `tests/integration/heights/heights.spec.ts`, add the import alongside the existing container imports: - -```ts -import { CATALOG_RECORDS_MAP, CATALOG_SYNC_MANAGER, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig'; -``` - -(That replaces the existing `import { CATALOG_RECORDS_MAP, DEM_TERRAIN_CACHE_MANAGER } from '../../../src/containerConfig';` line — add `CATALOG_SYNC_MANAGER` to it.) - -Then extend the `getApp` override in `beforeAll` so the real poll loop never runs in tests: - -```ts -const app = await getApp({ - override: [ - { token: SERVICES.LOGGER, provider: { useValue: jsLogger({ enabled: false }) } }, - { token: CATALOG_SYNC_MANAGER, provider: { useValue: { start: (): void => undefined, stop: (): void => undefined } } }, - ], -}); -``` - -- [ ] **Step 2: Rewrite `src/containerConfig.ts` to remove the worker and start the manager** - -Replace the entire contents of `src/containerConfig.ts` with: - -```ts -import config from 'config'; -import pino from 'pino'; -import client from 'prom-client'; -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 { getOtelMixin } from '@map-colonies/telemetry'; -import { SERVICES, SERVICE_NAME } from './common/constants'; -import { InjectionObject, registerDependencies } from './common/dependencyRegistration'; -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 { CatalogSyncManager } from './heights/models/catalogSyncManager'; - -export interface RegisterOptions { - override?: InjectionObject[]; - useChild?: boolean; -} - -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 productMetadataFields = config.get('productMetadataFields').split(','); - - tracing.start(); - const tracer = trace.getTracer(SERVICE_NAME); - - const dependencies: InjectionObject[] = [ - { token: SERVICES.CONFIG, provider: { useValue: config } }, - { token: SERVICES.LOGGER, provider: { useValue: logger } }, - { token: SERVICES.TRACER, provider: { useValue: tracer } }, - { - token: SERVICES.METRICS_REGISTRY, - provider: { - useFactory: instanceCachingFactory((container) => { - const config = container.resolve(SERVICES.CONFIG); - if (config.get('telemetry.metrics.enabled')) { - client.register.setDefaultLabels({ - app: SERVICE_NAME, - }); - return client.register; - } - }), - }, - }, - { 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()]); - }, - }, - }, - }, - ]; - - 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 registeredContainer; -}; -``` - -What changed vs. the #52 version: removed `import { Worker } from 'worker_threads'`, `import path`, the `PycswDemCatalogRecord` import, `import { isSame }`, and `import { WorkerEvent }`; deleted the entire `initCSWWorker` function and its call; added the `CatalogSyncManager` import, the `CATALOG_SYNC_MANAGER` symbol, its DI registration, the `.start(...)` call after registration, and the `.stop()` call in `onSignal`. The `container` import is retained (used by `onSignal` and the metrics factory). - -- [ ] **Step 3: Delete the worker source file** - -```bash -git rm src/workerCatalogRecords.ts -``` - -- [ ] **Step 4: Verify nothing else references the worker** - -Run: `git grep -n "workerCatalogRecords\|WorkerEvent\|initCSWWorker\|worker_threads" -- src tests` -Expected: no output (empty). If anything prints, fix that reference before proceeding. - -- [ ] **Step 5: Build to confirm no dangling imports / type errors** - -Run: `npm run build` -Expected: clean exit, no TypeScript errors. - -- [ ] **Step 6: Commit** - -```bash -git add src/containerConfig.ts tests/integration/heights/heights.spec.ts -git commit -m "refactor: replace CSW worker thread with in-process CatalogSyncManager" -``` - ---- - -## Task 4: Full verification - -**Files:** none (verification only) - -- [ ] **Step 1: Run the full unit suite with coverage** - -Run: `npm run test:unit` -Expected: all suites PASS, coverage gate met (branches ≥60, functions ≥80, lines ≥80). New `catalogSyncManager.ts` is under `src/heights/models/` (in the unit coverage set) and is covered by Task 2's tests; `containerConfig.ts` is excluded from unit coverage by the `!/src/*` glob. - -- [ ] **Step 2: Run the full integration suite** - -Run: `npm run test:integration` -Expected: 8/8 PASS, coverage gate met, and — unlike before — no `CatalogRecords ERROR` / worker-exit noise in the log (the sync manager is stubbed via the override; no real CSW fetch, no leaked timers, Jest exits cleanly). - -- [ ] **Step 3: Lint** - -Run: `npm run lint` -Expected: clean. - -- [ ] **Step 4: Confirm the worker is fully gone from the build output** - -Run: `ls dist/workerCatalogRecords.js 2>/dev/null && echo "STILL PRESENT — investigate" || echo "removed OK"` -Expected: `removed OK` (a stale `dist` from a previous build is fine to ignore; a fresh `npm run build` in Task 3 Step 5 will not emit it). - -- [ ] **Step 5: Push and open the PR against the #52 branch** - -```bash -git push -u origin refactor/catalog-sync-inprocess -gh pr create --base feat/geotiff-heights-migration --repo MapColonies/dem-heights \ - --title "refactor: in-process catalog sync + concurrent provider opens (MAPCO-11560)" \ - --body "Replaces the CSW worker thread with an in-process CatalogSyncManager poll loop and parallelizes GeoTIFF provider opening. Stacked on #52." -``` - -Note: base is `feat/geotiff-heights-migration`, so this PR merges into #52, not master. If #52 merges to master first, rebase this branch onto master before merging. - ---- - -## Self-Review - -**Spec coverage:** - -- Replace worker with in-process poll → Tasks 2 (manager) + 3 (wiring/removal). ✅ -- Fix silent-death → inherent: the `setTimeout` self-reschedule in `syncOnce`'s `finally` survives fetch errors; there is no separate process to die (Task 2, error test proves survival). ✅ -- Parallelize provider opens → Task 1. ✅ -- No API/behavior change → `/points`, response shape, `openapi3.yaml` untouched; startup remains non-blocking (loop kicked off, not awaited), matching prior behavior. ✅ -- No test regressions / no leaked timers → Task 3 Step 1 stub + Task 4. ✅ - -**Placeholder scan:** No TBD/TODO/"handle errors appropriately"; every code step contains full code. ✅ - -**Type consistency:** `initProviders` (not `initTerrainProviders`); `heightProviders` property; `HeightProviders` type; `start(catalogRecords, cacheManager)` signature matches both the `containerConfig` caller and the unit-test caller; `CATALOG_SYNC_MANAGER` symbol name identical across `containerConfig`, `onSignal`, and the integration override. ✅ - -**Decisions locked:** - -- Manager injects only CONFIG + LOGGER and receives singletons via `start(...)` — deliberately avoids importing the DI symbols from `containerConfig` (which imports the manager), preventing a circular-import decoration hazard. -- Provider-open concurrency reuses `samplingConcurrency` rather than adding a new config key (records are few; avoids config sprawl across `default.json` / `custom-environment-variables.json` / helm).