diff --git a/package-lock.json b/package-lock.json index 6d9cf78..f7821b4 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": "^6.1.0", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "9.0.0-alpha.1", + "@map-colonies/raster-shared": "9.0.0", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.18.0", "@map-colonies/shapefile-reader": "^1.0.1", @@ -6593,9 +6593,9 @@ } }, "node_modules/@map-colonies/raster-shared": { - "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==", + "version": "9.0.0", + "resolved": "https://registry.npmjs.org/@map-colonies/raster-shared/-/raster-shared-9.0.0.tgz", + "integrity": "sha512-RBrcPmVDm6i0/+odACjw/4HJYRuqWXwlJWnNXDm2sqamx74r+IxuTUYtRJHKktFwJc3r+g1YuT60VBeW9Lcaww==", "license": "ISC", "dependencies": { "@map-colonies/mc-priority-queue": "^9.1.0", diff --git a/package.json b/package.json index ecb606c..587850d 100644 --- a/package.json +++ b/package.json @@ -48,7 +48,7 @@ "@map-colonies/mc-priority-queue": "^9.1.2", "@map-colonies/mc-utils": "^6.1.0", "@map-colonies/prometheus": "^1.0.0", - "@map-colonies/raster-shared": "9.0.0-alpha.1", + "@map-colonies/raster-shared": "9.0.0", "@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/interfaces.ts b/src/common/interfaces.ts index 2b04c48..2571fc2 100644 --- a/src/common/interfaces.ts +++ b/src/common/interfaces.ts @@ -383,7 +383,7 @@ export type PolygonPartsProcessPayload = Pick { + public async create({ layerName, ingestionJob }: CreateCacheDeletionJobParams): Promise { await context.with(trace.setSpan(context.active(), this.tracer.startSpan(`${CacheDeletionJobCreator.name}.${this.create.name}`)), async () => { const activeSpan = trace.getActiveSpan(); const { jobType, buildTasks } = this.resolveStrategy(ingestionJob); @@ -168,6 +168,7 @@ export class CacheDeletionJobCreator { let taskBatch: CacheDeletionTask[] = []; let jobId: string | undefined; let taskCount = 0; + let isJobCreatedWithAllTasks = false; try { for await (const task of tasks) { @@ -175,13 +176,18 @@ export class CacheDeletionJobCreator { taskCount++; if (taskBatch.length === taskBatchSize) { - jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch); + jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch, false); taskBatch = []; } } if (taskBatch.length > 0) { - jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch); + isJobCreatedWithAllTasks = jobId === undefined; + jobId = await this.flushTaskBatch(job, jobType, catalogId, jobId, taskBatch, true); + } + + if (jobId !== undefined && !isJobCreatedWithAllTasks) { + await this.markTasksCreationCompleted(jobId, job.id); } } catch (err) { if (jobId !== undefined) { @@ -205,7 +211,8 @@ export class CacheDeletionJobCreator { jobType: string, catalogId: string, jobId: string | undefined, - tasks: CacheDeletionTask[] + tasks: CacheDeletionTask[], + isLastBatch: boolean ): Promise { if (jobId !== undefined) { await this.queueClient.jobManagerClient.createTaskForJob(jobId, tasks); @@ -213,12 +220,12 @@ export class CacheDeletionJobCreator { } const { resourceId, version, producerName, productName, productType, domain } = job; - const createJobRequest: ICreateJobBody = { + const createJobRequest: ICreateJobBody = { resourceId, internalId: catalogId, version, type: jobType, - parameters: { ingestionJobId: job.id, ingestionJobType: job.type }, + parameters: { ingestionJobId: job.id, tasksCreationCompleted: isLastBatch }, status: OperationStatus.IN_PROGRESS, producerName: producerName ?? undefined, productName, @@ -231,11 +238,11 @@ export class CacheDeletionJobCreator { 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 markTasksCreationCompleted(jobId: string, ingestionJobId: string): Promise { + const parameters: CacheDeletionJobParams = { ingestionJobId, tasksCreationCompleted: true }; + await this.queueClient.jobManagerClient.updateJob(jobId, { parameters }); + } + 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 }); diff --git a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts index 9348371..b6cd4ba 100644 --- a/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts +++ b/tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts @@ -2,8 +2,8 @@ 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, GetMapproxyCacheResponse } from '../../../../src/common/interfaces'; +import { StorageProvider, type CacheDeletionJobParams } from '@map-colonies/raster-shared'; +import type { CreateCacheDeletionJobParams, CacheDeletionTaskConfig, GetMapproxyCacheResponse } from '../../../../src/common/interfaces'; import { registerDefaultConfig, configMock, setValue } from '../../mocks/configMock'; import { createFakePolygonalGeometry } from '../../mocks/geometryMockData'; import { LayerCacheType } from '../../../../src/common/constants'; @@ -47,7 +47,7 @@ describe('CacheDeletionJobCreator', () => { mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); - const params: CacheDeletionJobParams = { layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }; + const params: CreateCacheDeletionJobParams = { layerName: 'layer-Orthophoto', ingestionJob: ingestionSwapUpdateFinalizeJob }; await cacheDeletionJobCreator.create(params); @@ -65,10 +65,12 @@ describe('CacheDeletionJobCreator', () => { type: jobType, status: OperationStatus.IN_PROGRESS, }); + // the single wipe task is the whole job, so tasks creation is completed with the job itself expect(request.parameters).toStrictEqual({ ingestionJobId: ingestionSwapUpdateFinalizeJob.id, - ingestionJobType: ingestionSwapUpdateFinalizeJob.type, - }); + tasksCreationCompleted: true, + } satisfies CacheDeletionJobParams); + expect(jobManagerClientMock.updateJob).not.toHaveBeenCalled(); // 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'); @@ -200,6 +202,58 @@ describe('CacheDeletionJobCreator', () => { } }); + it('should create the job with tasks creation completed and skip the close-out when all tasks fit in the first batch', async () => { + setValue('jobManagement.ingestion.tasks.cacheDeletion.taskBatchSize', 1000); + + const { cacheDeletionJobCreator, jobManagerClientMock, mapproxyClientMock, readProductGeometryMock } = await setupCacheDeletionJobCreatorTest(); + + mapproxyClientMock.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: randomUUID(), taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }); + + expect(jobManagerClientMock.createJob).toHaveBeenCalledTimes(1); + expect(jobManagerClientMock.createTaskForJob).not.toHaveBeenCalled(); + expect(jobManagerClientMock.createJob.mock.calls[0]![0].parameters).toStrictEqual({ + ingestionJobId: ingestionUpdateFinalizeJob.id, + tasksCreationCompleted: true, + } satisfies CacheDeletionJobParams); + expect(jobManagerClientMock.updateJob).not.toHaveBeenCalled(); + }); + + it('should flag tasks creation as completed only after the last batch is enqueued', async () => { + const jobId = randomUUID(); + + // Same config trick as the streaming test above - guarantees more than one flush. + 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.getRedisCache.mockResolvedValue(redisCache('layer-Orthophoto')); + readProductGeometryMock.mockResolvedValue(productGeometry); + jobManagerClientMock.createJob.mockResolvedValue({ id: jobId, taskIds: [randomUUID()] }); + + await cacheDeletionJobCreator.create({ layerName: 'layer-Orthophoto', ingestionJob: ingestionUpdateFinalizeJob }); + + expect(jobManagerClientMock.createJob.mock.calls[0]![0].parameters).toStrictEqual({ + ingestionJobId: ingestionUpdateFinalizeJob.id, + tasksCreationCompleted: false, + } satisfies CacheDeletionJobParams); + + // job-manager overwrites parameters on update, so the full object must be sent + expect(jobManagerClientMock.updateJob).toHaveBeenCalledTimes(1); + expect(jobManagerClientMock.updateJob).toHaveBeenCalledWith(jobId, { + parameters: { ingestionJobId: ingestionUpdateFinalizeJob.id, tasksCreationCompleted: true } satisfies CacheDeletionJobParams, + }); + + const lastEnqueueOrder = Math.max(...jobManagerClientMock.createTaskForJob.mock.invocationCallOrder); + + expect(jobManagerClientMock.updateJob.mock.invocationCallOrder[0]).toBeGreaterThan(lastEnqueueOrder); + }); + it('should fail the job explicitly when a mid-stream enqueue fails, rather than leaving a silent partial success', async () => { const jobId = randomUUID(); diff --git a/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts b/tests/unit/job/swapJobHandler/swapJobHandler.spec.ts index 02b5517..1dc1d1b 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 CacheDeletionJobParams } from '../../../../src/common/interfaces'; +import { Grid, type MergeTask, type CreateCacheDeletionJobParams } from '../../../../src/common/interfaces'; import { finalizeTaskForIngestionSwapUpdate, createTasksTaskForIngestionSwapUpdate } from '../../mocks/tasksMockData'; import { ingestionSwapUpdateFinalizeJob, ingestionSwapUpdateJob } from '../../mocks/jobsMockData'; import { jobTrackerClientMock } from '../../mocks/jobManagerMocks'; @@ -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 createCacheDeletionJobParams: CacheDeletionJobParams = { + const createCacheDeletionJobParams: CreateCacheDeletionJobParams = { ingestionJob: job, layerName, };