diff --git a/.vscode/settings.json b/.vscode/settings.json index 009aae3..8644b0a 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -1,5 +1,5 @@ { - "typescript.tsdk": "node_modules/typescript/lib", + "js/ts.tsdk.path": "node_modules/typescript/lib", "files.eol": "\n", "editor.defaultFormatter": "esbenp.prettier-vscode", "files.insertFinalNewline": true diff --git a/config/custom-environment-variables.json b/config/custom-environment-variables.json index 11aaa93..86c410e 100644 --- a/config/custom-environment-variables.json +++ b/config/custom-environment-variables.json @@ -73,24 +73,44 @@ "__format": "boolean" } }, - "s3": { - "endpoint": "S3_ENDPOINT", - "accessKeyId": "S3_ACCESS_KEY_ID", - "secretAccessKey": "S3_SECRET_ACCESS_KEY", - "sslEnabled": { - "__name": "S3_SSL_ENABLED", - "__format": "boolean" + "storage": { + "cleanupStorageProviders": { + "__name": "CLEANUP_STORAGE_PROVIDERS", + "__format": "json" }, - "forcePathStyle": { - "__name": "S3_FORCE_PATH_STYLE", - "__format": "boolean" + "s3": { + "delete": { + "batchSize": { + "__name": "S3_DELETE_BATCH_SIZE", + "__format": "number" + } + }, + "endpoint": "S3_ENDPOINT", + "accessKeyId": "S3_ACCESS_KEY_ID", + "secretAccessKey": "S3_SECRET_ACCESS_KEY", + "sslEnabled": { + "__name": "S3_SSL_ENABLED", + "__format": "boolean" + }, + "forcePathStyle": { + "__name": "S3_FORCE_PATH_STYLE", + "__format": "boolean" + }, + "region": "S3_REGION" }, - "region": "S3_REGION" + "fs": { + "delete": { + "batchSize": { + "__name": "FS_DELETE_BATCH_SIZE", + "__format": "number" + } + }, + "basePath": "FS_BASE_PATH" + } }, "strategies": { "tilesDeletion": { "s3Bucket": "TILES_DELETION_S3_BUCKET", - "fsBasePath": "TILES_DELETION_FS_BASE_PATH", "batchSize": { "__name": "TILES_DELETION_BATCH_SIZE", "__format": "number" @@ -103,6 +123,12 @@ "__name": "TILES_DELETION_FAILURE_SAMPLE_SIZE", "__format": "number" } + }, + "storedResourcesDeletion": { + "failureSampleSize": { + "__name": "STORED_RESOURCES_DELETION_FAILURE_SAMPLE_SIZE", + "__format": "number" + } } } } diff --git a/config/default.json b/config/default.json index 901d821..2dd1972 100644 --- a/config/default.json +++ b/config/default.json @@ -9,9 +9,7 @@ "level": "info", "prettyPrint": false, "opentelemetryOptions": { - "enabled": false, - "url": "", - "resourceAttributes": {} + "enabled": false } } }, @@ -28,18 +26,26 @@ { "job": "Ingestion_Swap_Update", "task": "tiles-deletion" + }, + { + "job": "Delete_Layer", + "task": "tiles-deletion" + }, + { + "job": "Delete_Layer", + "task": "artifacts-deletion" } ] } }, "queue": { - "jobManagerBaseUrl": "http://localhost:8080", - "heartbeatBaseUrl": "http://localhost:8081", + "jobManagerBaseUrl": "http://localhost:8081", + "heartbeatBaseUrl": "http://localhost:8082", "heartbeatIntervalMs": 1000, "dequeueIntervalMs": 3000 }, "servicesUrl": { - "jobTracker": "http://localhost:8082" + "jobTracker": "http://localhost:8083" }, "disableHttpClientLogs": false, "jobDefinitions": { @@ -49,12 +55,23 @@ }, "swapUpdate": { "type": "Ingestion_Swap_Update" + }, + "deleteLayer": { + "type": "Delete_Layer" } }, "tasks": { "tilesDeletion": { "type": "tiles-deletion", "maxAttempts": 3 + }, + "layerDeletion": { + "type": "tiles-deletion", + "maxAttempts": 3 + }, + "artifactsDeletion": { + "type": "artifacts-deletion", + "maxAttempts": 3 } } }, @@ -63,21 +80,35 @@ "delay": "exponential", "shouldResetTimeout": true }, - "s3": { - "endpoint": "http://localhost:9000", - "accessKeyId": "minioadmin", - "secretAccessKey": "minioadmin", - "sslEnabled": false, - "forcePathStyle": true, - "region": "us-east-1" + "storage": { + "cleanupStorageProviders": ["FS", "S3"], + "s3": { + "delete": { + "batchSize": 1000 + }, + "endpoint": "http://localhost:9000", + "accessKeyId": "minioadmin", + "secretAccessKey": "minioadmin", + "sslEnabled": false, + "forcePathStyle": true, + "region": "us-east-1" + }, + "fs": { + "delete": { + "batchSize": 1000 + }, + "basePath": "/tiles" + } }, "strategies": { "tilesDeletion": { "s3Bucket": "", - "fsBasePath": "/tiles", "batchSize": 1000, "concurrency": 10, "failureSampleSize": 3 + }, + "storedResourcesDeletion": { + "failureSampleSize": 3 } } } diff --git a/helm/templates/configmap.yaml b/helm/templates/configmap.yaml index 855dbaa..1860c39 100644 --- a/helm/templates/configmap.yaml +++ b/helm/templates/configmap.yaml @@ -2,7 +2,8 @@ {{- $serviceUrls := fromYaml (include "common.serviceUrls.merged" .) -}} {{- $storage := fromYaml (include "common.storage.merged" .) -}} {{- $s3 := ($storage.s3) | default dict -}} -{{- $internalPvc := (($storage.fs).internalPvc) | default dict -}} +{{- $fs := ($storage.fs) | default dict -}} +{{- $internalPvc := (($fs).internalPvc) | default dict -}} {{- $tilesFSBasePath := clean (printf "/%s/%s" $internalPvc.mountPath ($internalPvc.tilesSubPath | default "tiles")) -}} {{- if .Values.enabled -}} apiVersion: v1 @@ -54,15 +55,25 @@ data: {{- with .Values.env.jobnik.worker }} JOBNIK_WORKER_CONCURRENCY: {{ .concurrency | default 1 | quote }} {{- end }} + CLEANUP_STORAGE_PROVIDERS: {{ $storage.cleanupStorageProviders | toJson | quote }} + {{- if has "S3" $storage.cleanupStorageProviders }} + S3_DELETE_BATCH_SIZE: {{ $s3.delete.batchSize | default 1000 | quote }} S3_ENDPOINT: {{ if $s3.endpointUrl }}{{ printf "%s://%s" ($s3.sslEnabled | ternary "https" "http") $s3.endpointUrl | quote }}{{ else }}{{ "" | quote }}{{ end }} S3_FORCE_PATH_STYLE: {{ $s3.forcePathStyle | default false | quote }} S3_SSL_ENABLED: {{ $s3.sslEnabled | default false | quote }} S3_REGION: {{ $s3.region | default "us-east-1" | quote }} + {{- end }} + {{- if has "FS" $storage.cleanupStorageProviders }} + FS_DELETE_BATCH_SIZE: {{ $fs.delete.batchSize | default 1000 | quote }} + FS_BASE_PATH: {{ $tilesFSBasePath | quote }} + {{- end }} TILES_DELETION_S3_BUCKET: {{ $s3.tilesBucket | default "" | quote }} - TILES_DELETION_FS_BASE_PATH: {{ $tilesFSBasePath | quote }} {{- with .Values.env.strategies.tilesDeletion }} TILES_DELETION_BATCH_SIZE: {{ .batchSize | default 1000 | quote }} TILES_DELETION_CONCURRENCY: {{ .concurrency | default 10 | quote }} TILES_DELETION_FAILURE_SAMPLE_SIZE: {{ .failureSampleSize | default 3 | quote }} {{- end }} + {{- with .Values.env.strategies.storedResourcesDeletion }} + STORED_RESOURCES_DELETION_FAILURE_SAMPLE_SIZE: {{ .failureSampleSize | default 3 | quote }} + {{- end }} {{- end }} diff --git a/helm/templates/deployment.yaml b/helm/templates/deployment.yaml index ec4d9ff..0b8be5f 100644 --- a/helm/templates/deployment.yaml +++ b/helm/templates/deployment.yaml @@ -54,10 +54,10 @@ spec: imagePullPolicy: {{ .pullPolicy | default "IfNotPresent" }} {{- end }} {{- if .Values.command }} - command: + command: {{- toYaml .Values.command | nindent 12 }} {{- if .Values.args }} - args: + args: {{- toYaml .Values.args | nindent 12 }} {{- end }} {{- end }} @@ -122,7 +122,7 @@ spec: httpGet: path: {{ .Values.readinessProbe.path }} port: {{ .Values.env.targetPort }} - {{- end }} + {{- end }} {{- if .Values.resources.enabled }} resources: {{- toYaml .Values.resources.value | nindent 12 }} diff --git a/helm/values.yaml b/helm/values.yaml index be16b49..2203dec 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -13,7 +13,12 @@ serviceUrls: jobTracker: "" storage: + cleanupStorageProviders: + - FS + - S3 s3: + delete: + batchSize: 1000 endpointUrl: "" forcePathStyle: true secretName: "" @@ -21,6 +26,8 @@ storage: region: "" tilesBucket: "" fs: + delete: + batchSize: 1000 internalPvc: enabled: false name: "" @@ -128,6 +135,8 @@ env: batchSize: 1000 concurrency: 10 failureSampleSize: 3 + storedResourcesDeletion: + failureSampleSize: 3 resources: enabled: true diff --git a/package-lock.json b/package-lock.json index 4d16ca5..915d6fd 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +16,7 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-priority-queue": "^9.1.0", "@map-colonies/mc-utils": "^5.1.0", - "@map-colonies/raster-shared": "^8.1.0-alpha.3", + "@map-colonies/raster-shared": "^8.3.0-alpha.1", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.20.0", "@map-colonies/telemetry": "^10.0.1", @@ -2543,9 +2543,9 @@ "license": "ISC" }, "node_modules/@map-colonies/raster-shared": { - "version": "8.1.0-alpha.3", - "resolved": "https://registry.npmjs.org/@map-colonies/raster-shared/-/raster-shared-8.1.0-alpha.3.tgz", - "integrity": "sha512-igF8yhnSeXaUdpx6gW3C5gi4AM9J56jVwgqhh7+Zv0mYb/yBhxv0xsv4Mhyv+41UpFCVduYf90DJZRrY29kLpQ==", + "version": "8.3.0-alpha.1", + "resolved": "https://registry.npmjs.org/@map-colonies/raster-shared/-/raster-shared-8.3.0-alpha.1.tgz", + "integrity": "sha512-KP/n3ZPcG7ULHrtX3h9fMOTNnF5Qkos1N9ZPjawi/5yiDHit9S5V9sFI4Tgv5aMchgSS0zXu5gL0ybA4/CliYQ==", "license": "ISC", "dependencies": { "@map-colonies/mc-priority-queue": "^9.1.0", diff --git a/package.json b/package.json index d71fedb..d15fb70 100644 --- a/package.json +++ b/package.json @@ -16,7 +16,7 @@ "prebuild": "npm run clean", "build": "tsc --project tsconfig.build.json && tsc-alias -p tsconfig.build.json && npm run assets:copy", "start": "npm run build && cd dist && node --import ./instrumentation.mjs ./index.js", - "start:dev": "npm run build && cd dist && cross-env CONFIG_OFFLINE_MODE=true node --enable-source-maps --import ./instrumentation.mjs ./index.js", + "start:dev": "npm run build && cd dist && cross-env CONFIG_OFFLINE_MODE=true node --enable-source-maps --import ./instrumentation.mjs ./index.js", "assets:copy": "copyfiles -f ./config/* ./dist/config && copyfiles ./package.json dist", "clean": "rimraf dist", "prepare": "node .husky/install.mjs" @@ -37,7 +37,7 @@ "@map-colonies/js-logger": "^5.0.0", "@map-colonies/mc-priority-queue": "^9.1.0", "@map-colonies/mc-utils": "^5.1.0", - "@map-colonies/raster-shared": "^8.1.0-alpha.3", + "@map-colonies/raster-shared": "^8.3.0-alpha.1", "@map-colonies/read-pkg": "^1.0.0", "@map-colonies/schemas": "^1.20.0", "@map-colonies/telemetry": "^10.0.1", diff --git a/scripts/simulate-deletion.mjs b/scripts/simulate-deletion.mjs index 59ef463..5d8a19a 100644 --- a/scripts/simulate-deletion.mjs +++ b/scripts/simulate-deletion.mjs @@ -1,5 +1,5 @@ /** - * Simulation script for tiles-deletion tasks. + * Simulation script for tiles-deletion and artifacts-deletion tasks. * * Usage: * node scripts/simulate-deletion.mjs --provider S3 @@ -9,6 +9,15 @@ * node scripts/simulate-deletion.mjs --provider S3 --partial # multi-zoom partial deletion * node scripts/simulate-deletion.mjs --provider FS --real-tiles --source-tile scripts/tile_deletion_test.jpeg * node scripts/simulate-deletion.mjs --provider S3 --real-tiles --source-tile scripts/tile_deletion_test.jpeg --zooms 17,18,19,20 --tile-count 400 + * node scripts/simulate-deletion.mjs --provider S3 --artifacts-deletion + * node scripts/simulate-deletion.mjs --provider FS --artifacts-deletion --paths config,gpkg,reports + * + * --artifacts-deletion: simulate the Delete_Layer / artifacts-deletion task + * (DeleteStoredResourcesStrategy), which recursively deletes whole resource + * trees (e.g. gpkg, config, validation reports) rather than a tile range. + * --paths : comma-separated relative paths to seed & delete (default: config,gpkg,reports) + * Each path is seeded with a couple of nested fake files to prove the whole + * subtree — not just its top-level file — gets removed. * * --real-tiles: seed a real local tile file replicated across a multi-zoom grid. * --source-tile : local file to use as tile content for every seeded tile (required) @@ -42,7 +51,7 @@ * * Override defaults with env vars: * QUEUE_JOB_MANAGER_BASE_URL, S3_ENDPOINT, S3_ACCESS_KEY_ID, S3_SECRET_ACCESS_KEY, - * TILES_DELETION_S3_BUCKET, TILES_DELETION_FS_BASE_PATH, TILES_PATH, ZOOM, MIN_X, MAX_X, MIN_Y, MAX_Y, + * TILES_DELETION_S3_BUCKET, FS_BASE_PATH, TILES_PATH, ARTIFACTS_PATH, ZOOM, MIN_X, MAX_X, MIN_Y, MAX_Y, * SEED_CONCURRENCY (default: 50) */ @@ -69,6 +78,7 @@ Modes (mutually exclusive): --partial Seed tiles across 3 zoom levels; task only deletes a subset --real-tiles Use a real tile file instead of fake content --skip-seed Skip seeding; go straight to job creation (idempotent delete) + --artifacts-deletion Simulate the Delete_Layer / artifacts-deletion task (whole-resource deletion) Real-tiles options (require --real-tiles): --source-tile Local tile file to replicate across the grid (required) @@ -76,14 +86,18 @@ Real-tiles options (require --real-tiles): --zooms Zoom levels, comma-separated (default: 17,18,19,20) --tile-count Total tiles to seed across all zoom levels (default: 400) +Artifacts-deletion options (require --artifacts-deletion): + --paths Relative resource paths to seed & delete, comma-separated (default: config,gpkg,reports) + Env vars (all optional — pod ConfigMap values are used automatically): QUEUE_JOB_MANAGER_BASE_URL Job manager endpoint S3_ENDPOINT S3 endpoint URL S3_ACCESS_KEY_ID S3 access key S3_SECRET_ACCESS_KEY S3 secret key - TILES_DELETION_S3_BUCKET S3 bucket name - TILES_DELETION_FS_BASE_PATH FS base path for tile files + TILES_DELETION_S3_BUCKET S3 bucket name (tiles-deletion) / S3 bucket for artifacts-deletion + FS_BASE_PATH FS base path for stored files (tiles and artifacts) TILES_PATH Relative path prefix for tiles (default: simulate/layer/v1) + ARTIFACTS_PATH Relative base path for artifacts-deletion resources (default: simulate/artifacts) ZOOM Zoom level (default: 10) MIN_X, MAX_X X tile range (default: 0..3) MIN_Y, MAX_Y Y tile range (default: 0..3) @@ -94,6 +108,7 @@ Examples: node scripts/simulate-deletion.mjs --provider FS --partial node scripts/simulate-deletion.mjs --provider S3 --real-tiles --source-tile scripts/tile_deletion_test.jpeg node scripts/simulate-deletion.mjs --provider FS --real-tiles --source-tile scripts/tile_deletion_test.jpeg + node scripts/simulate-deletion.mjs --provider S3 --artifacts-deletion MAX_X=999 MAX_Y=999 ZOOM=18 node scripts/simulate-deletion.mjs --provider S3 --skip-seed `); process.exit(0); @@ -105,6 +120,7 @@ if (!providerFlag || !['S3', 'FS'].includes(providerFlag)) { console.error( ' node scripts/simulate-deletion.mjs --provider --real-tiles --source-tile [--zooms ] [--tile-count ]' ); + console.error(' node scripts/simulate-deletion.mjs --provider --artifacts-deletion [--paths ]'); console.error('\nRun with --help for full usage information.'); process.exit(1); } @@ -112,12 +128,16 @@ const PROVIDER = providerFlag; const SKIP_SEED = args.includes('--skip-seed'); const PARTIAL = args.includes('--partial'); const REAL_TILES = args.includes('--real-tiles'); +const ARTIFACTS_DELETION = args.includes('--artifacts-deletion'); // real-tiles specific args const SOURCE_TILE = args.includes('--source-tile') ? args[args.indexOf('--source-tile') + 1] : null; const ZOOMS_INPUT = args.includes('--zooms') ? args[args.indexOf('--zooms') + 1] : '17,18,19,20'; const TILE_COUNT = args.includes('--tile-count') ? Number(args[args.indexOf('--tile-count') + 1]) : 400; +// artifacts-deletion specific args +const ARTIFACT_PATHS_INPUT = args.includes('--paths') ? args[args.indexOf('--paths') + 1] : 'config,gpkg,reports'; + if (REAL_TILES) { if (!SOURCE_TILE) { console.error('--real-tiles requires --source-tile '); @@ -129,6 +149,11 @@ if (REAL_TILES) { } } +if (ARTIFACTS_DELETION && (PARTIAL || REAL_TILES)) { + console.error('--artifacts-deletion cannot be combined with --partial or --real-tiles'); + process.exit(1); +} + // ─── Read local.json as config base (env vars override) ────────────────────── let localConfig = {}; @@ -141,7 +166,8 @@ try { } const cfg = { - s3: localConfig.s3 ?? {}, + s3: localConfig.storage?.s3 ?? {}, + fs: localConfig.storage?.fs ?? {}, strategies: localConfig.strategies?.tilesDeletion ?? {}, queue: localConfig.queue ?? {}, }; @@ -154,7 +180,7 @@ const S3_ENDPOINT = process.env.S3_ENDPOINT ?? cfg.s3.endpoint ?? 'http://localh const S3_ACCESS_KEY_ID = process.env.S3_ACCESS_KEY_ID ?? cfg.s3.accessKeyId ?? 'minioadmin'; const S3_SECRET_ACCESS_KEY = process.env.S3_SECRET_ACCESS_KEY ?? cfg.s3.secretAccessKey ?? 'minioadmin'; const S3_BUCKET = process.env.TILES_DELETION_S3_BUCKET ?? cfg.strategies.s3Bucket ?? ''; -const FS_BASE_PATH = process.env.TILES_DELETION_FS_BASE_PATH ?? cfg.strategies.fsBasePath ?? '/tiles'; +const FS_BASE_PATH = process.env.FS_BASE_PATH ?? cfg.fs.basePath ?? '/tiles'; const SEED_CONCURRENCY = Number(process.env.SEED_CONCURRENCY ?? 200); // Tile range to seed + delete @@ -166,6 +192,11 @@ const MIN_Y = Number(process.env.MIN_Y ?? 0); const MAX_Y = Number(process.env.MAX_Y ?? 3); const FILE_EXTENSION = 'jpeg'; +// Artifacts-deletion resource paths to seed + delete +const ARTIFACTS_PATH = process.env.ARTIFACTS_PATH ?? 'simulate/artifacts'; +const ARTIFACT_SUBPATHS = ARTIFACT_PATHS_INPUT.split(',').map((s) => s.trim()); +const ARTIFACT_PATHS = ARTIFACT_SUBPATHS.map((subPath) => `${ARTIFACTS_PATH}/${subPath}`); + // ─── Partial-deletion scenario definition ──────────────────────────────────── // // Three zoom levels are seeded identically (4×4 grid by default). @@ -527,9 +558,125 @@ async function runRealTilesMode() { await createRealTilesJob(zooms, grid, TILES_PATH, ext); } +// ─── Artifacts-deletion mode (Delete_Layer / artifacts-deletion task) ───────── +// +// Unlike tiles-deletion (which deletes a tile-range grid), artifacts-deletion +// recursively removes whole resource trees given their root paths (e.g. a +// layer's gpkg, config and validation-report directories). Each seeded path +// gets a couple of nested files to prove the whole subtree is removed, not +// just its top-level entry. + +function* artifactFilesForPath(rootPath) { + yield `${rootPath}/file.dat`; + yield `${rootPath}/nested/file.dat`; +} + +async function seedS3Artifacts() { + if (!S3_BUCKET) { + throw new Error('S3_BUCKET env var is required for S3 provider (or set it in config/local.json)'); + } + + const client = new S3Client({ + endpoint: S3_ENDPOINT, + credentials: { accessKeyId: S3_ACCESS_KEY_ID, secretAccessKey: S3_SECRET_ACCESS_KEY }, + forcePathStyle: true, + region: 'us-east-1', + tls: false, + }); + + console.log(`[S3] Seeding artifacts under s3://${S3_BUCKET}/{${ARTIFACT_PATHS.join(', ')}}/...`); + for (const rootPath of ARTIFACT_PATHS) { + for (const key of artifactFilesForPath(rootPath)) { + await client.send( + new PutObjectCommand({ Bucket: S3_BUCKET, Key: key, Body: Buffer.from('fake-artifact'), ContentType: 'application/octet-stream' }) + ); + } + } + const firstKey = [...artifactFilesForPath(ARTIFACT_PATHS[0])][0]; + await client.send(new HeadObjectCommand({ Bucket: S3_BUCKET, Key: firstKey })); + console.log(`[S3] Verified: s3://${S3_BUCKET}/${firstKey} exists`); + console.log(`[S3] Done — seeded ${ARTIFACT_PATHS.length} artifact path(s).\n`); +} + +function seedFsArtifacts() { + console.log(`[FS] Seeding artifacts under ${FS_BASE_PATH}/{${ARTIFACT_PATHS.join(', ')}}/...`); + for (const rootPath of ARTIFACT_PATHS) { + for (const relativePath of artifactFilesForPath(rootPath)) { + const fullPath = join(FS_BASE_PATH, relativePath); + mkdirSync(join(fullPath, '..'), { recursive: true }); + writeFileSync(fullPath, 'fake-artifact'); + } + } + const firstPath = join(FS_BASE_PATH, [...artifactFilesForPath(ARTIFACT_PATHS[0])][0]); + if (!existsSync(firstPath)) throw new Error(`Seeding failed — ${firstPath} not found`); + console.log(`[FS] Verified: ${firstPath} exists`); + console.log(`[FS] Done — seeded ${ARTIFACT_PATHS.length} artifact path(s).\n`); +} + +async function createArtifactsDeletionJob() { + const taskParameters = { + storageProvider: PROVIDER, + paths: ARTIFACT_PATHS, + ...(PROVIDER === 'S3' && { bucket: S3_BUCKET }), + }; + + const body = { + resourceId: `simulate-artifacts-${PROVIDER.toLowerCase()}-${Date.now()}`, + version: '1.0.0', + type: 'Delete_Layer', + parameters: {}, + domain: 'RASTER', + tasks: [{ type: 'artifacts-deletion', parameters: taskParameters }], + }; + + console.log('[Job] Creating job with task parameters:'); + console.log(JSON.stringify(taskParameters, null, 2)); + + const res = await fetch(`${JOB_MANAGER_URL}/jobs`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(body), + }); + + if (!res.ok) { + const text = await res.text(); + throw new Error(`Job creation failed [${res.status}]: ${text}`); + } + + const { id: jobId, taskIds } = await res.json(); + + console.log(`\n[Job] Created successfully:`); + console.log(` Job ID : ${jobId}`); + console.log(` Task ID : ${taskIds[0]}`); + console.log(`\nStart the cleaner — it will pick up task "${taskIds[0]}" and delete ${ARTIFACT_PATHS.length} artifact path(s).`); + console.log(`Track progress: GET ${JOB_MANAGER_URL}/jobs/${jobId}?shouldReturnTasks=true`); +} + +async function runArtifactsDeletionMode() { + if (!SKIP_SEED) { + if (PROVIDER === 'S3') { + await seedS3Artifacts(); + } else { + seedFsArtifacts(); + } + } else { + console.log(`[seed] Skipped — missing paths are treated as success (idempotent delete).\n`); + } + + await createArtifactsDeletionJob(); +} + // ─── Main ───────────────────────────────────────────────────────────────────── async function main() { + if (ARTIFACTS_DELETION) { + console.log( + `\n=== Simulating artifacts-deletion | storageProvider: ${PROVIDER} | paths: ${ARTIFACT_PATHS.length} | seed: ${SKIP_SEED ? 'skipped' : 'yes'} ===\n` + ); + await runArtifactsDeletionMode(); + return; + } + if (REAL_TILES) { const zooms = parseZoomList(ZOOMS_INPUT); const ext = extname(SOURCE_TILE).slice(1).toLowerCase() || FILE_EXTENSION; diff --git a/src/cleaner/errors/errors.ts b/src/cleaner/errors/errors.ts index a3efe73..5c10557 100644 --- a/src/cleaner/errors/errors.ts +++ b/src/cleaner/errors/errors.ts @@ -89,8 +89,8 @@ export class ValidationError extends UnrecoverableError { * This is always unrecoverable since the strategy won't appear on retry. */ export class StrategyNotFoundError extends UnrecoverableError { - public constructor(taskType: string) { - super(`No strategy registered for task type: ${taskType}`); + public constructor({ jobType, taskType }: { jobType: string; taskType: string }) { + super(`No strategy registered for job type: ${jobType} and task type: ${taskType}`); this.name = StrategyNotFoundError.name; Error.captureStackTrace(this, this.constructor); } diff --git a/src/cleaner/errors/index.ts b/src/cleaner/errors/index.ts index 0315ef3..330ee2f 100644 --- a/src/cleaner/errors/index.ts +++ b/src/cleaner/errors/index.ts @@ -1,2 +1,2 @@ -export { toError, describeError, RecoverableError, UnrecoverableError, ConfigurationError, ValidationError, StrategyNotFoundError } from './errors'; export { ErrorHandler } from './errorHandler'; +export { ConfigurationError, describeError, RecoverableError, StrategyNotFoundError, toError, UnrecoverableError, ValidationError } from './errors'; diff --git a/src/cleaner/storageProviders/fsStorageProvider.ts b/src/cleaner/storageProviders/fsStorageProvider.ts index 87734d5..32b81d8 100644 --- a/src/cleaner/storageProviders/fsStorageProvider.ts +++ b/src/cleaner/storageProviders/fsStorageProvider.ts @@ -1,11 +1,36 @@ +import { accessSync, constants, statSync } from 'node:fs'; +import { rm, rmdir, stat, unlink } from 'node:fs/promises'; import { join } from 'node:path'; -import { stat, unlink, rmdir } from 'node:fs/promises'; import type { Logger } from '@map-colonies/js-logger'; -import { describeError } from '../errors'; -import type { DeleteFailure, IStorageProvider } from './iStorageProvider'; +import type { DeleteStoredResourcesParams } from '@map-colonies/raster-shared'; +import type { DeleteFailure, DeleteResourcesResult, IStorageProvider, StorageProvider } from '@src/cleaner/storageProviders'; +import { getChunk, normalizeFolderPath, resolveAbsolutePath } from '@src/cleaner/utils'; +import type { ConfigType } from '@src/common/config'; +import { ConfigurationError, describeError, UnrecoverableError } from '../errors'; -export class FsStorageProvider implements IStorageProvider { - public constructor(private readonly logger: Logger) {} +type FSStorageProviderType = Extract; + +export interface FsConfig { + delete: { + batchSize: number; + }; + basePath: string; +} + +export class FsStorageProvider implements IStorageProvider<'FS'> { + private readonly fsConfig: FsConfig; + private readonly basePath: string; + + public constructor( + private readonly config: ConfigType, + private readonly logger: Logger + ) { + this.fsConfig = this.config.get('storage.fs') as unknown as FsConfig; + if (this.fsConfig.delete.batchSize <= 0) throw new ConfigurationError('Deletion batch size must be greater than 0'); + this.basePath = resolveAbsolutePath(this.fsConfig.basePath); + this.canDeleteFromFolder(this.basePath); + this.logger.debug(`Using ${this.basePath} as base path for FS`); + } public async targetExists(storageTarget: string, relativePath: string): Promise { try { @@ -18,26 +43,24 @@ export class FsStorageProvider implements IStorageProvider { } public async delete(paths: string[], storageTarget: string): Promise { - if (paths.length === 0) { - return []; - } - this.logger.info({ msg: 'Deleting files from filesystem', basePath: storageTarget, count: paths.length }); - const results = await Promise.allSettled( - paths.map(async (relativePath) => { - await unlink(join(storageTarget, relativePath)); - }) - ); - const failures: DeleteFailure[] = []; - for (const [idx, result] of results.entries()) { - if (result.status === 'rejected') { - const relativePath = paths[idx]!; - const error: unknown = result.reason; - const reason = describeError(error); - this.logger.debug({ msg: 'Failed to delete file', path: join(storageTarget, relativePath), reason, error }); - failures.push({ path: relativePath, reason }); + for (const relativePaths of getChunk(paths, this.fsConfig.delete.batchSize)) { + const results = await Promise.allSettled( + relativePaths.map(async (relativePath) => { + await unlink(join(storageTarget, relativePath)); + }) + ); + + for (const [idx, result] of results.entries()) { + if (result.status === 'rejected') { + const relativePath = relativePaths[idx]!; + const error: unknown = result.reason; + const reason = describeError(error); + this.logger.debug({ msg: 'Failed to delete file', path: join(storageTarget, relativePath), reason, error }); + failures.push({ path: relativePath, reason }); + } } } @@ -46,6 +69,64 @@ export class FsStorageProvider implements IStorageProvider { return failures; } + public async deleteResources({ + paths, + }: Extract): Promise { + // Prevent path traversal (i.e. accessing folders above root folder) + if (!this.checkPathTraversal(paths)) throw new UnrecoverableError(`Cannot delete files/folders outside base path or base path itself`); + + const failures: DeleteFailure[] = []; + for (const relativePaths of getChunk(paths, this.fsConfig.delete.batchSize)) { + const results = await Promise.allSettled( + relativePaths.map(async (relativePath) => { + await rm(join(this.basePath, relativePath), { recursive: true, force: true }); + }) + ); + + for (const [idx, result] of results.entries()) { + if (result.status === 'rejected') { + const fullPath = join(this.basePath, relativePaths[idx]!); + const reason = describeError(result.reason); + this.logger.error({ msg: 'Failed to delete layer directory', fullPath, reason, err: result.reason }); + failures.push({ path: fullPath, reason }); + } + } + } + + return { failures }; + } + + private canDeleteFromFolder(path: string): void { + try { + accessSync(path, constants.F_OK | constants.R_OK | constants.W_OK); + } catch (err) { + if (err instanceof Error && 'code' in err && err.code === 'ENOENT') { + throw new ConfigurationError(`FS path does not exist: ${path}`); + } else if (err instanceof Error && 'code' in err && (err.code === 'EACCES' || err.code === 'EPERM')) { + throw new ConfigurationError(`FS path permission denied for path: ${path}`); + } else { + throw new ConfigurationError(`An unexpected error occurred on FS path accessibility check: ${describeError(err)}`); + } + } + + try { + const pathStat = statSync(path); + if (!pathStat.isDirectory()) { + throw new ConfigurationError(`FS path exists but it is a file, not a directory: ${path}`); + } + } catch (err) { + if (err instanceof ConfigurationError) throw err; + throw new ConfigurationError(`An unexpected error occurred on FS info check: ${describeError(err)}`); + } + } + + private checkPathTraversal(paths: string[]): boolean { + return paths.every((path) => { + const absolutePath = resolveAbsolutePath(join(this.basePath, path)); + return absolutePath.startsWith(normalizeFolderPath(this.basePath)); + }); + } + // Attempts to remove any directories that became empty after file deletion. // Strategy: collect every ancestor directory of every deleted file, grouped by // depth relative to the file (levelIdx 0 = direct parent, 1 = grandparent, …). diff --git a/src/cleaner/storageProviders/iStorageProvider.ts b/src/cleaner/storageProviders/iStorageProvider.ts index e392459..f798f12 100644 --- a/src/cleaner/storageProviders/iStorageProvider.ts +++ b/src/cleaner/storageProviders/iStorageProvider.ts @@ -1,3 +1,5 @@ +import type { DeleteStoredResourcesParams } from '@map-colonies/raster-shared'; + /** * A single failed deletion paired with a short reason string (e.g. 'ENOENT', * 'AccessDenied', 'NoSuchKey'). @@ -7,7 +9,13 @@ export interface DeleteFailure { reason: string; } -export interface IStorageProvider { +export interface DeleteResourcesResult { + failures: DeleteFailure[]; +} + +export type StorageProvider = DeleteStoredResourcesParams['storageProvider']; + +export interface IStorageProvider { /** * Deletes a batch of relative file paths within the given storage target. * - S3: storageTarget = bucket name; paths are object keys @@ -18,6 +26,13 @@ export interface IStorageProvider { */ delete: (paths: string[], storageTarget: string) => Promise; + /** + * Deletes ALL objects/files under the given paths. + * - S3: storageTarget = bucket name; paths are root paths to resource(s) + * - FS: storageTarget = base directory; paths are relative paths from mounted dir + */ + deleteResources: (deleteStoredResourcesParams: Extract) => Promise; + /** * Returns true if relativePath exists within storageTarget and contains data. * - S3: storageTarget = bucket, relativePath = key prefix — lists objects (KeyCount > 0) @@ -25,3 +40,7 @@ export interface IStorageProvider { */ targetExists: (storageTarget: string, relativePath: string) => Promise; } + +export type StorageProviders = { + [T in StorageProvider]?: IStorageProvider; +}; diff --git a/src/cleaner/storageProviders/index.ts b/src/cleaner/storageProviders/index.ts index 54b39c1..4fb7851 100644 --- a/src/cleaner/storageProviders/index.ts +++ b/src/cleaner/storageProviders/index.ts @@ -1,4 +1,4 @@ -export type { IStorageProvider, DeleteFailure } from './iStorageProvider'; -export { S3StorageProvider } from './s3StorageProvider'; -export { FsStorageProvider } from './fsStorageProvider'; export { summarizeDeleteFailures, type DeleteFailureSummary } from './deleteFailureSummary'; +export { FsStorageProvider, type FsConfig } from './fsStorageProvider'; +export type { DeleteFailure, DeleteResourcesResult, IStorageProvider, StorageProvider, StorageProviders } from './iStorageProvider'; +export { S3StorageProvider } from './s3StorageProvider'; diff --git a/src/cleaner/storageProviders/s3StorageProvider.ts b/src/cleaner/storageProviders/s3StorageProvider.ts index d554a20..e7e4318 100644 --- a/src/cleaner/storageProviders/s3StorageProvider.ts +++ b/src/cleaner/storageProviders/s3StorageProvider.ts @@ -1,68 +1,99 @@ /* eslint-disable @typescript-eslint/naming-convention */ -import { S3Client, DeleteObjectsCommand, ListObjectsV2Command, NoSuchBucket } from '@aws-sdk/client-s3'; +import { + DeleteObjectsCommand, + HeadBucketCommand, + HeadObjectCommand, + ListObjectsV2Command, + NoSuchBucket, + NotFound, + paginateListObjectsV2, + S3Client, + S3ServiceException, + type _Object, +} from '@aws-sdk/client-s3'; import type { Logger } from '@map-colonies/js-logger'; +import type { DeleteStoredResourcesParams } from '@map-colonies/raster-shared'; import type { ConfigType } from '@common/config'; -import { describeError } from '../errors'; -import type { DeleteFailure, IStorageProvider } from './iStorageProvider'; +import type { DeleteFailure, DeleteResourcesResult, IStorageProvider, StorageProvider } from '@src/cleaner/storageProviders'; +import { getChunk, normalizeFolderPath } from '@src/cleaner/utils'; +import { ConfigurationError, describeError, UnrecoverableError } from '../errors'; -const S3_MAX_DELETE_BATCH = 1000; +type S3StorageProviderType = Extract; -interface S3Config { +// S3/MinIO reject DeleteObjects requests with more than 1000 keys, regardless of configured batch size. +const S3_DELETE_OBJECTS_MAX_KEYS = 1000; + +export interface S3Config { + delete: { + batchSize: number; + }; endpoint: string; accessKeyId: string; secretAccessKey: string; - sslEnabled: boolean; - forcePathStyle: boolean; - region: string; + sslEnabled?: boolean; + forcePathStyle?: boolean; + region?: string; } -export class S3StorageProvider implements IStorageProvider { +export class S3StorageProvider implements IStorageProvider { private readonly s3Client: S3Client; + private readonly s3Config: S3Config; public constructor( config: ConfigType, private readonly logger: Logger ) { - const s3Config = config.get('s3') as S3Config; + this.s3Config = config.get('storage.s3') as unknown as S3Config; + if (this.s3Config.delete.batchSize <= 0) throw new ConfigurationError('Deletion batch size must be greater than 0'); this.s3Client = new S3Client({ - endpoint: s3Config.endpoint, + endpoint: this.s3Config.endpoint, credentials: { - accessKeyId: s3Config.accessKeyId, - secretAccessKey: s3Config.secretAccessKey, + accessKeyId: this.s3Config.accessKeyId, + secretAccessKey: this.s3Config.secretAccessKey, }, - forcePathStyle: s3Config.forcePathStyle, - region: s3Config.region, - tls: s3Config.sslEnabled, + forcePathStyle: this.s3Config.forcePathStyle, + region: this.s3Config.region, + tls: this.s3Config.sslEnabled, }); } - public async delete(paths: string[], storageTarget: string): Promise { - if (paths.length === 0) { - return []; + public async delete(paths: string[], bucket: string): Promise { + this.logger.debug({ msg: 'Deleting objects from S3', bucket, count: paths.length }); + + const failures: DeleteFailure[] = []; + + for (const chunk of getChunk(paths, this.s3Config.delete.batchSize)) { + const failed = await this.deleteChunk(chunk, bucket); + failures.push(...failed); } - this.logger.debug({ msg: 'Deleting objects from S3', bucket: storageTarget, count: paths.length }); + return failures; + } + public async deleteResources({ + bucket, + paths, + }: Extract): Promise { const failures: DeleteFailure[] = []; + this.logger.debug({ msg: `Starting S3 deletion`, bucket, paths }); - for (const chunk of this.chunk(paths, S3_MAX_DELETE_BATCH)) { - try { - const failed = await this.deleteChunk(chunk, storageTarget); - failures.push(...failed); - } catch (error) { - // Whole chunk failed (network/auth/etc) — mark every path in it with the same reason - // so the caller still gets a per-path failure list and a human-readable cause. - const reason = describeError(error); - this.logger.error({ msg: 'S3 batch request failed', bucket: storageTarget, reason, error }); - failures.push(...chunk.map((path) => ({ path, reason }))); - } + if (paths.some((path) => path.length === 0)) throw new UnrecoverableError('Cannot delete resources directly under root path of the bucket'); // Prevent root deletion + + const exists = await this.bucketExists(bucket); + if (!exists) { + throw new UnrecoverableError(`Bucket does not exist: ${bucket}`); } - return failures; + for (const path of paths) { + const pathFailures = await this.deleteResource({ bucket, path }); + failures.push(...pathFailures); + } + + return { failures }; } public async targetExists(bucket: string, relativePath: string): Promise { - const prefix = relativePath.endsWith('/') ? relativePath : `${relativePath}/`; + const prefix = normalizeFolderPath(relativePath); try { const result = await this.s3Client.send(new ListObjectsV2Command({ Bucket: bucket, Prefix: prefix, MaxKeys: 1 })); return (result.KeyCount ?? 0) > 0; @@ -73,21 +104,178 @@ export class S3StorageProvider implements IStorageProvider { } } + private async bucketExists(bucket: string): Promise { + try { + this.logger.debug({ msg: `Checking bucket exists`, bucket }); + const command = new HeadBucketCommand({ Bucket: bucket }); + await this.s3Client.send(command); // If it resolves, the bucket exists and you have permission to access it + this.logger.debug({ msg: `Bucket exists`, bucket }); + return true; + } catch (err) { + if (err instanceof NotFound) { + this.logger.error({ msg: `Bucket does not exist`, bucket, err }); + return false; + } + const reason = describeError(err); + this.logger.error({ msg: 'Failed to check if bucket exists', bucket, reason, err }); + throw err; + } + } + private async deleteChunk(paths: string[], bucket: string): Promise { - const command = new DeleteObjectsCommand({ - Bucket: bucket, - Delete: { Objects: paths.map((Key) => ({ Key })) }, + const failures: DeleteFailure[] = []; + + // The configured batch size may exceed the DeleteObjects API limit, so re-chunk defensively here. + for (const keys of getChunk(paths, S3_DELETE_OBJECTS_MAX_KEYS)) { + const chunkFailures = await this.deleteObjects(keys, bucket); + failures.push(...chunkFailures); + } + + return failures; + } + + private async deleteObjects(paths: string[], bucket: string): Promise { + try { + const command = new DeleteObjectsCommand({ + Bucket: bucket, + Delete: { Objects: paths.map((Key) => ({ Key })) }, + }); + + const response = await this.s3Client.send(command); + const errors = (response.Errors ?? []) + .filter((error): error is typeof error & { Key: string } => Boolean(error.Key)) // Only include entries with a Key so every failure maps to a specific path. + .map((error) => ({ path: error.Key, reason: error.Code ?? error.Message ?? 'Unknown' })); + if (errors.length > 0) this.logger.warn({ msg: `Failed to delete ${errors.length} out of ${paths.length} objects` }); + return errors; + } catch (err) { + const reason = describeError(err); + this.logger.error({ msg: 'S3 batch request failed', bucket, reason, err }); + return paths.map((path) => { + return { path, reason }; + }); + } + } + + private async deleteResource({ bucket, path }: { bucket: string; path: string }): Promise { + const failures: DeleteFailure[] = []; + let totalDeletedObjectsCount = 0, + totalFailedObjectsCount = 0; + + const s3Objects = this.getS3Objects({ + bucket, + prefix: path, + pageSize: this.s3Config.delete.batchSize, }); - const response = await this.s3Client.send(command); - return (response.Errors ?? []) - .filter((e): e is typeof e & { Key: string } => Boolean(e.Key)) // Only include entries with a Key so every failure maps to a specific path. - .map((e) => ({ path: e.Key, reason: e.Code ?? e.Message ?? 'Unknown' })); + try { + for await (const pageOfObjects of s3Objects) { + this.logger.debug({ + msg: `Received a batch of ${pageOfObjects.length} objects to delete`, + pageOfObjects, + path, + pageSize: this.s3Config.delete.batchSize, + }); + const keys = pageOfObjects.map((obj) => obj.Key).filter((key): key is string => key !== undefined && this.matchesTarget(key, path)); + + if (keys.length === 0) { + continue; + } + + const chunkFailures = await this.deleteChunk(keys, bucket); + failures.push(...chunkFailures); + + const failedObjectsCount = chunkFailures.length; + const deletedObjectsCount = keys.length - failedObjectsCount; + totalDeletedObjectsCount += deletedObjectsCount; + totalFailedObjectsCount += failedObjectsCount; + + this.logger.debug({ + msg: `Successfully deleted ${deletedObjectsCount} objects. Totally ${totalDeletedObjectsCount} successfully deleted objects`, + }); + + if (failedObjectsCount > 0) + this.logger.debug({ + msg: `Could not delete ${failedObjectsCount} objects. Totally ${totalFailedObjectsCount} objects could not be deleted`, + }); + } + this.logger.debug({ msg: 'Deletion completed', path, totalDeletedObjectsCount, totalFailedObjectsCount }); + return failures; + } catch (err) { + this.logger.error({ + msg: 'Stream of objects was interrupted by an error', + path, + totalDeletedObjectsCount, + totalFailedObjectsCount, + err, + }); + throw err; + } + } + + private async *getS3Objects({ + bucket, + prefix, + pageSize, + }: { + bucket: string; + prefix?: string; + pageSize?: number; + }): AsyncGenerator<_Object[], void, unknown> { + try { + // First, if object exists it is removed. This is to mitigate an issue in MinIO that shadows paths sharing common path with an object. + // Second, objects having this path are iterated and removed + if (prefix !== undefined && (await this.resourceExists({ bucket, path: prefix }))) { + yield [{ Key: prefix }]; + } + + const paginatorConfig = { + client: this.s3Client, + pageSize, + }; + + const commandInput = { + Bucket: bucket, + Prefix: prefix !== undefined ? normalizeFolderPath(prefix) : prefix, + }; + + const paginator = paginateListObjectsV2(paginatorConfig, commandInput); + for await (const page of paginator) { + yield page.Contents ?? []; + } + } catch (err) { + if (err instanceof NoSuchBucket) { + this.logger.error({ msg: `S3 Error [${err.name}] no such bucket: ${err.message} (Req ID: ${err.$metadata.requestId})` }); + } else if (err instanceof S3ServiceException) { + this.logger.error({ msg: `S3 Error [${err.name}] occured during pagination: ${err.message} (Req ID: ${err.$metadata.requestId})` }); + } else { + this.logger.error({ msg: 'Unexpected error occurred during pagination', err, bucket, prefix }); + } + throw err; + } + } + + // A key belongs to the target if it IS the target object or lives under 'target/'. + // The 'target/' guard prevents matching sibling keys that merely share the prefix + // (e.g. target 'photos' must not match 'photos_old/img.jpg', target 'metadata.txt' + // must not match 'metadata.txt.bak'). + private matchesTarget(key: string, target: string): boolean { + return key === target || key.startsWith(normalizeFolderPath(target)); } - private *chunk(paths: string[], size: number): Generator { - for (let i = 0; i < paths.length; i += size) { - yield paths.slice(i, i + size); + private async resourceExists({ bucket, path }: { bucket: string; path: string }): Promise { + try { + await this.s3Client.send(new HeadObjectCommand({ Bucket: bucket, Key: path })); + return true; + } catch (err) { + if (err instanceof NotFound) { + return false; + } else if (err instanceof NoSuchBucket) { + this.logger.warn({ msg: `S3 Error [${err.name}] no such bucket: ${err.message} (Req ID: ${err.$metadata.requestId})` }); + return false; + } else { + this.logger.error({ msg: 'resourceExists object check failed', err, bucket, path }); + throw err; + } } } } diff --git a/src/cleaner/strategies/deleteStoredResourcesStrategy.ts b/src/cleaner/strategies/deleteStoredResourcesStrategy.ts new file mode 100644 index 0000000..5b0aa09 --- /dev/null +++ b/src/cleaner/strategies/deleteStoredResourcesStrategy.ts @@ -0,0 +1,72 @@ +import type { Logger } from '@map-colonies/js-logger'; +import { deleteStoredResourcesParamsSchema, type DeleteStoredResourcesParams } from '@map-colonies/raster-shared'; +import { inject, injectable } from 'tsyringe'; +import type { ConfigType } from '@common/config'; +import { SERVICES } from '@common/constants'; +import { summarizeDeleteFailures, type IStorageProvider, type StorageProvider, type StorageProviders } from '@src/cleaner/storageProviders'; +import { RecoverableError, UnrecoverableError } from '../errors'; +import { validateSchema } from '../utils'; +import type { ITaskStrategy } from './taskStrategy'; + +@injectable() +export class DeleteStoredResourcesStrategy implements ITaskStrategy { + private readonly failureSampleSize: number; + + public constructor( + @inject(SERVICES.LOGGER) private readonly logger: Logger, + @inject(SERVICES.CONFIG) private readonly config: ConfigType, + @inject(SERVICES.STORAGE_PROVIDERS) private readonly storageProviders: StorageProviders + ) { + this.failureSampleSize = this.config.get('strategies.storedResourcesDeletion.failureSampleSize') as unknown as number; + } + + public validate(params: unknown): DeleteStoredResourcesParams { + this.logger.debug({ msg: `Validating input parameters` }); + return validateSchema(deleteStoredResourcesParamsSchema, params, this.logger); + } + + public async execute(params: DeleteStoredResourcesParams): Promise { + const { paths } = params; + const provider = this.resolveStorageProvider(params); + + this.logger.info({ + msg: 'Starting deletion', + provider: params.storageProvider, + count: paths.length, + paths, + ...(params.storageProvider === 'S3' && { bucket: params.bucket }), + }); + + const { failures } = await provider.deleteResources(params); + + if (failures.length > 0) { + const { counts, summary, sample } = summarizeDeleteFailures(failures, this.failureSampleSize); + this.logger.error({ + msg: 'Deletion failed', + provider: params.storageProvider, + paths, + failureCount: failures.length, + reasonCounts: counts, + sample, + }); + throw new RecoverableError(`Failed to delete resource(s). Reasons: ${summary}. Sample: ${sample.join(', ')}`); + } + + this.logger.info({ + msg: 'Deletion completed successfully', + provider: params.storageProvider, + paths, + }); + } + + private resolveStorageProvider( + params: Extract + ): IStorageProvider { + if (!(params.storageProvider in this.storageProviders)) throw new UnrecoverableError(`Unsupported storage provider ${params.storageProvider}`); + // eslint-disable-next-line @typescript-eslint/naming-convention + const storageProvider = this.storageProviders[params.storageProvider]; + if (storageProvider === undefined) throw new UnrecoverableError(`Unsupported storage provider ${params.storageProvider}`); + this.logger.debug({ msg: `Using ${params.storageProvider} provider` }); + return storageProvider; + } +} diff --git a/src/cleaner/strategies/index.ts b/src/cleaner/strategies/index.ts index f3895d4..2bcd2c4 100644 --- a/src/cleaner/strategies/index.ts +++ b/src/cleaner/strategies/index.ts @@ -1,3 +1,4 @@ -export { type ITaskStrategy } from './taskStrategy'; +export { DeleteStoredResourcesStrategy } from './deleteStoredResourcesStrategy'; export { StrategyFactory, type TaskContext } from './strategyFactory'; +export type { ITaskStrategy } from './taskStrategy'; export { TilesDeletionStrategy } from './tilesDeletionStrategy'; diff --git a/src/cleaner/strategies/strategyFactory.ts b/src/cleaner/strategies/strategyFactory.ts index 00da2be..c56f541 100644 --- a/src/cleaner/strategies/strategyFactory.ts +++ b/src/cleaner/strategies/strategyFactory.ts @@ -1,8 +1,9 @@ -import { container, inject, injectable } from 'tsyringe'; import type { Logger } from '@map-colonies/js-logger'; +import { container, inject, injectable } from 'tsyringe'; import { SERVICES } from '@common/constants'; -import { StrategyNotFoundError } from '../errors'; -import type { ITaskStrategy } from './taskStrategy'; +import { StrategyNotFoundError } from '@src/cleaner/errors'; +import type { ITaskStrategy } from '@src/cleaner/strategies'; +import { getJobAndTaskToken } from '@src/common/dependencyRegistration'; export interface TaskContext { jobId: string; @@ -27,8 +28,10 @@ export class StrategyFactory { public resolveWithContext(taskContext: TaskContext): ITaskStrategy { this.logger.debug({ msg: 'Resolving strategy with task context', ...taskContext }); - if (!container.isRegistered(taskContext.taskType)) { - throw new StrategyNotFoundError(taskContext.taskType); + const jobTaskToken = getJobAndTaskToken(taskContext); + + if (!container.isRegistered(jobTaskToken)) { + throw new StrategyNotFoundError({ jobType: taskContext.jobType, taskType: taskContext.taskType }); } const taskContainer = container.createChildContainer(); @@ -39,7 +42,7 @@ export class StrategyFactory { taskContainer.register(SERVICES.LOGGER, { useValue: taskLogger }); taskContainer.register(SERVICES.TASK_CONTEXT, { useValue: taskContext }); - const strategy = taskContainer.resolve(taskContext.taskType); + const strategy = taskContainer.resolve(jobTaskToken); taskLogger.debug({ msg: 'Strategy resolved successfully with task context' }); diff --git a/src/cleaner/strategies/tilesDeletionStrategy.ts b/src/cleaner/strategies/tilesDeletionStrategy.ts index c0b1b6c..e760e62 100644 --- a/src/cleaner/strategies/tilesDeletionStrategy.ts +++ b/src/cleaner/strategies/tilesDeletionStrategy.ts @@ -1,13 +1,20 @@ -import { inject, injectable } from 'tsyringe'; +import { NoSuchKey } from '@aws-sdk/client-s3'; import type { Logger } from '@map-colonies/js-logger'; -import { SourceType, TileRange, TilesDeletionParams, tilesDeletionParamsSchema } from '@map-colonies/raster-shared'; import type { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; -import { NoSuchKey } from '@aws-sdk/client-s3'; -import { PERCENTAGE_COMPLETE, SERVICES } from '@common/constants'; +import { SourceType, TileRange, TilesDeletionParams, tilesDeletionParamsSchema } from '@map-colonies/raster-shared'; +import { inject, injectable } from 'tsyringe'; import type { ConfigType } from '@common/config'; -import { validateSchema } from '../utils'; +import { PERCENTAGE_COMPLETE, SERVICES } from '@common/constants'; +import { + summarizeDeleteFailures, + type DeleteFailure, + type FsConfig, + type IStorageProvider, + type StorageProvider, + type StorageProviders, +} from '@src/cleaner/storageProviders'; import { RecoverableError, UnrecoverableError, describeError } from '../errors'; -import { summarizeDeleteFailures, type DeleteFailure, type IStorageProvider } from '../storageProviders'; +import { validateSchema } from '../utils'; import type { TaskContext } from './strategyFactory'; import type { ITaskStrategy } from './taskStrategy'; @@ -24,7 +31,7 @@ export class TilesDeletionStrategy implements ITaskStrategy public constructor( @inject(SERVICES.LOGGER) private readonly logger: Logger, @inject(SERVICES.CONFIG) config: ConfigType, - @inject(SERVICES.STORAGE_PROVIDERS) private readonly storageProviders: Map, + @inject(SERVICES.STORAGE_PROVIDERS) private readonly storageProviders: StorageProviders, @inject(SERVICES.QUEUE_CLIENT) private readonly queueClient: QueueClient, @inject(SERVICES.TASK_CONTEXT) private readonly taskContext: TaskContext ) { @@ -32,15 +39,16 @@ export class TilesDeletionStrategy implements ITaskStrategy this.concurrency = config.get('strategies.tilesDeletion.concurrency') as unknown as number; this.failureSampleSize = config.get('strategies.tilesDeletion.failureSampleSize') as unknown as number; this.s3Bucket = config.get('strategies.tilesDeletion.s3Bucket') as unknown as string; - this.fsBasePath = config.get('strategies.tilesDeletion.fsBasePath') as unknown as string; + this.fsBasePath = config.get('storage.fs.basePath') as unknown as FsConfig['basePath']; } public validate(params: unknown): TilesDeletionParams { + this.logger.debug({ msg: `Validating input parameters` }); return validateSchema(tilesDeletionParamsSchema, params, this.logger); } public async execute(params: TilesDeletionParams): Promise { - const { provider, storageTarget } = this.resolveProvider(params); + const { provider, storageTarget } = this.resolveStorageProvider(params); if (!(await provider.targetExists(storageTarget, params.tilesPath))) { throw new UnrecoverableError(`${params.sourceProvider} storage target does not exist: ${storageTarget}/${params.tilesPath}`); @@ -104,13 +112,16 @@ export class TilesDeletionStrategy implements ITaskStrategy this.logger.info({ msg: 'Tiles deletion completed successfully', deletedCount: totalTiles }); } - private resolveProvider(params: TilesDeletionParams): { provider: IStorageProvider; storageTarget: string } { - const provider = this.storageProviders.get(params.sourceProvider); - if (provider === undefined) { - throw new UnrecoverableError(`Unknown storage provider: ${params.sourceProvider}`); - } + private resolveStorageProvider( + params: Extract + ): { provider: IStorageProvider; storageTarget: string } { + if (!(params.sourceProvider in this.storageProviders)) throw new UnrecoverableError(`Unsupported storage provider ${params.sourceProvider}`); + // eslint-disable-next-line @typescript-eslint/naming-convention + const storageProvider = this.storageProviders[params.sourceProvider]; + if (storageProvider === undefined) throw new UnrecoverableError(`Unsupported storage provider ${params.sourceProvider}`); const storageTarget = params.sourceProvider === SourceType.S3 ? this.s3Bucket : this.fsBasePath; - return { provider, storageTarget }; + this.logger.debug({ msg: `Using ${params.sourceProvider} provider` }); + return { provider: storageProvider, storageTarget }; } private async deleteAllTiles( diff --git a/src/cleaner/utils/chunk.ts b/src/cleaner/utils/chunk.ts new file mode 100644 index 0000000..d419943 --- /dev/null +++ b/src/cleaner/utils/chunk.ts @@ -0,0 +1,5 @@ +export function* getChunk(items: T[], size: number): Generator { + for (let i = 0; i < items.length; i += size) { + yield items.slice(i, i + size); + } +} diff --git a/src/cleaner/utils/index.ts b/src/cleaner/utils/index.ts index fc203c9..9d2809c 100644 --- a/src/cleaner/utils/index.ts +++ b/src/cleaner/utils/index.ts @@ -1,2 +1,4 @@ +export { getChunk } from './chunk'; export { buildPollingPairs } from './pairBuilder'; export { validateSchema } from './validationHelper'; +export { normalizeFolderPath, resolveAbsolutePath } from './path'; diff --git a/src/cleaner/utils/path.ts b/src/cleaner/utils/path.ts new file mode 100644 index 0000000..c3049ee --- /dev/null +++ b/src/cleaner/utils/path.ts @@ -0,0 +1,16 @@ +import { resolve, sep } from 'node:path/posix'; + +export const normalizeFolderPath = (path: string): string => { + return path.endsWith(sep) ? path : `${path}${sep}`; +}; + +/** + * Resolves a file system path to an absolute path. + * Ensures the path is resolved as an absolute path and properly formatted + * with a leading separator if not already present. + * @param path - The input path string to normalize + * @returns An absolute path with proper path separators + */ +export const resolveAbsolutePath = (path: string): string => { + return resolve(`${path.startsWith(sep) ? '' : sep}${path}`); +}; diff --git a/src/common/dependencyRegistration.ts b/src/common/dependencyRegistration.ts index a591ff9..29f34e1 100644 --- a/src/common/dependencyRegistration.ts +++ b/src/common/dependencyRegistration.ts @@ -8,6 +8,8 @@ export interface InjectionObject { provider: Providers; } +export const getJobAndTaskToken = ({ jobType, taskType }: { jobType: string; taskType: string }): string => `${jobType}-${taskType}`; + export const registerDependencies = ( dependencies: InjectionObject[], override?: InjectionObject[], diff --git a/src/containerConfig.ts b/src/containerConfig.ts index 586613b..41ee7fb 100644 --- a/src/containerConfig.ts +++ b/src/containerConfig.ts @@ -1,22 +1,23 @@ -import { getOtelMixin } from '@map-colonies/telemetry'; +import { IWorker, JobnikSDK } from '@map-colonies/jobnik-sdk'; +import { jsLogger, type Logger } from '@map-colonies/js-logger'; +import { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; import { SourceType } from '@map-colonies/raster-shared'; +import { getOtelMixin } from '@map-colonies/telemetry'; import { trace } from '@opentelemetry/api'; import { Registry } from 'prom-client'; import { instancePerContainerCachingFactory } from 'tsyringe'; import { DependencyContainer } from 'tsyringe/dist/typings/types'; -import { jsLogger, type Logger } from '@map-colonies/js-logger'; -import { IWorker, JobnikSDK } from '@map-colonies/jobnik-sdk'; -import { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; -import { InjectionObject, registerDependencies } from '@common/dependencyRegistration'; -import { SERVICES, SERVICE_NAME } from '@common/constants'; +import { SERVICE_NAME, SERVICES } from '@common/constants'; +import { getJobAndTaskToken, InjectionObject, registerDependencies } from '@common/dependencyRegistration'; import { getTracing } from '@common/tracing'; +import type { StorageProviders } from '@src/cleaner/storageProviders'; +import { ErrorHandler } from './cleaner/errors'; +import { JobTrackerClient } from './cleaner/httpClients'; +import { FsStorageProvider, S3StorageProvider } from './cleaner/storageProviders'; +import { DeleteStoredResourcesStrategy, StrategyFactory, TilesDeletionStrategy } from './cleaner/strategies'; import type { QueueConfig } from './cleaner/types'; import { ConfigType, getConfig } from './common/config'; import { workerBuilder } from './worker'; -import { StrategyFactory, TilesDeletionStrategy } from './cleaner/strategies'; -import { ErrorHandler } from './cleaner/errors'; -import { S3StorageProvider, FsStorageProvider, type IStorageProvider } from './cleaner/storageProviders'; -import { JobTrackerClient } from './cleaner/httpClients'; export interface RegisterOptions { override?: InjectionObject[]; @@ -100,22 +101,58 @@ export const registerExternalValues = async (options?: RegisterOptions): Promise { token: SERVICES.STORAGE_PROVIDERS, provider: { - useFactory: instancePerContainerCachingFactory((container) => { + useFactory: instancePerContainerCachingFactory((container) => { const config = container.resolve(SERVICES.CONFIG); const logger = container.resolve(SERVICES.LOGGER); - return new Map([ - [SourceType.S3, new S3StorageProvider(config, logger)], - [SourceType.FS, new FsStorageProvider(logger)], - ]); + const cleanupStorageProviders = config.get('storage.cleanupStorageProviders') as unknown as string[]; + const providers = { + ...(cleanupStorageProviders.includes(SourceType.S3) && { [SourceType.S3]: new S3StorageProvider(config, logger) }), + ...(cleanupStorageProviders.includes(SourceType.FS) && { [SourceType.FS]: new FsStorageProvider(config, logger) }), + }; + return providers; }), }, }, { - token: configInstance.get('jobDefinitions.tasks.tilesDeletion.type') as unknown as string, //TODO: when we create worker config schema we can move this to a constant and remove the cast + token: getJobAndTaskToken({ + //TODO: when we create worker config schema we can move this to a constant and remove the cast + jobType: configInstance.get('jobDefinitions.jobs.update.type') as unknown as string, + taskType: configInstance.get('jobDefinitions.tasks.tilesDeletion.type') as unknown as string, + }), provider: { useClass: TilesDeletionStrategy, }, }, + { + token: getJobAndTaskToken({ + //TODO: when we create worker config schema we can move this to a constant and remove the cast + jobType: configInstance.get('jobDefinitions.jobs.swapUpdate.type') as unknown as string, + taskType: configInstance.get('jobDefinitions.tasks.tilesDeletion.type') as unknown as string, + }), + provider: { + useClass: TilesDeletionStrategy, + }, + }, + { + token: getJobAndTaskToken({ + //TODO: when we create worker config schema we can move this to a constant and remove the cast + jobType: configInstance.get('jobDefinitions.jobs.deleteLayer.type') as unknown as string, + taskType: configInstance.get('jobDefinitions.tasks.layerDeletion.type') as unknown as string, + }), + provider: { + useClass: DeleteStoredResourcesStrategy, + }, + }, + { + token: getJobAndTaskToken({ + //TODO: when we create worker config schema we can move this to a constant and remove the cast + jobType: configInstance.get('jobDefinitions.jobs.deleteLayer.type') as unknown as string, + taskType: configInstance.get('jobDefinitions.tasks.artifactsDeletion.type') as unknown as string, + }), + provider: { + useClass: DeleteStoredResourcesStrategy, + }, + }, { token: 'onSignal', provider: { diff --git a/tests/helpers/mocks.ts b/tests/helpers/mocks.ts index 079851f..13d31c7 100644 --- a/tests/helpers/mocks.ts +++ b/tests/helpers/mocks.ts @@ -1,12 +1,14 @@ -import { vi } from 'vitest'; import type { Logger } from '@map-colonies/js-logger'; import type { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; -import type { ConfigType } from '../../src/common/config'; -import type { ITaskStrategy, StrategyFactory } from '../../src/cleaner/strategies'; -import type { IStorageProvider } from '../../src/cleaner/storageProviders'; +import { vi } from 'vitest'; +import type { FsConfig } from '@src/cleaner/storageProviders/fsStorageProvider'; +import type { IStorageProvider, StorageProvider } from '@src/cleaner/storageProviders/iStorageProvider'; +import type { S3Config } from '@src/cleaner/storageProviders/s3StorageProvider'; import type { ErrorHandler } from '../../src/cleaner/errors'; -import type { ErrorDecision, PollingPairConfig } from '../../src/cleaner/types'; import type { JobTrackerClient } from '../../src/cleaner/httpClients'; +import type { ITaskStrategy, StrategyFactory } from '../../src/cleaner/strategies'; +import type { ErrorDecision, PollingPairConfig } from '../../src/cleaner/types'; +import type { ConfigType } from '../../src/common/config'; import { TaskPoller } from '../../src/worker/taskPoller'; // ─── Logger ────────────────────────────────────────────────────────────────── @@ -66,9 +68,10 @@ export function createMockErrorHandler(defaultDecision: ErrorDecision = { should // ─── StorageProvider ───────────────────────────────────────────────────────── -export function createMockStorageProvider(): IStorageProvider { +export function createMockStorageProvider(): IStorageProvider { return { delete: vi.fn().mockResolvedValue([]), + deleteResources: vi.fn().mockResolvedValue({ failures: [] }), targetExists: vi.fn().mockResolvedValue(true), }; } @@ -89,7 +92,21 @@ export function createMockStrategyConfig(overrides: Record = {} 'strategies.tilesDeletion.concurrency': TILES_DELETION_CONFIG_DEFAULTS.concurrency, 'strategies.tilesDeletion.failureSampleSize': TILES_DELETION_CONFIG_DEFAULTS.failureSampleSize, 'strategies.tilesDeletion.s3Bucket': TILES_DELETION_CONFIG_DEFAULTS.s3Bucket, - 'strategies.tilesDeletion.fsBasePath': TILES_DELETION_CONFIG_DEFAULTS.fsBasePath, + 'storage.fs.basePath': TILES_DELETION_CONFIG_DEFAULTS.fsBasePath, + ...overrides, + }; + return { get: vi.fn().mockImplementation((key: string) => values[key]) } as unknown as ConfigType; +} + +// ─── Strategy Config (DeleteStoredResourcesStrategy) ──────────────────────────── + +export const STORED_RESOURCES_DELETION_CONFIG_DEFAULTS = { + failureSampleSize: 3, +} as const; + +export function createMockStoredResourcesDeletionStrategyConfig(overrides: Record = {}): ConfigType { + const values: Record = { + 'strategies.storedResourcesDeletion.failureSampleSize': STORED_RESOURCES_DELETION_CONFIG_DEFAULTS.failureSampleSize, ...overrides, }; return { get: vi.fn().mockImplementation((key: string) => values[key]) } as unknown as ConfigType; @@ -98,17 +115,35 @@ export function createMockStrategyConfig(overrides: Record = {} // ─── S3 Storage Config (S3StorageProvider) ─────────────────────────────────── export const S3_STORAGE_CONFIG_DEFAULTS = { + delete: { + batchSize: 100, + }, endpoint: 'http://localhost:9000', accessKeyId: 'test-key', secretAccessKey: 'test-secret', sslEnabled: false, forcePathStyle: true, region: 'us-east-1', -} as const; +} as const satisfies S3Config; + +export function createMockS3Config(overrides: Record = {}): ConfigType { + return { + get: vi.fn().mockReturnValue({ ...S3_STORAGE_CONFIG_DEFAULTS, ...overrides }), + } as unknown as ConfigType; +} + +// ─── FS Storage Config (FsStorageProvider) ─────────────────────────────────── + +export const FS_STORAGE_CONFIG_DEFAULTS = { + delete: { + batchSize: 3, + }, + basePath: '/test/tiles', +} as const satisfies FsConfig; -export function createMockS3Config(): ConfigType { +export function createMockFsConfig(overrides: Record = {}): ConfigType { return { - get: vi.fn().mockReturnValue({ ...S3_STORAGE_CONFIG_DEFAULTS }), + get: vi.fn().mockReturnValue({ ...FS_STORAGE_CONFIG_DEFAULTS, ...overrides }), } as unknown as ConfigType; } diff --git a/tests/storageProviders/fsStorageProvider.spec.ts b/tests/storageProviders/fsStorageProvider.spec.ts index cdca3f4..ecdb749 100644 --- a/tests/storageProviders/fsStorageProvider.spec.ts +++ b/tests/storageProviders/fsStorageProvider.spec.ts @@ -1,30 +1,50 @@ -import { stat, unlink, rmdir } from 'node:fs/promises'; +import { accessSync, Stats, statSync } from 'node:fs'; +import { rm, rmdir, stat, unlink } from 'node:fs/promises'; import { join } from 'node:path'; -import { Stats } from 'node:fs'; -import { describe, it, expect, beforeEach, vi } from 'vitest'; -import { FsStorageProvider } from '@src/cleaner/storageProviders/fsStorageProvider'; -import { createMockLogger } from '../helpers/mocks'; +import type { Logger } from '@map-colonies/js-logger'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { faker } from '@faker-js/faker'; +import { ConfigurationError, UnrecoverableError } from '@src/cleaner/errors'; +import { FsStorageProvider, type FsConfig } from '@src/cleaner/storageProviders/fsStorageProvider'; +import type { ConfigType } from '@src/common/config'; +import { createMockFsConfig, createMockLogger, FS_STORAGE_CONFIG_DEFAULTS } from '../helpers/mocks'; vi.mock('node:fs/promises', () => ({ stat: vi.fn(), unlink: vi.fn(), rmdir: vi.fn(), + rm: vi.fn(), })); -const BASE_PATH = '/tiles/test'; +vi.mock(import('node:fs'), async (importOriginal) => { + const originModule = await importOriginal(); + return { + ...originModule, + accessSync: vi.fn(), + statSync: vi.fn(), + }; +}); +const BASE_PATH = FS_STORAGE_CONFIG_DEFAULTS.basePath; describe('FsStorageProvider', () => { let provider: FsStorageProvider; + let mockLogger: Logger; + let mockConfig: ConfigType; beforeEach(() => { vi.clearAllMocks(); vi.mocked(stat).mockResolvedValue({} as Stats); vi.mocked(unlink).mockResolvedValue(undefined); vi.mocked(rmdir).mockResolvedValue(undefined); - provider = new FsStorageProvider(createMockLogger()); + vi.mocked(rm).mockResolvedValue(undefined); + vi.mocked(accessSync).mockReturnValue(undefined); + vi.mocked(statSync).mockReturnValue({ isDirectory: () => true } as Stats); + mockLogger = createMockLogger(); + mockConfig = createMockFsConfig(); + provider = new FsStorageProvider(mockConfig, mockLogger); }); - describe('targetExists', () => { + describe('#targetExists', () => { const RELATIVE_PATH = 'layer/v1'; it('should call stat with full target path', async () => { @@ -50,7 +70,7 @@ describe('FsStorageProvider', () => { }); }); - describe('delete', () => { + describe('#delete', () => { it('should return empty array for empty input', async () => { const result = await provider.delete([], BASE_PATH); expect(result).toEqual([]); @@ -182,4 +202,114 @@ describe('FsStorageProvider', () => { }); }); }); + + describe('#deleteResources', () => { + const RELATIVE_PATH = 'layer/v1'; + + it('should successfully call delete all files and return without failures', async () => { + const result = await provider.deleteResources({ paths: [RELATIVE_PATH], storageProvider: 'FS' }); + + expect(rm).toHaveBeenCalledWith(join(BASE_PATH, RELATIVE_PATH), { recursive: true, force: true }); + expect(result).toEqual({ failures: [] }); + }); + + it('should throw UnrecoverableError when a path escapes the base path via traversal', async () => { + await expect(provider.deleteResources({ paths: ['../../../etc/passwd'], storageProvider: 'FS' })).rejects.toThrow(UnrecoverableError); + expect(rm).not.toHaveBeenCalled(); + }); + + it('should throw UnrecoverableError when only one of several paths escapes the base path', async () => { + await expect(provider.deleteResources({ paths: [RELATIVE_PATH, '../../escape'], storageProvider: 'FS' })).rejects.toThrow(UnrecoverableError); + expect(rm).not.toHaveBeenCalled(); + }); + + it('should throw UnrecoverableError when a path resolves to the base path itself', async () => { + await expect(provider.deleteResources({ paths: [''], storageProvider: 'FS' })).rejects.toThrow(UnrecoverableError); + expect(rm).not.toHaveBeenCalled(); + }); + + it('should throw UnrecoverableError when a path resolves to the base path itself via "."', async () => { + await expect(provider.deleteResources({ paths: ['.'], storageProvider: 'FS' })).rejects.toThrow(UnrecoverableError); + expect(rm).not.toHaveBeenCalled(); + }); + + it('should return failures entry when rm rejects', async () => { + vi.mocked(rm).mockRejectedValue(new Error('EACCES')); + + const result = await provider.deleteResources({ paths: [RELATIVE_PATH], storageProvider: 'FS' }); + + expect(result).toEqual({ + failures: [{ path: join(BASE_PATH, RELATIVE_PATH), reason: 'EACCES' }], + }); + }); + + it('should not throw when rm rejects', async () => { + vi.mocked(rm).mockRejectedValue(new Error('Permission denied')); + + const result = provider.deleteResources({ paths: [RELATIVE_PATH], storageProvider: 'FS' }); + + await expect(result).resolves.not.toThrow(); + }); + }); + + describe('#constructor', () => { + it('should construct successfully when the base path is accessible and is a directory', () => { + expect(() => new FsStorageProvider(mockConfig, mockLogger)).not.toThrow(); + expect(accessSync).toHaveBeenCalledWith(FS_STORAGE_CONFIG_DEFAULTS.basePath, expect.any(Number)); + expect(statSync).toHaveBeenCalledWith(FS_STORAGE_CONFIG_DEFAULTS.basePath); + }); + + it('should throw ConfigurationError when the base path does not exist', () => { + vi.mocked(accessSync).mockImplementationOnce(() => { + throw Object.assign(new Error('ENOENT'), { code: 'ENOENT' }); + }); + + expect(() => new FsStorageProvider(mockConfig, mockLogger)).toThrow(ConfigurationError); + }); + + it('should throw ConfigurationError when access to the base path is denied (EACCES)', () => { + vi.mocked(accessSync).mockImplementationOnce(() => { + throw Object.assign(new Error('EACCES'), { code: 'EACCES' }); + }); + + expect(() => new FsStorageProvider(mockConfig, mockLogger)).toThrow(ConfigurationError); + }); + + it('should throw ConfigurationError when access to the base path is denied (EPERM)', () => { + vi.mocked(accessSync).mockImplementationOnce(() => { + throw Object.assign(new Error('EPERM'), { code: 'EPERM' }); + }); + + expect(() => new FsStorageProvider(mockConfig, mockLogger)).toThrow(ConfigurationError); + }); + + it('should throw ConfigurationError when the base path exists but is not a directory', () => { + vi.mocked(statSync).mockReturnValueOnce({ isDirectory: () => false } as Stats); + + expect(() => new FsStorageProvider(mockConfig, mockLogger)).toThrow(ConfigurationError); + }); + + it('should resolve a relative base path (without a leading separator) to an absolute path', () => { + const relativeConfig = { + get: vi.fn().mockReturnValue({ basePath: 'relative/tiles', delete: { batchSize: 100 } } satisfies FsConfig), + } as unknown as ConfigType; + + expect(() => new FsStorageProvider(relativeConfig, createMockLogger())).not.toThrow(); + expect(accessSync).toHaveBeenCalledWith('/relative/tiles', expect.any(Number)); + }); + + it('should throw ConfigurationError for an unexpected accessibility error', () => { + vi.mocked(accessSync).mockImplementationOnce(() => { + throw new Error('disk exploded'); + }); + + expect(() => new FsStorageProvider(mockConfig, mockLogger)).toThrow(ConfigurationError); + }); + + it('should throw ConfigurationError when batchSize is less than or equal to 0', () => { + const zeroBatchConfig = createMockFsConfig({ delete: { batchSize: faker.number.int({ max: 0, min: -Number.MAX_SAFE_INTEGER }) } }); + + expect(() => new FsStorageProvider(zeroBatchConfig, mockLogger)).toThrow(ConfigurationError); + }); + }); }); diff --git a/tests/storageProviders/s3StorageProvider.spec.ts b/tests/storageProviders/s3StorageProvider.spec.ts index 1dde516..ade8544 100644 --- a/tests/storageProviders/s3StorageProvider.spec.ts +++ b/tests/storageProviders/s3StorageProvider.spec.ts @@ -1,23 +1,43 @@ /* eslint-disable @typescript-eslint/naming-convention */ -import { describe, it, expect, beforeEach, vi } from 'vitest'; -import { S3Client, DeleteObjectsCommand, ListObjectsV2Command, NoSuchBucket } from '@aws-sdk/client-s3'; +import { + DeleteObjectsCommand, + ListObjectsV2Command, + NoSuchBucket, + NotFound, + paginateListObjectsV2, + S3Client, + S3ServiceException, + type DeleteObjectsCommandOutput, +} from '@aws-sdk/client-s3'; +import { faker } from '@faker-js/faker'; +import type { Logger } from '@map-colonies/js-logger'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { ConfigurationError, UnrecoverableError } from '@src/cleaner/errors'; import { S3StorageProvider } from '@src/cleaner/storageProviders/s3StorageProvider'; +import type { ConfigType } from '@src/common/config'; import { createMockLogger, createMockS3Config, S3_STORAGE_CONFIG_DEFAULTS } from '../helpers/mocks'; const mockSend = vi.fn(); +const mockPaginateListObjectsV2Next = vi.fn(); +const mockPaginateListObjectsV2Return = vi.fn(); +const mockPaginateListObjectsV2Throw = vi.fn(); +const mockPaginateListObjectsV2Iterator = vi.fn().mockReturnValue({ + next: mockPaginateListObjectsV2Next, + return: mockPaginateListObjectsV2Return, + throw: mockPaginateListObjectsV2Throw, +}); +const mockPaginateListObjectsV2 = { + [Symbol.asyncIterator]: mockPaginateListObjectsV2Iterator, +}; -vi.mock('@aws-sdk/client-s3', () => { - class NoSuchBucket extends Error { - public constructor() { - super('NoSuchBucket'); - this.name = 'NoSuchBucket'; - } - } +vi.mock(import('@aws-sdk/client-s3'), async (importOriginal) => { + const originModule = await importOriginal(); return { - S3Client: vi.fn(() => ({ send: mockSend })), - DeleteObjectsCommand: vi.fn((input: unknown) => input), - ListObjectsV2Command: vi.fn((input: unknown) => input), - NoSuchBucket, + ...originModule, + S3Client: vi.fn(() => ({ send: mockSend })) as unknown as typeof S3Client, + DeleteObjectsCommand: vi.fn((input: unknown) => input) as unknown as typeof DeleteObjectsCommand, + ListObjectsV2Command: vi.fn((input: unknown) => input) as unknown as typeof ListObjectsV2Command, + paginateListObjectsV2: vi.fn(() => mockPaginateListObjectsV2) as unknown as typeof paginateListObjectsV2, }; }); @@ -25,14 +45,21 @@ const BUCKET = 'test-bucket'; describe('S3StorageProvider', () => { let provider: S3StorageProvider; + let mockConfig: ConfigType; + let mockLogger: Logger; beforeEach(() => { vi.clearAllMocks(); - mockSend.mockResolvedValue({ Errors: [] }); - provider = new S3StorageProvider(createMockS3Config(), createMockLogger()); + mockConfig = createMockS3Config(); + mockLogger = createMockLogger(); + provider = new S3StorageProvider(mockConfig, mockLogger); }); - describe('delete', () => { + describe('#delete', () => { + beforeEach(() => { + mockSend.mockResolvedValue({ Errors: [] }); + }); + it('should return empty array for empty input', async () => { const result = await provider.delete([], BUCKET); expect(result).toEqual([]); @@ -123,7 +150,12 @@ describe('S3StorageProvider', () => { it('should batch paths into chunks of 1000 (S3 limit)', async () => { const paths = Array.from({ length: 1500 }, (_, i) => `object-${i}.txt`); - + const mockConfig = createMockS3Config({ + delete: { + batchSize: 1000, + }, + }); + provider = new S3StorageProvider(mockConfig, mockLogger); await provider.delete(paths, BUCKET); expect(mockSend).toHaveBeenCalledTimes(2); @@ -137,6 +169,23 @@ describe('S3StorageProvider', () => { expect(secondCallInput.Delete.Objects).toHaveLength(500); }); + it('should cap DeleteObjects requests at S3 max keys limit (1000 keys) even when configured batchSize is larger', async () => { + const paths = Array.from({ length: 2500 }, (_, i) => `object-${i}.txt`); + const mockConfig = createMockS3Config({ + delete: { + batchSize: 2000, + }, + }); + provider = new S3StorageProvider(mockConfig, mockLogger); + + await provider.delete(paths, BUCKET); + + expect(mockSend).toHaveBeenCalledTimes(3); + const callInputs = vi.mocked(DeleteObjectsCommand).mock.calls.map((call) => call[0] as { Delete: { Objects: { Key: string }[] } }); + expect(callInputs.every((input) => input.Delete.Objects.length <= 1000)).toBe(true); + expect(callInputs.map((input) => input.Delete.Objects.length)).toEqual([1000, 1000, 500]); + }); + it('should accumulate failures across multiple chunks', async () => { const paths = Array.from({ length: 1500 }, (_, i) => `object-${i}.txt`); mockSend @@ -164,7 +213,11 @@ describe('S3StorageProvider', () => { }); }); - describe('targetExists', () => { + describe('#targetExists', () => { + beforeEach(() => { + mockSend.mockResolvedValue({ Errors: [] }); + }); + const PREFIX = 'some/prefix'; it('should list objects with bucket and relativePath prefix', async () => { @@ -215,7 +268,426 @@ describe('S3StorageProvider', () => { }); }); - describe('constructor', () => { + describe('#deleteResources', () => { + const PATH = 'layer/v1'; + const NORMALIZED_PATH = `${PATH}/`; + + it('should return empty result when listing returns no keys', async () => { + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found + mockPaginateListObjectsV2Next.mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should not double-add trailing slash when prefix already ends with one', async () => { + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found + mockPaginateListObjectsV2Next.mockResolvedValueOnce({ done: true, value: undefined }); + + await provider.deleteResources({ paths: [`${PATH}/`], bucket: BUCKET, storageProvider: 'S3' }); + + expect(paginateListObjectsV2).toHaveBeenCalledWith(expect.anything(), expect.objectContaining({ Prefix: NORMALIZED_PATH })); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should list and delete a single object', async () => { + const path = 'layer/v1/0/0.png'; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found for single object + mockPaginateListObjectsV2Next.mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [path], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should list and delete a single page of objects', async () => { + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ Errors: [] }); // delete page 1 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should list and delete a single page of objects - no matching single object (bucket does not exists)', async () => { + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NoSuchBucket({ $metadata: {}, message: 'no bucket' })) // no bucket found for single object + .mockResolvedValueOnce({ Errors: [] }); // delete page 1 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should list and delete multiple pages of objects', async () => { + const page1Keys = ['layer/v1/0/0.png']; + const page2Keys = ['layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ Errors: [] }) // delete page 1 + .mockResolvedValueOnce({ Errors: [] }); // delete page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(3); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(4); + }); + + it('should list and delete multiple pages of objects skipping deletion of empty keys', async () => { + const page1Keys = ['layer/v1/0/0.png']; + const page2Keys = ['layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ Errors: [] }); // delete page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: page1Keys.map(() => ({ Key: undefined })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(3); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should list and a single object skipping deletion of similar keys of objects', async () => { + const objectKey = 'layer/v1/0/0.png'; + const page1Keys = [`${objectKey}8`]; + const page2Keys = [objectKey]; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ Errors: [] }); // delete page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [objectKey], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(3); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should list and delete multiple paths and their objects', async () => { + const PATH1 = 'layer/v1'; + const PATH2 = 'layer/v2'; + const path1Page1Keys = ['layer/v1/0/0.png']; + const path1Page2Keys = ['layer/v1/0/1.png']; + const path2Page1Keys = ['layer/v2/0/0.png']; + const path2Page2Keys = ['layer/v2/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object for path 1 + .mockResolvedValueOnce({ Errors: [] }) // delete path 1 page 1 + .mockResolvedValueOnce({ Errors: [] }) // delete path 1 page 2 + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object for path 2 + .mockResolvedValueOnce({ Errors: [] }) // delete path 2 page 1 + .mockResolvedValueOnce({ Errors: [] }); // delete path 2 page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: path1Page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: path1Page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }) + .mockResolvedValueOnce({ done: false, value: { Contents: path2Page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: path2Page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH1, PATH2], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(6); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(7); + }); + + it('should list and delete a single object and matching paths and their objects', async () => { + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockResolvedValueOnce(undefined) // single object exists + .mockResolvedValueOnce({ Errors: [] }) // delete single object + .mockResolvedValueOnce({ Errors: [] }); // delete page 1 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(4); + }); + + it('should handle empty page Contents response', async () => { + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found + mockPaginateListObjectsV2Next.mockResolvedValueOnce({ done: false, value: {} }).mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should handle empty delete page Errors response', async () => { + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ $metadata: {} } satisfies DeleteObjectsCommandOutput); // delete page 1 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should handle delete page Errors response without elements', async () => { + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ $metadata: {}, Errors: [] } satisfies DeleteObjectsCommandOutput); // delete page 1 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + + it('should throw UnrecoverableError when bucket does not exist', async () => { + mockSend.mockRejectedValueOnce(new NotFound({ $metadata: {}, message: '' })); + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(UnrecoverableError); + expect(mockSend).toHaveBeenCalledTimes(1); + }); + + it('should throw an error when storage existence check failing', async () => { + const expectedError = new Error('error'); + mockSend.mockRejectedValueOnce(expectedError); + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockSend).toHaveBeenCalledTimes(1); + }); + + it('should throw UnrecoverableError when a path is empty (root deletion guard)', async () => { + const result = provider.deleteResources({ paths: [''], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(UnrecoverableError); + expect(paginateListObjectsV2).not.toHaveBeenCalled(); + expect(mockSend).toHaveBeenCalledTimes(0); + }); + + it('should throw UnrecoverableError when any path is empty even if others are valid', async () => { + const result = provider.deleteResources({ paths: [PATH, ''], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(UnrecoverableError); + expect(mockSend).toHaveBeenCalledTimes(0); + }); + + it('should throw an error when checking for matching single object is failing (generic error)', async () => { + const expectedError = new Error('NetworkError'); + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(expectedError); // error thrown for single object lookup + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should paginate and return all delete errors returned by S3 per object', async () => { + const page1Keys = ['layer/v1/0/0.png']; + const page2Keys = ['layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockResolvedValueOnce({ Errors: [{ Key: 'layer/v1/0/0.png', Code: 'AccessDenied' }] }) // delete page 1 + .mockResolvedValueOnce({ Errors: [] }); // delete page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [{ path: 'layer/v1/0/0.png', reason: 'AccessDenied' }] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(3); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(4); + }); + + it('should paginate and return all delete errors thrown and unhandled by S3', async () => { + const page1Keys = ['layer/v1/0/0.png']; + const page2Keys = ['layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })) // no listing found for single object + .mockRejectedValueOnce(new Error('NetworkError')) // delete page 1 + .mockResolvedValueOnce({ Errors: [] }); // delete page 2 + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: page1Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: false, value: { Contents: page2Keys.map((Key) => ({ Key })) } }) + .mockResolvedValueOnce({ done: true, value: undefined }); + + const result = await provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + expect(result).toEqual({ failures: [{ path: 'layer/v1/0/0.png', reason: 'NetworkError' }] }); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(3); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(4); + }); + + it('should stop pagination and throw an error when bucket does not exist', async () => { + const expectedError = new NoSuchBucket({ $metadata: { requestId: faker.string.uuid() }, message: 'msg' }); + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found for single object + mockPaginateListObjectsV2Next.mockRejectedValueOnce(expectedError); + + const result = provider.deleteResources({ paths: [PATH], bucket: '', storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should stop pagination and throw an error when S3 throws a service error', async () => { + const expectedError = new S3ServiceException({ $fault: 'server', $metadata: { requestId: faker.string.uuid() }, name: 'msg' }); + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found for single object + mockPaginateListObjectsV2Next.mockRejectedValueOnce(expectedError); + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should stop pagination and throw an error when S3 throws a bucket does not exists error', async () => { + const expectedError = new NoSuchBucket({ $metadata: {}, message: '' }); + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found for single object + mockPaginateListObjectsV2Next.mockRejectedValueOnce(expectedError); + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(1); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(2); + }); + + it('should stop pagination and return failure when list call throws', async () => { + const expectedError = new Error('NetworkError'); + const keys = ['layer/v1/0/0.png', 'layer/v1/0/1.png']; + mockSend + .mockResolvedValueOnce(undefined) // bucket exists + .mockRejectedValueOnce(new NotFound({ $metadata: {}, message: 'not found' })); // no listing found for single object + mockPaginateListObjectsV2Next + .mockResolvedValueOnce({ done: false, value: { Contents: keys.map((Key) => ({ Key })) } }) + .mockRejectedValueOnce(expectedError); + + const result = provider.deleteResources({ paths: [PATH], bucket: BUCKET, storageProvider: 'S3' }); + + await expect(result).rejects.toThrow(expectedError); + expect(mockPaginateListObjectsV2Next).toHaveBeenCalledTimes(2); + expect(mockPaginateListObjectsV2Return).toHaveBeenCalledTimes(0); + expect(mockPaginateListObjectsV2Throw).toHaveBeenCalledTimes(0); + expect(mockSend).toHaveBeenCalledTimes(3); + }); + }); + + describe('#constructor', () => { it('should construct S3Client with config values', () => { expect(S3Client).toHaveBeenCalledWith( expect.objectContaining({ @@ -226,5 +698,11 @@ describe('S3StorageProvider', () => { }) ); }); + + it('should throw ConfigurationError when batchSize is less than or equal to 0', () => { + const zeroBatchConfig = createMockS3Config({ delete: { batchSize: faker.number.int({ max: 0, min: -Number.MAX_SAFE_INTEGER }) } }); + + expect(() => new S3StorageProvider(zeroBatchConfig, mockLogger)).toThrow(ConfigurationError); + }); }); }); diff --git a/tests/strategies/deleteStoredResourcesStrategy.spec.ts b/tests/strategies/deleteStoredResourcesStrategy.spec.ts new file mode 100644 index 0000000..302b194 --- /dev/null +++ b/tests/strategies/deleteStoredResourcesStrategy.spec.ts @@ -0,0 +1,135 @@ +import type { Logger } from '@map-colonies/js-logger'; +import { SourceType, type DeleteStoredResourcesParams } from '@map-colonies/raster-shared'; +import { beforeEach, describe, expect, it, type vi } from 'vitest'; +import { RecoverableError, UnrecoverableError, ValidationError } from '@src/cleaner/errors'; +import type { IStorageProvider, StorageProviders } from '@src/cleaner/storageProviders'; +import { DeleteStoredResourcesStrategy } from '@src/cleaner/strategies/deleteStoredResourcesStrategy'; +import type { ConfigType } from '@src/common/config'; +import { createMockStoredResourcesDeletionStrategyConfig, createMockLogger, createMockStorageProvider } from '../helpers/mocks'; + +const S3_BUCKET = 'test-bucket'; +const s3Params = { storageProvider: SourceType.S3, paths: ['layer1'], bucket: S3_BUCKET }; +const fsParams = { storageProvider: SourceType.FS, paths: ['layer2'] }; + +describe('DeleteStoredResourcesStrategy', () => { + let strategy: DeleteStoredResourcesStrategy; + // eslint-disable-next-line @typescript-eslint/naming-convention + let mockS3Provider: IStorageProvider<'S3'>; + // eslint-disable-next-line @typescript-eslint/naming-convention + let mockFsProvider: IStorageProvider<'FS'>; + let mockLogger: Logger; + let mockConfig: ConfigType; + + beforeEach(() => { + mockS3Provider = createMockStorageProvider(); + mockFsProvider = createMockStorageProvider(); + mockLogger = createMockLogger(); + mockConfig = createMockStoredResourcesDeletionStrategyConfig(); + + const storageProviders: StorageProviders = { + [SourceType.FS]: mockFsProvider, + [SourceType.S3]: mockS3Provider, + }; + + strategy = new DeleteStoredResourcesStrategy(mockLogger, mockConfig, storageProviders); + }); + + describe('#validate', () => { + it('should validate and return S3 params', () => { + const result = strategy.validate(s3Params); + + expect(result).toEqual(s3Params); + }); + + it('should validate and return FS params', () => { + const result = strategy.validate(fsParams); + + expect(result).toEqual(fsParams); + }); + + it('should throw ValidationError when storageProvider is missing', () => { + expect(() => strategy.validate({ catalogId: 'layer1' })).toThrow(ValidationError); + }); + + it('should throw ValidationError when storageProvider is unknown', () => { + expect(() => strategy.validate({ storageProvider: 'GCS', catalogId: 'layer1' })).toThrow(ValidationError); + }); + + it('should throw ValidationError when tilesPath is empty string', () => { + expect(() => strategy.validate({ storageProvider: SourceType.S3, catalogId: '' })).toThrow(ValidationError); + }); + + it('should throw ValidationError when tilesPath is missing', () => { + expect(() => strategy.validate({ storageProvider: SourceType.S3 })).toThrow(ValidationError); + }); + + it('should throw ValidationError for null params', () => { + expect(() => strategy.validate(null)).toThrow(ValidationError); + }); + }); + + describe('#execute', () => { + it('should call delete all resources on S3 provider', async () => { + await strategy.execute(s3Params); + + expect(mockS3Provider.deleteResources).toHaveBeenCalledWith({ paths: s3Params.paths, bucket: S3_BUCKET, storageProvider: 'S3' }); + expect(mockFsProvider.deleteResources).not.toHaveBeenCalled(); + }); + + it('should call delete all resources on FS provider', async () => { + await strategy.execute(fsParams); + + expect(mockFsProvider.deleteResources).toHaveBeenCalledWith({ paths: fsParams.paths, storageProvider: 'FS' }); + expect(mockS3Provider.deleteResources).not.toHaveBeenCalled(); + }); + + it('should resolve without throwing when deleteResources returns no failures', async () => { + (mockS3Provider.deleteResources as ReturnType).mockResolvedValueOnce({ failures: [] }); + + const result = strategy.execute(s3Params); + + await expect(result).resolves.toBeUndefined(); + }); + + it('should throw UnrecoverableError for unknown provider', async () => { + const unknownParams = { ...s3Params, storageProvider: 'UNKNOWN' } as unknown as DeleteStoredResourcesParams; + + const result = strategy.execute(unknownParams); + + await expect(result).rejects.toThrow(UnrecoverableError); + expect(mockS3Provider.deleteResources).not.toHaveBeenCalled(); + }); + + it('should throw UnrecoverableError when the provider is registered but resolves to undefined', async () => { + const storageProviders = { + [SourceType.FS]: mockFsProvider, + [SourceType.S3]: undefined, + } as unknown as StorageProviders; + strategy = new DeleteStoredResourcesStrategy(mockLogger, mockConfig, storageProviders); + + const result = strategy.execute(s3Params); + + await expect(result).rejects.toThrow(UnrecoverableError); + expect(mockS3Provider.deleteResources).not.toHaveBeenCalled(); + }); + + it('should throw RecoverableError when deleteResources returns failures', async () => { + (mockS3Provider.deleteResources as ReturnType).mockResolvedValueOnce({ + failures: [{ path: 'layer1/0/0.png', reason: 'AccessDenied' }], + }); + + const result = strategy.execute(s3Params); + + await expect(result).rejects.toThrow(RecoverableError); + }); + + it('should rethrow error thrown by deleteResources', async () => { + const expectedError = new Error('Custom'); + (mockS3Provider.deleteResources as ReturnType).mockRejectedValueOnce(expectedError); + + const result = strategy.execute(s3Params); + + await expect(result).rejects.toThrow(expectedError); + }); + }); +}); diff --git a/tests/tilesDeletionStrategy.spec.ts b/tests/strategies/tilesDeletionStrategy.spec.ts similarity index 91% rename from tests/tilesDeletionStrategy.spec.ts rename to tests/strategies/tilesDeletionStrategy.spec.ts index 7669a80..161a2c7 100644 --- a/tests/tilesDeletionStrategy.spec.ts +++ b/tests/strategies/tilesDeletionStrategy.spec.ts @@ -1,12 +1,12 @@ -import { describe, it, expect, beforeEach, vi } from 'vitest'; -import { type TilesDeletionParams, SourceType } from '@map-colonies/raster-shared'; -import type { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; import { faker } from '@faker-js/faker'; -import { TilesDeletionStrategy } from '@src/cleaner/strategies/tilesDeletionStrategy'; +import type { TaskHandler as QueueClient } from '@map-colonies/mc-priority-queue'; +import { type TilesDeletionParams, SourceType } from '@map-colonies/raster-shared'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { RecoverableError, UnrecoverableError, ValidationError } from '@src/cleaner/errors'; +import type { IStorageProvider, StorageProviders } from '@src/cleaner/storageProviders'; import type { TaskContext } from '@src/cleaner/strategies/strategyFactory'; -import { ValidationError, RecoverableError, UnrecoverableError } from '@src/cleaner/errors'; -import type { IStorageProvider } from '@src/cleaner/storageProviders'; -import { createMockLogger, createMockStorageProvider, createMockStrategyConfig, TILES_DELETION_CONFIG_DEFAULTS } from './helpers/mocks'; +import { TilesDeletionStrategy } from '@src/cleaner/strategies/tilesDeletionStrategy'; +import { createMockLogger, createMockStorageProvider, createMockStrategyConfig, TILES_DELETION_CONFIG_DEFAULTS } from '../helpers/mocks'; const { s3Bucket: S3_BUCKET, fsBasePath: FS_BASE_PATH } = TILES_DELETION_CONFIG_DEFAULTS; @@ -28,8 +28,8 @@ const tilePath = (z: number, x: number, y: number): string => `${s3Params.tilesP describe('TilesDeletionStrategy', () => { let strategy: TilesDeletionStrategy; - let MockS3Provider: IStorageProvider; - let MockFsProvider: IStorageProvider; + let MockS3Provider: IStorageProvider<'S3'>; + let MockFsProvider: IStorageProvider<'FS'>; let mockUpdateProgress: ReturnType; beforeEach(() => { @@ -37,10 +37,10 @@ describe('TilesDeletionStrategy', () => { MockFsProvider = createMockStorageProvider(); mockUpdateProgress = vi.fn().mockResolvedValue(undefined); - const storageProviders = new Map([ - [SourceType.S3, MockS3Provider], - [SourceType.FS, MockFsProvider], - ]); + const storageProviders: StorageProviders = { + [SourceType.FS]: MockFsProvider, + [SourceType.S3]: MockS3Provider, + }; const queueClient = { updateProgress: mockUpdateProgress } as unknown as QueueClient; strategy = new TilesDeletionStrategy(createMockLogger(), createMockStrategyConfig(), storageProviders, queueClient, TASK_CONTEXT); @@ -138,6 +138,24 @@ describe('TilesDeletionStrategy', () => { await expect(strategy.execute(unknownParams)).rejects.toThrow(UnrecoverableError); }); + + it('should throw UnrecoverableError when the provider is registered but resolves to undefined', async () => { + const storageProviders = { + [SourceType.FS]: MockFsProvider, + [SourceType.S3]: undefined, + } satisfies StorageProviders; + strategy = new TilesDeletionStrategy( + createMockLogger(), + createMockStrategyConfig(), + storageProviders, + { updateProgress: mockUpdateProgress } as unknown as QueueClient, + TASK_CONTEXT + ); + + await expect(strategy.execute(s3Params)).rejects.toThrow(UnrecoverableError); + expect(MockS3Provider.targetExists).not.toHaveBeenCalled(); + expect(MockS3Provider.delete).not.toHaveBeenCalled(); + }); }); describe('tile path generation', () => { diff --git a/tests/strategyFactory.spec.ts b/tests/strategyFactory.spec.ts index 2daa995..027a1f1 100644 --- a/tests/strategyFactory.spec.ts +++ b/tests/strategyFactory.spec.ts @@ -1,11 +1,13 @@ -import { describe, it, expect, beforeEach, afterEach } from 'vitest'; -import { container } from 'tsyringe'; import { faker } from '@faker-js/faker'; import type { Logger } from '@map-colonies/js-logger'; -import { SERVICES } from '../src/common/constants'; -import { StrategyFactory, TilesDeletionStrategy, type ITaskStrategy, type TaskContext } from '../src/cleaner/strategies'; +import { SourceType } from '@map-colonies/raster-shared'; +import { container } from 'tsyringe'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { IStorageProvider, StorageProviders } from '@src/cleaner/storageProviders'; import { StrategyNotFoundError } from '../src/cleaner/errors'; -import { createMockLogger, createMockConfig, createMockQueueClient } from './helpers/mocks'; +import { StrategyFactory, TilesDeletionStrategy, type ITaskStrategy, type TaskContext } from '../src/cleaner/strategies'; +import { SERVICES } from '../src/common/constants'; +import { createMockConfig, createMockLogger, createMockQueueClient, createMockStorageProvider } from './helpers/mocks'; class MockStrategy implements ITaskStrategy { public validate(params: unknown): Record { @@ -20,13 +22,25 @@ class MockStrategy implements ITaskStrategy { describe('StrategyFactory', () => { let strategyFactory: StrategyFactory; let mockLogger: Logger; + // eslint-disable-next-line @typescript-eslint/naming-convention + let mockS3Provider: IStorageProvider<'S3'>; + // eslint-disable-next-line @typescript-eslint/naming-convention + let mockFsProvider: IStorageProvider<'FS'>; beforeEach(() => { mockLogger = createMockLogger(); + mockS3Provider = createMockStorageProvider(); + mockFsProvider = createMockStorageProvider(); + + const storageProviders: StorageProviders = { + [SourceType.FS]: mockFsProvider, + [SourceType.S3]: mockS3Provider, + }; + container.register(SERVICES.LOGGER, { useValue: mockLogger }); container.register(SERVICES.CONFIG, { useValue: createMockConfig() }); - container.register(SERVICES.STORAGE_PROVIDERS, { useValue: new Map() }); + container.register(SERVICES.STORAGE_PROVIDERS, { useValue: storageProviders }); container.register(SERVICES.QUEUE_CLIENT, { useValue: createMockQueueClient() }); strategyFactory = new StrategyFactory(mockLogger); @@ -37,15 +51,16 @@ describe('StrategyFactory', () => { container.clearInstances(); }); - describe('resolveWithContext', () => { + describe('#resolveWithContext', () => { it('should resolve registered strategy with enriched logger context', () => { + const jobType = 'Ingestion_Update'; const taskType = 'tiles-deletion'; - container.register(taskType, { useClass: TilesDeletionStrategy }); + container.register(`${jobType}-${taskType}`, { useClass: TilesDeletionStrategy }); const taskContext: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), - jobType: 'Ingestion_Update', + jobType, taskType, }; @@ -63,13 +78,14 @@ describe('StrategyFactory', () => { }); it('should create child logger with task context', () => { + const jobType = 'Ingestion_Swap_Update'; const taskType = 'tiles-deletion'; - container.register(taskType, { useClass: TilesDeletionStrategy }); + container.register(`${jobType}-${taskType}`, { useClass: TilesDeletionStrategy }); const taskContext: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), - jobType: 'Ingestion_Swap_Update', + jobType, taskType, }; @@ -96,11 +112,12 @@ describe('StrategyFactory', () => { }); it('should create separate instances for different tasks (child container isolation)', () => { + const jobType = 'Ingestion_Update'; const taskType = 'tiles-deletion'; - container.register(taskType, { useClass: TilesDeletionStrategy }); + container.register(`${jobType}-${taskType}`, { useClass: TilesDeletionStrategy }); - const context1: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType: 'Ingestion_Update', taskType }; - const context2: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType: 'Ingestion_Update', taskType }; + const context1: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType, taskType }; + const context2: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType, taskType }; const strategy1 = strategyFactory.resolveWithContext(context1); const strategy2 = strategyFactory.resolveWithContext(context2); @@ -109,14 +126,15 @@ describe('StrategyFactory', () => { }); it('should resolve different strategies for different task types', () => { + const jobType = 'Ingestion_Update'; const taskType1 = 'tiles-deletion'; const taskType2 = 'files-deletion'; - container.register(taskType1, { useClass: TilesDeletionStrategy }); - container.register(taskType2, { useClass: MockStrategy }); + container.register(`${jobType}-${taskType1}`, { useClass: TilesDeletionStrategy }); + container.register(`${jobType}-${taskType2}`, { useClass: MockStrategy }); - const context1: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType: 'Ingestion_Update', taskType: taskType1 }; - const context2: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType: 'Ingestion_Update', taskType: taskType2 }; + const context1: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType, taskType: taskType1 }; + const context2: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), jobType, taskType: taskType2 }; const strategy1 = strategyFactory.resolveWithContext(context1); const strategy2 = strategyFactory.resolveWithContext(context2); @@ -127,13 +145,14 @@ describe('StrategyFactory', () => { }); it('should handle special characters in task type', () => { + const jobType = 'CustomJob'; const taskType = 'task-with-special_chars.v2'; - container.register(taskType, { useClass: MockStrategy }); + container.register(`${jobType}-${taskType}`, { useClass: MockStrategy }); const taskContext: TaskContext = { jobId: faker.string.uuid(), taskId: faker.string.uuid(), - jobType: 'CustomJob', + jobType, taskType, };