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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
2 changes: 1 addition & 1 deletion src/common/interfaces.ts
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,7 @@ export type PolygonPartsProcessPayload = Pick<PolygonPartsPayload, 'productId' |

//#region cacheDeletionJobCreator

export interface CacheDeletionJobParams {
export interface CreateCacheDeletionJobParams {
layerName: LayerName;
ingestionJob: IngestionUpdateFinalizeJob | IngestionSwapUpdateFinalizeJob;
}
Expand Down
33 changes: 20 additions & 13 deletions src/job/models/ingestion/cacheDeletionJobCreator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,13 @@ 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 type { CacheDeletionJobParams, 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 { UnexpectedCacheGridsError, UnsupportedGridError } from '../../../common/errors';
import type { CacheDeletionJobParams, CacheDeletionTaskConfig, CacheDeletionTaskParams, IConfig } from '../../../common/interfaces';
import type { CreateCacheDeletionJobParams, 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';
Expand Down Expand Up @@ -57,7 +57,7 @@ export class CacheDeletionJobCreator {
this.wipeDelaySeconds = this.taskConfig.gracefulReloadMaxSeconds + this.taskConfig.reloadWindowMarginSeconds;
}

public async create({ layerName, ingestionJob }: CacheDeletionJobParams): Promise<void> {
public async create({ layerName, ingestionJob }: CreateCacheDeletionJobParams): Promise<void> {
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);
Expand Down Expand Up @@ -168,20 +168,26 @@ export class CacheDeletionJobCreator {
let taskBatch: CacheDeletionTask[] = [];
let jobId: string | undefined;
let taskCount = 0;
let isJobCreatedWithAllTasks = false;

try {
for await (const task of tasks) {
taskBatch.push(task);
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) {
Expand All @@ -205,20 +211,21 @@ export class CacheDeletionJobCreator {
jobType: string,
catalogId: string,
jobId: string | undefined,
tasks: CacheDeletionTask[]
tasks: CacheDeletionTask[],
isLastBatch: boolean
): Promise<string> {
if (jobId !== undefined) {
await this.queueClient.jobManagerClient.createTaskForJob(jobId, tasks);
return jobId;
}

const { resourceId, version, producerName, productName, productType, domain } = job;
const createJobRequest: ICreateJobBody<unknown, CacheDeletionTaskParams> = {
const createJobRequest: ICreateJobBody<CacheDeletionJobParams, CacheDeletionTaskParams> = {
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,
Expand All @@ -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<void> {
const parameters: CacheDeletionJobParams = { ingestionJobId, tasksCreationCompleted: true };
await this.queueClient.jobManagerClient.updateJob(jobId, { parameters });
}

private async failPartialJob(jobId: string, cause: unknown, logger: Logger): Promise<void> {
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 });
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,12 @@
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';
import { LayerCacheNotFoundError, UnexpectedCacheGridsError, UnsupportedGridError } from '../../../../src/common/errors';

Check warning on line 10 in tests/unit/job/cacheDeletionJobCreator/cacheDeletionJobCreator.spec.ts

View workflow job for this annotation

GitHub Actions / Run TS Project eslint (24.x)

'UnexpectedCacheGridsError' is defined but never used
import { ingestionSwapUpdateFinalizeJob, ingestionUpdateFinalizeJob } from '../../mocks/jobsMockData';
import type { CacheDeletionJobCreatorTestContext } from './cacheDeletionJobCreatorSetup';
import { setupCacheDeletionJobCreatorTest } from './cacheDeletionJobCreatorSetup';
Expand Down Expand Up @@ -47,7 +47,7 @@
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);

Expand All @@ -65,10 +65,12 @@
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');
Expand Down Expand Up @@ -200,6 +202,58 @@
}
});

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();

Expand Down
4 changes: 2 additions & 2 deletions tests/unit/job/swapJobHandler/swapJobHandler.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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,
};
Expand Down
Loading