From ee9683779f7d3a7873524bec90d3a0ec5df23c1b Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 12:20:09 +0300 Subject: [PATCH 01/15] feat: add range-capped tile batch generator (MAPCO-11265) --- package-lock.json | 6 +-- package.json | 2 +- src/utils/tileRangeBatcher.ts | 37 +++++++++++++ tests/unit/utils/tileRangeBatcher.spec.ts | 63 +++++++++++++++++++++++ 4 files changed, 104 insertions(+), 4 deletions(-) create mode 100644 src/utils/tileRangeBatcher.ts create mode 100644 tests/unit/utils/tileRangeBatcher.spec.ts diff --git a/package-lock.json b/package-lock.json index 66210eb..106629a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -20,7 +20,7 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-model-types": "^17.15.1", "@map-colonies/mc-priority-queue": "^9.1.2", - "@map-colonies/mc-utils": "^6.0.1", + "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", "@map-colonies/prometheus": "^1.0.0", "@map-colonies/raster-shared": "^9.0.0-alpha.0", "@map-colonies/read-pkg": "^1.0.0", @@ -4802,8 +4802,8 @@ }, "node_modules/@map-colonies/mc-utils": { "version": "6.0.1", - "resolved": "https://registry.npmjs.org/@map-colonies/mc-utils/-/mc-utils-6.0.1.tgz", - "integrity": "sha512-aHPeigSA/rl88LzXiah9CVd2gzVTXV6zlk2rt34rjwqgLpBa3vBadwOgBun/w9/1jOXrS/hGz9FhMKmQL7qYGQ==", + "resolved": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", + "integrity": "sha512-MKHh65bmqkW8Mc6Ky3Y3tKnV435gdJZiOIkeYjDB8vfsVowPsP4Rz2xq671zny/nzdKRrcWFxH2uXOqV3HnHFQ==", "license": "ISC", "dependencies": { "@map-colonies/types": "^1.9.0", diff --git a/package.json b/package.json index 4f6fc7e..aa9abf8 100644 --- a/package.json +++ b/package.json @@ -46,7 +46,7 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-model-types": "^17.15.1", "@map-colonies/mc-priority-queue": "^9.1.2", - "@map-colonies/mc-utils": "^6.0.1", + "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", "@map-colonies/prometheus": "^1.0.0", "@map-colonies/raster-shared": "^9.0.0-alpha.0", "@map-colonies/read-pkg": "^1.0.0", diff --git a/src/utils/tileRangeBatcher.ts b/src/utils/tileRangeBatcher.ts new file mode 100644 index 0000000..5fc2fe9 --- /dev/null +++ b/src/utils/tileRangeBatcher.ts @@ -0,0 +1,37 @@ +import { tileBatchGenerator } from '@map-colonies/mc-utils'; +import type { ITileRange } from '@map-colonies/mc-utils'; + +/** Adapts a sync generator to the async generator `tileBatchGenerator` expects. */ +// eslint-disable-next-line @typescript-eslint/require-await -- a pass-through async generator has nothing to await +async function* toAsyncGenerator(source: Generator): AsyncGenerator { + for (const item of source) { + yield item; + } +} + +/** + * Batches tile ranges bounded by tile count and range count per batch. + * + * Caps both tile count (via `tileBatchGenerator`) and ITileRange count, since complex + * footprints can yield thousands of ranges that exceed task params size limits. + * + * @param tileBatchSize maximum tiles per batch + * @param maxRangesPerTask maximum ITileRange objects per batch + * @param ranges sync generator of ranges, e.g. from `footprintToTileRanges` + */ +export async function* limitedTileBatchGenerator( + tileBatchSize: number, + maxRangesPerTask: number, + ranges: Generator +): AsyncGenerator { + if (tileBatchSize < 1 || maxRangesPerTask < 1) { + throw new RangeError(`tileBatchSize [${tileBatchSize}] and maxRangesPerTask [${maxRangesPerTask}] must both be at least 1`); + } + + for await (const batch of tileBatchGenerator(tileBatchSize, toAsyncGenerator(ranges))) { + // slicing preserves every range and their order; it only spreads them over more tasks + for (let index = 0; index < batch.length; index += maxRangesPerTask) { + yield batch.slice(index, index + maxRangesPerTask); + } + } +} diff --git a/tests/unit/utils/tileRangeBatcher.spec.ts b/tests/unit/utils/tileRangeBatcher.spec.ts new file mode 100644 index 0000000..11cde32 --- /dev/null +++ b/tests/unit/utils/tileRangeBatcher.spec.ts @@ -0,0 +1,63 @@ +import type { ITileRange } from '@map-colonies/mc-utils'; +import { limitedTileBatchGenerator } from '../../../src/utils/tileRangeBatcher'; + +/** One-tile range at z21, distinguishable by x so we can assert nothing is dropped. */ +const singleTile = (x: number): ITileRange => ({ zoom: 21, minX: x, maxX: x, minY: 0, maxY: 0 }); + +function* ranges(...items: ITileRange[]): Generator { + yield* items; +} + +const collect = async (source: AsyncGenerator): Promise => { + const batches: ITileRange[][] = []; + for await (const batch of source) { + batches.push(batch); + } + return batches; +}; + +describe('cappedTileBatchGenerator', () => { + it('should yield a single batch when the ranges fit both limits', async () => { + const batches = await collect(limitedTileBatchGenerator(100, 100, ranges(singleTile(0)))); + + expect(batches).toEqual([[singleTile(0)]]); + }); + + it('should split a batch that exceeds maxRangesPerTask, preserving every range in order', async () => { + const input = Array.from({ length: 10 }, (_, index) => singleTile(index)); + // `tileBatchGenerator` walks `range.minY` forward in place, so `input` is mutated as it is + // consumed while the yielded ranges are fresh objects. Build the expectation separately - + // asserting against `input` after the fact compares against the mutated originals and fails. + const expected = Array.from({ length: 10 }, (_, index) => singleTile(index)); + + // 10 tiles is well under the 1000-tile budget, so tileBatchGenerator produces one batch of + // 10 ranges; only the range cap splits it. + const batches = await collect(limitedTileBatchGenerator(1000, 4, ranges(...input))); + + expect(batches.map((batch) => batch.length)).toEqual([4, 4, 2]); + expect(batches.flat()).toEqual(expected); + }); + + it('should still split by tile count when the range cap is not reached', async () => { + // one 10-wide row, cut into 4 + 4 + 2 tiles by the tile budget + const batches = await collect(limitedTileBatchGenerator(4, 100, ranges({ zoom: 21, minX: 0, maxX: 9, minY: 0, maxY: 0 }))); + + expect(batches).toEqual([ + [{ zoom: 21, minX: 0, maxX: 3, minY: 0, maxY: 0 }], + [{ zoom: 21, minX: 4, maxX: 7, minY: 0, maxY: 0 }], + [{ zoom: 21, minX: 8, maxX: 9, minY: 0, maxY: 0 }], + ]); + }); + + it('should yield nothing for an empty range generator', async () => { + const batches = await collect(limitedTileBatchGenerator(100, 100, ranges())); + + expect(batches).toEqual([]); + }); + + it('should reject a maxRangesPerTask below 1, which would loop forever', async () => { + const action = collect(limitedTileBatchGenerator(100, 0, ranges(singleTile(0)))); + + await expect(action).rejects.toThrow(RangeError); + }); +}); From 7447389db34310d7389bf00d8aa5da09d36ce00a Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 14:47:31 +0300 Subject: [PATCH 02/15] feat: add cache-deletion job and task configuration (MAPCO-11265) Adds the Update_Delete_Cache / Swap_Delete_Cache job types and the cacheDeletion task config, wired through to the chart in the same change. Pins raster-shared to the build exporting GEODETIC_GRIDS. --- config/custom-environment-variables.json | 16 +++++++++ config/default.json | 16 +++++++++ helm/templates/configmap.yaml | 10 ++++++ helm/values.yaml | 12 +++++++ package-lock.json | 6 ++-- package.json | 2 +- src/common/errors.ts | 7 ++++ src/common/interfaces.ts | 41 ++++++++++++++++++++++++ tests/unit/mocks/configMock.ts | 16 +++++++++ 9 files changed, 122 insertions(+), 4 deletions(-) diff --git a/config/custom-environment-variables.json b/config/custom-environment-variables.json index 437cf0d..8d778c8 100644 --- a/config/custom-environment-variables.json +++ b/config/custom-environment-variables.json @@ -148,6 +148,12 @@ "jobs": { "seed": { "type": "INGESTION_SEED_JOB_TYPE" + }, + "updateCacheDeletion": { + "type": "UPDATE_CACHE_DELETION_JOB_TYPE" + }, + "swapCacheDeletion": { + "type": "SWAP_CACHE_DELETION_JOB_TYPE" } }, "tasks": { @@ -162,6 +168,16 @@ "__format": "number" } }, + "cacheDeletion": { + "type": "CACHE_DELETION_TASK_TYPE", + "grid": "CACHE_DELETION_GRID", + "maxZoom": { "__name": "CACHE_DELETION_MAX_ZOOM", "__format": "number" }, + "tileBatchSize": { "__name": "CACHE_DELETION_TILE_BATCH_SIZE", "__format": "number" }, + "maxRangesPerTask": { "__name": "CACHE_DELETION_MAX_RANGES_PER_TASK", "__format": "number" }, + "taskBatchSize": { "__name": "CACHE_DELETION_TASK_BATCH_SIZE", "__format": "number" }, + "gracefulReloadMaxSeconds": { "__name": "CACHE_DELETION_GRACEFUL_RELOAD_MAX_SECONDS", "__format": "number" }, + "reloadWindowMarginSeconds": { "__name": "CACHE_DELETION_RELOAD_WINDOW_MARGIN_SECONDS", "__format": "number" } + }, "tilesMerging": { "type": "TILES_MERGING_TASK_TYPE", "tileBatchSize": { diff --git a/config/default.json b/config/default.json index 4441ffd..d6236d4 100644 --- a/config/default.json +++ b/config/default.json @@ -115,6 +115,12 @@ "jobs": { "seed": { "type": "Ingestion_Seed" + }, + "updateCacheDeletion": { + "type": "Update_Delete_Cache" + }, + "swapCacheDeletion": { + "type": "Swap_Delete_Cache" } }, "tasks": { @@ -143,6 +149,16 @@ "type": "tiles-deletion", "tileBatchSize": 10000, "taskBatchSize": 5 + }, + "cacheDeletion": { + "type": "tiles-deletion", + "grid": "WorldCRS84", + "maxZoom": 21, + "tileBatchSize": 100000, + "maxRangesPerTask": 5000, + "taskBatchSize": 5, + "gracefulReloadMaxSeconds": 300, + "reloadWindowMarginSeconds": 8 } } }, diff --git a/helm/templates/configmap.yaml b/helm/templates/configmap.yaml index a2a9972..f379ee4 100644 --- a/helm/templates/configmap.yaml +++ b/helm/templates/configmap.yaml @@ -58,6 +58,8 @@ data: INGESTION_SWAP_UPDATE_JOB_TYPE: {{ $jobDefinitions.jobs.swapUpdate.type | quote }} INGESTION_DELETE_LAYER_JOB_TYPE: {{ $jobDefinitions.jobs.deleteLayer.type | quote }} INGESTION_SEED_JOB_TYPE : {{ $jobDefinitions.jobs.seed.type | quote }} + UPDATE_CACHE_DELETION_JOB_TYPE: {{ $jobDefinitions.jobs.updateCacheDeletion.type | quote }} + SWAP_CACHE_DELETION_JOB_TYPE: {{ $jobDefinitions.jobs.swapCacheDeletion.type | quote }} EXPORT_JOB_TYPE: {{ $jobDefinitions.jobs.export.type | quote }} EXPORT_CLEANUP_EXPIRATION_DAYS: {{ $jobDefinitions.jobs.export.cleanupExpirationDays | quote }} TILES_DELETION_TASK_TYPE: {{ $jobDefinitions.tasks.tilesDeletion.type | quote }} @@ -82,6 +84,14 @@ data: TILES_SEEDING_ZOOM_THRESHOLD: {{ $jobDefinitions.tasks.seed.zoomThreshold | quote }} TILES_SEEDING_MAX_TILES_PER_SEED_TASK: {{ $jobDefinitions.tasks.seed.maxTilesPerSeedTask | quote }} TILES_SEEDING_MAX_TILES_PER_CLEAN_TASK: {{ $jobDefinitions.tasks.seed.maxTilesPerCleanTask | quote }} + CACHE_DELETION_TASK_TYPE: {{ $jobDefinitions.tasks.cacheDeletion.type | quote }} + CACHE_DELETION_GRID: {{ $jobDefinitions.tasks.cacheDeletion.grid | quote }} + CACHE_DELETION_MAX_ZOOM: {{ $jobDefinitions.tasks.cacheDeletion.maxZoom | quote }} + CACHE_DELETION_TILE_BATCH_SIZE: {{ $jobDefinitions.tasks.cacheDeletion.tileBatchSize | quote }} + CACHE_DELETION_MAX_RANGES_PER_TASK: {{ $jobDefinitions.tasks.cacheDeletion.maxRangesPerTask | quote }} + CACHE_DELETION_TASK_BATCH_SIZE: {{ $jobDefinitions.tasks.cacheDeletion.taskBatchSize | quote }} + CACHE_DELETION_GRACEFUL_RELOAD_MAX_SECONDS: {{ .Values.global.gracefulReloadMaxSeconds | quote }} + CACHE_DELETION_RELOAD_WINDOW_MARGIN_SECONDS: {{ $jobDefinitions.tasks.cacheDeletion.reloadWindowMarginSeconds | quote }} TILES_EXPORTING_TASK_TYPE: {{ $jobDefinitions.tasks.export.type | quote }} MAPPROXY_API_URL: {{ $serviceUrls.mapproxyApi | quote }} GEOSERVER_API_URL: {{ $serviceUrls.geoserverApiUrl | quote }} diff --git a/helm/values.yaml b/helm/values.yaml index 3f12081..f0a1ee2 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -143,6 +143,10 @@ jobDefinitions: type: "" seed: type: "" + updateCacheDeletion: + type: "" + swapCacheDeletion: + type: "" export: type: "" cleanupExpirationDays: 14 @@ -178,6 +182,14 @@ jobDefinitions: type: "" tileBatchSize: 10000 taskBatchSize: 5 + cacheDeletion: + type: "" + grid: "WorldCRS84" # must be a geodetic grid - footprintToTileRanges supports no other + maxZoom: 21 + tileBatchSize: 100000 # max tiles per range-deletion task + maxRangesPerTask: 5000 # bounds the serialized task params for jagged footprints + taskBatchSize: 5 + reloadWindowMarginSeconds: 8 # added to global.gracefulReloadMaxSeconds: mapproxinator's poll plus drain config: dequeueIntervalMs: 3000 diff --git a/package-lock.json b/package-lock.json index 106629a..26336ab 100644 --- a/package-lock.json +++ b/package-lock.json @@ -22,7 +22,7 @@ "@map-colonies/mc-priority-queue": "^9.1.2", "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "^9.0.0-alpha.0", + "@map-colonies/raster-shared": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.18.0", "@map-colonies/shapefile-reader": "^1.0.1", @@ -6594,8 +6594,8 @@ }, "node_modules/@map-colonies/raster-shared": { "version": "9.0.0-alpha.0", - "resolved": "https://registry.npmjs.org/@map-colonies/raster-shared/-/raster-shared-9.0.0-alpha.0.tgz", - "integrity": "sha512-7NUVOMnlGaeSDYZu84wEXI25SyL/DatD/LHVBiGVP+CYWZ/8KjEe9WMH84EIPZS4VAlNpRbNZ770582AJNRjxQ==", + "resolved": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", + "integrity": "sha512-GZyFW9Dk8lD87qbWlGMFU7KxTTNMrr2fiOwnes7Wjy+ePdFU2maKUoWfDBitCEGderkSIz4iOLRD+yC8MS/M3Q==", "license": "ISC", "dependencies": { "@map-colonies/mc-priority-queue": "^9.1.0", diff --git a/package.json b/package.json index aa9abf8..0cf9bb7 100644 --- a/package.json +++ b/package.json @@ -48,7 +48,7 @@ "@map-colonies/mc-priority-queue": "^9.1.2", "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "^9.0.0-alpha.0", + "@map-colonies/raster-shared": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.18.0", "@map-colonies/shapefile-reader": "^1.0.1", diff --git a/src/common/errors.ts b/src/common/errors.ts index 27f5af7..783dbcd 100644 --- a/src/common/errors.ts +++ b/src/common/errors.ts @@ -80,6 +80,13 @@ export class LayerCacheNotFoundError extends Error { } } +export class UnsupportedGridError extends Error { + public constructor(grid: string, supported: readonly string[]) { + super(`Unsupported grid (${grid}) for cache deletion, supported grids: ${supported.join(', ')}`); + this.name = UnsupportedGridError.name; + } +} + export class SeedJobCreationError extends Error { public constructor(msg: string, err: Error) { super(msg); diff --git a/src/common/interfaces.ts b/src/common/interfaces.ts index 3eba00b..22ae15a 100644 --- a/src/common/interfaces.ts +++ b/src/common/interfaces.ts @@ -10,6 +10,8 @@ import type { InputFiles, LayerName, RasterLayerMetadata, + RedisDeleteStoredResourcesParams, + RedisTilesDeletionParams, TileFormatStrategy, TileOutputFormat, } from '@map-colonies/raster-shared'; @@ -60,6 +62,10 @@ export interface JobConfig { export interface IngestionJobsConfig { seed: JobConfig | undefined; + /** job type for the update flow's cache deletion; selects cleaner's range-deletion strategy */ + updateCacheDeletion: JobConfig | undefined; + /** job type for the swap flow's cache deletion; selects cleaner's prefix-wipe strategy */ + swapCacheDeletion: JobConfig | undefined; } export interface IngestionPollingJobsConfig { @@ -79,6 +85,7 @@ export interface IngestionTasksConfig { tilesMerging: TilesMergingTaskConfig; tilesSeeding: TilesSeedingTaskConfig; tilesDeletion: TilesDeletionTaskConfig; + cacheDeletion: CacheDeletionTaskConfig; } export interface ExportTasksConfig { @@ -122,6 +129,24 @@ export interface TilesSeedingTaskConfig { skipUncached: boolean; } +export interface CacheDeletionTaskConfig { + /** `tiles-deletion` - shared with the S3/FS path, told apart by job type */ + type: string; + /** must be a member of GEODETIC_GRIDS; the sole input to the composed `${cacheName}_${grid}` key prefix */ + grid: string; + maxZoom: number; + /** max tiles a single range-deletion task covers */ + tileBatchSize: number; + /** max ITileRange objects a single task carries, bounding the serialized params size */ + maxRangesPerTask: number; + /** tasks pushed per createTaskForJob call */ + taskBatchSize: number; + /** mirrors helm's global.gracefulReloadMaxSeconds */ + gracefulReloadMaxSeconds: number; + /** added on top, covering mapproxinator's <=5s poll plus pod drain */ + reloadWindowMarginSeconds: number; +} + export interface TilesExportingTaskConfig { type: string; } @@ -399,6 +424,22 @@ export interface SeedTaskParams { //#endregion seedingJobCreator +//#region cacheDeletionJobCreator + +export interface CacheDeletionJobParams { + layerName: LayerName; + ingestionJob: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob; +} + +/** + * Params of a `tiles-deletion` cache-deletion task. Both shapes share that task type; which one a + * task carries follows from the job type - Update_Delete_Cache means `ranges`, Swap_Delete_Cache + * means prefix-only - and that is also how cleaner picks the strategy to run. + */ +export type CacheDeletionTaskParams = RedisTilesDeletionParams | RedisDeleteStoredResourcesParams; + +//#endregion cacheDeletionJobCreator + //#region telemetry export interface TraceParentContext { traceparent?: string; diff --git a/tests/unit/mocks/configMock.ts b/tests/unit/mocks/configMock.ts index b3b6b3b..4cbe50e 100644 --- a/tests/unit/mocks/configMock.ts +++ b/tests/unit/mocks/configMock.ts @@ -159,6 +159,12 @@ const registerDefaultConfig = (): void => { seed: { type: 'Ingestion_Seed', }, + updateCacheDeletion: { + type: 'Update_Delete_Cache', + }, + swapCacheDeletion: { + type: 'Swap_Delete_Cache', + }, }, tasks: { tilesMerging: { @@ -179,6 +185,16 @@ const registerDefaultConfig = (): void => { tileBatchSize: 10000, taskBatchSize: 5, }, + cacheDeletion: { + type: 'tiles-deletion', + grid: 'WorldCRS84', + maxZoom: 21, + tileBatchSize: 100000, + maxRangesPerTask: 5000, + taskBatchSize: 5, + gracefulReloadMaxSeconds: 300, + reloadWindowMarginSeconds: 8, + }, }, }, export: { From ab2e16c021db9196f9e34f978c9d116d6dd173f8 Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 18:01:59 +0300 Subject: [PATCH 03/15] feat: add cache deletion job creator (MAPCO-11265) Composes the redis key prefix from the mapproxy cache name and the configured grid, emits a prefix-wipe task for swap-update and streamed range tasks for update, and fails a partially-populated job explicitly. --- .../ingestion/cacheDeletionJobCreator.ts | 242 ++++++++++++++ .../cacheDeletionJobCreator.spec.ts | 301 ++++++++++++++++++ .../cacheDeletionJobCreatorSetup.ts | 45 +++ tests/unit/mocks/geometryMockData.ts | 7 +- 4 files changed, 594 insertions(+), 1 deletion(-) create mode 100644 src/job/models/ingestion/cacheDeletionJobCreator.ts create mode 100644 tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts create mode 100644 tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts diff --git a/src/job/models/ingestion/cacheDeletionJobCreator.ts b/src/job/models/ingestion/cacheDeletionJobCreator.ts new file mode 100644 index 0000000..eb737a8 --- /dev/null +++ b/src/job/models/ingestion/cacheDeletionJobCreator.ts @@ -0,0 +1,242 @@ +import type { Logger } from '@map-colonies/js-logger'; +import { footprintToTileRanges } from '@map-colonies/mc-utils'; +import type { ICreateJobBody, ICreateTaskBody } from '@map-colonies/mc-priority-queue'; +import { OperationStatus, TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; +import type { LayerName } from '@map-colonies/raster-shared'; +import { GEODETIC_GRIDS, StorageProvider } from '@map-colonies/raster-shared'; +import { inject, injectable } from 'tsyringe'; +import { context, SpanStatusCode, trace, type Span, type Tracer } from '@opentelemetry/api'; +import { LayerCacheType, SERVICES } from '../../../common/constants'; +import { UnsupportedGridError } from '../../../common/errors'; +import type { CacheDeletionJobParams, CacheDeletionTaskConfig, CacheDeletionTaskParams, IConfig } from '../../../common/interfaces'; +import { MapproxyApiClient } from '../../../httpClients/mapproxyClient'; +import { internalIdSchema } from '../../../utils/zod/schemas/jobParameters.schema'; +import { IngestionSwapUpdateFinalizeJob, IngestionUpdateFinalizeJob } from '../../../utils/zod/schemas/job.schema'; +import { limitedTileBatchGenerator } from '../../../utils/tileRangeBatcher'; +import type { ReadProductGeometry } from '../../../utils/storage/productReader'; + +type IngestionFinalizeJob = IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob; +type CacheDeletionTask = ICreateTaskBody; + +/** + * Creates the job that deletes a layer's MapProxy redis tile cache after an ingestion. + * + * Two strategies: + * - swap-update: wipes entire prefix (layer redirects to new tiles path) + * - update: deletes specific tile ranges over ingested footprint across all zooms + */ +@injectable() +export class CacheDeletionJobCreator { + private readonly taskConfig: CacheDeletionTaskConfig; + /** the ingestion job type we are reacting to, used only to pick the flow */ + private readonly swapUpdateJobType: string; + /** cache-deletion job type per flow - cleaner picks its strategy from this */ + private readonly updateCacheDeletionJobType: string; + private readonly swapCacheDeletionJobType: string; + private readonly wipeDelaySeconds: number; + + public constructor( + @inject(SERVICES.LOGGER) private readonly logger: Logger, + @inject(SERVICES.TRACER) private readonly tracer: Tracer, + @inject(SERVICES.CONFIG) private readonly config: IConfig, + @inject(SERVICES.QUEUE_CLIENT) private readonly queueClient: QueueClient, + @inject(MapproxyApiClient) private readonly mapproxyClient: MapproxyApiClient, + @inject(SERVICES.PRODUCT_READER) private readonly readProductGeometry: ReadProductGeometry + ) { + this.taskConfig = this.config.get('jobManagement.ingestion.tasks.cacheDeletion'); + this.swapUpdateJobType = this.config.get('jobManagement.ingestion.pollingJobs.swapUpdate.type'); + this.updateCacheDeletionJobType = this.config.get('jobManagement.ingestion.jobs.updateCacheDeletion.type'); + this.swapCacheDeletionJobType = this.config.get('jobManagement.ingestion.jobs.swapCacheDeletion.type'); + + // Refuse to start rather than delete the wrong keys and report success. + if (!GEODETIC_GRIDS.includes(this.taskConfig.grid)) { + throw new UnsupportedGridError(this.taskConfig.grid, GEODETIC_GRIDS); + } + // Serving pods reload config on their own schedule (gracefulReloadMaxSeconds). + // Delay wiping cache keys until all pods have reloaded to prevent stale pods from re-caching + // old tiles under removed keys. + this.wipeDelaySeconds = this.taskConfig.gracefulReloadMaxSeconds + this.taskConfig.reloadWindowMarginSeconds; + } + + public async create({ layerName, ingestionJob }: CacheDeletionJobParams): Promise { + await context.with(trace.setSpan(context.active(), this.tracer.startSpan(`${CacheDeletionJobCreator.name}.${this.create.name}`)), async () => { + const activeSpan = trace.getActiveSpan(); + const isSwapUpdate = ingestionJob.type === this.swapUpdateJobType; + + // The job type is what tells cleaner which strategy to run, so it is chosen per flow rather + // than fixed at construction. + const jobType = isSwapUpdate ? this.swapCacheDeletionJobType : this.updateCacheDeletionJobType; + + const logger = this.logger.child({ + ingestionJobId: ingestionJob.id, + jobType, + taskType: this.taskConfig.type, + layerName, + }); + + try { + logger.info({ msg: 'Starting cache deletion job creation process' }); + activeSpan?.setAttributes({ + ingestionJobId: ingestionJob.id, + cacheDeletionJobType: jobType, + cacheDeletionShape: isSwapUpdate ? 'prefix-wipe' : 'range-deletion', + layerName, + }); + + const prefix = await this.resolvePrefix(layerName, logger); + const catalogId = internalIdSchema.parse(ingestionJob).internalId; + + const tasks = isSwapUpdate ? this.buildWipeTask(prefix) : this.buildRangeTasks(ingestionJob, prefix, logger); + + await this.createJobWithStreamedTasks(ingestionJob, jobType, catalogId, tasks, logger, activeSpan); + } catch (err) { + if (err instanceof Error) { + activeSpan?.recordException(err); + activeSpan?.setStatus({ code: SpanStatusCode.ERROR }); + logger.error({ msg: `Failed to create cache deletion job: ${err.message}`, err }); + } + } finally { + activeSpan?.end(); + } + }); + } + + /** + * Composes the redis key prefix in format `${cacheName}_${gridName}`, matching mapproxy's loader.py behavior. + * The grid comes from config since mapproxy-api does not report it. + */ + private async resolvePrefix(layerName: LayerName, logger: Logger): Promise { + const cacheName = await this.mapproxyClient.getRedisCacheName({ layerName, cacheType: LayerCacheType.REDIS }); + const prefix = `${cacheName}_${this.taskConfig.grid}`; + + logger.info({ msg: 'Composed redis key prefix', cacheName, grid: this.taskConfig.grid, prefix }); + trace.getActiveSpan()?.setAttributes({ cacheName, redisPrefix: prefix }); + + return prefix; + } + + // eslint-disable-next-line @typescript-eslint/require-await + private async *buildWipeTask(prefix: string): AsyncGenerator { + yield { + type: this.taskConfig.type, + description: 'redis cache prefix wipe', + parameters: { storageProvider: StorageProvider.REDIS, prefix, delaySeconds: this.wipeDelaySeconds }, + }; + } + + private async *buildRangeTasks(job: IngestionFinalizeJob, prefix: string, logger: Logger): AsyncGenerator { + const { grid, maxZoom, tileBatchSize, maxRangesPerTask } = this.taskConfig; + + const geometry = await this.readProductGeometry(job.parameters.inputFiles.productShapefilePath); + logger.info({ msg: 'Computing tile ranges over the updated footprint', grid, maxZoom, tileBatchSize, maxRangesPerTask }); + + const ranges = footprintToTileRanges(geometry, { minZoom: 0, maxZoom }); + + for await (const batch of limitedTileBatchGenerator(tileBatchSize, maxRangesPerTask, ranges)) { + yield { + type: this.taskConfig.type, + description: 'redis cache range deletion', + parameters: { storageProvider: StorageProvider.REDIS, prefix, ranges: batch }, + }; + } + } + + /** + * Tasks are streamed to avoid memory overhead. The job is created with its first batch + * so job-tracker never sees an empty job. + */ + private async createJobWithStreamedTasks( + job: IngestionFinalizeJob, + jobType: string, + catalogId: string, + tasks: AsyncGenerator, + logger: Logger, + activeSpan: Span | undefined + ): Promise { + const { taskBatchSize } = this.taskConfig; + let taskBatch: CacheDeletionTask[] = []; + let jobId: string | undefined; + let taskCount = 0; + + try { + for await (const task of tasks) { + taskBatch.push(task); + taskCount++; + + if (taskBatch.length === taskBatchSize) { + jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch); + taskBatch = []; + } + } + + if (taskBatch.length > 0) { + jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch); + } + } catch (err) { + if (jobId !== undefined) { + await this.failPartialJob(jobId, err, logger); + } + throw err; // create()'s catch logs and swallows, so ingestion is unaffected + } + + if (jobId === undefined) { + logger.warn({ msg: 'No cache deletion tasks produced, skipping job creation' }); + activeSpan?.addEvent('createJob.skipped', { reason: 'No tasks produced' }); + return; + } + + logger.info({ msg: 'Cache deletion job created successfully', cacheDeletionJobId: jobId, taskCount }); + activeSpan?.setAttributes({ cacheDeletionJobId: jobId, taskCount }); + } + + private async flushTaskBatch( + job: IngestionFinalizeJob, + jobType: string, + catalogId: string, + jobId: string | undefined, + tasks: CacheDeletionTask[] + ): Promise { + if (jobId !== undefined) { + await this.queueClient.jobManagerClient.createTaskForJob(jobId, tasks); + return jobId; + } + + const { resourceId, version, producerName, productName, productType, domain } = job; + const createJobRequest: ICreateJobBody = { + resourceId, + internalId: catalogId, + version, + type: jobType, + parameters: { ingestionJobId: job.id, ingestionJobType: job.type }, + status: OperationStatus.IN_PROGRESS, + producerName: producerName ?? undefined, + productName, + productType, + domain, + tasks, + }; + + const res = await this.queueClient.jobManagerClient.createJob(createJobRequest); + return res.id; + } + + /** + * A job created with its first batch but missing later ones would run its enqueued tasks to + * completion, and job-tracker would then report the whole cache deletion as successful. Fail it + * explicitly instead - silent partial success is the failure mode this creator exists to avoid. + */ + private async failPartialJob(jobId: string, cause: unknown, logger: Logger): Promise { + const reason = `cache deletion task enqueue failed: ${cause instanceof Error ? cause.message : String(cause)}`; + logger.error({ msg: 'Cache deletion job is missing tasks after an enqueue failure, marking it failed', cacheDeletionJobId: jobId, reason }); + + try { + await this.queueClient.jobManagerClient.updateJob(jobId, { status: OperationStatus.FAILED, reason }); + } catch (updateErr) { + logger.error({ + msg: 'Could not mark the partial cache deletion job as failed, it may be reported as successful', + cacheDeletionJobId: jobId, + err: updateErr, + }); + } + } +} diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts new file mode 100644 index 0000000..3820cd2 --- /dev/null +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -0,0 +1,301 @@ +/* eslint-disable @typescript-eslint/unbound-method */ +import { randomUUID } from 'node:crypto'; +import type { MultiPolygon, Polygon } from 'geojson'; +import { OperationStatus } from '@map-colonies/mc-priority-queue'; +import { StorageProvider } from '@map-colonies/raster-shared'; +import type { CacheDeletionJobParams, CacheDeletionTaskConfig } from '../../../../src/common/interfaces'; +import { registerDefaultConfig, configMock, setValue } from '../../mocks/configMock'; +import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; +import { LayerCacheType } from '../../../../src/common/constants'; +import { LayerCacheNotFoundError } from '../../../../src/common/errors'; +import { ingestionSwapUpdateFinalizeJob, ingestionUpdateFinalizeJob } from '../../mocks/jobsMockData'; +import type { CacheDeletionJobCreatorTestContext } from './cacheDeletionJobCreatorSetup'; +import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup'; + +describe('CacheDeletionJobCreator', () => { + let ctx: CacheDeletionJobCreatorTestContext; + let productGeometry: Polygon | MultiPolygon; + + beforeEach(async () => { + vi.resetAllMocks(); + registerDefaultConfig(); + ctx = await setupCacheDeletionJobCreatorTest(); + productGeometry = createFakePolygonalGeometry({ radiusInMeters: 50 }); + }); + + afterEach(() => { + vi.resetAllMocks(); + vi.restoreAllMocks(); + }); + + describe('swap-update', () => { + it('should create a Swap_Delete_Cache job with a single prefix-wipe task carrying the delay', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; + const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); + const jobType = configMock.get('jobManagement.ingestion.jobs.swapCacheDeletion.type'); + const jobId = randomUUID(); + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); + + const params: CacheDeletionJobParams = { layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }; + + await cacheDeletionJobCreator.create(params); + + // must ask for the REDIS cache explicitly - the layer's tiles cache is file or s3 + expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: 'layer-Orthophoto', cacheType: LayerCacheType.REDIS }); + expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); + expect(jobManagerClientMock.createTaskForJob).not.toHaveBeenCalled(); + + const request = jobManagerClientMock.createJob.mock.calls[0]![0]; + + expect(request).toMatchObject({ + resourceId: ingestionSwapUpdateFinalizeJob.resourceId, + internalId: ingestionSwapUpdateFinalizeJob.internalId, + version: ingestionSwapUpdateFinalizeJob.version, + type: jobType, + parameters: {}, + status: OperationStatus.IN_PROGRESS, + }); + // cleaner resolves DeleteStoredResourcesStrategy from this job type; the update job type + // would route the same params to the range strategy and fail validation + expect(request.type).toBe('Swap_Delete_Cache'); + expect(request.tasks).toStrictEqual([ + { + type: taskConfig.type, + description: 'redis cache prefix wipe', + parameters: { + storageProvider: StorageProvider.REDIS, + // no explicit prefix from mapproxy, so loader.py's fallback applies + prefix: `layer-Orthophoto-redis_${taskConfig.grid}`, + delaySeconds: taskConfig.gracefulReloadMaxSeconds + taskConfig.reloadWindowMarginSeconds, + }, + }, + ]); + }); + + it('should compose the prefix from the cache name and the configured grid', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; + const grid = configMock.get('jobManagement.ingestion.tasks.cacheDeletion').grid; + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('sss-Orthophoto-redis'); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'sss-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); + + const request = jobManagerClientMock.createJob.mock.calls[0]![0]; + + // mapproxy reports no prefix of its own, so the creator composes loader.py's shape + expect(request.tasks![0]!.parameters).toMatchObject({ prefix: `sss-Orthophoto-redis_${grid}` }); + expect(request.tasks![0]!.parameters).toMatchObject({ prefix: 'sss-Orthophoto-redis_WorldCRS84' }); + }); + + it('should not read the product shapefile, since a wipe needs no geometry', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); + + expect(readProductGeometryMock).not.toHaveBeenCalled(); + }); + }); + + describe('update', () => { + it('should create an Update_Delete_Cache job with range tasks and no delay', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; + const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }); + + expect(readProductGeometryMock).toHaveBeenCalledWith(ingestionUpdateFinalizeJob.parameters.inputFiles.productShapefilePath); + expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); + // cleaner resolves TilesDeletionStrategy from this job type + expect(jobManagerClientMock.createJob.mock.calls[0]![0].type).toBe('Update_Delete_Cache'); + + const tasks = jobManagerClientMock.createJob.mock.calls[0]![0].tasks!; + + expect(tasks.length).toBeGreaterThan(0); + + for (const task of tasks) { + expect(task.type).toBe(taskConfig.type); + + const { ranges } = task.parameters as { ranges: unknown[] }; + + expect(Array.isArray(ranges)).toBe(true); + expect(ranges.length).toBeLessThanOrEqual(taskConfig.maxRangesPerTask); + // toStrictEqual (rather than toMatchObject) pins the exact field set: no delaySeconds, and + // no other field the worker's .strict() params schema would reject at runtime. + expect(task.parameters).toStrictEqual({ + storageProvider: StorageProvider.REDIS, + prefix: `layer-Orthophoto-redis_${taskConfig.grid}`, + ranges, + }); + } + }); + + it('should cover every zoom from 0 to maxZoom', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; + const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + jobManagerClientMock.createTaskForJob.mockResolvedValue(undefined); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }); + + // createTaskForJob's param type is `ICreateTaskBody | ICreateTaskBody[]`; the creator + // always passes an array, so flatten through a cast. + const streamedTasks = jobManagerClientMock.createTaskForJob.mock.calls.flatMap((call) => call[1] as { parameters: unknown }[]); + const allTasks = [...jobManagerClientMock.createJob.mock.calls[0]![0].tasks!, ...streamedTasks]; + const zooms = new Set(allTasks.flatMap((task) => (task.parameters as { ranges: { zoom: number }[] }).ranges.map((range) => range.zoom))); + + expect(Math.min(...zooms)).toBe(0); + expect(Math.max(...zooms)).toBe(taskConfig.maxZoom); + }); + + it('should stream the remaining tasks with createTaskForJob once the first batch fills the job', async () => { + const jobId = randomUUID(); + + // A 1-tile budget makes every range its own task, and a shallow maxZoom keeps the count + // small. The creator reads config in its constructor, so the context must be rebuilt after + // changing these - the `ctx` from beforeEach still holds the defaults. + setValue('jobManagement.ingestion.tasks.cacheDeletion.tileBatchSize', 1); + setValue('jobManagement.ingestion.tasks.cacheDeletion.taskBatchSize', 2); + setValue('jobManagement.ingestion.tasks.cacheDeletion.maxZoom', 3); + + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }); + + expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); + expect(jobManagerClientMock.createJob.mock.calls[0]![0].tasks).toHaveLength(2); + expect(jobManagerClientMock.createTaskForJob).toHaveBeenCalled(); + + for (const call of jobManagerClientMock.createTaskForJob.mock.calls) { + expect(call[0]).toBe(jobId); + // createTaskForJob's param is `ICreateTaskBody | ICreateTaskBody[]`; we always pass an + // array. The final flush may be a partial batch, hence the inequality. + expect((call[1] as unknown[]).length).toBeLessThanOrEqual(2); + } + }); + + it('should fail the job explicitly when a mid-stream enqueue fails, rather than leaving a silent partial success', async () => { + const jobId = randomUUID(); + + // Same config trick as the streaming test above: a 1-tile budget and a shallow maxZoom + // guarantee more than one flush, so the job-creating flush (createJob) succeeds before the + // failure hits a later flush (createTaskForJob) - the scenario where a job already exists + // but is missing tasks. + setValue('jobManagement.ingestion.tasks.cacheDeletion.tileBatchSize', 1); + setValue('jobManagement.ingestion.tasks.cacheDeletion.taskBatchSize', 2); + setValue('jobManagement.ingestion.tasks.cacheDeletion.maxZoom', 3); + + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); + + mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); + jobManagerClientMock.createTaskForJob.mockRejectedValue(new Error('job-manager unreachable')); + + // create() must still resolve - the caller (updateJobHandler) runs this after its own task + // already completed, so throwing here would fail an ingestion finalize whose work is done. + await expect( + cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) + ).resolves.toBeUndefined(); + + expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); + expect(jobManagerClientMock.createTaskForJob).toHaveBeenCalled(); + // the job already exists with only its first batch enqueued - left alone it would run to + // completion and be reported as a fully successful cache deletion, so it must be failed + // explicitly, against the id of the job that was actually created. + expect(jobManagerClientMock.updateJob).toHaveBeenCalledWith(jobId, expect.objectContaining({ status: OperationStatus.FAILED })); + }); + }); + + describe('production key fixture', () => { + // Anchors the two assumptions that fail silently AND successfully if wrong - a bad prefix + // deletes nothing and the task still acks, bad grid math deletes the wrong tiles. Both are + // pinned here against a key observed in a deployed environment: + // bluemarble_swap_test-RasterVectorBest-redis_WorldCRS84-6-76-43 + const CACHE_NAME = 'bluemarble_swap_test-RasterVectorBest-redis'; + const LAYER_NAME = 'bluemarble_swap_test-RasterVectorBest'; + const KNOWN_KEY = 'bluemarble_swap_test-RasterVectorBest-redis_WorldCRS84-6-76-43'; + const TILE = { zoom: 6, x: 76, y: 43 }; + + it('should emit params that reconstruct the known production key', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; + + // a tiny footprint at the centre of tile z6/76/43 - lon [33.75, 36.5625], lat [30.9375, 33.75] + const span = 180 / 2 ** TILE.zoom; + const lon = -180 + TILE.x * span + span / 2; + const lat = -90 + TILE.y * span + span / 2; + const delta = 0.001; + const footprint: Polygon = { + type: 'Polygon', + coordinates: [ + [ + [lon - delta, lat - delta], + [lon + delta, lat - delta], + [lon + delta, lat + delta], + [lon - delta, lat + delta], + [lon - delta, lat - delta], + ], + ], + }; + + // no explicit prefix, mirroring the deployed configuration this key came from + mapproxyClientMock.getRedisCacheName.mockResolvedValue(CACHE_NAME); + readProductGeometryMock.mockResolvedValue(footprint); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: LAYER_NAME, ingestionJob: ingestionUpdateFinalizeJob }); + + const params = jobManagerClientMock.createJob.mock.calls[0]![0].tasks!.map( + (task) => task.parameters as { prefix: string; ranges: { zoom: number; minX: number; maxX: number; minY: number; maxY: number }[] } + ); + + const { prefix } = params[0]!; + + expect(prefix).toBe('bluemarble_swap_test-RasterVectorBest-redis_WorldCRS84'); + + const rangeForTile = params + .flatMap((param) => param.ranges) + .find((range) => range.zoom === TILE.zoom && TILE.x >= range.minX && TILE.x <= range.maxX && TILE.y >= range.minY && TILE.y <= range.maxY); + + expect(rangeForTile).toBeDefined(); + + // the key cleaner will build from these params must match what mapproxy actually wrote + expect(`${prefix}-${TILE.zoom}-${TILE.x}-${TILE.y}`).toBe(KNOWN_KEY); + }); + }); + + describe('failure handling', () => { + it('should swallow and log a mapproxy failure without creating a job', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; + + mapproxyClientMock.getRedisCacheName.mockRejectedValue(new LayerCacheNotFoundError('layer-Orthophoto', 'redis')); + + await expect( + cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) + ).resolves.toBeUndefined(); + expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); + }); + + it('should refuse to construct against a grid footprintToTileRanges does not support', async () => { + setValue('jobManagement.ingestion.tasks.cacheDeletion.grid', 'webmercator'); + + await expect(setupCacheDeletionJobCreatorTest()).rejects.toThrow('Unsupported grid (webmercator)'); + }); + }); +}); diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts new file mode 100644 index 0000000..36b601f --- /dev/null +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts @@ -0,0 +1,45 @@ +import type { Mocked, MockedFunction } from 'vitest'; +import type { JobManagerClient, TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; +import { getTestLogger } from '../../../configurations/testLogger'; +import { CacheDeletionJobCreator } from '../../../../src/job/models/ingestion/cacheDeletionJobCreator'; +import type { MapproxyApiClient } from '../../../../src/httpClients/mapproxyClient'; +import { configMock } from '../../mocks/configMock'; +import { tracerMock } from '../../mocks/tracerMock'; +import { readProductGeometryMock } from '../../mocks/productReaderMock'; + +export interface CacheDeletionJobCreatorTestContext { + cacheDeletionJobCreator: CacheDeletionJobCreator; + jobManagerClientMock: Mocked; + mapproxyClientMock: Mocked; + readProductGeometryMock: MockedFunction; +} + +export const setupCacheDeletionJobCreatorTest = async (): Promise => { + const jobManagerClientMock = { + createJob: vi.fn(), + createTaskForJob: vi.fn(), + updateJob: vi.fn(), + } as unknown as Mocked; + + const queueClientMock = { + jobManagerClient: jobManagerClientMock, + } as unknown as Mocked; + + const mapproxyClientMock = { getRedisCacheName: vi.fn() } as unknown as Mocked; + + const cacheDeletionJobCreator = new CacheDeletionJobCreator( + await getTestLogger(), + tracerMock, + configMock, + queueClientMock, + mapproxyClientMock, + readProductGeometryMock + ); + + return { + cacheDeletionJobCreator, + jobManagerClientMock, + mapproxyClientMock, + readProductGeometryMock, + }; +}; diff --git a/tests/unit/mocks/geometryMockData.ts b/tests/unit/mocks/geometryMockData.ts index f62d344..ddcd94d 100644 --- a/tests/unit/mocks/geometryMockData.ts +++ b/tests/unit/mocks/geometryMockData.ts @@ -17,7 +17,12 @@ export function createFakePolygon(options?: { radiusInMeters?: number }): Polygo const radius = options?.radiusInMeters ?? 1000; // Center point const centerLon = faker.location.longitude({ min: -180, max: 180 }); - const centerLat = faker.location.latitude({ min: -90, max: 90 }); + // The centre latitude is clamped to +/-85 rather than the full +/-90. Longitude degrees shrink by + // cos(lat), so at extreme latitudes a fixed metric radius legitimately spans tens of degrees of + // longitude - which is correct geodesy, but it makes the polygon enormous in tile terms for any + // consumer that walks it down to high zoom (e.g. footprintToTileRanges). Clamping the centre keeps + // the span bounded while leaving radiusInMeters honest at every latitude the fixture can produce. + const centerLat = faker.location.latitude({ min: -85, max: 85 }); // Convert meters to approximate degrees (1 degree ≈ 111320 meters at equator) const radiusInDegrees = radius / 111320; From 6a276518afde7e1e106b4519771f75f2d30f6023 Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 19:16:16 +0300 Subject: [PATCH 04/15] feat!: create the cache deletion job instead of the seeding job (MAPCO-11265) Both ingestion finalize handlers now create an Update_Delete_Cache or Swap_Delete_Cache job instead of a seeding job, and SeedingJobCreator with its SeedMode and Seed* param types is removed. BREAKING CHANGE: overseer no longer creates seeding jobs. --- src/common/constants.ts | 7 - src/common/errors.ts | 8 - src/common/interfaces.ts | 28 +- .../ingestion/cacheDeletionJobCreator.ts | 2 - src/job/models/ingestion/seedingJobCreator.ts | 296 ------------ src/job/models/ingestion/swapJobHandler.ts | 8 +- src/job/models/ingestion/updateJobHandler.ts | 8 +- .../SeedingJobCreator.spec.ts | 453 ------------------ .../seedingJobCreatorSetup.ts | 59 --- .../job/swapJobHandler/swapJobHandler.spec.ts | 8 +- .../job/swapJobHandler/swapJobHandlerSetup.ts | 10 +- .../updateJobHandler/updateJobHandlerSetup.ts | 10 +- tests/unit/mocks/jobsMockData.ts | 70 --- 13 files changed, 23 insertions(+), 944 deletions(-) delete mode 100644 src/job/models/ingestion/seedingJobCreator.ts delete mode 100644 tests/unit/job/seedingJobCreator/SeedingJobCreator.spec.ts delete mode 100644 tests/unit/job/seedingJobCreator/seedingJobCreatorSetup.ts diff --git a/src/common/constants.ts b/src/common/constants.ts index bda0753..03acf36 100644 --- a/src/common/constants.ts +++ b/src/common/constants.ts @@ -44,10 +44,3 @@ export const storageProviderToCacheTypeMap = new Map([ [StorageProvider.S3, LayerCacheType.S3], [StorageProvider.REDIS, LayerCacheType.REDIS], ]); - -export const SeedMode = { - SEED: 'seed', - CLEAN: 'clean', -} as const; - -export type SeedMode = (typeof SeedMode)[keyof typeof SeedMode]; diff --git a/src/common/errors.ts b/src/common/errors.ts index 783dbcd..b2bfe2d 100644 --- a/src/common/errors.ts +++ b/src/common/errors.ts @@ -87,14 +87,6 @@ export class UnsupportedGridError extends Error { } } -export class SeedJobCreationError extends Error { - public constructor(msg: string, err: Error) { - super(msg); - this.name = SeedJobCreationError.name; - this.stack = err.stack; - } -} - export class S3Error extends Error { public constructor(err: unknown, customMessage?: string) { const message = `S3 Error(${customMessage}): ${err instanceof Error ? err.message : 'unknown'}`; diff --git a/src/common/interfaces.ts b/src/common/interfaces.ts index 22ae15a..23d384e 100644 --- a/src/common/interfaces.ts +++ b/src/common/interfaces.ts @@ -26,7 +26,7 @@ import type { ingestionSwapUpdateFinalizeJobParamsSchema, ingestionUpdateFinalizeJobParamsSchema, } from '../utils/zod/schemas/jobParameters.schema'; -import type { LayerCacheType, SeedMode } from './constants'; +import type { LayerCacheType } from './constants'; export type StepKey = keyof T & { [K in keyof T]: T[K] extends boolean ? K : never }[keyof T]; // this is a utility type that extracts the keys of T that are of type boolean @@ -398,32 +398,6 @@ export type CatalogUpdateMetadata = Partial; export type PolygonPartsProcessPayload = Pick & { shouldClearEntities?: boolean }; //#endregion PolygonPartsManagerClient -//#region seedingJobCreator - -export interface SeedJobParams { - layerName: LayerName; - ingestionJob: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob; -} -export interface SeedTaskOptions { - mode: SeedMode; - grid: string; - fromZoomLevel: number; - toZoomLevel: number; - geometry: Footprint; - skipUncached: boolean; - layerId: string; // cache name as configured in mapproxy - refreshBefore: string; -} - -export interface SeedTaskParams { - seedTasks: SeedTaskOptions[]; - catalogId: string; - traceParentContext?: TraceParentContext; - cacheType: LayerCacheType; -} - -//#endregion seedingJobCreator - //#region cacheDeletionJobCreator export interface CacheDeletionJobParams { diff --git a/src/job/models/ingestion/cacheDeletionJobCreator.ts b/src/job/models/ingestion/cacheDeletionJobCreator.ts index eb737a8..d8c4c7f 100644 --- a/src/job/models/ingestion/cacheDeletionJobCreator.ts +++ b/src/job/models/ingestion/cacheDeletionJobCreator.ts @@ -28,9 +28,7 @@ type CacheDeletionTask = ICreateTaskBody; @injectable() export class CacheDeletionJobCreator { private readonly taskConfig: CacheDeletionTaskConfig; - /** the ingestion job type we are reacting to, used only to pick the flow */ private readonly swapUpdateJobType: string; - /** cache-deletion job type per flow - cleaner picks its strategy from this */ private readonly updateCacheDeletionJobType: string; private readonly swapCacheDeletionJobType: string; private readonly wipeDelaySeconds: number; diff --git a/src/job/models/ingestion/seedingJobCreator.ts b/src/job/models/ingestion/seedingJobCreator.ts deleted file mode 100644 index 6eab3b9..0000000 --- a/src/job/models/ingestion/seedingJobCreator.ts +++ /dev/null @@ -1,296 +0,0 @@ -import type { Logger } from '@map-colonies/js-logger'; -import { degreesPerPixelToZoomLevel, getUTCDate, featureToTilesCount } from '@map-colonies/mc-utils'; -import { feature } from '@turf/turf'; -import type { MultiPolygon, Polygon } from 'geojson'; -import { ICreateJobBody, ICreateTaskBody, OperationStatus, TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; -import { inject, injectable } from 'tsyringe'; -import { context, SpanStatusCode, trace, type Tracer } from '@opentelemetry/api'; -import { LayerCacheType, SeedMode, SERVICES } from '../../../common/constants'; -import type { Footprint, IConfig, SeedJobParams, SeedTaskOptions, SeedTaskParams, TilesSeedingTaskConfig } from '../../../common/interfaces'; -import { MapproxyApiClient } from '../../../httpClients/mapproxyClient'; -import { internalIdSchema } from '../../../utils/zod/schemas/jobParameters.schema'; -import { IngestionSwapUpdateFinalizeJob, IngestionUpdateFinalizeJob } from '../../../utils/zod/schemas/job.schema'; -import { splitGeometryByTileCount } from '../../../utils/geoUtils'; -import type { ReadProductGeometry } from '../../../utils/storage/productReader'; -import { CatalogClient } from '../../../httpClients/catalogClient'; - -@injectable() -export class SeedingJobCreator { - private readonly tilesSeedingConfig: TilesSeedingTaskConfig; - private readonly seedJobType: string; - private readonly zoomThreshold: number; - private readonly maxZoom: number; - private readonly maxTilesPerSeedTask: number; - private readonly maxTilesPerCleanTask: number; - private readonly updateJobType: string; - private readonly swapUpdateJobType: string; - - public constructor( - @inject(SERVICES.LOGGER) private readonly logger: Logger, - @inject(SERVICES.TRACER) private readonly tracer: Tracer, - @inject(SERVICES.CONFIG) private readonly config: IConfig, - @inject(SERVICES.QUEUE_CLIENT) protected queueClient: QueueClient, - @inject(MapproxyApiClient) private readonly mapproxyClient: MapproxyApiClient, - @inject(SERVICES.PRODUCT_READER) private readonly readProductGeometry: ReadProductGeometry, - @inject(CatalogClient) private readonly catalogClient: CatalogClient - ) { - this.tilesSeedingConfig = this.config.get('jobManagement.ingestion.tasks.tilesSeeding'); - this.seedJobType = this.config.get('jobManagement.ingestion.jobs.seed.type'); - this.zoomThreshold = this.config.get('jobManagement.ingestion.tasks.tilesSeeding.zoomThreshold'); - this.maxTilesPerSeedTask = this.config.get('jobManagement.ingestion.tasks.tilesSeeding.maxTilesPerSeedTask'); - this.maxTilesPerCleanTask = this.config.get('jobManagement.ingestion.tasks.tilesSeeding.maxTilesPerCleanTask'); - this.updateJobType = this.config.get('jobManagement.ingestion.pollingJobs.update.type'); - this.swapUpdateJobType = this.config.get('jobManagement.ingestion.pollingJobs.swapUpdate.type'); - this.maxZoom = this.config.get('jobManagement.ingestion.tasks.tilesSeeding.maxZoom'); - } - - public async create({ layerName, ingestionJob }: SeedJobParams): Promise { - await context.with(trace.setSpan(context.active(), this.tracer.startSpan(`${SeedingJobCreator.name}.${this.create.name}`)), async () => { - const activeSpan = trace.getActiveSpan(); - try { - const { type: seedTaskType } = this.tilesSeedingConfig; - - const logger = this.logger.child({ ingestionJobId: ingestionJob.id, jobType: this.seedJobType, taskType: seedTaskType }); - logger.info({ msg: 'Starting seeding job creation process' }); - - activeSpan?.setAttributes({ - ingestionJobId: ingestionJob.id, - seedJobType: this.seedJobType, - seedTaskType, - layerName, - }); - - logger.debug({ msg: 'Getting cache name for layer', layerName }); - - const cacheName = await this.mapproxyClient.getRedisCacheName({ layerName, cacheType: LayerCacheType.REDIS }); - activeSpan?.addEvent('getRedisCacheName.success', { cacheName }); - - const validCatalogId = internalIdSchema.parse(ingestionJob).internalId; - - const seedTasks: ICreateTaskBody[] = []; - - // Handle different modes - const cleanModeTasks = await this.handleCleanMode(ingestionJob, cacheName, validCatalogId); - const seedModeTasks = await this.handleSeedMode(ingestionJob, cacheName, validCatalogId); - - seedTasks.push(...cleanModeTasks); - seedTasks.push(...seedModeTasks); - - if (seedTasks.length === 0) { - logger.warn({ msg: 'No tasks created, skipping job creation' }); - activeSpan?.addEvent('createJob.skipped', { reason: 'No tasks created' }); - return; - } - - const { resourceId, version, producerName, productType, domain, productName } = ingestionJob; - const createJobRequest: ICreateJobBody = { - resourceId, - internalId: validCatalogId, - version, - type: this.seedJobType, - parameters: {}, - status: OperationStatus.IN_PROGRESS, - producerName: producerName ?? undefined, - productName, - productType, - domain, - tasks: seedTasks, - }; - - const res = await this.queueClient.jobManagerClient.createJob(createJobRequest); - activeSpan?.addEvent('createJob.success', { seedJobId: res.id }); - logger.info({ msg: 'Seeding job created successfully', seedJobId: res.id, seedTaskIds: res.taskIds }); - } catch (err) { - if (err instanceof Error) { - activeSpan?.recordException(err); - activeSpan?.setStatus({ code: SpanStatusCode.ERROR }); - this.logger.error({ msg: `Failed to create seeding job: ${err.message}`, err }); - } - } finally { - activeSpan?.end(); - } - }); - } - - private async handleCleanMode( - job: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob, - cacheName: string, - catalogId: string - ): Promise[]> { - const activeSpan = trace.getActiveSpan(); - const seedTasks: ICreateTaskBody[] = []; - const seedTaskType = this.tilesSeedingConfig.type; - const logger = this.logger.child({ mode: SeedMode.CLEAN, jobId: job.id, catalogId: job.internalId }); - - const cleanGeometry = await this.calculateGeometryByMode(SeedMode.CLEAN, job, catalogId); - if (!cleanGeometry) { - activeSpan?.addEvent('calculateCleanGeometry.empty'); - logger.warn({ msg: 'No geometry found for CLEAN mode' }); - return []; - } - - activeSpan?.addEvent('calculateCleanGeometry.success', { geometry: JSON.stringify(cleanGeometry) }); - - if (job.type === this.swapUpdateJobType) { - const cleanOptions = this.createSeedOptions(SeedMode.CLEAN, cleanGeometry, cacheName); - activeSpan?.addEvent('createSeedOptions.success', { seedOptions: JSON.stringify(cleanOptions) }); - - const taskParams = this.createTaskParams(catalogId, cleanOptions); - activeSpan?.addEvent('createTaskParams.success', { taskParams: JSON.stringify(taskParams) }); - - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } else if (job.type === this.updateJobType) { - const maxUpdateZoomLevel = degreesPerPixelToZoomLevel(job.parameters.ingestionResolution); //TODO: When we support multi parts resolution, we need to extract maxUpdateZoomLevel from polygon parts table - if (maxUpdateZoomLevel + 1 <= this.maxZoom) { - const cleanOptions = this.createSeedOptions(SeedMode.CLEAN, cleanGeometry, cacheName, maxUpdateZoomLevel + 1); - activeSpan?.addEvent('createSeedOptions.success', { seedOptions: JSON.stringify(cleanOptions) }); - - const taskParams = this.createTaskParams(catalogId, cleanOptions); - activeSpan?.addEvent('createTaskParams.success', { taskParams: JSON.stringify(taskParams) }); - - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } - } - return seedTasks; - - //TODO: add this when we would want to seperate both clean tasks by zoom and maxTilesCount - /** - * const maxUpdatedZoom = extractMaxUpdateZoomLevel(job); - const startCleanZoom: number = job.type === this.swapUpdateJobType ? this.zoomThreshold + 1 : maxUpdatedZoom + 1; - - // For swap update, create a base task from 0 to zoomThreshold - if (job.type === this.swapUpdateJobType) { - const cleanOptions = this.createSeedOptions(SeedMode.CLEAN, cleanGeometry, cacheName, 0, this.zoomThreshold); - activeSpan?.addEvent('createSeedOptions.success', { seedOptions: JSON.stringify(cleanOptions) }); - - const taskParams = this.createTaskParams(catalogId, cleanOptions); - activeSpan?.addEvent('createTaskParams.success', { taskParams: JSON.stringify(taskParams) }); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } - - // For both swap and update, create tasks for high-resolution zoom levels if needed - for (let zoom = startCleanZoom; zoom <= this.maxZoom; zoom++) { - const estimatedTiles = featureToTilesCount(feature(cleanGeometry), zoom); - - if (estimatedTiles <= this.maxTilesPerCleanTask) { - // If tiles count is within limit, create a single task - const options = this.createSeedOptions(SeedMode.CLEAN, cleanGeometry, cacheName, zoom, zoom); - const taskParams = this.createTaskParams(catalogId, options); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } else { - // If tiles count exceeds limit, split the geometry - const splitGeometries = this.splitGeometryByTileCount(cleanGeometry, zoom, this.maxTilesPerCleanTask); - for (const geometry of splitGeometries) { - const options = this.createSeedOptions(SeedMode.CLEAN, geometry, cacheName, zoom, zoom); - const taskParams = this.createTaskParams(catalogId, options); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } - } - } - - return seedTasks; - */ - } - - private async handleSeedMode( - job: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob, - cacheName: string, - catalogId: string - ): Promise[]> { - const activeSpan = trace.getActiveSpan(); - const seedTasks: ICreateTaskBody[] = []; - const seedTaskType = this.tilesSeedingConfig.type; - - if (job.type !== this.updateJobType) { - this.logger.debug({ msg: 'Ingestion Job is not of update type, skipping SEED creation' }); - activeSpan?.addEvent('handleSeedMode.skipped'); - return []; - } - - const seedGeometry = await this.calculateGeometryByMode(SeedMode.SEED, job, catalogId); - if (!seedGeometry) { - activeSpan?.addEvent('calculateSeedGeometry.empty'); - this.logger.warn({ msg: 'No geometry found for SEED mode' }); - return []; - } - - const maxUpdatedZoom = degreesPerPixelToZoomLevel(job.parameters.ingestionResolution); //TODO: When we support multi parts resolution, we need to extract maxUpdateZoomLevel from polygon parts table - - // Step 1: Handle all res from 0 to zoomThreshold in one seed task - const baseZoomLevel = 0; - const seedOptions = this.createSeedOptions(SeedMode.SEED, seedGeometry, cacheName, baseZoomLevel, Math.min(maxUpdatedZoom, this.zoomThreshold)); - const taskParams = this.createTaskParams(catalogId, seedOptions); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - - // Step 2: Handle high-res parts individually by zoom level - for (let zoom = this.zoomThreshold + 1; zoom <= Math.min(maxUpdatedZoom, this.maxZoom); zoom++) { - // the min between the parts updated zoom level and the max allowed zoom- currently we are using the product, later we will use parts table to extract the zoom level - const estimatedTiles = featureToTilesCount(feature(seedGeometry), zoom); - - if (estimatedTiles <= this.maxTilesPerSeedTask) { - // If tiles count is within limit, create a single task - const options = this.createSeedOptions(SeedMode.SEED, seedGeometry, cacheName, zoom, zoom); - const taskParams = this.createTaskParams(catalogId, options); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } else { - // If tiles count exceeds limit, split the geometry - const splitGeometries = splitGeometryByTileCount(seedGeometry, zoom, this.maxTilesPerSeedTask); - for (const geometry of splitGeometries) { - const options = this.createSeedOptions(SeedMode.SEED, geometry, cacheName, zoom, zoom); - const taskParams = this.createTaskParams(catalogId, options); - seedTasks.push({ type: seedTaskType, parameters: taskParams }); - } - } - } - - return seedTasks; - } - - private createSeedOptions( - mode: SeedMode, - geometry: Footprint, - cacheName: string, - fromZoomLevel: number = 0, - toZoomLevel: number = this.maxZoom - ): SeedTaskOptions { - const { grid, skipUncached } = this.tilesSeedingConfig; - const refreshBefore = getUTCDate().toISOString().replace(/\..+/, ''); - return { - mode, - grid, - fromZoomLevel: fromZoomLevel, - toZoomLevel: toZoomLevel, - geometry, - skipUncached, - layerId: cacheName, - refreshBefore, - }; - } - - private createTaskParams(catalogId: string, seedOptions: SeedTaskOptions): SeedTaskParams { - return { - seedTasks: [seedOptions], - catalogId, - traceParentContext: undefined, // todo - add tracing - cacheType: LayerCacheType.REDIS, - }; - } - - private async calculateGeometryByMode( - mode: SeedMode, - job: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob, - catalogId: string - ): Promise { - const logger = this.logger.child({ mode }); - logger.debug({ msg: 'Getting geometry for seeding job' }); - if (mode === SeedMode.CLEAN && job.type === this.swapUpdateJobType) { - const layer = await this.catalogClient.findLayer(catalogId); - const footprint = layer.metadata.footprint; - const geometry = footprint.type === 'Polygon' || footprint.type === 'MultiPolygon' ? footprint : undefined; - return geometry; - } - - const geometry = await this.readProductGeometry(job.parameters.inputFiles.productShapefilePath); - return geometry; - } -} diff --git a/src/job/models/ingestion/swapJobHandler.ts b/src/job/models/ingestion/swapJobHandler.ts index d7d175e..007c951 100644 --- a/src/job/models/ingestion/swapJobHandler.ts +++ b/src/job/models/ingestion/swapJobHandler.ts @@ -22,7 +22,7 @@ import { CatalogClient } from '../../../httpClients/catalogClient'; import { TaskMetrics } from '../../../utils/metrics/taskMetrics'; import { SERVICES } from '../../../common/constants'; import { JobHandler } from '../jobHandler'; -import { SeedingJobCreator } from './seedingJobCreator'; +import { CacheDeletionJobCreator } from './cacheDeletionJobCreator'; @injectable() export class SwapJobHandler @@ -37,7 +37,7 @@ export class SwapJobHandler @inject(TileMergeTaskManager) private readonly taskBuilder: TileMergeTaskManager, @inject(MapproxyApiClient) private readonly mapproxyClient: MapproxyApiClient, @inject(CatalogClient) private readonly catalogClient: CatalogClient, - @inject(SeedingJobCreator) private readonly seedingJobCreator: SeedingJobCreator, + @inject(CacheDeletionJobCreator) private readonly cacheDeletionJobCreator: CacheDeletionJobCreator, @inject(JobTrackerClient) jobTrackerClient: JobTrackerClient, @inject(PolygonPartsMangerClient) private readonly polygonPartsMangerClient: PolygonPartsMangerClient, @inject(SERVICES.PRODUCT_READER) private readonly readProductGeometry: ReadProductGeometry, @@ -142,8 +142,8 @@ export class SwapJobHandler logger.info({ msg: 'All finalize steps completed successfully', ...finalizeTaskParams }); await this.completeTask(job, task, { taskTracker: taskProcessTracking, tracingSpan: activeSpan }); - activeSpan?.addEvent('createSeedingJob.start', { layerName }); - await this.seedingJobCreator.create({ layerName, ingestionJob: job }); + activeSpan?.addEvent('createCacheDeletionJob.start', { layerName }); + await this.cacheDeletionJobCreator.create({ layerName, ingestionJob: job }); } } catch (err) { await this.handleError(err, job, task, { taskTracker: taskProcessTracking, tracingSpan: activeSpan }); diff --git a/src/job/models/ingestion/updateJobHandler.ts b/src/job/models/ingestion/updateJobHandler.ts index d64137b..b414a0d 100644 --- a/src/job/models/ingestion/updateJobHandler.ts +++ b/src/job/models/ingestion/updateJobHandler.ts @@ -20,7 +20,7 @@ import { TaskMetrics } from '../../../utils/metrics/taskMetrics'; import { TileMergeTaskManager } from '../../../task/models/tileMergeTaskManager'; import { TileDeletionTaskManager } from '../../../task/models/deletionTaskManager'; import { JobHandler } from '../jobHandler'; -import { SeedingJobCreator } from './seedingJobCreator'; +import { CacheDeletionJobCreator } from './cacheDeletionJobCreator'; @injectable() export class UpdateJobHandler @@ -35,7 +35,7 @@ export class UpdateJobHandler @inject(TileDeletionTaskManager) private readonly tileDeletionTaskManager: TileDeletionTaskManager, @inject(SERVICES.QUEUE_CLIENT) protected override queueClient: QueueClient, @inject(CatalogClient) private readonly catalogClient: CatalogClient, - @inject(SeedingJobCreator) private readonly seedingJobCreator: SeedingJobCreator, + @inject(CacheDeletionJobCreator) private readonly cacheDeletionJobCreator: CacheDeletionJobCreator, @inject(JobTrackerClient) jobTrackerClient: JobTrackerClient, @inject(PolygonPartsMangerClient) private readonly polygonPartsMangerClient: PolygonPartsMangerClient, @inject(SERVICES.PRODUCT_READER) private readonly readProductGeometry: ReadProductGeometry, @@ -122,8 +122,8 @@ export class UpdateJobHandler logger.info({ msg: 'All finalize steps completed successfully', ...finalizeTaskParams }); await this.completeTask(job, task, { taskTracker: taskProcessTracking, tracingSpan: activeSpan }); - activeSpan?.addEvent('createSeedingJob.start', { layerName }); - await this.seedingJobCreator.create({ layerName, ingestionJob: job }); + activeSpan?.addEvent('createCacheDeletionJob.start', { layerName }); + await this.cacheDeletionJobCreator.create({ layerName, ingestionJob: job }); } } catch (err) { await this.handleError(err, job, task, { taskTracker: taskProcessTracking, tracingSpan: activeSpan }); diff --git a/tests/unit/job/seedingJobCreator/SeedingJobCreator.spec.ts b/tests/unit/job/seedingJobCreator/SeedingJobCreator.spec.ts deleted file mode 100644 index 08b1715..0000000 --- a/tests/unit/job/seedingJobCreator/SeedingJobCreator.spec.ts +++ /dev/null @@ -1,453 +0,0 @@ -/* eslint-disable @typescript-eslint/unbound-method */ -import { randomUUID } from 'node:crypto'; -import type { MultiPolygon, Polygon } from 'geojson'; -import nock from 'nock'; -import { degreesPerPixelToZoomLevel } from '@map-colonies/mc-utils'; -import { OperationStatus } from '@map-colonies/mc-priority-queue'; -import { LayerCacheType, SeedMode } from '../../../../src/common/constants'; -import type { FindLayerResponse, SeedJobParams, SeedTaskOptions, SeedTaskParams, TilesSeedingTaskConfig } from '../../../../src/common/interfaces'; -import { registerDefaultConfig } from '../../mocks/configMock'; -import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; -import { LayerCacheNotFoundError } from '../../../../src/common/errors'; -import { - ingestionSwapUpdateFinalizeJob, - ingestionUpdateJob, - ingestionUpdateJobHighRes, - ingestionUpdateJobHighResMaxTiles, - createBaseSeedTaskOptions, - createHighResSeedTaskOptions, - createCleanTaskOptions, - createSeedJob, -} from '../../mocks/jobsMockData'; -import { splitGeometryByTileCount } from '../../../../src/utils/geoUtils'; -import { layerRecord } from '../../mocks/catalogClientMockData'; -import type { SeedingJobCreatorTestContext } from './seedingJobCreatorSetup'; -import { seedJobParameters, setupSeedingJobCreatorTest } from './seedingJobCreatorSetup'; - -describe('SeedingJobCreator', () => { - let seedingJobCreatorContext: SeedingJobCreatorTestContext; - let productGeometry: Polygon | MultiPolygon; - - beforeEach(async () => { - vi.resetAllMocks(); - registerDefaultConfig(); - seedingJobCreatorContext = await setupSeedingJobCreatorTest(); - productGeometry = createFakePolygonalGeometry(); - }); - - afterEach(() => { - vi.useRealTimers(); - vi.resetAllMocks(); - vi.restoreAllMocks(); - }); - - describe('createSeedingJob', () => { - it('should create seeding job successfully with clean task mode on swapUpdate', async () => { - const { seedingJobCreator, queueClientMock, jobManagerClientMock, mapproxyClientMock, configMock, readProductGeometryMock, catalogClientMock } = - seedingJobCreatorContext; - const baseUrl = configMock.get('jobManagement.config.jobManagerBaseUrl'); - const seedJobType = configMock.get('jobManagement.ingestion.jobs.seed.type'); - const tilesSeedingConfig = configMock.get('jobManagement.ingestion.tasks.tilesSeeding'); - const layerCacheName = 'cache-Name-s3'; - const layer: FindLayerResponse = { ...layerRecord, metadata: { ...layerRecord.metadata, footprint: productGeometry } }; - const taskId = randomUUID(); - const seedJobId = randomUUID(); - - const seedJobParams: SeedJobParams = { - layerName: 'layer-Orthophoto', - ingestionJob: ingestionSwapUpdateFinalizeJob, - }; - - const seedTaskOptions: SeedTaskOptions = { - fromZoomLevel: 0, - toZoomLevel: tilesSeedingConfig.maxZoom, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.CLEAN, - }; - - const seedJob = { - resourceId: ingestionSwapUpdateFinalizeJob.resourceId, - internalId: ingestionSwapUpdateFinalizeJob.internalId, - version: ingestionSwapUpdateFinalizeJob.version, - type: seedJobType, - parameters: {}, - status: OperationStatus.IN_PROGRESS, - producerName: ingestionSwapUpdateFinalizeJob.producerName, - productName: ingestionSwapUpdateFinalizeJob.productName, - productType: ingestionSwapUpdateFinalizeJob.productType, - domain: ingestionSwapUpdateFinalizeJob.domain, - tasks: [ - { - type: tilesSeedingConfig.type, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionSwapUpdateFinalizeJob.internalId, - seedTasks: [seedTaskOptions], - traceParentContext: undefined, - } as SeedTaskParams, - }, - ], - }; - - readProductGeometryMock.mockResolvedValue(productGeometry); - catalogClientMock.findLayer.mockResolvedValue(layer); - vi.useFakeTimers().setSystemTime(new Date('2024-11-05T13:50:27Z')); - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - jobManagerClientMock.createJob.mockResolvedValue({ id: seedJobId, taskIds: [taskId] }); - - nock(baseUrl) - .post('/jobs') - .reply(200, { id: seedJobId, taskIds: [taskId] }); - - const res = await seedingJobCreator.create(seedJobParams); - - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: seedJobParams.layerName, cacheType: LayerCacheType.REDIS }); - expect(queueClientMock.jobManagerClient.createJob).toHaveBeenCalledWith(seedJob); - expect(res).toBeUndefined(); - }); - - it('should create seeding job successfully with seed task and clean task on update', async () => { - const { seedingJobCreator, queueClientMock, jobManagerClientMock, mapproxyClientMock, configMock, catalogClientMock, readProductGeometryMock } = - seedingJobCreatorContext; - const baseUrl = configMock.get('jobManagement.config.jobManagerBaseUrl'); - const seedJobType = configMock.get('jobManagement.ingestion.jobs.seed.type'); - const tilesSeedingConfig = configMock.get('jobManagement.ingestion.tasks.tilesSeeding'); - const layerCacheName = 'cache-Name-s3'; - const layer: FindLayerResponse = { ...layerRecord, metadata: { ...layerRecord.metadata, footprint: productGeometry } }; - const ingestionZoomLevel = degreesPerPixelToZoomLevel(ingestionUpdateJob.parameters.ingestionResolution); - const taskId = randomUUID(); - const seedJobId = randomUUID(); - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - }; - - const seedTaskOptions: SeedTaskOptions = { - fromZoomLevel: 0, - toZoomLevel: ingestionZoomLevel, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.SEED, - }; - - const cleanTaskOptions: SeedTaskOptions = { - fromZoomLevel: ingestionZoomLevel + 1, - toZoomLevel: tilesSeedingConfig.maxZoom, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.CLEAN, - }; - - const seedJob = { - resourceId: ingestionUpdateJob.resourceId, - internalId: ingestionUpdateJob.internalId, - version: ingestionUpdateJob.version, - type: seedJobType, - parameters: {}, - status: OperationStatus.IN_PROGRESS, - producerName: ingestionUpdateJob.producerName, - productName: ingestionUpdateJob.productName, - productType: ingestionUpdateJob.productType, - domain: ingestionUpdateJob.domain, - tasks: [ - { - type: tilesSeedingConfig.type, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionUpdateJob.internalId, - seedTasks: [cleanTaskOptions], - traceParentContext: undefined, - } as SeedTaskParams, - }, - { - type: tilesSeedingConfig.type, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionUpdateJob.internalId, - seedTasks: [seedTaskOptions], - traceParentContext: undefined, - } as SeedTaskParams, - }, - ], - }; - - readProductGeometryMock.mockResolvedValue(productGeometry); - catalogClientMock.findLayer.mockResolvedValue(layer); - vi.useFakeTimers().setSystemTime(new Date('2024-11-05T13:50:27Z')); - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - jobManagerClientMock.createJob.mockResolvedValue({ id: seedJobId, taskIds: [taskId] }); - - nock(baseUrl) - .post('/jobs') - .reply(200, { id: seedJobId, taskIds: [taskId] }); - - const res = await seedingJobCreator.create(seedJobParams); - - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: seedJobParams.layerName, cacheType: LayerCacheType.REDIS }); - expect(queueClientMock.jobManagerClient.createJob).toHaveBeenCalledWith(seedJob); - expect(res).toBeUndefined(); - }); - - it('should create seeding job successfully with seed task mode', async () => { - const { seedingJobCreator, queueClientMock, jobManagerClientMock, mapproxyClientMock, configMock, readProductGeometryMock } = - seedingJobCreatorContext; - const baseUrl = configMock.get('jobManagement.config.jobManagerBaseUrl'); - const seedJobType = configMock.get('jobManagement.ingestion.jobs.seed.type'); - const tilesSeedingConfig = configMock.get('jobManagement.ingestion.tasks.tilesSeeding'); - const layerCacheName = 'cache-Name-s3'; - const taskId = randomUUID(); - const seedJobId = randomUUID(); - const ingestionZoomLevel = degreesPerPixelToZoomLevel(ingestionUpdateJob.parameters.ingestionResolution); - - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - }; - - const seedTaskOptions: SeedTaskOptions = { - fromZoomLevel: 0, - toZoomLevel: ingestionZoomLevel, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.SEED, - }; - - const cleanTaskOptions: SeedTaskOptions = { - fromZoomLevel: ingestionZoomLevel + 1, - toZoomLevel: tilesSeedingConfig.maxZoom, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.CLEAN, - }; - - const seedJob = { - resourceId: ingestionUpdateJob.resourceId, - internalId: ingestionUpdateJob.internalId, - version: ingestionUpdateJob.version, - type: seedJobType, - parameters: {}, - status: OperationStatus.IN_PROGRESS, - producerName: ingestionUpdateJob.producerName, - productType: ingestionUpdateJob.productType, - productName: ingestionUpdateJob.productName, - domain: ingestionUpdateJob.domain, - tasks: [ - { - type: tilesSeedingConfig.type, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionUpdateJob.internalId, - seedTasks: [cleanTaskOptions], - traceParentContext: undefined, - } as SeedTaskParams, - }, - { - type: tilesSeedingConfig.type, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionUpdateJob.internalId, - seedTasks: [seedTaskOptions], - traceParentContext: undefined, - } as SeedTaskParams, - }, - ], - }; - - readProductGeometryMock.mockResolvedValue(productGeometry); - vi.useFakeTimers().setSystemTime(new Date('2024-11-05T13:50:27Z')); - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - jobManagerClientMock.createJob.mockResolvedValue({ id: seedJobId, taskIds: [taskId] }); - - nock(baseUrl) - .post('/jobs') - .reply(200, { id: seedJobId, taskIds: [taskId] }); - - const res = await seedingJobCreator.create(seedJobParams); - - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: seedJobParams.layerName, cacheType: LayerCacheType.REDIS }); - expect(queueClientMock.jobManagerClient.createJob).toHaveBeenCalledWith(seedJob); - expect(res).toBeUndefined(); - }); - - it('should skip creating seeding job - no cache name found', async () => { - const { seedingJobCreator, mapproxyClientMock, jobManagerClientMock } = seedingJobCreatorContext; - const layerCacheName = 'not-exist-s3'; - mapproxyClientMock.getRedisCacheName.mockRejectedValue(new LayerCacheNotFoundError(layerCacheName, LayerCacheType.REDIS)); - - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - }; - - await seedingJobCreator.create(seedJobParams); - - expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); - }); - - it('should skip creating seeding job - no footprint found', async () => { - const { seedingJobCreator, mapproxyClientMock, jobManagerClientMock } = seedingJobCreatorContext; - const layerCacheName = 'layer-name'; - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - }; - - vi.spyOn( - seedingJobCreator as unknown as { calculateGeometryByMode: (...args: unknown[]) => unknown }, - 'calculateGeometryByMode' - ).mockReturnValue(undefined); - await seedingJobCreator.create(seedJobParams); - - expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); - }); - - describe('multiple seed tasks', () => { - it(`should create multiple seed tasks when high-res parts doesn't exceed maxTilesPerSeedTask`, async () => { - const { seedingJobCreator, queueClientMock, jobManagerClientMock, mapproxyClientMock, configMock, readProductGeometryMock } = - seedingJobCreatorContext; - const baseUrl = configMock.get('jobManagement.config.jobManagerBaseUrl'); - const seedJobType = configMock.get('jobManagement.ingestion.jobs.seed.type'); - const tilesSeedingConfig = configMock.get('jobManagement.ingestion.tasks.tilesSeeding'); - productGeometry = createFakePolygonalGeometry({ geometryType: 'Polygon', radiusInMeters: 1000 }); // currently we have bug(MAPCO-9456) calculating tiles count in multipolygon- later we can use the random option to cover more scenarios. - const layerCacheName = 'cache-Name-s3'; - const taskId = randomUUID(); - const seedJobId = randomUUID(); - - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - ingestionJob: ingestionUpdateJobHighRes, - }; - - // Base task from 0 to zoomThreshold - const baseSeedTaskOptions: SeedTaskOptions = { - fromZoomLevel: 0, - toZoomLevel: 16, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.SEED, - }; - - // High-res task that will be split - const highResSeedTaskOptions: SeedTaskOptions = { - fromZoomLevel: 17, - toZoomLevel: 17, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.SEED, - }; - - const cleanTaskOptions: SeedTaskOptions = { - fromZoomLevel: 18, - toZoomLevel: tilesSeedingConfig.maxZoom, - skipUncached: tilesSeedingConfig.skipUncached, - geometry: productGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: tilesSeedingConfig.grid, - mode: SeedMode.CLEAN, - }; - - const seedJob = createSeedJob(ingestionUpdateJobHighRes, seedJobType, tilesSeedingConfig.type, [ - cleanTaskOptions, - baseSeedTaskOptions, - highResSeedTaskOptions, - ]); - - readProductGeometryMock.mockResolvedValue(productGeometry); - vi.useFakeTimers().setSystemTime(new Date('2024-11-05T13:50:27Z')); - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - jobManagerClientMock.createJob.mockResolvedValue({ id: seedJobId, taskIds: [taskId] }); - - nock(baseUrl) - .post('/jobs') - .reply(200, { id: seedJobId, taskIds: [taskId] }); - - const action = async () => { - await seedingJobCreator.create(seedJobParams); - }; - - await expect(action()).resolves.not.toThrow(); - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: seedJobParams.layerName, cacheType: LayerCacheType.REDIS }); - expect(queueClientMock.jobManagerClient.createJob).toHaveBeenCalledWith(seedJob); - }); - - it('should create multiple seed tasks when high-res parts exceed maxTilesPerSeedTask', async () => { - const { seedingJobCreator, queueClientMock, jobManagerClientMock, mapproxyClientMock, configMock, readProductGeometryMock } = - seedingJobCreatorContext; - const baseUrl = configMock.get('jobManagement.config.jobManagerBaseUrl'); - const seedJobType = configMock.get('jobManagement.ingestion.jobs.seed.type'); - const tilesSeedingConfig = configMock.get('jobManagement.ingestion.tasks.tilesSeeding'); - productGeometry = createFakePolygonalGeometry({ geometryType: 'Polygon', radiusInMeters: 100000 }); // Large geometry to ensure high tile count // currently we have bug(MAPCO-9456) calculating tiles count in multipolygon- later we can use the random option to cover more scenarios. - const layerCacheName = 'cache-Name-s3'; - const taskId = randomUUID(); - const seedJobId = randomUUID(); - - const seedJobParams: SeedJobParams = { - ...seedJobParameters, - ingestionJob: ingestionUpdateJobHighResMaxTiles, - }; - - // Get the first part's geometry for high-res splitting - - // Split the first part geometry by tile count (zoom level 17, maxTiles from config) - const maxTilesPerSeedTask = configMock.get('jobManagement.ingestion.tasks.tilesSeeding.maxTilesPerSeedTask'); - const splitGeometries = splitGeometryByTileCount(productGeometry, 17, maxTilesPerSeedTask); - - // Base task from 0 to zoomThreshold - const baseSeedTaskOptions = createBaseSeedTaskOptions(productGeometry, layerCacheName, 16); - - // High-res tasks using the split geometries - const highResSeedTaskOptions: SeedTaskOptions[] = []; - for (const geom of splitGeometries) { - const options = createHighResSeedTaskOptions(geom, layerCacheName); - highResSeedTaskOptions.push(options); - } - - const cleanTaskOptions = createCleanTaskOptions(productGeometry, layerCacheName, 18, tilesSeedingConfig.maxZoom); - - const seedJob = createSeedJob(ingestionUpdateJobHighResMaxTiles, seedJobType, tilesSeedingConfig.type, [ - cleanTaskOptions, - baseSeedTaskOptions, - ...highResSeedTaskOptions, - ]); - - readProductGeometryMock.mockResolvedValue(productGeometry); - vi.useFakeTimers().setSystemTime(new Date('2024-11-05T13:50:27Z')); - mapproxyClientMock.getRedisCacheName.mockResolvedValue(layerCacheName); - jobManagerClientMock.createJob.mockResolvedValue({ id: seedJobId, taskIds: [taskId] }); - - nock(baseUrl) - .post('/jobs') - .reply(200, { id: seedJobId, taskIds: [taskId] }); - - const action = async () => { - await seedingJobCreator.create(seedJobParams); - }; - - await expect(action()).resolves.not.toThrow(); - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: seedJobParams.layerName, cacheType: LayerCacheType.REDIS }); - expect(queueClientMock.jobManagerClient.createJob).toHaveBeenCalledWith(seedJob); - }); - }); - }); -}); diff --git a/tests/unit/job/seedingJobCreator/seedingJobCreatorSetup.ts b/tests/unit/job/seedingJobCreator/seedingJobCreatorSetup.ts deleted file mode 100644 index 9b1d838..0000000 --- a/tests/unit/job/seedingJobCreator/seedingJobCreatorSetup.ts +++ /dev/null @@ -1,59 +0,0 @@ -import type { Mocked, MockedFunction } from 'vitest'; -import type { JobManagerClient, TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; -import { getTestLogger } from '../../../configurations/testLogger'; -import { SeedingJobCreator } from '../../../../src/job/models/ingestion/seedingJobCreator'; -import type { MapproxyApiClient } from '../../../../src/httpClients/mapproxyClient'; -import { configMock } from '../../mocks/configMock'; -import type { SeedJobParams } from '../../../../src/common/interfaces'; -import { ingestionUpdateFinalizeJob } from '../../mocks/jobsMockData'; -import { tracerMock } from '../../mocks/tracerMock'; -import { readProductGeometryMock } from '../../mocks/productReaderMock'; -import type { CatalogClient } from '../../../../src/httpClients/catalogClient'; - -export interface SeedingJobCreatorTestContext { - seedingJobCreator: SeedingJobCreator; - queueClientMock: Mocked; - jobManagerClientMock: Mocked; - mapproxyClientMock: Mocked; - configMock: typeof configMock; - readProductGeometryMock: MockedFunction; - catalogClientMock: Mocked; -} - -export const setupSeedingJobCreatorTest = async (): Promise => { - const jobManagerClientMock = { - createJob: vi.fn(), - } as unknown as Mocked; - - const queueClientMock = { - jobManagerClient: jobManagerClientMock, - } as unknown as Mocked; - - const mapproxyClientMock = { getRedisCacheName: vi.fn() } as unknown as Mocked; - const catalogClientMock = { update: vi.fn(), findLayer: vi.fn() } as unknown as Mocked; - - const seedingJobCreator = new SeedingJobCreator( - await getTestLogger(), - tracerMock, - configMock, - queueClientMock, - mapproxyClientMock, - readProductGeometryMock, - catalogClientMock - ); - - return { - seedingJobCreator, - queueClientMock, - jobManagerClientMock, - mapproxyClientMock, - configMock, - readProductGeometryMock, - catalogClientMock, - }; -}; - -export const seedJobParameters: SeedJobParams = { - layerName: 'layer-Orthophoto', - ingestionJob: ingestionUpdateFinalizeJob, -}; diff --git a/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts b/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts index c4da7a9..02b5517 100644 --- a/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts +++ b/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts @@ -2,7 +2,7 @@ import { getEntityName, getMapServingLayerName, type LayerName, swapUpdateAdditionalParamsSchema } from '@map-colonies/raster-shared'; import { registerDefaultConfig } from '../../mocks/configMock'; import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; -import { Grid, type MergeTask, type SeedJobParams } from '../../../../src/common/interfaces'; +import { Grid, type MergeTask, type CacheDeletionJobParams } from '../../../../src/common/interfaces'; import { finalizeTaskForIngestionSwapUpdate, createTasksTaskForIngestionSwapUpdate } from '../../mocks/tasksMockData'; import { ingestionSwapUpdateFinalizeJob, ingestionSwapUpdateJob } from '../../mocks/jobsMockData'; import { jobTrackerClientMock } from '../../mocks/jobManagerMocks'; @@ -81,7 +81,7 @@ describe('swapJobHandler', () => { jobManagerClientMock, mapproxyClientMock, catalogClientMock, - seedingJobCreatorMock, + cacheDeletionJobCreatorMock, polygonPartsManagerClientMock, } = await setupSwapJobHandlerTest(); const job = structuredClone(ingestionSwapUpdateFinalizeJob); @@ -93,7 +93,7 @@ describe('swapJobHandler', () => { const layerName: LayerName = getMapServingLayerName(job.resourceId, productType); const entityName = getEntityName(job.resourceId, productType); const layerRelativePath = `${job.internalId}/${displayPath}`; - const createSeedingJobParams: SeedJobParams = { + const createCacheDeletionJobParams: CacheDeletionJobParams = { ingestionJob: job, layerName, }; @@ -119,7 +119,7 @@ describe('swapJobHandler', () => { }); expect(queueClientMock.ack).toHaveBeenCalledWith(job.id, task.id); expect(jobTrackerClientMock.notify).toHaveBeenCalledWith(task); - expect(seedingJobCreatorMock.create).toHaveBeenCalledWith(createSeedingJobParams); + expect(cacheDeletionJobCreatorMock.create).toHaveBeenCalledWith(createCacheDeletionJobParams); }); it('should handle job finalize failure and reject the task', async () => { diff --git a/tests/unit/job/swapJobHandler/swapJobHandlerSetup.ts b/tests/unit/job/swapJobHandler/swapJobHandlerSetup.ts index d3e2840..d1877de 100644 --- a/tests/unit/job/swapJobHandler/swapJobHandlerSetup.ts +++ b/tests/unit/job/swapJobHandler/swapJobHandlerSetup.ts @@ -5,7 +5,7 @@ import type { TileMergeTaskManager } from '../../../../src/task/models/tileMerge import type { MapproxyApiClient } from '../../../../src/httpClients/mapproxyClient'; import type { CatalogClient } from '../../../../src/httpClients/catalogClient'; import { SwapJobHandler } from '../../../../src/job/models/ingestion/swapJobHandler'; -import type { SeedingJobCreator } from '../../../../src/job/models/ingestion/seedingJobCreator'; +import type { CacheDeletionJobCreator } from '../../../../src/job/models/ingestion/cacheDeletionJobCreator'; import { taskMetricsMock } from '../../mocks/metricsMock'; import { jobManagerClientMock, jobTrackerClientMock, queueClientMock } from '../../mocks/jobManagerMocks'; import { tracerMock } from '../../mocks/tracerMock'; @@ -22,7 +22,7 @@ export interface SwapJobHandlerTestContext { jobManagerClientMock: Mocked; mapproxyClientMock: Mocked; catalogClientMock: Mocked; - seedingJobCreatorMock: Mocked; + cacheDeletionJobCreatorMock: Mocked; jobTrackerClientMock: Mocked; polygonPartsManagerClientMock: Mocked; readProductGeometryMock: MockedFunction; @@ -36,7 +36,7 @@ export const setupSwapJobHandlerTest = async (): Promise; const catalogClientMock = { update: vi.fn() } as unknown as Mocked; - const seedingJobCreatorMock = { create: vi.fn() } as unknown as Mocked; + const cacheDeletionJobCreatorMock = { create: vi.fn() } as unknown as Mocked; const swapJobHandler = new SwapJobHandler( await getTestLogger(), @@ -46,7 +46,7 @@ export const setupSwapJobHandlerTest = async (): Promise; mapproxyClientMock: Mocked; catalogClientMock: Mocked; - seedingJobCreatorMock: Mocked; + cacheDeletionJobCreatorMock: Mocked; jobTrackerClientMock: Mocked; polygonPartsManagerClientMock: Mocked; readProductGeometryMock: MockedFunction; @@ -42,7 +42,7 @@ export const setupUpdateJobHandlerTest = async (): Promise; const catalogClientMock = { publish: vi.fn(), update: vi.fn() } as unknown as Mocked; - const seedingJobCreatorMock = { create: vi.fn() } as unknown as Mocked; + const cacheDeletionJobCreatorMock = { create: vi.fn() } as unknown as Mocked; const updateJobHandler = new UpdateJobHandler( await getTestLogger(), configMock, @@ -51,7 +51,7 @@ export const setupUpdateJobHandlerTest = async (): Promise ({ - fromZoomLevel: 0, - toZoomLevel, - skipUncached: true, - geometry: seedGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: 'WorldCRS84', - mode: SeedMode.SEED, -}); - -export const createHighResSeedTaskOptions = (geometry: Footprint, layerCacheName: string, zoomLevel: number = 17): SeedTaskOptions => ({ - fromZoomLevel: zoomLevel, - toZoomLevel: zoomLevel, - skipUncached: true, - geometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: 'WorldCRS84', - mode: SeedMode.SEED, -}); - -export const createCleanTaskOptions = ( - seedGeometry: Footprint, - layerCacheName: string, - fromZoomLevel: number = 18, - maxZoom: number = 18 -): SeedTaskOptions => ({ - fromZoomLevel, - toZoomLevel: maxZoom, - skipUncached: true, - geometry: seedGeometry, - refreshBefore: '2024-11-05T13:50:27', - layerId: layerCacheName, - grid: 'WorldCRS84', - mode: SeedMode.CLEAN, -}); - -// Seed Job Mock Function -export const createSeedJob = ( - ingestionJob: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob, - seedJobType: string, - tilesSeedingConfigType: string, - seedTaskOptionsList: SeedTaskOptions[] -): ICreateJobBody => ({ - resourceId: ingestionJob.resourceId, - internalId: ingestionJob.internalId, - version: ingestionJob.version, - type: seedJobType, - parameters: {}, - status: OperationStatus.IN_PROGRESS, - producerName: ingestionJob.producerName, - productName: ingestionJob.productName, - productType: ingestionJob.productType, - domain: ingestionJob.domain, - tasks: seedTaskOptionsList.map((seedTask) => ({ - type: tilesSeedingConfigType, - parameters: { - cacheType: LayerCacheType.REDIS, - catalogId: ingestionJob.internalId, - seedTasks: [seedTask], - traceParentContext: undefined, - } as SeedTaskParams, - })), -}); From 40a49048163b149c0b1e3f06df57b5ea66696514 Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 19:16:37 +0300 Subject: [PATCH 05/15] test: assert the cache deletion job's ingestion trace params (MAPCO-11265) The previous toMatchObject assertion used an empty parameters object, which matches anything, so ingestionJobId and ingestionJobType were unasserted. Verified the new assertion fails when they are removed. --- .../cacheDeletionJobCreator.spec.ts | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 3820cd2..78ee142 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -54,9 +54,14 @@ describe('CacheDeletionJobCreator', () => { internalId: ingestionSwapUpdateFinalizeJob.internalId, version: ingestionSwapUpdateFinalizeJob.version, type: jobType, - parameters: {}, status: OperationStatus.IN_PROGRESS, }); + // job params are how an operator traces a cache deletion back to the ingestion that caused it. + // Asserted strictly rather than via toMatchObject, whose `{}` would match any object at all. + expect(request.parameters).toStrictEqual({ + ingestionJobId: ingestionSwapUpdateFinalizeJob.id, + ingestionJobType: ingestionSwapUpdateFinalizeJob.type, + }); // cleaner resolves DeleteStoredResourcesStrategy from this job type; the update job type // would route the same params to the range strategy and fail validation expect(request.type).toBe('Swap_Delete_Cache'); From 9adaa8d3f5d4df1cddee2bf92fd73bf28d56007a Mon Sep 17 00:00:00 2001 From: almog8k Date: Wed, 12 Aug 2026 19:22:51 +0300 Subject: [PATCH 06/15] refactor: remove redundant comments in CacheDeletionJobCreator tests --- .../job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts | 2 -- 1 file changed, 2 deletions(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 78ee142..2aa2e99 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -56,8 +56,6 @@ describe('CacheDeletionJobCreator', () => { type: jobType, status: OperationStatus.IN_PROGRESS, }); - // job params are how an operator traces a cache deletion back to the ingestion that caused it. - // Asserted strictly rather than via toMatchObject, whose `{}` would match any object at all. expect(request.parameters).toStrictEqual({ ingestionJobId: ingestionSwapUpdateFinalizeJob.id, ingestionJobType: ingestionSwapUpdateFinalizeJob.type, From 3dcb29d3408ae348f47c3b78b9c88d6b1491ba14 Mon Sep 17 00:00:00 2001 From: almog8k Date: Thu, 3 Sep 2026 15:37:29 +0300 Subject: [PATCH 07/15] feat: remove grid configuration and update related logic to fetch from mapproxy --- config/custom-environment-variables.json | 1 - config/default.json | 1 - helm/templates/configmap.yaml | 1 - helm/values.yaml | 1 - src/common/errors.ts | 8 +++ src/common/interfaces.ts | 6 +- src/httpClients/mapproxyClient.ts | 4 +- .../ingestion/cacheDeletionJobCreator.ts | 33 +++++---- tests/unit/httpClients/mapproxyClient.spec.ts | 16 +++-- .../cacheDeletionJobCreator.spec.ts | 69 ++++++++++++++----- .../cacheDeletionJobCreatorSetup.ts | 2 +- tests/unit/mocks/configMock.ts | 1 - 12 files changed, 93 insertions(+), 50 deletions(-) diff --git a/config/custom-environment-variables.json b/config/custom-environment-variables.json index 8d778c8..93093c0 100644 --- a/config/custom-environment-variables.json +++ b/config/custom-environment-variables.json @@ -170,7 +170,6 @@ }, "cacheDeletion": { "type": "CACHE_DELETION_TASK_TYPE", - "grid": "CACHE_DELETION_GRID", "maxZoom": { "__name": "CACHE_DELETION_MAX_ZOOM", "__format": "number" }, "tileBatchSize": { "__name": "CACHE_DELETION_TILE_BATCH_SIZE", "__format": "number" }, "maxRangesPerTask": { "__name": "CACHE_DELETION_MAX_RANGES_PER_TASK", "__format": "number" }, diff --git a/config/default.json b/config/default.json index d6236d4..cc0aa5a 100644 --- a/config/default.json +++ b/config/default.json @@ -152,7 +152,6 @@ }, "cacheDeletion": { "type": "tiles-deletion", - "grid": "WorldCRS84", "maxZoom": 21, "tileBatchSize": 100000, "maxRangesPerTask": 5000, diff --git a/helm/templates/configmap.yaml b/helm/templates/configmap.yaml index f379ee4..c349fb2 100644 --- a/helm/templates/configmap.yaml +++ b/helm/templates/configmap.yaml @@ -85,7 +85,6 @@ data: TILES_SEEDING_MAX_TILES_PER_SEED_TASK: {{ $jobDefinitions.tasks.seed.maxTilesPerSeedTask | quote }} TILES_SEEDING_MAX_TILES_PER_CLEAN_TASK: {{ $jobDefinitions.tasks.seed.maxTilesPerCleanTask | quote }} CACHE_DELETION_TASK_TYPE: {{ $jobDefinitions.tasks.cacheDeletion.type | quote }} - CACHE_DELETION_GRID: {{ $jobDefinitions.tasks.cacheDeletion.grid | quote }} CACHE_DELETION_MAX_ZOOM: {{ $jobDefinitions.tasks.cacheDeletion.maxZoom | quote }} CACHE_DELETION_TILE_BATCH_SIZE: {{ $jobDefinitions.tasks.cacheDeletion.tileBatchSize | quote }} CACHE_DELETION_MAX_RANGES_PER_TASK: {{ $jobDefinitions.tasks.cacheDeletion.maxRangesPerTask | quote }} diff --git a/helm/values.yaml b/helm/values.yaml index f0a1ee2..ab0c1cd 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -184,7 +184,6 @@ jobDefinitions: taskBatchSize: 5 cacheDeletion: type: "" - grid: "WorldCRS84" # must be a geodetic grid - footprintToTileRanges supports no other maxZoom: 21 tileBatchSize: 100000 # max tiles per range-deletion task maxRangesPerTask: 5000 # bounds the serialized task params for jagged footprints diff --git a/src/common/errors.ts b/src/common/errors.ts index b2bfe2d..72f91b6 100644 --- a/src/common/errors.ts +++ b/src/common/errors.ts @@ -87,6 +87,14 @@ export class UnsupportedGridError extends Error { } } +export class UnexpectedCacheGridsError extends Error { + public constructor(cacheName: string, grids: readonly string[] | undefined) { + const reported = grids === undefined ? 'none' : grids.join(', '); + super(`Expected exactly one grid for cache ${cacheName}, mapproxy reported: ${reported}`); + this.name = UnexpectedCacheGridsError.name; + } +} + export class S3Error extends Error { public constructor(err: unknown, customMessage?: string) { const message = `S3 Error(${customMessage}): ${err instanceof Error ? err.message : 'unknown'}`; diff --git a/src/common/interfaces.ts b/src/common/interfaces.ts index 23d384e..e15ce97 100644 --- a/src/common/interfaces.ts +++ b/src/common/interfaces.ts @@ -62,9 +62,7 @@ export interface JobConfig { export interface IngestionJobsConfig { seed: JobConfig | undefined; - /** job type for the update flow's cache deletion; selects cleaner's range-deletion strategy */ updateCacheDeletion: JobConfig | undefined; - /** job type for the swap flow's cache deletion; selects cleaner's prefix-wipe strategy */ swapCacheDeletion: JobConfig | undefined; } @@ -132,8 +130,6 @@ export interface TilesSeedingTaskConfig { export interface CacheDeletionTaskConfig { /** `tiles-deletion` - shared with the S3/FS path, told apart by job type */ type: string; - /** must be a member of GEODETIC_GRIDS; the sole input to the composed `${cacheName}_${grid}` key prefix */ - grid: string; maxZoom: number; /** max tiles a single range-deletion task covers */ tileBatchSize: number; @@ -364,6 +360,8 @@ export interface GetMapproxyCacheResponse { cacheName: string; // eslint-disable-next-line @typescript-eslint/naming-convention cache: { type: LayerCacheType; directory?: string; directory_layout?: string; bucket_name?: string }; + /** the grids the cache is served on; mapproxy derives one redis key prefix per grid */ + grids?: string[]; } //#endregion mapproxyApi diff --git a/src/httpClients/mapproxyClient.ts b/src/httpClients/mapproxyClient.ts index 0c4136f..6b7ba94 100644 --- a/src/httpClients/mapproxyClient.ts +++ b/src/httpClients/mapproxyClient.ts @@ -125,7 +125,7 @@ export class MapproxyApiClient extends HttpClient { return this.fetchLayerCache(layerName, this.layerCacheType); } - public async getRedisCacheName(getCacheReq: GetMapproxyCacheRequest): Promise { + public async getRedisCache(getCacheReq: GetMapproxyCacheRequest): Promise { const { layerName, cacheType } = getCacheReq; const res = await this.fetchLayerCache(layerName, cacheType); if (res === undefined) { @@ -134,7 +134,7 @@ export class MapproxyApiClient extends HttpClient { if (res.cache.type !== LayerCacheType.REDIS) { throw new UnsupportedLayerCacheError(layerName, cacheType); } - return res.cacheName; + return res; } private async fetchLayerCache(layerName: LayerName, cacheType: LayerCacheType): Promise { diff --git a/src/job/models/ingestion/cacheDeletionJobCreator.ts b/src/job/models/ingestion/cacheDeletionJobCreator.ts index d8c4c7f..02473d2 100644 --- a/src/job/models/ingestion/cacheDeletionJobCreator.ts +++ b/src/job/models/ingestion/cacheDeletionJobCreator.ts @@ -7,7 +7,7 @@ import { GEODETIC_GRIDS, StorageProvider } from '@map-colonies/raster-shared'; import { inject, injectable } from 'tsyringe'; import { context, SpanStatusCode, trace, type Span, type Tracer } from '@opentelemetry/api'; import { LayerCacheType, SERVICES } from '../../../common/constants'; -import { UnsupportedGridError } from '../../../common/errors'; +import { UnexpectedCacheGridsError, UnsupportedGridError } from '../../../common/errors'; import type { CacheDeletionJobParams, CacheDeletionTaskConfig, CacheDeletionTaskParams, IConfig } from '../../../common/interfaces'; import { MapproxyApiClient } from '../../../httpClients/mapproxyClient'; import { internalIdSchema } from '../../../utils/zod/schemas/jobParameters.schema'; @@ -46,10 +46,6 @@ export class CacheDeletionJobCreator { this.updateCacheDeletionJobType = this.config.get('jobManagement.ingestion.jobs.updateCacheDeletion.type'); this.swapCacheDeletionJobType = this.config.get('jobManagement.ingestion.jobs.swapCacheDeletion.type'); - // Refuse to start rather than delete the wrong keys and report success. - if (!GEODETIC_GRIDS.includes(this.taskConfig.grid)) { - throw new UnsupportedGridError(this.taskConfig.grid, GEODETIC_GRIDS); - } // Serving pods reload config on their own schedule (gracefulReloadMaxSeconds). // Delay wiping cache keys until all pods have reloaded to prevent stale pods from re-caching // old tiles under removed keys. @@ -101,18 +97,31 @@ export class CacheDeletionJobCreator { /** * Composes the redis key prefix in format `${cacheName}_${gridName}`, matching mapproxy's loader.py behavior. - * The grid comes from config since mapproxy-api does not report it. + * Both parts come from mapproxy-api's cache response, so the grid is not a second source of truth here. */ private async resolvePrefix(layerName: LayerName, logger: Logger): Promise { - const cacheName = await this.mapproxyClient.getRedisCacheName({ layerName, cacheType: LayerCacheType.REDIS }); - const prefix = `${cacheName}_${this.taskConfig.grid}`; + const { cacheName, grids } = await this.mapproxyClient.getRedisCache({ layerName, cacheType: LayerCacheType.REDIS }); + const grid = this.resolveGrid(cacheName, grids); + const prefix = `${cacheName}_${grid}`; - logger.info({ msg: 'Composed redis key prefix', cacheName, grid: this.taskConfig.grid, prefix }); - trace.getActiveSpan()?.setAttributes({ cacheName, redisPrefix: prefix }); + logger.info({ msg: 'Composed redis key prefix', cacheName, grid, prefix }); + trace.getActiveSpan()?.setAttributes({ cacheName, grid, redisPrefix: prefix }); return prefix; } + private resolveGrid(cacheName: string, grids: string[] | undefined): string { + const grid = grids?.length === 1 ? grids[0] : undefined; + if (grid === undefined) { + throw new UnexpectedCacheGridsError(cacheName, grids); + } + if (!GEODETIC_GRIDS.includes(grid)) { + throw new UnsupportedGridError(grid, GEODETIC_GRIDS); + } + + return grid; + } + // eslint-disable-next-line @typescript-eslint/require-await private async *buildWipeTask(prefix: string): AsyncGenerator { yield { @@ -123,10 +132,10 @@ export class CacheDeletionJobCreator { } private async *buildRangeTasks(job: IngestionFinalizeJob, prefix: string, logger: Logger): AsyncGenerator { - const { grid, maxZoom, tileBatchSize, maxRangesPerTask } = this.taskConfig; + const { maxZoom, tileBatchSize, maxRangesPerTask } = this.taskConfig; const geometry = await this.readProductGeometry(job.parameters.inputFiles.productShapefilePath); - logger.info({ msg: 'Computing tile ranges over the updated footprint', grid, maxZoom, tileBatchSize, maxRangesPerTask }); + logger.info({ msg: 'Computing tile ranges over the updated footprint', maxZoom, tileBatchSize, maxRangesPerTask }); const ranges = footprintToTileRanges(geometry, { minZoom: 0, maxZoom }); diff --git a/tests/unit/httpClients/mapproxyClient.spec.ts b/tests/unit/httpClients/mapproxyClient.spec.ts index a7195cf..7cfe971 100644 --- a/tests/unit/httpClients/mapproxyClient.spec.ts +++ b/tests/unit/httpClients/mapproxyClient.spec.ts @@ -115,23 +115,25 @@ describe('mapproxyClient', () => { }); }); - describe('getRedisCacheName', () => { - it('should get cache name from mapproxy', async () => { + describe('getRedisCache', () => { + it('should get the whole redis cache, grids included, from mapproxy', async () => { const baseUrl = configMock.get('servicesUrl.mapproxyApi'); const layerName: LayerName = 'test-Orthophoto'; const cacheType = LayerCacheType.REDIS; const cacheName = 'cacheName'; + const grids = ['WorldCRS84']; nock(baseUrl) .get(`/layer/${layerName}/${cacheType}`) - .reply(200, { cacheName, cache: { type: cacheType } }); + .reply(200, { cacheName, grids, cache: { type: cacheType } }); - const action = mapproxyApiClient.getRedisCacheName({ layerName, cacheType }); + const action = mapproxyApiClient.getRedisCache({ layerName, cacheType }); await expect(action).resolves.not.toThrow(); // eslint-disable-next-line import-x/no-named-as-default-member expect(nock.isDone()).toBe(true); - await expect(action).resolves.toBe(cacheName); + // the grid is the second half of the redis key prefix, so it must survive the client + await expect(action).resolves.toMatchObject({ cacheName, grids }); }); it('should throw an error for unsupported layer cache type', async () => { @@ -144,7 +146,7 @@ describe('mapproxyClient', () => { .get(`/layer/${layerName}/${cacheType}`) .reply(200, { cacheName, cache: { type: cacheType } }); - const action = mapproxyApiClient.getRedisCacheName({ layerName, cacheType }); + const action = mapproxyApiClient.getRedisCache({ layerName, cacheType }); await expect(action).rejects.toThrow(UnsupportedLayerCacheError); }); @@ -156,7 +158,7 @@ describe('mapproxyClient', () => { nock(baseUrl).get(`/layer/${layerName}/${cacheType}`).reply(404); - const action = mapproxyApiClient.getRedisCacheName({ layerName, cacheType }); + const action = mapproxyApiClient.getRedisCache({ layerName, cacheType }); await expect(action).rejects.toThrow(LayerCacheNotFoundError); }); diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 2aa2e99..5d0a247 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -3,7 +3,7 @@ import { randomUUID } from 'node:crypto'; import type { MultiPolygon, Polygon } from 'geojson'; import { OperationStatus } from '@map-colonies/mc-priority-queue'; import { StorageProvider } from '@map-colonies/raster-shared'; -import type { CacheDeletionJobParams, CacheDeletionTaskConfig } from '../../../../src/common/interfaces'; +import type { CacheDeletionJobParams, CacheDeletionTaskConfig, GetMapproxyCacheResponse } from '../../../../src/common/interfaces'; import { registerDefaultConfig, configMock, setValue } from '../../mocks/configMock'; import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; import { LayerCacheType } from '../../../../src/common/constants'; @@ -12,6 +12,15 @@ import { ingestionSwapUpdateFinalizeJob, ingestionUpdateFinalizeJob } from '../. import type { CacheDeletionJobCreatorTestContext } from './cacheDeletionJobCreatorSetup'; import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup'; +/** the only grid GEODETIC_GRIDS admits, and the one the deployed mapproxy serves these caches on */ +const GRID = 'WorldCRS84'; + +const redisCache = (cacheName: string, grids: string[] = [GRID]): GetMapproxyCacheResponse => ({ + cacheName, + cache: { type: LayerCacheType.REDIS }, + grids, +}); + describe('CacheDeletionJobCreator', () => { let ctx: CacheDeletionJobCreatorTestContext; let productGeometry: Polygon | MultiPolygon; @@ -35,7 +44,7 @@ describe('CacheDeletionJobCreator', () => { const jobType = configMock.get('jobManagement.ingestion.jobs.swapCacheDeletion.type'); const jobId = randomUUID(); - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); const params: CacheDeletionJobParams = { layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }; @@ -43,7 +52,7 @@ describe('CacheDeletionJobCreator', () => { await cacheDeletionJobCreator.create(params); // must ask for the REDIS cache explicitly - the layer's tiles cache is file or s3 - expect(mapproxyClientMock.getRedisCacheName).toHaveBeenCalledWith({ layerName: 'layer-Orthophoto', cacheType: LayerCacheType.REDIS }); + expect(mapproxyClientMock.getRedisCache).toHaveBeenCalledWith({ layerName: 'layer-Orthophoto', cacheType: LayerCacheType.REDIS }); expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); expect(jobManagerClientMock.createTaskForJob).not.toHaveBeenCalled(); @@ -70,18 +79,17 @@ describe('CacheDeletionJobCreator', () => { parameters: { storageProvider: StorageProvider.REDIS, // no explicit prefix from mapproxy, so loader.py's fallback applies - prefix: `layer-Orthophoto-redis_${taskConfig.grid}`, + prefix: `layer-Orthophoto-redis_${GRID}`, delaySeconds: taskConfig.gracefulReloadMaxSeconds + taskConfig.reloadWindowMarginSeconds, }, }, ]); }); - it('should compose the prefix from the cache name and the configured grid', async () => { + it('should compose the prefix from the cache name and grid mapproxy reports, not from config', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; - const grid = configMock.get('jobManagement.ingestion.tasks.cacheDeletion').grid; - mapproxyClientMock.getRedisCacheName.mockResolvedValue('sss-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('sss-Orthophoto-redis', ['WorldCRS84'])); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); await cacheDeletionJobCreator.create({ layerName: 'sss-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); @@ -89,14 +97,13 @@ describe('CacheDeletionJobCreator', () => { const request = jobManagerClientMock.createJob.mock.calls[0]![0]; // mapproxy reports no prefix of its own, so the creator composes loader.py's shape - expect(request.tasks![0]!.parameters).toMatchObject({ prefix: `sss-Orthophoto-redis_${grid}` }); expect(request.tasks![0]!.parameters).toMatchObject({ prefix: 'sss-Orthophoto-redis_WorldCRS84' }); }); it('should not read the product shapefile, since a wipe needs no geometry', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); @@ -110,7 +117,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); @@ -136,7 +143,7 @@ describe('CacheDeletionJobCreator', () => { // no other field the worker's .strict() params schema would reject at runtime. expect(task.parameters).toStrictEqual({ storageProvider: StorageProvider.REDIS, - prefix: `layer-Orthophoto-redis_${taskConfig.grid}`, + prefix: `layer-Orthophoto-redis_${GRID}`, ranges, }); } @@ -146,7 +153,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); jobManagerClientMock.createTaskForJob.mockResolvedValue(undefined); @@ -175,7 +182,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); @@ -206,7 +213,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); - mapproxyClientMock.getRedisCacheName.mockResolvedValue('layer-Orthophoto-redis'); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); jobManagerClientMock.createTaskForJob.mockRejectedValue(new Error('job-manager unreachable')); @@ -258,7 +265,7 @@ describe('CacheDeletionJobCreator', () => { }; // no explicit prefix, mirroring the deployed configuration this key came from - mapproxyClientMock.getRedisCacheName.mockResolvedValue(CACHE_NAME); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache(CACHE_NAME)); readProductGeometryMock.mockResolvedValue(footprint); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); @@ -287,7 +294,19 @@ describe('CacheDeletionJobCreator', () => { it('should swallow and log a mapproxy failure without creating a job', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; - mapproxyClientMock.getRedisCacheName.mockRejectedValue(new LayerCacheNotFoundError('layer-Orthophoto', 'redis')); + mapproxyClientMock.getRedisCache.mockRejectedValue(new LayerCacheNotFoundError('layer-Orthophoto', 'redis')); + + await expect( + cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) + ).resolves.toBeUndefined(); + expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); + }); + + it('should create no job when mapproxy reports a grid footprintToTileRanges does not support', async () => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; + + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis', ['webmercator'])); + readProductGeometryMock.mockResolvedValue(productGeometry); await expect( cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) @@ -295,10 +314,22 @@ describe('CacheDeletionJobCreator', () => { expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); }); - it('should refuse to construct against a grid footprintToTileRanges does not support', async () => { - setValue('jobManagement.ingestion.tasks.cacheDeletion.grid', 'webmercator'); + // loader.py derives a prefix per grid, so there is no single prefix that covers the cache - + // creating a job anyway would delete one grid's keys and report the whole cache as deleted + it.each<{ name: string; cache: GetMapproxyCacheResponse }>([ + { name: 'several grids', cache: redisCache('layer-Orthophoto-redis', [GRID, 'epsg3857']) }, + { name: 'no grids', cache: redisCache('layer-Orthophoto-redis', []) }, + // mapproxy-api returns the cache verbatim, so a cache without grids has no grids field + { name: 'no grids field at all', cache: { cacheName: 'layer-Orthophoto-redis', cache: { type: LayerCacheType.REDIS } } }, + ])('should create no job when mapproxy reports $name', async ({ cache }) => { + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; - await expect(setupCacheDeletionJobCreatorTest()).rejects.toThrow('Unsupported grid (webmercator)'); + mapproxyClientMock.getRedisCache.mockResolvedValue(cache); + + await expect( + cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }) + ).resolves.toBeUndefined(); + expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); }); }); }); diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts index 36b601f..0637b03 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreatorSetup.ts @@ -25,7 +25,7 @@ export const setupCacheDeletionJobCreatorTest = async (): Promise; - const mapproxyClientMock = { getRedisCacheName: vi.fn() } as unknown as Mocked; + const mapproxyClientMock = { getRedisCache: vi.fn() } as unknown as Mocked; const cacheDeletionJobCreator = new CacheDeletionJobCreator( await getTestLogger(), diff --git a/tests/unit/mocks/configMock.ts b/tests/unit/mocks/configMock.ts index 4cbe50e..f7f5645 100644 --- a/tests/unit/mocks/configMock.ts +++ b/tests/unit/mocks/configMock.ts @@ -187,7 +187,6 @@ const registerDefaultConfig = (): void => { }, cacheDeletion: { type: 'tiles-deletion', - grid: 'WorldCRS84', maxZoom: 21, tileBatchSize: 100000, maxRangesPerTask: 5000, From d20c15f49c4c662f85eeaca939a7f945df68f152 Mon Sep 17 00:00:00 2001 From: almog8k Date: Thu, 3 Sep 2026 16:00:57 +0300 Subject: [PATCH 08/15] chore: remove redundant comments --- src/common/interfaces.ts | 6 ------ 1 file changed, 6 deletions(-) diff --git a/src/common/interfaces.ts b/src/common/interfaces.ts index e15ce97..7278bb1 100644 --- a/src/common/interfaces.ts +++ b/src/common/interfaces.ts @@ -128,18 +128,12 @@ export interface TilesSeedingTaskConfig { } export interface CacheDeletionTaskConfig { - /** `tiles-deletion` - shared with the S3/FS path, told apart by job type */ type: string; maxZoom: number; - /** max tiles a single range-deletion task covers */ tileBatchSize: number; - /** max ITileRange objects a single task carries, bounding the serialized params size */ maxRangesPerTask: number; - /** tasks pushed per createTaskForJob call */ taskBatchSize: number; - /** mirrors helm's global.gracefulReloadMaxSeconds */ gracefulReloadMaxSeconds: number; - /** added on top, covering mapproxinator's <=5s poll plus pod drain */ reloadWindowMarginSeconds: number; } From 64508f676d8ff3c41c8198df7d364c154ba8565b Mon Sep 17 00:00:00 2001 From: almog8k Date: Thu, 3 Sep 2026 16:10:54 +0300 Subject: [PATCH 09/15] chore: remove redundant comments in CacheDeletionJobCreator --- src/job/models/ingestion/cacheDeletionJobCreator.ts | 8 -------- 1 file changed, 8 deletions(-) diff --git a/src/job/models/ingestion/cacheDeletionJobCreator.ts b/src/job/models/ingestion/cacheDeletionJobCreator.ts index 02473d2..9decbec 100644 --- a/src/job/models/ingestion/cacheDeletionJobCreator.ts +++ b/src/job/models/ingestion/cacheDeletionJobCreator.ts @@ -95,10 +95,6 @@ export class CacheDeletionJobCreator { }); } - /** - * Composes the redis key prefix in format `${cacheName}_${gridName}`, matching mapproxy's loader.py behavior. - * Both parts come from mapproxy-api's cache response, so the grid is not a second source of truth here. - */ private async resolvePrefix(layerName: LayerName, logger: Logger): Promise { const { cacheName, grids } = await this.mapproxyClient.getRedisCache({ layerName, cacheType: LayerCacheType.REDIS }); const grid = this.resolveGrid(cacheName, grids); @@ -148,10 +144,6 @@ export class CacheDeletionJobCreator { } } - /** - * Tasks are streamed to avoid memory overhead. The job is created with its first batch - * so job-tracker never sees an empty job. - */ private async createJobWithStreamedTasks( job: IngestionFinalizeJob, jobType: string, From dc09d0838b1ea223d16d98832093ef5c515f6729 Mon Sep 17 00:00:00 2001 From: almog8k Date: Thu, 3 Sep 2026 16:14:18 +0300 Subject: [PATCH 10/15] feat: update mc-utils and raster-shared dependencies to latest versions --- package-lock.json | 16 ++++++++-------- package.json | 4 ++-- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/package-lock.json b/package-lock.json index 26336ab..b7ce5fa 100644 --- a/package-lock.json +++ b/package-lock.json @@ -20,9 +20,9 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-model-types": "^17.15.1", "@map-colonies/mc-priority-queue": "^9.1.2", - "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", + "@map-colonies/mc-utils": "^6.1.0", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", + "@map-colonies/raster-shared": "9.0.0-alpha.1", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.18.0", "@map-colonies/shapefile-reader": "^1.0.1", @@ -4801,9 +4801,9 @@ } }, "node_modules/@map-colonies/mc-utils": { - "version": "6.0.1", - "resolved": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", - "integrity": "sha512-MKHh65bmqkW8Mc6Ky3Y3tKnV435gdJZiOIkeYjDB8vfsVowPsP4Rz2xq671zny/nzdKRrcWFxH2uXOqV3HnHFQ==", + "version": "6.1.0", + "resolved": "https://registry.npmjs.org/@map-colonies/mc-utils/-/mc-utils-6.1.0.tgz", + "integrity": "sha512-dnt1y76YLAAdsoxsZ2I3KdSJHsho9EnoOJL+G174sqJdgVSgvN+EiDWJ8p1H0GDDcPiDhgJ371+xoW35eE4opg==", "license": "ISC", "dependencies": { "@map-colonies/types": "^1.9.0", @@ -6593,9 +6593,9 @@ } }, "node_modules/@map-colonies/raster-shared": { - "version": "9.0.0-alpha.0", - "resolved": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", - "integrity": "sha512-GZyFW9Dk8lD87qbWlGMFU7KxTTNMrr2fiOwnes7Wjy+ePdFU2maKUoWfDBitCEGderkSIz4iOLRD+yC8MS/M3Q==", + "version": "9.0.0-alpha.1", + "resolved": "https://registry.npmjs.org/@map-colonies/raster-shared/-/raster-shared-9.0.0-alpha.1.tgz", + "integrity": "sha512-D2DUyafRleXEQsmZzIsd6zN5HZrO1E0LHzZAaaSwGlD5JO5VqGe+eUq09woW3MLAZTW8UXzczzqBSZG4ae3VfQ==", "license": "ISC", "dependencies": { "@map-colonies/mc-priority-queue": "^9.1.0", diff --git a/package.json b/package.json index 0cf9bb7..3d7f250 100644 --- a/package.json +++ b/package.json @@ -46,9 +46,9 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-model-types": "^17.15.1", "@map-colonies/mc-priority-queue": "^9.1.2", - "@map-colonies/mc-utils": "https://ghatmpstorage.blob.core.windows.net/npm-packages/mc-utils-594e87195f23d153b8b6619e3e69c3f3bf247b0f.tgz", + "@map-colonies/mc-utils": "^6.1.0", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "https://ghatmpstorage.blob.core.windows.net/npm-packages/raster-shared-b84c16c3baf208c842c66c54c8868b1177d4afb3.tgz", + "@map-colonies/raster-shared": "9.0.0-alpha.1", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.18.0", "@map-colonies/shapefile-reader": "^1.0.1", From 5a5e3f5a437d787c2e461e8ad5c8ee140ab6a4f9 Mon Sep 17 00:00:00 2001 From: almog8k Date: Thu, 3 Sep 2026 17:29:39 +0300 Subject: [PATCH 11/15] refactor: use strategy pattern for task creation --- .../ingestion/cacheDeletionJobCreator.ts | 30 +++++++++++++------ 1 file changed, 21 insertions(+), 9 deletions(-) diff --git a/src/job/models/ingestion/cacheDeletionJobCreator.ts b/src/job/models/ingestion/cacheDeletionJobCreator.ts index 9decbec..5769c36 100644 --- a/src/job/models/ingestion/cacheDeletionJobCreator.ts +++ b/src/job/models/ingestion/cacheDeletionJobCreator.ts @@ -18,6 +18,11 @@ import type { ReadProductGeometry } from '../../../utils/storage/productReader'; type IngestionFinalizeJob = IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob; type CacheDeletionTask = ICreateTaskBody; +interface CacheDeletionStrategy { + jobType: string; + buildTasks: (prefix: string, logger: Logger) => AsyncGenerator; +} + /** * Creates the job that deletes a layer's MapProxy redis tile cache after an ingestion. * @@ -55,11 +60,7 @@ export class CacheDeletionJobCreator { public async create({ layerName, ingestionJob }: CacheDeletionJobParams): Promise { await context.with(trace.setSpan(context.active(), this.tracer.startSpan(`${CacheDeletionJobCreator.name}.${this.create.name}`)), async () => { const activeSpan = trace.getActiveSpan(); - const isSwapUpdate = ingestionJob.type === this.swapUpdateJobType; - - // The job type is what tells cleaner which strategy to run, so it is chosen per flow rather - // than fixed at construction. - const jobType = isSwapUpdate ? this.swapCacheDeletionJobType : this.updateCacheDeletionJobType; + const { jobType, buildTasks } = this.resolveStrategy(ingestionJob); const logger = this.logger.child({ ingestionJobId: ingestionJob.id, @@ -73,16 +74,13 @@ export class CacheDeletionJobCreator { activeSpan?.setAttributes({ ingestionJobId: ingestionJob.id, cacheDeletionJobType: jobType, - cacheDeletionShape: isSwapUpdate ? 'prefix-wipe' : 'range-deletion', layerName, }); const prefix = await this.resolvePrefix(layerName, logger); const catalogId = internalIdSchema.parse(ingestionJob).internalId; - const tasks = isSwapUpdate ? this.buildWipeTask(prefix) : this.buildRangeTasks(ingestionJob, prefix, logger); - - await this.createJobWithStreamedTasks(ingestionJob, jobType, catalogId, tasks, logger, activeSpan); + await this.createJobWithStreamedTasks(ingestionJob, jobType, catalogId, buildTasks(prefix, logger), logger, activeSpan); } catch (err) { if (err instanceof Error) { activeSpan?.recordException(err); @@ -95,6 +93,20 @@ export class CacheDeletionJobCreator { }); } + private resolveStrategy(ingestionJob: IngestionFinalizeJob): CacheDeletionStrategy { + if (ingestionJob.type === this.swapUpdateJobType) { + return { + jobType: this.swapCacheDeletionJobType, + buildTasks: (prefix) => this.buildWipeTask(prefix), + }; + } + + return { + jobType: this.updateCacheDeletionJobType, + buildTasks: (prefix, logger) => this.buildRangeTasks(ingestionJob, prefix, logger), + }; + } + private async resolvePrefix(layerName: LayerName, logger: Logger): Promise { const { cacheName, grids } = await this.mapproxyClient.getRedisCache({ layerName, cacheType: LayerCacheType.REDIS }); const grid = this.resolveGrid(cacheName, grids); From a96a6de95a06090fc210f1d3c86420b8a153e287 Mon Sep 17 00:00:00 2001 From: razbroc Date: Sun, 6 Sep 2026 11:04:15 +0300 Subject: [PATCH 12/15] test: update structure for concating 'redis' for redis cache --- .../cacheDeletionJobCreator.spec.ts | 26 +++++++++---------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 5d0a247..990b335 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -16,7 +16,7 @@ import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup const GRID = 'WorldCRS84'; const redisCache = (cacheName: string, grids: string[] = [GRID]): GetMapproxyCacheResponse => ({ - cacheName, + cacheName: `${cacheName}-redis`, cache: { type: LayerCacheType.REDIS }, grids, }); @@ -44,7 +44,7 @@ describe('CacheDeletionJobCreator', () => { const jobType = configMock.get('jobManagement.ingestion.jobs.swapCacheDeletion.type'); const jobId = randomUUID(); - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); const params: CacheDeletionJobParams = { layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }; @@ -89,7 +89,7 @@ describe('CacheDeletionJobCreator', () => { it('should compose the prefix from the cache name and grid mapproxy reports, not from config', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('sss-Orthophoto-redis', ['WorldCRS84'])); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('sss-Orthophoto', ['WorldCRS84'])); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); await cacheDeletionJobCreator.create({ layerName: 'sss-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); @@ -103,7 +103,7 @@ describe('CacheDeletionJobCreator', () => { it('should not read the product shapefile, since a wipe needs no geometry', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }); @@ -117,7 +117,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); @@ -153,7 +153,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; const taskConfig = configMock.get('jobManagement.ingestion.tasks.cacheDeletion'); - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); jobManagerClientMock.createTaskForJob.mockResolvedValue(undefined); @@ -182,7 +182,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); @@ -213,7 +213,7 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis')); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); readProductGeometryMock.mockResolvedValue(productGeometry); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); jobManagerClientMock.createTaskForJob.mockRejectedValue(new Error('job-manager unreachable')); @@ -238,7 +238,7 @@ describe('CacheDeletionJobCreator', () => { // deletes nothing and the task still acks, bad grid math deletes the wrong tiles. Both are // pinned here against a key observed in a deployed environment: // bluemarble_swap_test-RasterVectorBest-redis_WorldCRS84-6-76-43 - const CACHE_NAME = 'bluemarble_swap_test-RasterVectorBest-redis'; + const CACHE_NAME = 'bluemarble_swap_test-RasterVectorBest'; const LAYER_NAME = 'bluemarble_swap_test-RasterVectorBest'; const KNOWN_KEY = 'bluemarble_swap_test-RasterVectorBest-redis_WorldCRS84-6-76-43'; const TILE = { zoom: 6, x: 76, y: 43 }; @@ -305,7 +305,7 @@ describe('CacheDeletionJobCreator', () => { it('should create no job when mapproxy reports a grid footprintToTileRanges does not support', async () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; - mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto-redis', ['webmercator'])); + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto', ['webmercator'])); readProductGeometryMock.mockResolvedValue(productGeometry); await expect( @@ -317,10 +317,10 @@ describe('CacheDeletionJobCreator', () => { // loader.py derives a prefix per grid, so there is no single prefix that covers the cache - // creating a job anyway would delete one grid's keys and report the whole cache as deleted it.each<{ name: string; cache: GetMapproxyCacheResponse }>([ - { name: 'several grids', cache: redisCache('layer-Orthophoto-redis', [GRID, 'epsg3857']) }, - { name: 'no grids', cache: redisCache('layer-Orthophoto-redis', []) }, + { name: 'several grids', cache: redisCache('layer-Orthophoto', [GRID, 'epsg3857']) }, + { name: 'no grids', cache: redisCache('layer-Orthophoto', []) }, // mapproxy-api returns the cache verbatim, so a cache without grids has no grids field - { name: 'no grids field at all', cache: { cacheName: 'layer-Orthophoto-redis', cache: { type: LayerCacheType.REDIS } } }, + { name: 'no grids field at all', cache: { cacheName: 'layer-Orthophoto', cache: { type: LayerCacheType.REDIS } } }, ])('should create no job when mapproxy reports $name', async ({ cache }) => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock } = ctx; From e7aada07a22245d8d4489b7826987d117ef8fcac Mon Sep 17 00:00:00 2001 From: razbroc Date: Sun, 6 Sep 2026 11:12:11 +0300 Subject: [PATCH 13/15] chore: remove outdated comments in CacheDeletionJobCreator tests --- .../job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts | 2 -- 1 file changed, 2 deletions(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 990b335..6ee2793 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -314,8 +314,6 @@ describe('CacheDeletionJobCreator', () => { expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); }); - // loader.py derives a prefix per grid, so there is no single prefix that covers the cache - - // creating a job anyway would delete one grid's keys and report the whole cache as deleted it.each<{ name: string; cache: GetMapproxyCacheResponse }>([ { name: 'several grids', cache: redisCache('layer-Orthophoto', [GRID, 'epsg3857']) }, { name: 'no grids', cache: redisCache('layer-Orthophoto', []) }, From af3897891a5df95d9535441b8e1cadd409481896 Mon Sep 17 00:00:00 2001 From: razbroc Date: Sun, 6 Sep 2026 11:26:58 +0300 Subject: [PATCH 14/15] test: add UnsupportedGridError handling in CacheDeletionJobCreator tests --- .../cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 6ee2793..e407184 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -7,7 +7,7 @@ import type { CacheDeletionJobParams, CacheDeletionTaskConfig, GetMapproxyCacheR import { registerDefaultConfig, configMock, setValue } from '../../mocks/configMock'; import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; import { LayerCacheType } from '../../../../src/common/constants'; -import { LayerCacheNotFoundError } from '../../../../src/common/errors'; +import { LayerCacheNotFoundError, UnsupportedGridError } from '../../../../src/common/errors'; import { ingestionSwapUpdateFinalizeJob, ingestionUpdateFinalizeJob } from '../../mocks/jobsMockData'; import type { CacheDeletionJobCreatorTestContext } from './cacheDeletionJobCreatorSetup'; import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup'; @@ -306,11 +306,13 @@ describe('CacheDeletionJobCreator', () => { const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = ctx; mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto', ['webmercator'])); + const resolveGridSpy = vi.spyOn(cacheDeletionJobCreator as unknown as { resolveGrid: () => void }, 'resolveGrid'); readProductGeometryMock.mockResolvedValue(productGeometry); await expect( cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) ).resolves.toBeUndefined(); + expect(resolveGridSpy).toThrowError(UnsupportedGridError); expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); }); From 694aa25a0bf5cf9d36146aa1e251bb6e1ea96d59 Mon Sep 17 00:00:00 2001 From: razbroc Date: Sun, 6 Sep 2026 12:18:11 +0300 Subject: [PATCH 15/15] fix: add UnexpectedCacheGridsError handling in CacheDeletionJobCreator tests --- .../cacheDeletionJobCreator.spec.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index e407184..9348371 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -7,7 +7,7 @@ import type { CacheDeletionJobParams, CacheDeletionTaskConfig, GetMapproxyCacheR import { registerDefaultConfig, configMock, setValue } from '../../mocks/configMock'; import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; import { LayerCacheType } from '../../../../src/common/constants'; -import { LayerCacheNotFoundError, UnsupportedGridError } from '../../../../src/common/errors'; +import { LayerCacheNotFoundError, UnexpectedCacheGridsError, UnsupportedGridError } from '../../../../src/common/errors'; import { ingestionSwapUpdateFinalizeJob, ingestionUpdateFinalizeJob } from '../../mocks/jobsMockData'; import type { CacheDeletionJobCreatorTestContext } from './cacheDeletionJobCreatorSetup'; import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup'; @@ -312,8 +312,12 @@ describe('CacheDeletionJobCreator', () => { await expect( cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }) ).resolves.toBeUndefined(); - expect(resolveGridSpy).toThrowError(UnsupportedGridError); expect(jobManagerClientMock.createJob).not.toHaveBeenCalled(); + expect(resolveGridSpy.mock.results[0]).toMatchObject({ + type: 'throw', + // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment + value: expect.any(UnsupportedGridError), + }); }); it.each<{ name: string; cache: GetMapproxyCacheResponse }>([