Skip to content
Merged
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.0.1",
"@map-colonies/prometheus": "^1.0.0",
"@map-colonies/raster-shared": "^8.3.0",
"@map-colonies/raster-shared": "^9.0.0-alpha.0",
"@map-colonies/read-pkg": "^1.0.0",
"@map-colonies/schemas": "^1.18.0",
"@map-colonies/shapefile-reader": "^1.0.1",
Expand Down
9 changes: 4 additions & 5 deletions src/common/constants.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,6 @@
import { SourceType } from '@map-colonies/raster-shared';
import { StorageProvider } from '@map-colonies/raster-shared';
import { readPackageJsonSync } from '@map-colonies/read-pkg';

export type StorageProvider = Exclude<SourceType, 'GPKG'>;

export const SERVICE_NAME = readPackageJsonSync().name ?? 'unknown_service';
export const SERVICE_VERSION = readPackageJsonSync().version ?? 'unknown_version';
export const DEFAULT_SERVER_PORT = 80;
Expand Down Expand Up @@ -42,8 +40,9 @@ export const LayerCacheType = {
export type LayerCacheType = (typeof LayerCacheType)[keyof typeof LayerCacheType];

export const storageProviderToCacheTypeMap = new Map([
[SourceType.FS, LayerCacheType.FS],
[SourceType.S3, LayerCacheType.S3],
[StorageProvider.FS, LayerCacheType.FS],
[StorageProvider.S3, LayerCacheType.S3],
[StorageProvider.REDIS, LayerCacheType.REDIS],
]);

export const SeedMode = {
Expand Down
4 changes: 2 additions & 2 deletions src/httpClients/mapproxyClient.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
import type { Logger } from '@map-colonies/js-logger';
import type { LayerName, TileOutputFormat } from '@map-colonies/raster-shared';
import type { LayerName, StorageProvider, TileOutputFormat } from '@map-colonies/raster-shared';
import { context, SpanStatusCode, trace, type Tracer } from '@opentelemetry/api';
import { HttpClient, type IHttpRetryConfig } from '@map-colonies/mc-utils';
import { inject, injectable } from 'tsyringe';
import { NotFoundError } from '@map-colonies/error-types';
import type { IConfig, GetMapproxyCacheRequest, GetMapproxyCacheResponse, PublishMapLayerRequest } from '../common/interfaces';
import { LayerCacheType, SERVICES, storageProviderToCacheTypeMap, StorageProvider } from '../common/constants';
import { LayerCacheType, SERVICES, storageProviderToCacheTypeMap } from '../common/constants';
import {
DeleteLayerError,
LayerCacheNotFoundError,
Expand Down
15 changes: 8 additions & 7 deletions src/job/models/deletion/deleteLayerHandler.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
import type { Logger } from '@map-colonies/js-logger';
import { context, trace, type Tracer } from '@opentelemetry/api';
import { TaskHandler as QueueClient, type ICreateTaskBody } from '@map-colonies/mc-priority-queue';
import { SourceType, type DeleteTaskParams, type DeleteStoredResourcesParams, type LayerName, Storage } from '@map-colonies/raster-shared';
import { type DeleteTaskParams, type DeleteStoredResourcesParams, type LayerName, StorageProvider, Storage } from '@map-colonies/raster-shared';
import { inject, injectable } from 'tsyringe';
import type { IConfig, IJobHandler, JobAndTaskTelemetry, StepKey } from '../../../common/interfaces';
import { SERVICES, type StorageProvider } from '../../../common/constants';
import { SERVICES } from '../../../common/constants';
import { LayerCacheNotFoundError } from '../../../common/errors';
import { CatalogClient } from '../../../httpClients/catalogClient';
import { GeoserverClient } from '../../../httpClients/geoserverClient';
Expand All @@ -30,7 +30,8 @@ interface DeletionStep {
@injectable()
export class DeleteLayerHandler extends JobHandler implements IJobHandler<never, never, never, never, DeleteLayerJob, DeleteTask> {
private readonly tilesDeletionType: string;
private readonly tilesStorageProvider: StorageProvider;
//TODO: when we support redis tiles deletion, change the type to StorageProvider and remove the Exclude<> wrapper
private readonly tilesStorageProvider: Exclude<StorageProvider, 'REDIS'>;
Comment thread
almog8k marked this conversation as resolved.
private readonly tilesBucketConfig: string;
private readonly tilesSubPathConfig: string;

Expand All @@ -49,7 +50,7 @@ export class DeleteLayerHandler extends JobHandler implements IJobHandler<never,
super(logger, config, queueClient, jobTrackerClient);
// whole-layer tiles deletion shares the 'tiles-deletion' task type with the range-based ingestion flow (raster-shared DeletionTaskTypes.LayerTilesDeletion); the Cleaner distinguishes them by params shape
this.tilesDeletionType = this.config.get<string>('jobManagement.ingestion.tasks.tilesDeletion.type');
this.tilesStorageProvider = this.config.get<StorageProvider>('tilesStorageProvider');
this.tilesStorageProvider = this.config.get<Exclude<StorageProvider, 'REDIS'>>('tilesStorageProvider');
this.tilesBucketConfig = this.config.get<string>('S3.tilesBucket');
this.tilesSubPathConfig = this.config.get<string>('storage.internalPvc.tilesSubPath');
}
Expand Down Expand Up @@ -136,7 +137,7 @@ export class DeleteLayerHandler extends JobHandler implements IJobHandler<never,
throw new LayerCacheNotFoundError(layerName, this.tilesStorageProvider);
}

if (this.tilesStorageProvider !== SourceType.S3) {
if (this.tilesStorageProvider !== StorageProvider.S3) {
return { path };
}
return { path, bucket: this.resolveTilesBucket(cache?.cache.bucket_name, layerName) };
Expand All @@ -148,7 +149,7 @@ export class DeleteLayerHandler extends JobHandler implements IJobHandler<never,
* - FS: the first segment is mapproxy's tiles-PVC mount path; the Cleaner remounts the same PVC at its own base, so it is dropped.
*/
private toRelativeTilesPath(directory: string): string {
return this.tilesStorageProvider === SourceType.S3 ? directory.replace(/^\/+/, '') : directory.replace(/^\/?[^/]+\//, '');
return this.tilesStorageProvider === StorageProvider.S3 ? directory.replace(/^\/+/, '') : directory.replace(/^\/?[^/]+\//, '');
}

private resolveTilesBucket(cacheBucket: string | undefined, layerName: LayerName): string {
Expand All @@ -164,7 +165,7 @@ export class DeleteLayerHandler extends JobHandler implements IJobHandler<never,

// bucket is always resolved for S3 in resolveTilesLocation; the fallback only satisfies the optional TilesLocation.bucket type
const storage: Storage =
this.tilesStorageProvider === SourceType.S3
this.tilesStorageProvider === StorageProvider.S3
? { storageProvider: this.tilesStorageProvider, bucket: tilesLocation.bucket ?? this.tilesBucketConfig }
: { storageProvider: this.tilesStorageProvider, subPath: this.tilesSubPathConfig };

Expand Down
10 changes: 2 additions & 8 deletions src/job/models/export/exportJobHandler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,19 +11,13 @@ import {
ExportFinalizeType,
RasterLayerMetadata,
SourceType,
StorageProvider,
} from '@map-colonies/raster-shared';
import { type Logger } from '@map-colonies/js-logger';
import { context, trace, type Tracer } from '@opentelemetry/api';
import { ArtifactRasterType } from '@map-colonies/types';
import { OperationStatus, TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue';
import {
EXPORT_FAILURE_MESSAGE,
EXPORT_SUCCESS_MESSAGE,
GPKG_CONTENT_TYPE,
JSON_CONTENT_TYPE,
SERVICES,
StorageProvider,
} from '../../../common/constants';
import { EXPORT_FAILURE_MESSAGE, EXPORT_SUCCESS_MESSAGE, GPKG_CONTENT_TYPE, JSON_CONTENT_TYPE, SERVICES } from '../../../common/constants';
import { JobHandler } from '../jobHandler';
import { TaskMetrics } from '../../../utils/metrics/taskMetrics';
import type {
Expand Down
31 changes: 25 additions & 6 deletions src/task/models/deletionTaskManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,20 @@ import { feature as turfFeature, featureCollection as turfFeatureCollection, uni
import { inject, injectable } from 'tsyringe';
import type { Logger } from '@map-colonies/js-logger';
import { ShapefileChunkReader } from '@map-colonies/shapefile-reader';
import type { IntersectionFeatureCollection, IngestionValidationTaskParams, TilesDeletionParams } from '@map-colonies/raster-shared';
import type {
IntersectionFeatureCollection,
IngestionValidationTaskParams,
TilesDeletionParams,
FsStorage,
S3Storage,
} from '@map-colonies/raster-shared';
import { StorageProvider } from '@map-colonies/raster-shared';
import type { ICreateTaskBody, ITaskResponse } from '@map-colonies/mc-priority-queue';
import { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue';
import type { MultiPolygon, Polygon } from 'geojson';
import { NotFoundError, UnprocessableEntityError } from '@map-colonies/error-types';
import type { IConfig, BuildDeletionTaskParams } from '../../common/interfaces';
import { SERVICES, StorageProvider } from '../../common/constants';
import { SERVICES } from '../../common/constants';
import { TaskMetrics } from '../../utils/metrics/taskMetrics';
import { createChildSpan } from '../../common/tracing';
import { IngestionCreateTasksTask, IngestionUpdateCreateTasksJob } from '../../utils/zod/schemas/job.schema';
Expand All @@ -23,7 +30,7 @@ export class TileDeletionTaskManager {
private readonly tileBatchSize: number;
private readonly taskBatchSize: number;
private readonly taskType: string;
private readonly sourceProvider: StorageProvider;
private readonly tilesStorage: S3Storage | FsStorage;
private readonly shapefileReader: ShapefileChunkReader;

public constructor(
Expand All @@ -38,7 +45,7 @@ export class TileDeletionTaskManager {
this.tileBatchSize = this.config.get<number>('jobManagement.ingestion.tasks.tilesDeletion.tileBatchSize');
this.taskBatchSize = this.config.get<number>('jobManagement.ingestion.tasks.tilesDeletion.taskBatchSize');
this.taskType = this.config.get<string>('jobManagement.ingestion.tasks.tilesDeletion.type');
this.sourceProvider = this.config.get<StorageProvider>('tilesStorageProvider');
this.tilesStorage = this.resolveTilesStorage();
this.shapefileReader = new ShapefileChunkReader({
maxVerticesPerChunk: this.config.get<number>('shapefileReader.maxVerticesPerChunk'),
});
Expand Down Expand Up @@ -226,10 +233,10 @@ export class TileDeletionTaskManager {

for await (const batch of batches) {
const taskParameters: TilesDeletionParams = {
tilesPath: layerRelativePath,
...this.tilesStorage,
tilesRelativePath: layerRelativePath,
ranges: batch,
fileExtension: tileOutputFormat.toLowerCase(),
sourceProvider: this.sourceProvider,
};

yield {
Expand All @@ -243,6 +250,18 @@ export class TileDeletionTaskManager {
}
}

private resolveTilesStorage(): S3Storage | FsStorage {
const provider = this.config.get<StorageProvider>('tilesStorageProvider');

if (provider === StorageProvider.REDIS) {
throw new Error('Redis storage provider is not supported for tile deletion tasks');
}

return provider === StorageProvider.S3
? { storageProvider: provider, bucket: this.config.get<string>('S3.tilesBucket') }
: { storageProvider: provider, subPath: this.config.get<string>('storage.internalPvc.tilesSubPath') };
}

private async fetchValidationTask(jobId: string): Promise<ITaskResponse<IngestionValidationTaskParams>> {
const validationTaskType = this.config.get<string>('jobManagement.ingestion.tasks.validation.type');
const validationTasks = await this.queueClient.jobManagerClient.findTasks<IngestionValidationTaskParams>({
Expand Down
5 changes: 3 additions & 2 deletions src/task/models/exportTaskManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,20 +11,21 @@ import {
polygonSchema,
RasterLayerMetadata,
SourceType,
StorageProvider,
type RoiFeature,
type RoiFeatureCollection,
} from '@map-colonies/raster-shared';
import { BBox2d, bboxToTileRange, degreesPerPixelToZoomLevel, type ITileRange } from '@map-colonies/mc-utils';
import type { BBox, Feature, MultiPolygon, Polygon } from 'geojson';
import { SERVICES, StorageProvider } from '../../common/constants';
import { SERVICES } from '../../common/constants';
import type { IConfig, TaskSources, ZoomBoundsParameters } from '../../common/interfaces';
import type { ExportJob } from '../../utils/zod/schemas/job.schema';
import { createChildSpan } from '../../common/tracing';

@injectable()
export class ExportTaskManager {
private readonly allWorldBounds: BBox;
private readonly tilesProvider: SourceType;
private readonly tilesProvider: StorageProvider;
public constructor(
@inject(SERVICES.LOGGER) private readonly logger: Logger,
@inject(SERVICES.CONFIG) private readonly config: IConfig,
Expand Down
6 changes: 3 additions & 3 deletions src/task/models/tileMergeTaskManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { degreesPerPixelToZoomLevel, tileBatchGenerator, TileRanger } from '@map
import { bbox, feature } from '@turf/turf';
import { inject, injectable } from 'tsyringe';
import type { Logger } from '@map-colonies/js-logger';
import { type InputFiles } from '@map-colonies/raster-shared';
import { StorageProvider, type InputFiles } from '@map-colonies/raster-shared';
import type { ICreateTaskBody } from '@map-colonies/mc-priority-queue';
import { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue';
import type {
Expand All @@ -20,7 +20,7 @@ import type {
ZoomDefinitions,
FeatureTask,
} from '../../common/interfaces';
import { SERVICES, type StorageProvider } from '../../common/constants';
import { SERVICES } from '../../common/constants';
import { fileExtensionExtractor } from '../../utils/fileUtil';
import { TaskMetrics } from '../../utils/metrics/taskMetrics';
import { createChildSpan } from '../../common/tracing';
Expand All @@ -29,7 +29,7 @@ import { Grid } from '../../common/interfaces';

@injectable()
export class TileMergeTaskManager {
private readonly tilesStorageProvider: string;
private readonly tilesStorageProvider: StorageProvider;
private readonly tileBatchSize: number;
private readonly taskBatchSize: number;
private readonly taskType: string;
Expand Down
4 changes: 4 additions & 0 deletions tests/unit/mocks/configMock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@ const registerDefaultConfig = (): void => {
maxVerticesPerChunk: 2500,
},
ingestionSourcesDirPath: '/layerSources',
// eslint-disable-next-line @typescript-eslint/naming-convention
S3: {
tilesBucket: 'tiles-bucket',
},
tilesStorageProvider: 'FS',
gpkgStorageProvider: 'FS',
storage: {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,8 +1,15 @@
/* eslint-disable @typescript-eslint/unbound-method */
import { NotFoundError } from '@map-colonies/error-types';
import type { TaskBlockDuplicationParam } from '@map-colonies/raster-shared';
import type {
FsTilesDeletionParams,
S3TilesDeletionParams,
TaskBlockDuplicationParam,
TileRange,
TilesDeletionParams,
} from '@map-colonies/raster-shared';
import type { ICreateTaskBody } from '@map-colonies/mc-priority-queue';
import type { Polygon } from 'geojson';
import { configMock, registerDefaultConfig } from '../../mocks/configMock';
import { configMock, registerDefaultConfig, setValue } from '../../mocks/configMock';
import { ingestionUpdateJob } from '../../mocks/jobsMockData';
import {
createFakeTask,
Expand Down Expand Up @@ -191,6 +198,85 @@ describe('TileDeletionTaskManager', () => {
tileDeletionTaskManager.buildAndPushTasks(ingestionUpdateJob, task, polygonPartsEntityName, layerRelativePath)
).resolves.not.toThrow();
});

describe('tiles storage on the built task parameters', () => {
// a single conflict feature that yields at least one tile batch
const conflictGeometry: Polygon = {
type: 'Polygon',
coordinates: [
[
[0, 0],
[1, 0],
[1, 1],
[0, 1],
[0, 0],
],
],
};

const arrangeSingleIntersection = (): void => {
jobManagerClientMock.findTasks.mockResolvedValue([validationTaskWithResolutionErrors]);
vi.spyOn(reportUtils, 'readConflictFeatures').mockResolvedValue([
// eslint-disable-next-line @typescript-eslint/naming-convention
{ type: 'Feature', geometry: conflictGeometry, properties: { e_res: 'Resolution Conflict' } },
]);
// one zoom level with an intersection, then empty so the upward iteration stops
polygonPartsMangerClientMock.getIntersection
.mockResolvedValueOnce({ type: 'FeatureCollection', features: [{ type: 'Feature', geometry: conflictGeometry, properties: {} }] })
.mockResolvedValue({ type: 'FeatureCollection', features: [] });
};

const getTasksParameters = (): TilesDeletionParams[] => {
const batches = jobManagerClientMock.createTaskForJob.mock.calls.map(([, batch]) => batch) as ICreateTaskBody<TilesDeletionParams>[][];
return batches.flat().map((task) => task.parameters);
};

it('should carry the configured subPath on every task when tiles are stored on FS', async () => {
setValue('tilesStorageProvider', 'FS');
setValue('storage.internalPvc.tilesSubPath', 'raster/artifacts/tiles');
const { tileDeletionTaskManager } = await setupTileDeletionTaskManagerTest();
arrangeSingleIntersection();

await tileDeletionTaskManager.buildAndPushTasks(ingestionUpdateJob, task, polygonPartsEntityName, layerRelativePath);

const tasksParameters = getTasksParameters();

expect(tasksParameters.length).toBeGreaterThan(0);

tasksParameters.forEach((parameters) => {
expect(parameters).toStrictEqual({
storageProvider: 'FS',
subPath: 'raster/artifacts/tiles',
tilesRelativePath: layerRelativePath,
fileExtension: ingestionUpdateJob.parameters.additionalParams.tileOutputFormat.toLowerCase(),
ranges: expect.any(Array) as TileRange[],
} satisfies FsTilesDeletionParams);
});
});

it('should carry the configured bucket on every task when tiles are stored on S3', async () => {
setValue('tilesStorageProvider', 'S3');
setValue('S3.tilesBucket', 'tiles-bucket');
const { tileDeletionTaskManager } = await setupTileDeletionTaskManagerTest();
arrangeSingleIntersection();

await tileDeletionTaskManager.buildAndPushTasks(ingestionUpdateJob, task, polygonPartsEntityName, layerRelativePath);

const tasksParameters = getTasksParameters();

expect(tasksParameters.length).toBeGreaterThan(0);

tasksParameters.forEach((parameters) => {
expect(parameters).toStrictEqual({
storageProvider: 'S3',
bucket: 'tiles-bucket',
tilesRelativePath: layerRelativePath,
fileExtension: ingestionUpdateJob.parameters.additionalParams.tileOutputFormat.toLowerCase(),
ranges: expect.any(Array) as TileRange[],
} satisfies S3TilesDeletionParams);
});
});
});
});

describe('fetchValidationTask (via buildAndPushTasks)', () => {
Expand Down
Loading