Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
75 changes: 14 additions & 61 deletions src/containerConfig.ts
Original file line number Diff line number Diff line change
@@ -1,24 +1,19 @@
import { Worker } from 'worker_threads';
import path from 'path';
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 { 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';
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<unknown>[];
Expand All @@ -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<DependencyContainer> => {
const loggerConfig = config.get<LoggerOptions>('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<CatalogRecords>(CATALOG_RECORDS_MAP);
demTerrainCacheManager = container.resolve<DEMTerrainCacheManager>(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<string>('productMetadataFields').split(',');

tracing.start();
Expand All @@ -114,18 +56,29 @@ 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<void> => {
container.resolve<CatalogSyncManager>(CATALOG_SYNC_MANAGER).stop();
await Promise.all([tracing.stop()]);
},
},
},
},
];

return Promise.resolve(registerDependencies(dependencies, options?.override, options?.useChild));
const registeredContainer = registerDependencies(dependencies, options?.override, options?.useChild);

registeredContainer
.resolve<CatalogSyncManager>(CATALOG_SYNC_MANAGER)
.start(
registeredContainer.resolve<CatalogRecords>(CATALOG_RECORDS_MAP),
registeredContainer.resolve<DEMTerrainCacheManager>(DEM_TERRAIN_CACHE_MANAGER)
);

return Promise.resolve(registeredContainer);
};
29 changes: 16 additions & 13 deletions src/heights/models/DEMTerrainCacheManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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<number>('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;
}
Expand Down
95 changes: 95 additions & 0 deletions src/heights/models/catalogSyncManager.ts
Original file line number Diff line number Diff line change
@@ -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<IService>('csw')
);
this.intervalMs = this.config.get<number>('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<void> {
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<PycswDemCatalogRecord[]> {
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,
});
}
}
60 changes: 0 additions & 60 deletions src/workerCatalogRecords.ts

This file was deleted.

7 changes: 5 additions & 2 deletions tests/integration/heights/heights.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -40,7 +40,10 @@ describe('heights', function () {
productMetadataFields = config.get<string>('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);
Expand Down
15 changes: 15 additions & 0 deletions tests/unit/heights/models/DEMTerrainCacheManager.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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']);
});
});
Loading
Loading