From b1fc7ae599debfbccf404100ce042ec688939992 Mon Sep 17 00:00:00 2001 From: TULCHINSKI LIRAN Date: Tue, 29 Sep 2026 14:48:00 +0300 Subject: [PATCH 1/2] feat: create Jobnik ingestion + delete jobs (MAPCO-11590) Add JobnikClient wrapping the @map-colonies/jobnik-sdk Producer. - Ingestion: create 3D-Ingestion job + validation stage/task (+ data-extraction stage/task only for 3tz input). - Delete: validate deletability (record exists, productType not QuantizedMeshDTMBest, productStatus UNPUBLISHED) then create a 3D-Delete job. RecordManager ingest/delete now return real Jobnik jobs. Add externalServices.jobManager config (schema, default, custom-environment-variables, helm). Unit + integration tests. --- config/custom-environment-variables.json | 38 ++ config/default.json | 26 + helm/templates/configmap.yaml | 22 + helm/values.yaml | 21 + package-lock.json | 534 ++++++++++++++++++ package.json | 1 + src/common/config.ts | 90 ++- src/common/constants.ts | 15 + src/common/util.ts | 4 + src/containerConfig.ts | 2 + src/externalServices/catalog/interfaces.ts | 2 + src/externalServices/jobnik/jobnikClient.ts | 105 ++++ src/providers/getProvider.ts | 11 + src/providers/interfaces.ts | 3 + src/providers/nfsProvider.ts | 43 ++ src/providers/s3Provider.ts | 56 ++ src/record/controllers/recordController.ts | 4 +- src/record/models/recordManager.ts | 13 +- src/validator/validationManager.ts | 48 +- tests/integration/record/record.spec.ts | 23 +- .../jobnik/jobnikClient.spec.ts | 89 +++ tests/unit/providers/nfsProvider.spec.ts | 38 ++ tests/unit/providers/s3Provider.spec.ts | 59 ++ .../unit/record/models/recordManager.spec.ts | 16 +- .../unit/validator/validationManager.spec.ts | 89 ++- 25 files changed, 1316 insertions(+), 36 deletions(-) create mode 100644 src/common/util.ts create mode 100644 src/externalServices/jobnik/jobnikClient.ts create mode 100644 src/providers/getProvider.ts create mode 100644 src/providers/interfaces.ts create mode 100644 src/providers/nfsProvider.ts create mode 100644 src/providers/s3Provider.ts create mode 100644 tests/unit/externalServices/jobnik/jobnikClient.spec.ts create mode 100644 tests/unit/providers/nfsProvider.spec.ts create mode 100644 tests/unit/providers/s3Provider.spec.ts diff --git a/config/custom-environment-variables.json b/config/custom-environment-variables.json index c956cd1..6c1060f 100644 --- a/config/custom-environment-variables.json +++ b/config/custom-environment-variables.json @@ -5,5 +5,43 @@ "url": "LOOKUP_TABLES_URL", "subUrl": "LOOKUP_TABLES_SUB_URL" } + }, + "jobManager": { + "url": "JOB_MANAGER_URL", + "ingestion": { + "jobType": "JOB_INGESTION_TYPE", + "taskType": "TASK_INGESTION_TYPE", + "batches": { + "__name": "INGESTION_TASK_BATCHES", + "__format": "number" + } + }, + "delete": { + "jobType": "JOB_DELETE_TYPE", + "taskType": "TASK_DELETE_TYPE" + } + }, + "provider": "PROVIDER_FROM", + "NFS": { + "pvPath": "PV_SOURCE_PATH" + }, + "S3": { + "accessKeyId": "S3_ACCESS_KEY_ID", + "secretAccessKey": "S3_SECRET_ACCESS_KEY", + "endpointUrl": "S3_ENDPOINT_URL", + "bucket": "S3_BUCKET", + "region": "S3_REGION", + "forcePathStyle": { + "__name": "S3_FORCE_PATH_STYLE", + "__format": "boolean" + }, + "sslEnabled": { + "__name": "S3_SSL_ENABLED", + "__format": "boolean" + }, + "maxAttempts": { + "__name": "S3_MAX_ATTEMPTS", + "__format": "number" + } } } diff --git a/config/default.json b/config/default.json index 0f5589d..319977f 100644 --- a/config/default.json +++ b/config/default.json @@ -39,5 +39,31 @@ "url": "http://127.0.0.1:8080", "subUrl": "lookup-tables/lookupData" } + }, + "jobManager": { + "url": "http://127.0.0.1:8080", + "ingestion": { + "jobType": "Ingestion_New_3D", + "taskType": "tilesCopying", + "batches": 100 + }, + "delete": { + "jobType": "Delete_3D", + "taskType": "deleteModel" + } + }, + "provider": "NFS", + "NFS": { + "pvPath": "/app/models" + }, + "S3": { + "accessKeyId": "minio", + "secretAccessKey": "minio", + "endpointUrl": "http://127.0.0.1:9000", + "bucket": "3dtiles", + "region": "us-east-1", + "forcePathStyle": true, + "sslEnabled": false, + "maxAttempts": 3 } } diff --git a/helm/templates/configmap.yaml b/helm/templates/configmap.yaml index 2635724..f150bf7 100644 --- a/helm/templates/configmap.yaml +++ b/helm/templates/configmap.yaml @@ -37,6 +37,28 @@ data: LOOKUP_TABLES_SUB_URL: {{ .subUrl | quote }} {{- end }} {{- end }} + {{- with .Values.jobManager }} + JOB_MANAGER_URL: {{ .url | quote }} + JOB_INGESTION_TYPE: {{ .ingestion.jobType | quote }} + TASK_INGESTION_TYPE: {{ .ingestion.taskType | quote }} + INGESTION_TASK_BATCHES: {{ .ingestion.batches | quote }} + JOB_DELETE_TYPE: {{ .delete.jobType | quote }} + TASK_DELETE_TYPE: {{ .delete.taskType | quote }} + {{- end }} + PROVIDER_FROM: {{ .Values.provider | quote }} + {{- if eq .Values.provider "NFS" }} + PV_SOURCE_PATH: {{ .Values.NFS.pvPath | quote }} + {{- end }} + {{- if eq .Values.provider "S3" }} + {{- with .Values.S3 }} + S3_ENDPOINT_URL: {{ .endpointUrl | quote }} + S3_BUCKET: {{ .bucket | quote }} + S3_REGION: {{ .region | quote }} + S3_FORCE_PATH_STYLE: {{ .forcePathStyle | quote }} + S3_SSL_ENABLED: {{ .sslEnabled | quote }} + S3_MAX_ATTEMPTS: {{ .maxAttempts | quote }} + {{- end }} + {{- end }} {{- with .Values.configManagement }} CONFIG_NAME: {{ .name | quote }} CONFIG_VERSION: {{ .version | quote }} diff --git a/helm/values.yaml b/helm/values.yaml index a6b7973..c3cea6d 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -84,6 +84,27 @@ externalServices: url: '' subUrl: '' +jobManager: + url: + ingestion: + jobType: '' + taskType: '' + batches: + delete: + jobType: '' + taskType: '' + +provider: 'NFS' +NFS: + pvPath: '' +S3: + endpointUrl: '' + bucket: '' + region: '' + forcePathStyle: true + sslEnabled: false + maxAttempts: 3 + telemetry: logger: level: info diff --git a/package-lock.json b/package-lock.json index def67c3..f23dc6e 100644 --- a/package-lock.json +++ b/package-lock.json @@ -9,6 +9,7 @@ "version": "1.0.0", "license": "ISC", "dependencies": { + "@aws-sdk/client-s3": "^3.1142.0", "@godaddy/terminus": "^4.12.1", "@map-colonies/3d-shared": "file:../../3d-general/3d-shared", "@map-colonies/config": "^4.0.1", @@ -141,6 +142,416 @@ "url": "https://github.com/sponsors/philsturgeon" } }, + "node_modules/@aws-sdk/checksums": { + "version": "3.1001.1", + "resolved": "https://registry.npmjs.org/@aws-sdk/checksums/-/checksums-3.1001.1.tgz", + "integrity": "sha512-x12Q17KYlJAd3nKf8LV5LV0vt8sh8/6YfQLGPtrGnQf/tW4jqxPGq5GPpuVitpQYM3eUR4XB7CbxZf751NMbLw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/checksums/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/client-s3": { + "version": "3.1142.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/client-s3/-/client-s3-3.1142.0.tgz", + "integrity": "sha512-OC9AcGMFOsBc95YSPWH3O7DE6RUTy7jqeOpVvzO2Dv3rw3w4vQRK+YdLXIbO2yXfmx3N59Y55yAJdDRJFMQK3Q==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/checksums": "^3.1001.1", + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/credential-provider-node": "^3.972.84", + "@aws-sdk/middleware-sdk-s3": "^3.972.77", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/client-s3/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/core": { + "version": "3.978.1", + "resolved": "https://registry.npmjs.org/@aws-sdk/core/-/core-3.978.1.tgz", + "integrity": "sha512-LbY9aGsEiznDWmUc30Nwv3aIX/+dbwTx8KfS0yOC3NPYMO+O91e6jkT1azf34FwjOndq8/Q+RcVVZz5xnerwdg==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/types": "^3.974.6", + "@aws-sdk/xml-builder": "^3.972.41", + "@aws/lambda-invoke-store": "^0.3.0", + "@smithy/core": "^3.35.0", + "@smithy/signature-v4": "^5.7.3", + "@smithy/types": "^4.19.0", + "bowser": "^2.11.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/core/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-env": { + "version": "3.972.72", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-env/-/credential-provider-env-3.972.72.tgz", + "integrity": "sha512-xTKO/FWJPozTIXbozVnVGoNBhaGba8TBcx+KyUjRVeOlXE+dUc7GTR1cLvu0uTdIdmemzaFbqqCshXeZA1fZew==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-env/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-http": { + "version": "3.972.74", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-http/-/credential-provider-http-3.972.74.tgz", + "integrity": "sha512-u91E/hT8f4d1xy0Jl7VG4nVKJ3lxbrZkoBTeSVoJdWBiSEUMwMS/9+e0H/aJVQV//Lt5wuzP+E69v4aRSsNTmw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-http/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-ini": { + "version": "3.973.17", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-ini/-/credential-provider-ini-3.973.17.tgz", + "integrity": "sha512-ged4KXdBkvIC81bLvNHHuQKdKak/VXhQTR1NWYTTqW0474nlmsxy9O/vlgTIohDDWH3xpBdtVMZRyjb+DnocDA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/credential-provider-env": "^3.972.72", + "@aws-sdk/credential-provider-http": "^3.972.74", + "@aws-sdk/credential-provider-login": "^3.972.79", + "@aws-sdk/credential-provider-process": "^3.972.72", + "@aws-sdk/credential-provider-sso": "^3.973.16", + "@aws-sdk/credential-provider-web-identity": "^3.972.78", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/credential-provider-imds": "^4.5.2", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-ini/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-login": { + "version": "3.972.79", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-login/-/credential-provider-login-3.972.79.tgz", + "integrity": "sha512-L+Z85anONJd8MaiuraO4wRxATCdEejBZ3K3eymzWI5JPXa9sOS9CkIm72PBKqXKX+Z9p9NGMX5AIMXm0LEflgw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-login/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-node": { + "version": "3.972.84", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-node/-/credential-provider-node-3.972.84.tgz", + "integrity": "sha512-oHt854odINVwzwsh+c5x69j0ajm4DbqqqVJ+O1ECsCIZeMDAbzFpXItaqP7UZstJj/ATdTk/KFSH0LaNAgV+kA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/credential-provider-env": "^3.972.72", + "@aws-sdk/credential-provider-http": "^3.972.74", + "@aws-sdk/credential-provider-ini": "^3.973.17", + "@aws-sdk/credential-provider-process": "^3.972.72", + "@aws-sdk/credential-provider-sso": "^3.973.16", + "@aws-sdk/credential-provider-web-identity": "^3.972.78", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/credential-provider-imds": "^4.5.2", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-node/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-process": { + "version": "3.972.72", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-process/-/credential-provider-process-3.972.72.tgz", + "integrity": "sha512-rLIp2xbMjX/k9/od7APpqq1ZgXXnV0pOL1Th3ZsL8Wu0TRtBsDTVS8iPqcfRFcHakFxPvR04OSTv2ka2qOb/2A==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-process/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-sso": { + "version": "3.973.16", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-sso/-/credential-provider-sso-3.973.16.tgz", + "integrity": "sha512-IGihaJfFZYacJJr/odqILCoK7W/mvrZ7cuK7ECn3sAu4vLC6u0V8bS7mCGbdugJ8Aum2tnvqmx0F2MRFp2rn9g==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/token-providers": "3.1138.0", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-sso/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/credential-provider-web-identity": { + "version": "3.972.78", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-web-identity/-/credential-provider-web-identity-3.972.78.tgz", + "integrity": "sha512-/y9WvNtlcPBGLR0qc1a+9J/xtYZfVczvLUOuXaVWylzttH7ewsxwHtjmiJSolNrVSDorIxHGHMU61CbonRkmwA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-web-identity/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/middleware-sdk-s3": { + "version": "3.972.77", + "resolved": "https://registry.npmjs.org/@aws-sdk/middleware-sdk-s3/-/middleware-sdk-s3-3.972.77.tgz", + "integrity": "sha512-E7W2UOeUoc+lg3uIfR/dM7ZwusHwhBQrKMnlkRv4EXRR+C0YtV1pg25xC7GdZIhXH+NAMgZPCbE7o5to2cjFiw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/middleware-sdk-s3/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/nested-clients": { + "version": "3.997.46", + "resolved": "https://registry.npmjs.org/@aws-sdk/nested-clients/-/nested-clients-3.997.46.tgz", + "integrity": "sha512-oRxtBcka/JGHGs9l9p9IVajGoTP8vTPmoAzdHGy4Qcy9P5vPnDf6nhIeM/COQNY9k/OahImTRaLkHftoXvfcmQ==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/nested-clients/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/signature-v4-multi-region": { + "version": "3.996.47", + "resolved": "https://registry.npmjs.org/@aws-sdk/signature-v4-multi-region/-/signature-v4-multi-region-3.996.47.tgz", + "integrity": "sha512-Zk08macMvQTHzQJCLJVkOlviVoqwYMrpXv4lmLN7b7sAbiMoOK7Go0NYdR5UeF+MW8LIbRmwrNy9u/5VvX1U5g==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/types": "^3.974.6", + "@smithy/signature-v4": "^5.7.3", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/signature-v4-multi-region/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/token-providers": { + "version": "3.1138.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/token-providers/-/token-providers-3.1138.0.tgz", + "integrity": "sha512-GpyAr0DD63YOEmYFM6Df+gJuIgC92MMTiBK4FTKfxii5MJ9ge20epR7LyroulscYlG89J+ZB2ivFDPjvfQhzdw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/token-providers/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/types": { + "version": "3.974.6", + "resolved": "https://registry.npmjs.org/@aws-sdk/types/-/types-3.974.6.tgz", + "integrity": "sha512-v/clNZzZnDxGyvpHMOGpJKVXFAExJzUNAAjaWGdcx8QAcXLGwTaOkw33p5SHAi0YAioK32xB3hWwOekRVfmfKg==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/types/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws-sdk/xml-builder": { + "version": "3.972.41", + "resolved": "https://registry.npmjs.org/@aws-sdk/xml-builder/-/xml-builder-3.972.41.tgz", + "integrity": "sha512-ctjVSyCMegrWfXlx6VqzSBFI6UqmQ5ZlnfMhdLIiWmhoH8UAQxSCP5N3OpG7X3k4LnS7ou74C4mt20+bfTW2aQ==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/xml-builder/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@aws/lambda-invoke-store": { + "version": "0.3.0", + "resolved": "https://registry.npmjs.org/@aws/lambda-invoke-store/-/lambda-invoke-store-0.3.0.tgz", + "integrity": "sha512-sl4Bm6yiMNYrZKkqqDFWN0UfnWhlS8ivKxrYl+6t0gCLrqr8y3B2IqZZbFRkfaVVp7C/baApyh71P+LeE1A2sQ==", + "license": "Apache-2.0", + "engines": { + "node": ">=18.0.0" + } + }, "node_modules/@babel/code-frame": { "version": "7.29.7", "resolved": "https://registry.npmjs.org/@babel/code-frame/-/code-frame-7.29.7.tgz", @@ -4224,6 +4635,123 @@ "url": "https://ko-fi.com/dangreen" } }, + "node_modules/@smithy/core": { + "version": "3.35.0", + "resolved": "https://registry.npmjs.org/@smithy/core/-/core-3.35.0.tgz", + "integrity": "sha512-zRMhfkByhT2snNdr1si24vJitU6Cr9ix2MikUfWmkAgp4jrNP0GcKSP5YvwQ+TlI8AZXER5QOGJn3JsVtSD9/A==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/core/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@smithy/credential-provider-imds": { + "version": "4.5.2", + "resolved": "https://registry.npmjs.org/@smithy/credential-provider-imds/-/credential-provider-imds-4.5.2.tgz", + "integrity": "sha512-A9uSdn72ozbRUSit0eib0TW7nXuNPlaeM0zcGkJ+nE6tFcSDbnmtwoxbTCFBukVQcszDAyvsd7+rTduPTXpygg==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.2", + "@smithy/types": "^4.17.2", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/credential-provider-imds/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@smithy/fetch-http-handler": { + "version": "5.8.0", + "resolved": "https://registry.npmjs.org/@smithy/fetch-http-handler/-/fetch-http-handler-5.8.0.tgz", + "integrity": "sha512-ycSJu3tFAQ4v04CBB0agqFMVsSQ1iG3yw+SpgxRqKfaURpQD4CZ8Wn0zPMmSnOuTpTh65Vz+EA0rMrw089wvkA==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.3", + "@smithy/types": "^4.18.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/fetch-http-handler/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@smithy/node-http-handler": { + "version": "4.12.1", + "resolved": "https://registry.npmjs.org/@smithy/node-http-handler/-/node-http-handler-4.12.1.tgz", + "integrity": "sha512-ThMkboGeONWXAelq9FvGsuJC4rOi+qyC4/zhUF58xYpxUg5sQKx2VXZYJmtNjr4dSuBJ1HeJXETQILCz3wOHvw==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.3", + "@smithy/types": "^4.18.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/node-http-handler/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@smithy/signature-v4": { + "version": "5.7.4", + "resolved": "https://registry.npmjs.org/@smithy/signature-v4/-/signature-v4-5.7.4.tgz", + "integrity": "sha512-tHy0K0VtqNd5Y7Y41h0a0Lhh0L1GzC08dTWg0F7vRJWFtTENg7IZikf3wQkanYIRdb7ngoIPMTmqgUi401fEeQ==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/signature-v4/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, + "node_modules/@smithy/types": { + "version": "4.19.0", + "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.19.0.tgz", + "integrity": "sha512-r7jh49VJxGerfAcTQA6gXcKc+98zOp/tqRwzYjgOE+iSQsP6cEU1hq2QzbuipmP68QtYdY9wKEhiCQZIzHgZ4Q==", + "license": "Apache-2.0", + "dependencies": { + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/types/node_modules/tslib": { + "version": "2.8.1", + "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.8.1.tgz", + "integrity": "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w==", + "license": "0BSD" + }, "node_modules/@standard-schema/spec": { "version": "1.1.0", "resolved": "https://registry.npmjs.org/@standard-schema/spec/-/spec-1.1.0.tgz", @@ -5748,6 +6276,12 @@ "integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==", "license": "MIT" }, + "node_modules/bowser": { + "version": "2.14.1", + "resolved": "https://registry.npmjs.org/bowser/-/bowser-2.14.1.tgz", + "integrity": "sha512-tzPjzCxygAKWFOJP011oxFHs57HzIhOEracIgAePE4pqB3LikALKnSzUyU4MGs9/iCEUuHlAJTjTc5M+u7YEGg==", + "license": "MIT" + }, "node_modules/brace-expansion": { "version": "5.0.7", "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.7.tgz", diff --git a/package.json b/package.json index 6fb2382..4d9db38 100644 --- a/package.json +++ b/package.json @@ -34,6 +34,7 @@ "node": ">=24.0.0" }, "dependencies": { + "@aws-sdk/client-s3": "^3.1142.0", "@godaddy/terminus": "^4.12.1", "@map-colonies/3d-shared": "file:../../3d-general/3d-shared", "@map-colonies/config": "^4.0.1", diff --git a/src/common/config.ts b/src/common/config.ts index 47a6218..7dfad2a 100644 --- a/src/common/config.ts +++ b/src/common/config.ts @@ -11,7 +11,45 @@ interface ExternalServicesConfig { catalog: string; } -type OpsTriggerConfigType = commonBoilerplateV3Type & { externalServices: ExternalServicesConfig }; +interface JobManagerConfig { + url: string; + ingestion: { + jobType: string; + taskType: string; + batches: number; + }; + delete: { + jobType: string; + taskType: string; + }; +} + +interface NFSConfig { + pvPath: string; +} + +interface S3Config { + accessKeyId: string; + secretAccessKey: string; + endpointUrl: string; + bucket: string; + region: string; + forcePathStyle: boolean; + sslEnabled: boolean; + maxAttempts: number; +} + +type ProviderSource = 'NFS' | 'S3'; + +type OpsTriggerConfigType = commonBoilerplateV3Type & { + externalServices: ExternalServicesConfig; + jobManager: JobManagerConfig; + provider: ProviderSource; + // eslint-disable-next-line @typescript-eslint/naming-convention + NFS: NFSConfig; + // eslint-disable-next-line @typescript-eslint/naming-convention + S3: S3Config; +}; type ConfigType = ConfigInstance; @@ -22,8 +60,30 @@ const opsTriggerConfigSchema = { { $ref: commonBoilerplateV3.$id }, { type: 'object', - required: ['externalServices'], + required: ['externalServices', 'jobManager', 'provider'], properties: { + provider: { type: 'string', enum: ['NFS', 'S3'] }, + // eslint-disable-next-line @typescript-eslint/naming-convention + NFS: { + type: 'object', + required: ['pvPath'], + properties: { pvPath: { type: 'string' } }, + }, + // eslint-disable-next-line @typescript-eslint/naming-convention + S3: { + type: 'object', + required: ['accessKeyId', 'secretAccessKey', 'endpointUrl', 'bucket', 'region', 'forcePathStyle', 'sslEnabled', 'maxAttempts'], + properties: { + accessKeyId: { type: 'string' }, + secretAccessKey: { type: 'string' }, + endpointUrl: { type: 'string' }, + bucket: { type: 'string' }, + region: { type: 'string' }, + forcePathStyle: { type: 'boolean' }, + sslEnabled: { type: 'boolean' }, + maxAttempts: { type: 'number' }, + }, + }, externalServices: { type: 'object', required: ['lookupTables', 'catalog'], @@ -39,6 +99,30 @@ const opsTriggerConfigSchema = { catalog: { type: 'string' }, }, }, + jobManager: { + type: 'object', + required: ['url', 'ingestion', 'delete'], + properties: { + url: { type: 'string' }, + ingestion: { + type: 'object', + required: ['jobType'], + properties: { + jobType: { type: 'string' }, + taskType: { type: 'string' }, + batches: { type: 'number' }, + }, + }, + delete: { + type: 'object', + required: ['jobType'], + properties: { + jobType: { type: 'string' }, + taskType: { type: 'string' }, + }, + }, + }, + }, }, }, ], @@ -61,4 +145,4 @@ function getConfig(): ConfigType { } export { getConfig, initConfig }; -export type { ConfigType, LookupTablesConfig, ExternalServicesConfig }; +export type { ConfigType, LookupTablesConfig, ExternalServicesConfig, JobManagerConfig, NFSConfig, S3Config, ProviderSource }; diff --git a/src/common/constants.ts b/src/common/constants.ts index 6e7c100..492751a 100644 --- a/src/common/constants.ts +++ b/src/common/constants.ts @@ -6,11 +6,26 @@ export const DEFAULT_SERVER_PORT = 80; export const IGNORED_OUTGOING_TRACE_ROUTES = [/^.*\/v1\/metrics.*$/]; export const IGNORED_INCOMING_TRACE_ROUTES = [/^.*\/docs.*$/]; +export const DOMAIN = '3D'; + +/* eslint-disable @typescript-eslint/naming-convention */ +export const STAGE_TYPES = { + DATA_EXTRACTION: 'data-extraction', + VALIDATION: 'validation', + CREATE_UPLOAD_MODEL_TASKS: 'create-upload-model-tasks', + UPLOAD_MODEL_DATA: 'upload-model-data', + CLEAR_DATA: 'clear-data', + UPLOAD_MODEL_PARTS: 'upload-model-parts', + INGESTION_FINALIZER: 'ingestion-finalizer', +} satisfies Record; +/* eslint-enable @typescript-eslint/naming-convention */ + /* eslint-disable @typescript-eslint/naming-convention */ export const SERVICES = { LOGGER: Symbol('Logger'), CONFIG: Symbol('Config'), TRACER: Symbol('Tracer'), METRICS: Symbol('METRICS'), + PROVIDER: Symbol('Provider'), } satisfies Record; /* eslint-enable @typescript-eslint/naming-convention */ diff --git a/src/common/util.ts b/src/common/util.ts new file mode 100644 index 0000000..542af20 --- /dev/null +++ b/src/common/util.ts @@ -0,0 +1,4 @@ +export const is3tz = (modelPath: string): boolean => modelPath.toLowerCase().endsWith('.3tz'); + +export const buildModelFilePath = (modelPath: string, tilesetFilename: string): string => + is3tz(modelPath) ? modelPath : `${modelPath}/${tilesetFilename}`; diff --git a/src/containerConfig.ts b/src/containerConfig.ts index afe665b..3dbe41f 100644 --- a/src/containerConfig.ts +++ b/src/containerConfig.ts @@ -7,6 +7,7 @@ import { type InjectionObject, registerDependencies } from '@common/dependencyRe import { SERVICES, SERVICE_NAME } from '@common/constants'; import { getTracing } from '@common/tracing'; import { recordRouterFactory, RECORD_ROUTER_SYMBOL } from './record/routes/recordRouter'; +import { providerFactory } from './providers/getProvider'; import { getConfig } from './common/config'; export interface RegisterOptions { @@ -30,6 +31,7 @@ export const registerExternalValues = async (options?: RegisterOptions): Promise { token: SERVICES.LOGGER, provider: { useValue: logger } }, { token: SERVICES.TRACER, provider: { useValue: tracer } }, { token: SERVICES.METRICS, provider: { useValue: metricsRegistry } }, + { token: SERVICES.PROVIDER, provider: { useFactory: providerFactory } }, { token: RECORD_ROUTER_SYMBOL, provider: { useFactory: recordRouterFactory } }, { token: 'onSignal', diff --git a/src/externalServices/catalog/interfaces.ts b/src/externalServices/catalog/interfaces.ts index 7aa1b75..0b0609b 100644 --- a/src/externalServices/catalog/interfaces.ts +++ b/src/externalServices/catalog/interfaces.ts @@ -2,7 +2,9 @@ export interface Record3D { id: string; productId?: string; productName?: string; + productType?: string; productVersion?: number; + producerName?: string; productStatus?: string; links?: string; } diff --git a/src/externalServices/jobnik/jobnikClient.ts b/src/externalServices/jobnik/jobnikClient.ts new file mode 100644 index 0000000..8315a3c --- /dev/null +++ b/src/externalServices/jobnik/jobnikClient.ts @@ -0,0 +1,105 @@ +import { inject, injectable } from 'tsyringe'; +import type { Logger } from '@map-colonies/js-logger'; +import type { Registry } from 'prom-client'; +import { JobnikSDK } from '@map-colonies/jobnik-sdk'; +import { SERVICES, STAGE_TYPES } from '@common/constants'; +import { is3tz } from '@common/util'; +import type { ConfigType, JobManagerConfig } from '@common/config'; +import type { LogContext } from '@common/interfaces'; +import type { IngestionPayload, JobResponse } from '../../record/models/recordManager'; +import type { Record3D } from '../catalog/interfaces'; + +interface StageDescriptor { + type: string; + data: Record; + task?: Record; + only3tz?: boolean; +} + +@injectable() +export class JobnikClient { + private readonly logContext: LogContext; + private readonly producer: ReturnType; + private readonly jobManager: JobManagerConfig; + + public constructor( + @inject(SERVICES.CONFIG) private readonly config: ConfigType, + @inject(SERVICES.LOGGER) private readonly logger: Logger, + @inject(SERVICES.METRICS) private readonly metricsRegistry: Registry + ) { + this.jobManager = this.config.get('jobManager'); + const sdk = new JobnikSDK({ baseUrl: this.jobManager.url, metricsRegistry: this.metricsRegistry }); + this.producer = sdk.getProducer(); + this.logContext = { + fileName: __filename, + class: JobnikClient.name, + }; + } + + public async createIngestionJob(payload: IngestionPayload): Promise { + const logContext = { ...this.logContext, function: this.createIngestionJob.name }; + const isArchive = is3tz(payload.modelPath); + + const job = await this.producer.createJob({ + name: this.jobManager.ingestion.jobType, + data: { modelPath: payload.modelPath, tilesetFilename: payload.tilesetFilename, metadata: payload.metadata }, + }); + + const stages: StageDescriptor[] = [ + { + type: STAGE_TYPES.DATA_EXTRACTION, + data: { modelPath: payload.modelPath }, + task: { modelPath: payload.modelPath, tilesetFilename: payload.tilesetFilename }, + only3tz: true, + }, + { + type: STAGE_TYPES.VALIDATION, + data: {}, + task: { modelPath: payload.modelPath, tilesetFilename: payload.tilesetFilename, metadata: payload.metadata }, + }, + { + type: STAGE_TYPES.CREATE_UPLOAD_MODEL_TASKS, + data: {}, + task: { modelPath: payload.modelPath, tilesetFilename: payload.tilesetFilename }, + }, + { type: STAGE_TYPES.UPLOAD_MODEL_DATA, data: {} }, + { type: STAGE_TYPES.CLEAR_DATA, data: { modelPath: payload.modelPath }, task: { modelPath: payload.modelPath }, only3tz: true }, + { type: STAGE_TYPES.UPLOAD_MODEL_PARTS, data: {}, task: { metadata: payload.metadata } }, + { type: STAGE_TYPES.INGESTION_FINALIZER, data: { metadata: payload.metadata }, task: { metadata: payload.metadata } }, + ]; + + for (const descriptor of stages) { + if (descriptor.only3tz === true && !isArchive) { + continue; + } + const stage = await this.producer.createStage(job.id, { type: descriptor.type, data: descriptor.data }); + if (descriptor.task !== undefined) { + await this.producer.createTasks(stage.id, descriptor.type, [{ data: descriptor.task }]); + } + } + + this.logger.info({ msg: 'ingestion job created', logContext, jobId: job.id, is3tz: isArchive }); + + return { jobId: job.id, status: job.status }; + } + + public async createDeleteJob(record: Record3D): Promise { + const logContext = { ...this.logContext, function: this.createDeleteJob.name }; + + const job = await this.producer.createJob({ + name: this.jobManager.delete.jobType, + data: { + modelId: record.id, + productId: record.productId, + productVersion: record.productVersion, + productName: record.productName, + productType: record.productType, + producerName: record.producerName, + }, + }); + + this.logger.info({ msg: 'delete job created', logContext, jobId: job.id, recordId: record.id }); + + return { jobId: job.id, status: job.status }; + } +} diff --git a/src/providers/getProvider.ts b/src/providers/getProvider.ts new file mode 100644 index 0000000..26e6d30 --- /dev/null +++ b/src/providers/getProvider.ts @@ -0,0 +1,11 @@ +import type { FactoryFunction } from 'tsyringe'; +import { SERVICES } from '@common/constants'; +import type { ConfigType } from '@common/config'; +import { NFSProvider } from './nfsProvider'; +import { S3Provider } from './s3Provider'; +import type { Provider } from './interfaces'; + +export const providerFactory: FactoryFunction = (dependencyContainer) => { + const config = dependencyContainer.resolve(SERVICES.CONFIG); + return config.get('provider') === 'S3' ? dependencyContainer.resolve(S3Provider) : dependencyContainer.resolve(NFSProvider); +}; diff --git a/src/providers/interfaces.ts b/src/providers/interfaces.ts new file mode 100644 index 0000000..87ad465 --- /dev/null +++ b/src/providers/interfaces.ts @@ -0,0 +1,3 @@ +export interface Provider { + fileExists: (relativePath: string) => Promise; +} diff --git a/src/providers/nfsProvider.ts b/src/providers/nfsProvider.ts new file mode 100644 index 0000000..a4cb296 --- /dev/null +++ b/src/providers/nfsProvider.ts @@ -0,0 +1,43 @@ +import { access } from 'node:fs/promises'; +import { join } from 'node:path'; +import { inject, injectable } from 'tsyringe'; +import type { Logger } from '@map-colonies/js-logger'; +import { StatusCodes } from 'http-status-codes'; +import { SERVICES } from '@common/constants'; +import type { ConfigType } from '@common/config'; +import type { LogContext } from '@common/interfaces'; +import { AppError } from '@common/appError'; +import type { Provider } from './interfaces'; + +@injectable() +export class NFSProvider implements Provider { + private readonly logContext: LogContext; + private readonly pvPath: string; + + public constructor( + @inject(SERVICES.CONFIG) private readonly config: ConfigType, + @inject(SERVICES.LOGGER) private readonly logger: Logger + ) { + this.pvPath = this.config.get('NFS').pvPath; + this.logContext = { + fileName: __filename, + class: NFSProvider.name, + }; + } + + public async fileExists(relativePath: string): Promise { + const logContext = { ...this.logContext, function: this.fileExists.name }; + const fullPath = join(this.pvPath, relativePath); + try { + await access(fullPath); + return true; + } catch (err) { + if ((err as NodeJS.ErrnoException).code === 'ENOENT') { + this.logger.debug({ msg: `path does not exist: ${fullPath}`, logContext }); + return false; + } + this.logger.error({ msg: 'something went wrong with the NFS provider', logContext, fullPath, err }); + throw new AppError('nfs', StatusCodes.INTERNAL_SERVER_ERROR, 'there is a problem with the NFS provider', true); + } + } +} diff --git a/src/providers/s3Provider.ts b/src/providers/s3Provider.ts new file mode 100644 index 0000000..7f95301 --- /dev/null +++ b/src/providers/s3Provider.ts @@ -0,0 +1,56 @@ +import { inject, injectable } from 'tsyringe'; +import type { Logger } from '@map-colonies/js-logger'; +import { S3Client, HeadObjectCommand, type S3ClientConfig } from '@aws-sdk/client-s3'; +import { StatusCodes } from 'http-status-codes'; +import { SERVICES } from '@common/constants'; +import type { ConfigType, S3Config } from '@common/config'; +import type { LogContext } from '@common/interfaces'; +import { AppError } from '@common/appError'; +import type { Provider } from './interfaces'; + +@injectable() +export class S3Provider implements Provider { + private readonly logContext: LogContext; + private readonly s3Config: S3Config; + private readonly s3: S3Client; + + public constructor( + @inject(SERVICES.CONFIG) private readonly config: ConfigType, + @inject(SERVICES.LOGGER) private readonly logger: Logger + ) { + this.s3Config = this.config.get('S3'); + const s3ClientConfig: S3ClientConfig = { + endpoint: this.s3Config.endpointUrl, + forcePathStyle: this.s3Config.forcePathStyle, + credentials: { + accessKeyId: this.s3Config.accessKeyId, + secretAccessKey: this.s3Config.secretAccessKey, + }, + region: this.s3Config.region, + maxAttempts: this.s3Config.maxAttempts, + tls: this.s3Config.sslEnabled, + }; + this.s3 = new S3Client(s3ClientConfig); + this.logContext = { + fileName: __filename, + class: S3Provider.name, + }; + } + + public async fileExists(key: string): Promise { + const logContext = { ...this.logContext, function: this.fileExists.name }; + try { + // eslint-disable-next-line @typescript-eslint/naming-convention + await this.s3.send(new HeadObjectCommand({ Bucket: this.s3Config.bucket, Key: key })); + return true; + } catch (err) { + const statusCode = (err as { $metadata?: { httpStatusCode?: number } }).$metadata?.httpStatusCode; + if (statusCode === StatusCodes.NOT_FOUND) { + this.logger.debug({ msg: `key does not exist: ${key}`, logContext }); + return false; + } + this.logger.error({ msg: 'something went wrong with S3', logContext, key, err }); + throw new AppError('s3', StatusCodes.INTERNAL_SERVER_ERROR, 'there is a problem with S3', true); + } + } +} diff --git a/src/record/controllers/recordController.ts b/src/record/controllers/recordController.ts index 84fed2a..af0b6ee 100644 --- a/src/record/controllers/recordController.ts +++ b/src/record/controllers/recordController.ts @@ -40,11 +40,11 @@ export class RecordController { } }; - public deleteRecord: TypedRequestHandlers['deleteRecord'] = (req, res, next) => { + public deleteRecord: TypedRequestHandlers['deleteRecord'] = async (req, res, next) => { const logContext = { ...this.logContext, function: this.deleteRecord.name }; const { id } = req.params; try { - const job = this.manager.deleteRecord(id); + const job = await this.manager.deleteRecord(id); return res.status(StatusCodes.OK).json(job); } catch (err) { this.logger.error({ msg: 'failed to create delete job', logContext, err, recordId: id }); diff --git a/src/record/models/recordManager.ts b/src/record/models/recordManager.ts index 7f907ae..37867d8 100644 --- a/src/record/models/recordManager.ts +++ b/src/record/models/recordManager.ts @@ -4,6 +4,7 @@ import type { components } from '@openapi'; import { SERVICES } from '@common/constants'; import type { LogContext } from '@common/interfaces'; import { ValidationManager } from '../../validator/validationManager'; +import { JobnikClient } from '../../externalServices/jobnik/jobnikClient'; export type IngestionPayload = components['schemas']['ingestionPayload']; export type UpdatePayload = components['schemas']['updatePayload']; @@ -17,7 +18,8 @@ export class RecordManager { public constructor( @inject(SERVICES.LOGGER) private readonly logger: Logger, - @inject(ValidationManager) private readonly validator: ValidationManager + @inject(ValidationManager) private readonly validator: ValidationManager, + @inject(JobnikClient) private readonly jobnik: JobnikClient ) { this.logContext = { fileName: __filename, @@ -28,14 +30,15 @@ export class RecordManager { public async createIngestion(payload: IngestionPayload): Promise { const logContext = { ...this.logContext, function: this.createIngestion.name }; this.logger.info({ msg: 'creating ingestion job', logContext, modelPath: payload.modelPath, tilesetFilename: payload.tilesetFilename }); - await this.validator.validateIngestion(payload.metadata); - return { jobId: 'stub-ingestion-job-id', status: 'PENDING' }; + await this.validator.validateIngestion(payload); + return this.jobnik.createIngestionJob(payload); } - public deleteRecord(id: string): JobResponse { + public async deleteRecord(id: string): Promise { const logContext = { ...this.logContext, function: this.deleteRecord.name }; this.logger.info({ msg: 'creating delete job', logContext, recordId: id }); - return { jobId: 'stub-delete-job-id', status: 'PENDING' }; + const record = await this.validator.validateDelete(id); + return this.jobnik.createDeleteJob(record); } public updateMetadata(id: string, update: UpdatePayload): AckResponse { diff --git a/src/validator/validationManager.ts b/src/validator/validationManager.ts index bd97e2d..1ce2a50 100644 --- a/src/validator/validationManager.ts +++ b/src/validator/validationManager.ts @@ -5,14 +5,25 @@ import { new3DLayerMetadataSchema, geometrySchema } from '@map-colonies/3d-share import { SERVICES } from '@common/constants'; import type { LogContext } from '@common/interfaces'; import { AppError } from '@common/appError'; +import { buildModelFilePath } from '@common/util'; import { LookupTablesCall } from '../externalServices/lookupTables/lookupTablesCall'; import { CatalogCall } from '../externalServices/catalog/catalogCall'; +import type { Record3D } from '../externalServices/catalog/interfaces'; +import type { Provider } from '../providers/interfaces'; +import type { IngestionPayload } from '../record/models/recordManager'; + +const BLOCKED_DELETE_PRODUCT_TYPE = 'QuantizedMeshDTMBest'; +const RECORD_STATUS_UNPUBLISHED = 'UNPUBLISHED'; export const ERROR_METADATA_DATE = 'imagingTimeBeginUTC must not be later than imagingTimeEndUTC'; export const ERROR_METADATA_MISSING_DATE = 'imagingTimeBeginUTC and imagingTimeEndUTC are required'; export const ERROR_METADATA_INVALID_DATE = 'imagingTimeBeginUTC and imagingTimeEndUTC must be valid dates'; export const ERROR_METADATA_FOOTPRINT = 'Invalid footprint! Must be a GeoJSON Polygon or MultiPolygon with all-2D or all-3D coordinates'; export const ERROR_METADATA_PRODUCT_NAME_UNIQUE = 'product name is not unique!'; +export const ERROR_DELETE_RECORD_NOT_FOUND = "recordId doesn't match exactly one existing record"; +export const ERROR_DELETE_PRODUCT_TYPE = 'Cannot delete a record whose productType is "QuantizedMeshDTMBest"'; +export const ERROR_DELETE_STATUS = 'Cannot delete a record whose productStatus is not "UNPUBLISHED"'; +export const ERROR_FILE_NOT_FOUND = 'The model files do not exist in the agreed storage'; @injectable() export class ValidationManager { @@ -21,7 +32,8 @@ export class ValidationManager { public constructor( @inject(SERVICES.LOGGER) private readonly logger: Logger, @inject(LookupTablesCall) private readonly lookupTables: LookupTablesCall, - @inject(CatalogCall) private readonly catalog: CatalogCall + @inject(CatalogCall) private readonly catalog: CatalogCall, + @inject(SERVICES.PROVIDER) private readonly provider: Provider ) { this.logContext = { fileName: __filename, @@ -29,10 +41,11 @@ export class ValidationManager { }; } - public async validateIngestion(metadata: Record): Promise { + public async validateIngestion(payload: IngestionPayload): Promise { const logContext = { ...this.logContext, function: this.validateIngestion.name }; this.logger.info({ msg: 'ingestion validation start', logContext }); + const metadata = payload.metadata; const parsed = new3DLayerMetadataSchema.safeParse(metadata); if (!parsed.success) { const message = parsed.error.issues.map((issue) => `${issue.path.join('.')}: ${issue.message}`).join('; '); @@ -48,6 +61,37 @@ export class ValidationManager { await this.validateClassification(parsed.data.classification); await this.validateProductNameUnique(parsed.data.productName); + await this.validateFileExists(payload.modelPath, payload.tilesetFilename); + } + + public async validateDelete(recordId: string): Promise { + const logContext = { ...this.logContext, function: this.validateDelete.name }; + const records = await this.catalog.findRecords({ id: recordId }); + this.logger.debug({ msg: 'delete validation', logContext, recordId, matches: records.length }); + + const [record] = records; + if (records.length !== 1 || record === undefined) { + throw new AppError('badRequest', StatusCodes.BAD_REQUEST, ERROR_DELETE_RECORD_NOT_FOUND, true); + } + + if (record.productType === BLOCKED_DELETE_PRODUCT_TYPE) { + throw new AppError('badRequest', StatusCodes.BAD_REQUEST, ERROR_DELETE_PRODUCT_TYPE, true); + } + if (record.productStatus !== RECORD_STATUS_UNPUBLISHED) { + throw new AppError('badRequest', StatusCodes.BAD_REQUEST, ERROR_DELETE_STATUS, true); + } + + return record; + } + + private async validateFileExists(modelPath: string, tilesetFilename: string): Promise { + const logContext = { ...this.logContext, function: this.validateFileExists.name }; + const path = buildModelFilePath(modelPath, tilesetFilename); + const exists = await this.provider.fileExists(path); + this.logger.debug({ msg: 'file existence validation', logContext, path, exists }); + if (!exists) { + throw new AppError('badRequest', StatusCodes.BAD_REQUEST, ERROR_FILE_NOT_FOUND, true); + } } private async validateProductNameUnique(productName: string): Promise { diff --git a/tests/integration/record/record.spec.ts b/tests/integration/record/record.spec.ts index 0a1ccfb..0275902 100644 --- a/tests/integration/record/record.spec.ts +++ b/tests/integration/record/record.spec.ts @@ -9,10 +9,29 @@ import { SERVICES } from '@common/constants'; import { initConfig } from '@src/common/config'; import { LookupTablesCall } from '@src/externalServices/lookupTables/lookupTablesCall'; import { CatalogCall } from '@src/externalServices/catalog/catalogCall'; +import { JobnikClient } from '@src/externalServices/jobnik/jobnikClient'; import { buildValidMetadata } from '@tests/helpers/metadata'; +const deletableRecord = { + id: 'rec-1', + productId: 'p-1', + productName: 'afula', + productType: '3DPhotoRealistic', + productVersion: 1, + producerName: 'IDFMU', + productStatus: 'UNPUBLISHED', +}; + const lookupStub = { getClassifications: vi.fn().mockResolvedValue(['abc123']) } as unknown as LookupTablesCall; -const catalogStub = { findRecords: vi.fn().mockResolvedValue([]) } as unknown as CatalogCall; +// findRecords by id (delete) → an existing deletable record; by productName (uniqueness) → none +const catalogStub = { + findRecords: vi.fn().mockImplementation((payload: { id?: string }) => (payload.id !== undefined ? [deletableRecord] : [])), +} as unknown as CatalogCall; +const jobnikStub = { + createIngestionJob: vi.fn().mockResolvedValue({ jobId: 'job-1', status: 'PENDING' }), + createDeleteJob: vi.fn().mockResolvedValue({ jobId: 'del-1', status: 'PENDING' }), +} as unknown as JobnikClient; +const providerStub = { fileExists: vi.fn().mockResolvedValue(true) }; const validMetadata = buildValidMetadata(); @@ -36,6 +55,8 @@ describe('record', function () { { token: SERVICES.TRACER, provider: { useValue: trace.getTracer('testTracer') } }, { token: LookupTablesCall, provider: { useValue: lookupStub } }, { token: CatalogCall, provider: { useValue: catalogStub } }, + { token: JobnikClient, provider: { useValue: jobnikStub } }, + { token: SERVICES.PROVIDER, provider: { useValue: providerStub } }, ], useChild: true, }); diff --git a/tests/unit/externalServices/jobnik/jobnikClient.spec.ts b/tests/unit/externalServices/jobnik/jobnikClient.spec.ts new file mode 100644 index 0000000..c34ae47 --- /dev/null +++ b/tests/unit/externalServices/jobnik/jobnikClient.spec.ts @@ -0,0 +1,89 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { jsLogger } from '@map-colonies/js-logger'; +import { Registry } from 'prom-client'; +import { JobnikClient } from '@src/externalServices/jobnik/jobnikClient'; +import { STAGE_TYPES } from '@src/common/constants'; +import type { ConfigType } from '@src/common/config'; +import type { IngestionPayload } from '@src/record/models/recordManager'; + +const jobManagerConfig = { + url: 'http://job-manager', + ingestion: { jobType: 'Ingestion_New_3D', taskType: 'tilesCopying', batches: 100 }, + delete: { jobType: 'Delete_3D', taskType: 'deleteModel' }, +}; + +const producerMock = vi.hoisted(() => ({ + createJob: vi.fn(), + createStage: vi.fn(), + createTasks: vi.fn(), +})); + +vi.mock('@map-colonies/jobnik-sdk', () => ({ + // eslint-disable-next-line @typescript-eslint/naming-convention + JobnikSDK: class { + public getProducer(): typeof producerMock { + return producerMock; + } + }, +})); + +const configStub = { get: (): unknown => jobManagerConfig } as unknown as ConfigType; + +const buildPayload = (modelPath: string): IngestionPayload => ({ + modelPath, + tilesetFilename: 'tileset.json', + metadata: { productName: 'afula' }, +}); + +describe('JobnikClient', function () { + let client: JobnikClient; + + beforeEach(async function () { + vi.clearAllMocks(); + producerMock.createJob.mockResolvedValue({ id: 'job-1', status: 'PENDING' }); + producerMock.createStage.mockResolvedValue({ id: 'stage-1' }); + producerMock.createTasks.mockResolvedValue([]); + client = new JobnikClient(configStub, await jsLogger({ enabled: false }), new Registry()); + }); + + it('should create the ingestion job and a validation stage for a folder input', async function () { + const result = await client.createIngestionJob(buildPayload('/shared/models/afula')); + + expect(producerMock.createJob).toHaveBeenCalledWith(expect.objectContaining({ name: jobManagerConfig.ingestion.jobType })); + const folderStageTypes = (producerMock.createStage.mock.calls as [string, { type: string }][]).map((call) => call[1].type); + expect(folderStageTypes).toEqual([ + STAGE_TYPES.VALIDATION, + STAGE_TYPES.CREATE_UPLOAD_MODEL_TASKS, + STAGE_TYPES.UPLOAD_MODEL_DATA, + STAGE_TYPES.UPLOAD_MODEL_PARTS, + STAGE_TYPES.INGESTION_FINALIZER, + ]); + const validationTaskCall = (producerMock.createTasks.mock.calls as [string, string, { data: { metadata?: unknown } }[]][]).find( + (call) => call[1] === STAGE_TYPES.VALIDATION + ); + expect(validationTaskCall?.[2][0]?.data.metadata).toEqual({ productName: 'afula' }); + expect(result).toEqual({ jobId: 'job-1', status: 'PENDING' }); + }); + + it('should also create the data-extraction and clear-data stages for a 3tz input', async function () { + await client.createIngestionJob(buildPayload('/shared/models/afula.3tz')); + + const stageTypes = (producerMock.createStage.mock.calls as [string, { type: string }][]).map((call) => call[1].type); + expect(stageTypes).toEqual([ + STAGE_TYPES.DATA_EXTRACTION, + STAGE_TYPES.VALIDATION, + STAGE_TYPES.CREATE_UPLOAD_MODEL_TASKS, + STAGE_TYPES.UPLOAD_MODEL_DATA, + STAGE_TYPES.CLEAR_DATA, + STAGE_TYPES.UPLOAD_MODEL_PARTS, + STAGE_TYPES.INGESTION_FINALIZER, + ]); + }); + + it('should create a delete job with the record details', async function () { + const result = await client.createDeleteJob({ id: 'rec-1', productId: 'p-1', productName: 'afula', productType: '3DPhotoRealistic' }); + + expect(producerMock.createJob).toHaveBeenCalledWith(expect.objectContaining({ name: jobManagerConfig.delete.jobType })); + expect(result).toEqual({ jobId: 'job-1', status: 'PENDING' }); + }); +}); diff --git a/tests/unit/providers/nfsProvider.spec.ts b/tests/unit/providers/nfsProvider.spec.ts new file mode 100644 index 0000000..a37add4 --- /dev/null +++ b/tests/unit/providers/nfsProvider.spec.ts @@ -0,0 +1,38 @@ +import { access } from 'node:fs/promises'; +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { jsLogger } from '@map-colonies/js-logger'; +import { NFSProvider } from '@src/providers/nfsProvider'; +import { AppError } from '@src/common/appError'; +import type { ConfigType } from '@src/common/config'; + +vi.mock('node:fs/promises', () => ({ access: vi.fn() })); +const mockedAccess = vi.mocked(access); + +const configStub = { get: (): unknown => ({ pvPath: '/data/models' }) } as unknown as ConfigType; + +describe('NFSProvider', function () { + let provider: NFSProvider; + + beforeEach(async function () { + vi.clearAllMocks(); + provider = new NFSProvider(configStub, await jsLogger({ enabled: false })); + }); + + it('should return true when the path is accessible', async function () { + mockedAccess.mockResolvedValue(undefined); + + await expect(provider.fileExists('afula/tileset.json')).resolves.toBe(true); + }); + + it('should return false when the path does not exist (ENOENT)', async function () { + mockedAccess.mockRejectedValue(Object.assign(new Error('not found'), { code: 'ENOENT' })); + + await expect(provider.fileExists('missing/tileset.json')).resolves.toBe(false); + }); + + it('should throw an AppError on non-ENOENT errors', async function () { + mockedAccess.mockRejectedValue(Object.assign(new Error('permission denied'), { code: 'EACCES' })); + + await expect(provider.fileExists('locked/tileset.json')).rejects.toThrow(AppError); + }); +}); diff --git a/tests/unit/providers/s3Provider.spec.ts b/tests/unit/providers/s3Provider.spec.ts new file mode 100644 index 0000000..4e0f1e9 --- /dev/null +++ b/tests/unit/providers/s3Provider.spec.ts @@ -0,0 +1,59 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { jsLogger } from '@map-colonies/js-logger'; +import { StatusCodes } from 'http-status-codes'; +import { S3Provider } from '@src/providers/s3Provider'; +import { AppError } from '@src/common/appError'; +import type { ConfigType } from '@src/common/config'; + +const sendMock = vi.hoisted(() => vi.fn()); + +vi.mock('@aws-sdk/client-s3', () => ({ + // eslint-disable-next-line @typescript-eslint/naming-convention + S3Client: class { + public send = sendMock; + }, + // eslint-disable-next-line @typescript-eslint/naming-convention + HeadObjectCommand: class { + public constructor(public readonly input: unknown) {} + }, +})); + +const configStub = { + get: (): unknown => ({ + accessKeyId: 'k', + secretAccessKey: 's', + endpointUrl: 'http://s3', + bucket: 'b', + region: 'us-east-1', + forcePathStyle: true, + sslEnabled: false, + maxAttempts: 3, + }), +} as unknown as ConfigType; + +describe('S3Provider', function () { + let provider: S3Provider; + + beforeEach(async function () { + vi.clearAllMocks(); + provider = new S3Provider(configStub, await jsLogger({ enabled: false })); + }); + + it('should return true when the object exists', async function () { + sendMock.mockResolvedValue({}); + + await expect(provider.fileExists('afula/tileset.json')).resolves.toBe(true); + }); + + it('should return false when the object is not found (404)', async function () { + sendMock.mockRejectedValue({ $metadata: { httpStatusCode: StatusCodes.NOT_FOUND } }); + + await expect(provider.fileExists('missing')).resolves.toBe(false); + }); + + it('should throw an AppError on other S3 errors', async function () { + sendMock.mockRejectedValue({ $metadata: { httpStatusCode: StatusCodes.INTERNAL_SERVER_ERROR } }); + + await expect(provider.fileExists('boom')).rejects.toThrow(AppError); + }); +}); diff --git a/tests/unit/record/models/recordManager.spec.ts b/tests/unit/record/models/recordManager.spec.ts index fde0ed2..403dc42 100644 --- a/tests/unit/record/models/recordManager.spec.ts +++ b/tests/unit/record/models/recordManager.spec.ts @@ -2,14 +2,22 @@ import { jsLogger } from '@map-colonies/js-logger'; import { describe, it, expect, beforeEach, vi } from 'vitest'; import { RecordManager, type IngestionPayload } from '@src/record/models/recordManager'; import type { ValidationManager } from '@src/validator/validationManager'; +import type { JobnikClient } from '@src/externalServices/jobnik/jobnikClient'; -const noopValidator = { validateIngestion: vi.fn().mockResolvedValue(undefined) } as unknown as ValidationManager; +const noopValidator = { + validateIngestion: vi.fn().mockResolvedValue(undefined), + validateDelete: vi.fn().mockResolvedValue({ id: 'rec-1', productName: 'afula' }), +} as unknown as ValidationManager; +const jobnikStub = { + createIngestionJob: vi.fn().mockResolvedValue({ jobId: 'job-1', status: 'PENDING' }), + createDeleteJob: vi.fn().mockResolvedValue({ jobId: 'del-1', status: 'PENDING' }), +} as unknown as JobnikClient; describe('RecordManager', function () { let manager: RecordManager; beforeEach(async function () { - manager = new RecordManager(await jsLogger({ enabled: false }), noopValidator); + manager = new RecordManager(await jsLogger({ enabled: false }), noopValidator, jobnikStub); }); describe('createIngestion', function () { @@ -28,8 +36,8 @@ describe('RecordManager', function () { }); describe('deleteRecord', function () { - it('should return a job response', function () { - const result = manager.deleteRecord('rec-1'); + it('should return a job response', async function () { + const result = await manager.deleteRecord('rec-1'); expect(result.jobId).toBeTypeOf('string'); expect(result.status).toBeTypeOf('string'); diff --git a/tests/unit/validator/validationManager.spec.ts b/tests/unit/validator/validationManager.spec.ts index c2638a8..c9231fc 100644 --- a/tests/unit/validator/validationManager.spec.ts +++ b/tests/unit/validator/validationManager.spec.ts @@ -7,40 +7,54 @@ import { ERROR_METADATA_FOOTPRINT, ERROR_METADATA_MISSING_DATE, ERROR_METADATA_INVALID_DATE, + ERROR_DELETE_RECORD_NOT_FOUND, + ERROR_DELETE_PRODUCT_TYPE, + ERROR_DELETE_STATUS, + ERROR_FILE_NOT_FOUND, } from '@src/validator/validationManager'; import { AppError } from '@src/common/appError'; import type { LookupTablesCall } from '@src/externalServices/lookupTables/lookupTablesCall'; import type { CatalogCall } from '@src/externalServices/catalog/catalogCall'; -import { buildValidMetadata as validMetadata } from '@tests/helpers/metadata'; +import type { Provider } from '@src/providers/interfaces'; +import type { IngestionPayload } from '@src/record/models/recordManager'; +import { buildValidMetadata } from '@tests/helpers/metadata'; const lookupStub = { getClassifications: vi.fn().mockResolvedValue(['abc123']) } as unknown as LookupTablesCall; +const ingest = (metadata: Record): IngestionPayload => ({ + modelPath: '/shared/models/afula', + tilesetFilename: 'tileset.json', + metadata, +}); + describe('ValidationManager', function () { let validator: ValidationManager; let catalogStub: CatalogCall; + let providerStub: Provider; beforeEach(async function () { catalogStub = { findRecords: vi.fn().mockResolvedValue([]) } as unknown as CatalogCall; - validator = new ValidationManager(await jsLogger({ enabled: false }), lookupStub, catalogStub); + providerStub = { fileExists: vi.fn().mockResolvedValue(true) }; + validator = new ValidationManager(await jsLogger({ enabled: false }), lookupStub, catalogStub, providerStub); }); - it('should pass a fully valid metadata object', async function () { - await expect(validator.validateIngestion(validMetadata())).resolves.toBeUndefined(); + it('should pass a fully valid ingestion payload', async function () { + await expect(validator.validateIngestion(ingest(buildValidMetadata()))).resolves.toBeUndefined(); }); it('should throw 400 when a required core field is missing', async function () { - const metadata = validMetadata(); + const metadata = buildValidMetadata(); delete metadata.productName; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(AppError); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(AppError); }); it('should throw 400 for an invalid productType', async function () { - const metadata = { ...validMetadata(), productType: 'NOT_A_3D_TYPE' }; + const metadata = { ...buildValidMetadata(), productType: 'NOT_A_3D_TYPE' }; let thrown: unknown; try { - await validator.validateIngestion(metadata); + await validator.validateIngestion(ingest(metadata)); } catch (err) { thrown = err; } @@ -50,39 +64,76 @@ describe('ValidationManager', function () { }); it('should throw the footprint error for a malformed geometry', async function () { - const metadata = { ...validMetadata(), footprint: { type: 'Point', coordinates: [34.45, 31.48] } }; + const metadata = { ...buildValidMetadata(), footprint: { type: 'Point', coordinates: [34.45, 31.48] } }; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(ERROR_METADATA_FOOTPRINT); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(ERROR_METADATA_FOOTPRINT); }); it('should throw the date error when imagingTimeBeginUTC is after imagingTimeEndUTC', async function () { - const metadata = { ...validMetadata(), imagingTimeBeginUTC: '2025-07-11T00:00:00.000Z' }; + const metadata = { ...buildValidMetadata(), imagingTimeBeginUTC: '2025-07-11T00:00:00.000Z' }; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(ERROR_METADATA_DATE); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(ERROR_METADATA_DATE); }); it('should throw the missing-date error when a source date is absent', async function () { - const metadata = validMetadata(); + const metadata = buildValidMetadata(); delete metadata.imagingTimeEndUTC; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(ERROR_METADATA_MISSING_DATE); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(ERROR_METADATA_MISSING_DATE); }); it('should throw the invalid-date error for an unparseable date', async function () { - const metadata = { ...validMetadata(), imagingTimeBeginUTC: 'not-a-date' }; + const metadata = { ...buildValidMetadata(), imagingTimeBeginUTC: 'not-a-date' }; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(ERROR_METADATA_INVALID_DATE); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(ERROR_METADATA_INVALID_DATE); }); it('should throw 400 when the classification is not in the lookup table', async function () { - const metadata = { ...validMetadata(), classification: 'not-a-real-classification' }; + const metadata = { ...buildValidMetadata(), classification: 'not-a-real-classification' }; - await expect(validator.validateIngestion(metadata)).rejects.toThrow(AppError); + await expect(validator.validateIngestion(ingest(metadata))).rejects.toThrow(AppError); }); it('should throw 400 when the product name already exists in the catalog', async function () { (catalogStub.findRecords as ReturnType).mockResolvedValueOnce([{ id: 'existing', productName: 'afula' }]); - await expect(validator.validateIngestion(validMetadata())).rejects.toThrow(AppError); + await expect(validator.validateIngestion(ingest(buildValidMetadata()))).rejects.toThrow(AppError); + }); + + it('should throw 400 when the model files do not exist', async function () { + (providerStub.fileExists as ReturnType).mockResolvedValueOnce(false); + + await expect(validator.validateIngestion(ingest(buildValidMetadata()))).rejects.toThrow(ERROR_FILE_NOT_FOUND); + }); + + describe('validateDelete', function () { + it('should return the record when it is deletable', async function () { + const record = { id: 'rec-1', productType: '3DPhotoRealistic', productStatus: 'UNPUBLISHED' }; + (catalogStub.findRecords as ReturnType).mockResolvedValueOnce([record]); + + await expect(validator.validateDelete('rec-1')).resolves.toEqual(record); + }); + + it('should throw when the record is not found', async function () { + (catalogStub.findRecords as ReturnType).mockResolvedValueOnce([]); + + await expect(validator.validateDelete('missing')).rejects.toThrow(ERROR_DELETE_RECORD_NOT_FOUND); + }); + + it('should throw when the productType is blocked', async function () { + (catalogStub.findRecords as ReturnType).mockResolvedValueOnce([ + { id: 'rec-1', productType: 'QuantizedMeshDTMBest', productStatus: 'UNPUBLISHED' }, + ]); + + await expect(validator.validateDelete('rec-1')).rejects.toThrow(ERROR_DELETE_PRODUCT_TYPE); + }); + + it('should throw when the record is not unpublished', async function () { + (catalogStub.findRecords as ReturnType).mockResolvedValueOnce([ + { id: 'rec-1', productType: '3DPhotoRealistic', productStatus: 'PUBLISHED' }, + ]); + + await expect(validator.validateDelete('rec-1')).rejects.toThrow(ERROR_DELETE_STATUS); + }); }); }); From 104fb380677585928f6e93b6b3a25bb4ec745c85 Mon Sep 17 00:00:00 2001 From: TULCHINSKI LIRAN Date: Tue, 29 Sep 2026 18:36:12 +0300 Subject: [PATCH 2/2] feat: reject duplicate in-flight ingestion jobs (MAPCO-11589) Add ValidationManager check that rejects an ingestion when an in-flight Jobnik ingestion job already exists for the same productName, complementing the catalog unique-name check (which only sees finished records). JobnikClient.hasInFlightIngestionJob queries GET /v1/jobs by job name and filters client-side on non-terminal status + data.metadata.productName, paginating through all pages. Query failures raise AppError. JobnikClient is now a singleton so RecordManager and ValidationManager share one SDK instance instead of double-registering its metrics. Interim until Jobnik ships native duplicate support (epic MAPCO-11833). --- src/common/constants.ts | 3 + src/externalServices/jobnik/jobnikClient.ts | 39 +++++++++- src/validator/validationManager.ts | 14 ++++ tests/integration/record/record.spec.ts | 1 + .../jobnik/jobnikClient.spec.ts | 76 +++++++++++++++++++ .../unit/validator/validationManager.spec.ts | 21 ++++- 6 files changed, 150 insertions(+), 4 deletions(-) diff --git a/src/common/constants.ts b/src/common/constants.ts index 492751a..9715684 100644 --- a/src/common/constants.ts +++ b/src/common/constants.ts @@ -1,4 +1,5 @@ import { readPackageJsonSync } from '@map-colonies/read-pkg'; +import type { Job } from '@map-colonies/jobnik-sdk'; export const SERVICE_NAME = readPackageJsonSync().name ?? 'unknown_service'; export const DEFAULT_SERVER_PORT = 80; @@ -20,6 +21,8 @@ export const STAGE_TYPES = { } satisfies Record; /* eslint-enable @typescript-eslint/naming-convention */ +export const IN_FLIGHT_JOB_STATUSES: readonly Job['status'][] = ['PENDING', 'IN_PROGRESS', 'PAUSED', 'CREATED']; + /* eslint-disable @typescript-eslint/naming-convention */ export const SERVICES = { LOGGER: Symbol('Logger'), diff --git a/src/externalServices/jobnik/jobnikClient.ts b/src/externalServices/jobnik/jobnikClient.ts index 8315a3c..593ce4e 100644 --- a/src/externalServices/jobnik/jobnikClient.ts +++ b/src/externalServices/jobnik/jobnikClient.ts @@ -1,9 +1,11 @@ -import { inject, injectable } from 'tsyringe'; +import { inject, singleton } from 'tsyringe'; +import { StatusCodes } from 'http-status-codes'; import type { Logger } from '@map-colonies/js-logger'; import type { Registry } from 'prom-client'; import { JobnikSDK } from '@map-colonies/jobnik-sdk'; -import { SERVICES, STAGE_TYPES } from '@common/constants'; +import { IN_FLIGHT_JOB_STATUSES, SERVICES, STAGE_TYPES } from '@common/constants'; import { is3tz } from '@common/util'; +import { AppError } from '@common/appError'; import type { ConfigType, JobManagerConfig } from '@common/config'; import type { LogContext } from '@common/interfaces'; import type { IngestionPayload, JobResponse } from '../../record/models/recordManager'; @@ -16,10 +18,11 @@ interface StageDescriptor { only3tz?: boolean; } -@injectable() +@singleton() export class JobnikClient { private readonly logContext: LogContext; private readonly producer: ReturnType; + private readonly apiClient: ReturnType; private readonly jobManager: JobManagerConfig; public constructor( @@ -30,12 +33,42 @@ export class JobnikClient { this.jobManager = this.config.get('jobManager'); const sdk = new JobnikSDK({ baseUrl: this.jobManager.url, metricsRegistry: this.metricsRegistry }); this.producer = sdk.getProducer(); + this.apiClient = sdk.getApiClient(); this.logContext = { fileName: __filename, class: JobnikClient.name, }; } + public async hasInFlightIngestionJob(productName: string): Promise { + const logContext = { ...this.logContext, function: this.hasInFlightIngestionJob.name }; + const pageSize = 100; + + for (let page = 1; ; page++) { + const { data, error } = await this.apiClient.GET('/v1/jobs', { + // eslint-disable-next-line @typescript-eslint/naming-convention -- Jobnik query params are snake_case + params: { query: { job_name: this.jobManager.ingestion.jobType, page, page_size: pageSize } }, + }); + if (error !== undefined) { + this.logger.error({ msg: 'failed querying Jobnik for in-flight jobs', logContext, productName, err: error }); + throw new AppError('jobnik', StatusCodes.INTERNAL_SERVER_ERROR, 'failed querying Jobnik for in-flight jobs', false); + } + + const jobs = data.items; + const match = jobs.some((job) => { + const metadata = (job.data as { metadata?: { productName?: string } }).metadata; + return IN_FLIGHT_JOB_STATUSES.includes(job.status) && metadata?.productName === productName; + }); + if (match) { + this.logger.debug({ msg: 'in-flight ingestion job found', logContext, productName }); + return true; + } + if (jobs.length < pageSize) { + return false; + } + } + } + public async createIngestionJob(payload: IngestionPayload): Promise { const logContext = { ...this.logContext, function: this.createIngestionJob.name }; const isArchive = is3tz(payload.modelPath); diff --git a/src/validator/validationManager.ts b/src/validator/validationManager.ts index 1ce2a50..07edbf0 100644 --- a/src/validator/validationManager.ts +++ b/src/validator/validationManager.ts @@ -8,6 +8,7 @@ import { AppError } from '@common/appError'; import { buildModelFilePath } from '@common/util'; import { LookupTablesCall } from '../externalServices/lookupTables/lookupTablesCall'; import { CatalogCall } from '../externalServices/catalog/catalogCall'; +import { JobnikClient } from '../externalServices/jobnik/jobnikClient'; import type { Record3D } from '../externalServices/catalog/interfaces'; import type { Provider } from '../providers/interfaces'; import type { IngestionPayload } from '../record/models/recordManager'; @@ -20,6 +21,7 @@ export const ERROR_METADATA_MISSING_DATE = 'imagingTimeBeginUTC and imagingTimeE export const ERROR_METADATA_INVALID_DATE = 'imagingTimeBeginUTC and imagingTimeEndUTC must be valid dates'; export const ERROR_METADATA_FOOTPRINT = 'Invalid footprint! Must be a GeoJSON Polygon or MultiPolygon with all-2D or all-3D coordinates'; export const ERROR_METADATA_PRODUCT_NAME_UNIQUE = 'product name is not unique!'; +export const ERROR_PRODUCT_NAME_IN_FLIGHT = 'product name already has an in-flight ingestion job'; export const ERROR_DELETE_RECORD_NOT_FOUND = "recordId doesn't match exactly one existing record"; export const ERROR_DELETE_PRODUCT_TYPE = 'Cannot delete a record whose productType is "QuantizedMeshDTMBest"'; export const ERROR_DELETE_STATUS = 'Cannot delete a record whose productStatus is not "UNPUBLISHED"'; @@ -33,6 +35,7 @@ export class ValidationManager { @inject(SERVICES.LOGGER) private readonly logger: Logger, @inject(LookupTablesCall) private readonly lookupTables: LookupTablesCall, @inject(CatalogCall) private readonly catalog: CatalogCall, + @inject(JobnikClient) private readonly jobnik: JobnikClient, @inject(SERVICES.PROVIDER) private readonly provider: Provider ) { this.logContext = { @@ -61,6 +64,7 @@ export class ValidationManager { await this.validateClassification(parsed.data.classification); await this.validateProductNameUnique(parsed.data.productName); + await this.validateProductNameNotInFlight(parsed.data.productName); await this.validateFileExists(payload.modelPath, payload.tilesetFilename); } @@ -104,6 +108,16 @@ export class ValidationManager { } } + private async validateProductNameNotInFlight(productName: string): Promise { + const logContext = { ...this.logContext, function: this.validateProductNameNotInFlight.name }; + const inFlight = await this.jobnik.hasInFlightIngestionJob(productName); + this.logger.debug({ msg: 'in-flight duplicate validation', logContext, productName, inFlight }); + + if (inFlight) { + throw new AppError('conflict', StatusCodes.CONFLICT, ERROR_PRODUCT_NAME_IN_FLIGHT, true); + } + } + private validateDates(start: unknown, end: unknown): void { if (start === undefined || end === undefined) { throw new AppError('badRequest', StatusCodes.BAD_REQUEST, ERROR_METADATA_MISSING_DATE, true); diff --git a/tests/integration/record/record.spec.ts b/tests/integration/record/record.spec.ts index 0275902..f640d2c 100644 --- a/tests/integration/record/record.spec.ts +++ b/tests/integration/record/record.spec.ts @@ -30,6 +30,7 @@ const catalogStub = { const jobnikStub = { createIngestionJob: vi.fn().mockResolvedValue({ jobId: 'job-1', status: 'PENDING' }), createDeleteJob: vi.fn().mockResolvedValue({ jobId: 'del-1', status: 'PENDING' }), + hasInFlightIngestionJob: vi.fn().mockResolvedValue(false), } as unknown as JobnikClient; const providerStub = { fileExists: vi.fn().mockResolvedValue(true) }; diff --git a/tests/unit/externalServices/jobnik/jobnikClient.spec.ts b/tests/unit/externalServices/jobnik/jobnikClient.spec.ts index c34ae47..395cc7e 100644 --- a/tests/unit/externalServices/jobnik/jobnikClient.spec.ts +++ b/tests/unit/externalServices/jobnik/jobnikClient.spec.ts @@ -18,12 +18,20 @@ const producerMock = vi.hoisted(() => ({ createTasks: vi.fn(), })); +const apiClientMock = vi.hoisted(() => ({ + // eslint-disable-next-line @typescript-eslint/naming-convention -- mirrors the SDK ApiClient.GET method + GET: vi.fn(), +})); + vi.mock('@map-colonies/jobnik-sdk', () => ({ // eslint-disable-next-line @typescript-eslint/naming-convention JobnikSDK: class { public getProducer(): typeof producerMock { return producerMock; } + public getApiClient(): typeof apiClientMock { + return apiClientMock; + } }, })); @@ -43,6 +51,7 @@ describe('JobnikClient', function () { producerMock.createJob.mockResolvedValue({ id: 'job-1', status: 'PENDING' }); producerMock.createStage.mockResolvedValue({ id: 'stage-1' }); producerMock.createTasks.mockResolvedValue([]); + apiClientMock.GET.mockResolvedValue({ data: { total: 0, items: [] } }); client = new JobnikClient(configStub, await jsLogger({ enabled: false }), new Registry()); }); @@ -50,7 +59,9 @@ describe('JobnikClient', function () { const result = await client.createIngestionJob(buildPayload('/shared/models/afula')); expect(producerMock.createJob).toHaveBeenCalledWith(expect.objectContaining({ name: jobManagerConfig.ingestion.jobType })); + const folderStageTypes = (producerMock.createStage.mock.calls as [string, { type: string }][]).map((call) => call[1].type); + expect(folderStageTypes).toEqual([ STAGE_TYPES.VALIDATION, STAGE_TYPES.CREATE_UPLOAD_MODEL_TASKS, @@ -58,9 +69,11 @@ describe('JobnikClient', function () { STAGE_TYPES.UPLOAD_MODEL_PARTS, STAGE_TYPES.INGESTION_FINALIZER, ]); + const validationTaskCall = (producerMock.createTasks.mock.calls as [string, string, { data: { metadata?: unknown } }[]][]).find( (call) => call[1] === STAGE_TYPES.VALIDATION ); + expect(validationTaskCall?.[2][0]?.data.metadata).toEqual({ productName: 'afula' }); expect(result).toEqual({ jobId: 'job-1', status: 'PENDING' }); }); @@ -69,6 +82,7 @@ describe('JobnikClient', function () { await client.createIngestionJob(buildPayload('/shared/models/afula.3tz')); const stageTypes = (producerMock.createStage.mock.calls as [string, { type: string }][]).map((call) => call[1].type); + expect(stageTypes).toEqual([ STAGE_TYPES.DATA_EXTRACTION, STAGE_TYPES.VALIDATION, @@ -86,4 +100,66 @@ describe('JobnikClient', function () { expect(producerMock.createJob).toHaveBeenCalledWith(expect.objectContaining({ name: jobManagerConfig.delete.jobType })); expect(result).toEqual({ jobId: 'job-1', status: 'PENDING' }); }); + + describe('hasInFlightIngestionJob', function () { + it('should query jobs by the ingestion job name', async function () { + await client.hasInFlightIngestionJob('afula'); + + expect(apiClientMock.GET).toHaveBeenCalledWith('/v1/jobs', { + // eslint-disable-next-line @typescript-eslint/naming-convention -- Jobnik query params are snake_case + params: { query: { job_name: jobManagerConfig.ingestion.jobType, page: 1, page_size: 100 } }, + }); + }); + + it('should return true when a non-terminal job matches the product name', async function () { + apiClientMock.GET.mockResolvedValueOnce({ + data: { total: 1, items: [{ status: 'IN_PROGRESS', data: { metadata: { productName: 'afula' } } }] }, + }); + + await expect(client.hasInFlightIngestionJob('afula')).resolves.toBe(true); + }); + + it('should return false when the only matching job is in a terminal status', async function () { + apiClientMock.GET.mockResolvedValueOnce({ + data: { total: 1, items: [{ status: 'COMPLETED', data: { metadata: { productName: 'afula' } } }] }, + }); + + await expect(client.hasInFlightIngestionJob('afula')).resolves.toBe(false); + }); + + it('should return false when no in-flight job matches the product name', async function () { + apiClientMock.GET.mockResolvedValueOnce({ + data: { total: 1, items: [{ status: 'PENDING', data: { metadata: { productName: 'haifa' } } }] }, + }); + + await expect(client.hasInFlightIngestionJob('afula')).resolves.toBe(false); + }); + + it('should page through results and match a job on a later page', async function () { + const fullPage = Array.from({ length: 100 }, () => ({ status: 'IN_PROGRESS', data: { metadata: { productName: 'haifa' } } })); + apiClientMock.GET.mockResolvedValueOnce({ data: { total: 101, items: fullPage } }); + apiClientMock.GET.mockResolvedValueOnce({ + data: { total: 101, items: [{ status: 'PENDING', data: { metadata: { productName: 'afula' } } }] }, + }); + + await expect(client.hasInFlightIngestionJob('afula')).resolves.toBe(true); + expect(apiClientMock.GET).toHaveBeenNthCalledWith(2, '/v1/jobs', { + // eslint-disable-next-line @typescript-eslint/naming-convention -- Jobnik query params are snake_case + params: { query: { job_name: jobManagerConfig.ingestion.jobType, page: 2, page_size: 100 } }, + }); + }); + + it('should stop paging once a page is shorter than the page size', async function () { + apiClientMock.GET.mockResolvedValueOnce({ data: { total: 1, items: [{ status: 'PENDING', data: { metadata: { productName: 'haifa' } } }] } }); + + await expect(client.hasInFlightIngestionJob('afula')).resolves.toBe(false); + expect(apiClientMock.GET).toHaveBeenCalledTimes(1); + }); + + it('should throw an AppError when the jobs query returns an error', async function () { + apiClientMock.GET.mockResolvedValueOnce({ error: { message: 'boom' } }); + + await expect(client.hasInFlightIngestionJob('afula')).rejects.toThrow(/failed querying Jobnik/); + }); + }); }); diff --git a/tests/unit/validator/validationManager.spec.ts b/tests/unit/validator/validationManager.spec.ts index c9231fc..d88ae0c 100644 --- a/tests/unit/validator/validationManager.spec.ts +++ b/tests/unit/validator/validationManager.spec.ts @@ -11,10 +11,12 @@ import { ERROR_DELETE_PRODUCT_TYPE, ERROR_DELETE_STATUS, ERROR_FILE_NOT_FOUND, + ERROR_PRODUCT_NAME_IN_FLIGHT, } from '@src/validator/validationManager'; import { AppError } from '@src/common/appError'; import type { LookupTablesCall } from '@src/externalServices/lookupTables/lookupTablesCall'; import type { CatalogCall } from '@src/externalServices/catalog/catalogCall'; +import type { JobnikClient } from '@src/externalServices/jobnik/jobnikClient'; import type { Provider } from '@src/providers/interfaces'; import type { IngestionPayload } from '@src/record/models/recordManager'; import { buildValidMetadata } from '@tests/helpers/metadata'; @@ -30,12 +32,14 @@ const ingest = (metadata: Record): IngestionPayload => ({ describe('ValidationManager', function () { let validator: ValidationManager; let catalogStub: CatalogCall; + let jobnikStub: JobnikClient; let providerStub: Provider; beforeEach(async function () { catalogStub = { findRecords: vi.fn().mockResolvedValue([]) } as unknown as CatalogCall; + jobnikStub = { hasInFlightIngestionJob: vi.fn().mockResolvedValue(false) } as unknown as JobnikClient; providerStub = { fileExists: vi.fn().mockResolvedValue(true) }; - validator = new ValidationManager(await jsLogger({ enabled: false }), lookupStub, catalogStub, providerStub); + validator = new ValidationManager(await jsLogger({ enabled: false }), lookupStub, catalogStub, jobnikStub, providerStub); }); it('should pass a fully valid ingestion payload', async function () { @@ -100,6 +104,21 @@ describe('ValidationManager', function () { await expect(validator.validateIngestion(ingest(buildValidMetadata()))).rejects.toThrow(AppError); }); + it('should throw 409 when the product name has an in-flight ingestion job', async function () { + (jobnikStub.hasInFlightIngestionJob as ReturnType).mockResolvedValueOnce(true); + + let thrown: unknown; + try { + await validator.validateIngestion(ingest(buildValidMetadata())); + } catch (err) { + thrown = err; + } + + expect(thrown).toBeInstanceOf(AppError); + expect((thrown as AppError).status).toBe(StatusCodes.CONFLICT); + expect((thrown as AppError).message).toBe(ERROR_PRODUCT_NAME_IN_FLIGHT); + }); + it('should throw 400 when the model files do not exist', async function () { (providerStub.fileExists as ReturnType).mockResolvedValueOnce(false);