Compare commits

...
Author SHA1 Message Date
marcos-silva-rodrigues 901518d26c FIX: remove full table 2026-04-09 17:38:28 -03:00
marcos-silva-rodrigues 4ab8147149 FEAT: remove endpoints 2026-04-09 17:23:58 -03:00
marcos-silva-rodrigues 0086528d51 FEAT: update connectors 2026-04-07 18:05:30 -03:00
vinicius gadea bf12598915 Merge pull request #476 from dadosfera/feat/beta-delete-pipeline-job
Feat/beta delete pipeline job
2026-04-02 17:35:06 -03:00
viniciusgadea efb0b58648 FEAT: update deleteTable endpoint to use pipelineId in path and refactor parameters 2026-04-02 17:28:42 -03:00
viniciusgadea 8a9b58b612 FIX: rename method normalizeJobId to normalizePipelineId for clarity 2026-04-02 16:40:35 -03:00
viniciusgadea 837a9d7265 FEAT: update job deletion logic to mark tables as deleted and adjust API endpoints accordingly 2026-04-02 15:40:07 -03:00
viniciusgadea d0122a9c20 FEAT: remove unused status field from InputDocument and updateInputTable method 2026-04-01 14:43:51 -03:00
viniciusgadea 968f75b688 FEAT: update job deletion endpoint to include inputId in path and implement rollback for table deletion 2026-04-01 14:26:55 -03:00
viniciusgadea 877cb9d281 FEAT: refine PipelineTablesConfig type definition and update parsing logic in PipelinesController 2026-03-31 18:38:25 -03:00
viniciusgadea 59efdb6272 FEAT: update job deletion and mark associated table as deleted in DynamoDB. Update protospack-v2 2026-03-31 16:36:59 -03:00
Marcos Rodrigues Silva 7b1049224c Merge pull request #477 from dadosfera/hotfix/axios-vulnerability
FIX: preventing the axios vulnerability
2026-03-31 15:01:03 -03:00
marcos-silva-rodrigues 0369f10b5c FIX: preventing the axios vulnerability 2026-03-31 14:52:52 -03:00
viniciusgadea ecca9106f0 FEAT: remove PipelinesV2Module from PlatformApiModule imports 2026-03-31 11:17:32 -03:00
viniciusgadea efc1f49d92 FEAT: update delete job endpoint summary and remove unused InputsModule from platform-api module 2026-03-31 10:55:31 -03:00
viniciusgadea 2d46ac3213 FEAT: add delete job endpoint and update related services; update protospack-v2 version to 3.40.0-beta.3 2026-03-31 10:54:50 -03:00
Marcos Rodrigues Silva 271174176b Merge pull request #473 from dadosfera/feat/qualify
FEAT: update input
2026-03-30 12:06:08 -03:00
marcos-silva-rodrigues dcf7aed51c MERGE: resolve package conflicts 2026-03-30 12:05:36 -03:00
marcos-silva-rodrigues 7fbce8b5be FEAT: update input 2026-03-30 12:02:27 -03:00
vinicius gadea 049eca9700 Merge pull request #472 from dadosfera/feat/cancel-pipeline-run-beta
FEAT: update protospack-v2 to version 3.40.0-beta.1 and adjust relate…
2026-03-26 14:54:32 -03:00
viniciusgadea 3bcd78bb32 FEAT: update protospack-v2 to version 3.40.0-beta.1 and adjust related scripts 2026-03-26 14:53:33 -03:00
vinicius gadea f603b64679 Merge pull request #471 from dadosfera/feat/cancel-pipeline-run-beta
Feat/cancel pipeline run beta
2026-03-26 11:29:09 -03:00
viniciusgadea b640465624 Merge remote-tracking branch 'origin/beta' into feat/cancel-pipeline-run-beta 2026-03-26 11:19:30 -03:00
viniciusgadea 98282471a9 FIX: remove unused data asset methods from ElasticsearchService 2026-03-26 11:08:40 -03:00
viniciusgadea 8983430889 REF: remove last_run_status tracking and related Elasticsearch update logic from pipeline cancellation 2026-03-26 10:32:09 -03:00
viniciusgadea d3aeca12cd FIX: remove last_run_canceled_at from PipelineDocument and update last_run_status handling 2026-03-26 10:32:09 -03:00
viniciusgadea ce619942a3 FEAT: add endpoint to cancel running pipeline runs and update Elasticsearch status 2026-03-26 10:32:09 -03:00
Marcos Rodrigues Silva 71b3f278a5 Merge pull request #469 from dadosfera/feat/override-menu2
Feat/override menu2
2026-03-25 18:02:14 -03:00
viniciusgadea c92542ed90 FEAT: update protospack-v2 version to 3.39.0 in package.json and package-lock.json 2026-03-25 15:53:02 -03:00
viniciusgadeaandClaude Sonnet 4.6 69a9d78642 FIX: add multer@2.0.2 and update package-lock.json — missing transitive dep of protospack-v2@3.38.0-beta.30
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-03-25 15:41:02 -03:00
viniciusgadea bf4f3cfd8c FEAT: update CustomerSidebarSection and related DTOs to use object type for title; update protospack-v2 version; update from id to customerId 2026-03-25 15:40:57 -03:00
viniciusgadeaandClaude Sonnet 4.6 01afa69fcb FIX: remove duplicate identifier_columns declaration in TableColumns DTO
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-03-25 15:37:15 -03:00
viniciusgadea 80f3913fd2 FEAT: update CustomerSidebarSection to support both menu and link items 2026-03-25 15:24:59 -03:00
viniciusgadea 10293ad6a9 FEAT: update customer links handling and DTOs for improved structure sidebar. update protospack-v2 2026-03-25 15:24:59 -03:00
Marcos Rodrigues Silva e4d0c9c3e6 Merge pull request #465 from dadosfera/feat/qualify
FEAT: update nimbus after input
2026-03-24 14:10:42 -03:00
marcos-silva-rodrigues 21c57e5620 FEAT: update nimbus after input 2026-03-24 12:28:20 -03:00
marcos-silva-rodrigues 959210e354 FIX: npm ci 2026-03-19 14:56:58 -03:00
15 changed files with 559 additions and 338 deletions
+2 -2
View File
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
# run aws cli without mounting secret, because CI already has AWS credentials
FROM build_base AS ci_image
RUN aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
RUN npm ci
RUN npm ci --ignore-scripts
COPY . .
@@ -37,7 +37,7 @@ FROM build_base AS dev
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
# flag --build-from-source is required to force-build sqlite3
RUN npm ci
RUN npm ci --ignore-scripts
COPY . .
ENTRYPOINT npm run start:dev
+1 -1
View File
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
FROM build_base AS build
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
RUN npm ci
RUN npm ci --ignore-scripts
COPY . .
RUN npm run build
+83 -16
View File
@@ -3001,14 +3001,7 @@
"parameters": [],
"responses": {
"200": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "object"
}
}
}
"description": ""
}
},
"tags": [
@@ -3583,14 +3576,7 @@
],
"responses": {
"200": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "object"
}
}
}
"description": ""
}
},
"tags": [
@@ -4533,6 +4519,50 @@
]
}
},
"/platform/pipeline/{pipelineId}/pipeline_run/{runId}/cancel": {
"post": {
"operationId": "PlatformApiController_cancelPipelineRun",
"summary": "Cancel a running pipeline run",
"parameters": [
{
"name": "pipelineId",
"required": true,
"in": "path",
"schema": {
"type": "string"
}
},
{
"name": "runId",
"required": true,
"in": "path",
"schema": {
"type": "string"
}
}
],
"responses": {
"201": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "object"
}
}
}
}
},
"tags": [
"Platform API"
],
"security": [
{
"access-token": []
}
]
}
},
"/platform/jobs/{jobId}/input": {
"put": {
"operationId": "PlatformApiController_updateJobInput",
@@ -4675,6 +4705,43 @@
]
}
},
"/platform/pipelines/{pipelineId}/inputs/{inputId}": {
"delete": {
"operationId": "PlatformApiController_deleteTable",
"summary": "Mark a table as deleted and delete its associated job via platform-api",
"parameters": [
{
"name": "pipelineId",
"required": true,
"in": "path",
"schema": {
"type": "string"
}
},
{
"name": "inputId",
"required": true,
"in": "path",
"schema": {
"type": "string"
}
}
],
"responses": {
"200": {
"description": ""
}
},
"tags": [
"Platform API"
],
"security": [
{
"access-token": []
}
]
}
},
"/platform/jobs/jdbc/{jobId}": {
"get": {
"operationId": "PlatformApiController_getJdbcJob",
+114 -8
View File
@@ -17,7 +17,7 @@
"@aws-sdk/signature-v4": "^3.370.0",
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
"@dadosfera/protospack": "2.5.3",
"@dadosfera/protospack-v2": "^3.38.0-beta.30",
"@dadosfera/protospack-v2": "^3.40.0-beta.5",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
@@ -31,7 +31,7 @@
"@nestjs/schematics": "^9.2.0",
"@nestjs/swagger": "^6.3.0",
"@nestjs/testing": "^9.4.3",
"axios": "^0.30.2",
"axios": "0.30.3",
"cache-manager": "^5.1.4",
"cache-manager-ioredis-yet": "^1.1.0",
"class-transformer": "^0.5.1",
@@ -47,6 +47,7 @@
"jwk-to-pem": "^2.0.5",
"mixpanel": "^0.17.0",
"ms": "^3.0.0-canary.1",
"multer": "^2.0.2",
"openid-client": "^5.7.1",
"passport": "^0.6.0",
"passport-facebook": "^3.0.0",
@@ -1744,9 +1745,9 @@
}
},
"node_modules/@dadosfera/protospack-v2": {
"version": "3.38.0-beta.30",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.38.0-beta.30.tgz",
"integrity": "sha512-EePtEV4Bjr47BCtQgeZP8WzappDv1mRWMro15ms+Vdev/FIVBydMew/BtK3spzI7KHWJniB0Ctsvpcdj2GiTDA==",
"version": "3.40.0-beta.5",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.5.tgz",
"integrity": "sha512-iocKv/XXp2jKAasO5ONgm31cKfLgNsU4pEKZMr6YnR7nQaH11WcW7rnuagNxWoik++wLUqbYyf0bZWRDzMlCPA==",
"dependencies": {
"@grpc/grpc-js": "^1.9.3",
"rxjs": "^7.5.5"
@@ -2927,6 +2928,20 @@
"node": ">= 0.6"
}
},
"node_modules/@nestjs/platform-express/node_modules/concat-stream": {
"version": "1.6.2",
"resolved": "https://registry.npmjs.org/concat-stream/-/concat-stream-1.6.2.tgz",
"integrity": "sha512-27HBghJxjiZtIk3Ycvn/4kbJk/1uZuJFfuPEns6LaEvpvG1f0hTea8lilrouyo9mVc2GWdcEZ8OLoGmSADlrCw==",
"engines": [
"node >= 0.8"
],
"dependencies": {
"buffer-from": "^1.0.0",
"inherits": "^2.0.3",
"readable-stream": "^2.2.2",
"typedarray": "^0.0.6"
}
},
"node_modules/@nestjs/platform-express/node_modules/content-disposition": {
"version": "0.5.4",
"resolved": "https://registry.npmjs.org/content-disposition/-/content-disposition-0.5.4.tgz",
@@ -3096,6 +3111,24 @@
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
"license": "MIT"
},
"node_modules/@nestjs/platform-express/node_modules/multer": {
"version": "1.4.4-lts.1",
"resolved": "https://registry.npmjs.org/multer/-/multer-1.4.4-lts.1.tgz",
"integrity": "sha512-WeSGziVj6+Z2/MwQo3GvqzgR+9Uc+qt8SwHKh3gvNPiISKfsMfG4SvCOFYlxxgkXt7yIV2i1yczehm0EOKIxIg==",
"deprecated": "Multer 1.x is impacted by a number of vulnerabilities, which have been patched in 2.x. You should upgrade to the latest 2.x version.",
"dependencies": {
"append-field": "^1.0.0",
"busboy": "^1.0.0",
"concat-stream": "^1.5.2",
"mkdirp": "^0.5.4",
"object-assign": "^4.1.1",
"type-is": "^1.6.4",
"xtend": "^4.0.0"
},
"engines": {
"node": ">= 6.0.0"
}
},
"node_modules/@nestjs/platform-express/node_modules/negotiator": {
"version": "0.6.3",
"resolved": "https://registry.npmjs.org/negotiator/-/negotiator-0.6.3.tgz",
@@ -3120,6 +3153,25 @@
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/@nestjs/platform-express/node_modules/readable-stream": {
"version": "2.3.8",
"resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-2.3.8.tgz",
"integrity": "sha512-8p0AUk4XODgIewSi0l8Epjs+EVnWiK7NoDIEGU0HhE7+ZyY8D1IMY7odu5lRrFXGg71L15KG8QrPmum45RTtdA==",
"dependencies": {
"core-util-is": "~1.0.0",
"inherits": "~2.0.3",
"isarray": "~1.0.0",
"process-nextick-args": "~2.0.0",
"safe-buffer": "~5.1.1",
"string_decoder": "~1.1.1",
"util-deprecate": "~1.0.1"
}
},
"node_modules/@nestjs/platform-express/node_modules/readable-stream/node_modules/safe-buffer": {
"version": "5.1.2",
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.1.2.tgz",
"integrity": "sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g=="
},
"node_modules/@nestjs/platform-express/node_modules/safe-buffer": {
"version": "5.2.1",
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.2.1.tgz",
@@ -3194,6 +3246,19 @@
"node": ">= 0.8"
}
},
"node_modules/@nestjs/platform-express/node_modules/string_decoder": {
"version": "1.1.1",
"resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-1.1.1.tgz",
"integrity": "sha512-n/ShnvDi6FHbbVfviro+WojiFzv+s8MPMHBczVePfUpDJLwoLT0ht1l4YwBCbi8pJAveEEdnkHyPyTP/mzRfwg==",
"dependencies": {
"safe-buffer": "~5.1.0"
}
},
"node_modules/@nestjs/platform-express/node_modules/string_decoder/node_modules/safe-buffer": {
"version": "5.1.2",
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.1.2.tgz",
"integrity": "sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g=="
},
"node_modules/@nestjs/platform-express/node_modules/tslib": {
"version": "2.5.3",
"resolved": "https://registry.npmjs.org/tslib/-/tslib-2.5.3.tgz",
@@ -5533,9 +5598,9 @@
}
},
"node_modules/axios": {
"version": "0.30.2",
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.2.tgz",
"integrity": "sha512-0pE4RQ4UQi1jKY6p7u6i1Tkzqmu+d+/tHS7Q7rKunWLB9WyilBTpHHpXzPNMDj5hTbK0B0PTLSz07yqMBiF6xg==",
"version": "0.30.3",
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.3.tgz",
"integrity": "sha512-5/tmEb6TmE/ax3mdXBc/Mi6YdPGxQsv+0p5YlciXWt3PHIn0VamqCXhRMtScnwY3lbgSXLneOuXAKUhgmSRpwg==",
"license": "MIT",
"dependencies": {
"follow-redirects": "^1.15.4",
@@ -6501,6 +6566,20 @@
"integrity": "sha512-/Srv4dswyQNBfohGpz9o6Yb3Gz3SrUDqBH5rTuhGR7ahtlbYKnVxw2bCFMRljaA7EXHaXZ8wsHdodFvbkhKmqg==",
"license": "MIT"
},
"node_modules/concat-stream": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/concat-stream/-/concat-stream-2.0.0.tgz",
"integrity": "sha512-MWufYdFw53ccGjCA+Ol7XJYpAlW6/prSMzuPOTRnJGcGzuhLn4Scrz7qf6o8bROZ514ltazcIFJZevcfbo0x7A==",
"engines": [
"node >= 6.0"
],
"dependencies": {
"buffer-from": "^1.0.0",
"inherits": "^2.0.3",
"readable-stream": "^3.0.2",
"typedarray": "^0.0.6"
}
},
"node_modules/consola": {
"version": "2.15.3",
"resolved": "https://registry.npmjs.org/consola/-/consola-2.15.3.tgz",
@@ -9112,6 +9191,11 @@
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/isarray": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/isarray/-/isarray-1.0.0.tgz",
"integrity": "sha512-VLghIWNM6ELQzo7zwmcg0NmTVyWKYjvIeM83yjp0wRDTmUnrM678fQbcKBo6n2CJEF0szoG//ytg+TKla89ALQ=="
},
"node_modules/isexe": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/isexe/-/isexe-2.0.0.tgz",
@@ -10705,6 +10789,23 @@
"node": ">=18"
}
},
"node_modules/multer": {
"version": "2.0.2",
"resolved": "https://registry.npmjs.org/multer/-/multer-2.0.2.tgz",
"integrity": "sha512-u7f2xaZ/UG8oLXHvtF/oWTRvT44p9ecwBBqTwgJVq0+4BW1g8OW01TyMEGWBHbyMOYVHXslaut7qEQ1meATXgw==",
"dependencies": {
"append-field": "^1.0.0",
"busboy": "^1.6.0",
"concat-stream": "^2.0.0",
"mkdirp": "^0.5.6",
"object-assign": "^4.1.1",
"type-is": "^1.6.18",
"xtend": "^4.0.2"
},
"engines": {
"node": ">= 10.16.0"
}
},
"node_modules/mute-stream": {
"version": "0.0.8",
"resolved": "https://registry.npmjs.org/mute-stream/-/mute-stream-0.0.8.tgz",
@@ -11694,6 +11795,11 @@
"url": "https://github.com/chalk/ansi-styles?sponsor=1"
}
},
"node_modules/process-nextick-args": {
"version": "2.0.1",
"resolved": "https://registry.npmjs.org/process-nextick-args/-/process-nextick-args-2.0.1.tgz",
"integrity": "sha512-3ouUOpQhtgrbOa17J7+uxOTpITYWaGP7/AhoR3+A+/1e9skrzelGi/dXzEYyvbxubEF6Wn2ypscTKiKJFFn1ag=="
},
"node_modules/process-warning": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/process-warning/-/process-warning-1.0.0.tgz",
+8 -4
View File
@@ -10,7 +10,7 @@
},
"scripts": {
"co:login": "aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1",
"proto-update": "npm i @dadosfera/protospack-v2@latest --save-exact",
"proto-update": "npm i @dadosfera/protospack-v2@v3.40.0-beta.1 --save-exact",
"prebuild": "rimraf dist",
"build": "nest build",
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
@@ -35,7 +35,7 @@
"@aws-sdk/signature-v4": "^3.370.0",
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
"@dadosfera/protospack": "2.5.3",
"@dadosfera/protospack-v2": "^3.38.0-beta.30",
"@dadosfera/protospack-v2": "^3.40.0-beta.5",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
@@ -49,7 +49,7 @@
"@nestjs/schematics": "^9.2.0",
"@nestjs/swagger": "^6.3.0",
"@nestjs/testing": "^9.4.3",
"axios": "^0.30.2",
"axios": "0.30.3",
"cache-manager": "^5.1.4",
"cache-manager-ioredis-yet": "^1.1.0",
"class-transformer": "^0.5.1",
@@ -65,6 +65,7 @@
"jwk-to-pem": "^2.0.5",
"mixpanel": "^0.17.0",
"ms": "^3.0.0-canary.1",
"multer": "^2.0.2",
"openid-client": "^5.7.1",
"passport": "^0.6.0",
"passport-facebook": "^3.0.0",
@@ -80,7 +81,7 @@
"swagger-ui-express": "^4.6.3"
},
"overrides": {
"multer": "2.0.2",
"axios": "0.30.3",
"form-data": "^4.0.4",
"body-parser": "^1.20.3",
"cross-spawn": "^7.0.5",
@@ -117,5 +118,8 @@
"ts-node": "^10.9.1",
"tsconfig-paths": "^3.14.2",
"typescript": "^4.9.5"
},
"resolutions": {
"axios": "0.30.3"
}
}
+40 -7
View File
@@ -17,6 +17,8 @@ import {
InputCreateGenericRequest,
InputCreateS3Request,
InputNewCreateRequest,
InputUpdateResponse,
RollbackInputRequest,
TestConnectionRequest,
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
@@ -71,8 +73,8 @@ export class InputsService {
objectCamelToSnake(createInputResponse);
return createInputResponse;
},
update: async (updateInputDTO: UpdateInputRequest) => {
this.logger.info('InputClientService - Update');
update: async (updateInputDTO: UpdateInputRequest): Promise<InputUpdateResponse> => {
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
const updateInputResponse = await lastValueFrom(
this.inputWriteService.InputUpdate(updateInputDTO),
);
@@ -207,21 +209,44 @@ export class InputsService {
async update(id: string, data, info: Info) {
// this.validateCron({ ...data, info });
try {
const updateInputResponse: any = await this.OLD_inputClient.update({
const {
tablesUpdate,
dataAssetUpdate,
input
} = await this.OLD_inputClient.update({
id,
info,
...data,
info,
});
updateInputResponse.input = this.adjustInputPayload(
updateInputResponse?.input,
const updateInputResponse = this.adjustInputPayload(
input,
);
return updateInputResponse;
return {
input: updateInputResponse,
tablesUpdate,
dataAssetUpdate
};
} catch (err) {
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
}
}
async rollbackUpdate(
data: RollbackInputRequest
) {
this.logger.info('PipelinesClientService - rollbackUpdate');
this.logger.info('Rolling back input update with data: ' + JSON.stringify(data));
const updatePipelineResponse = await lastValueFrom(
this.inputWriteService.RollbackInputUpdate(
data
),
);
this.logger.info('Done');
return updatePipelineResponse;
}
async remove(idRequest: IIdRequest) {
return lastValueFrom(this.inputWriteService.InputRemove(idRequest));
}
@@ -263,4 +288,12 @@ export class InputsService {
};
return formatedPayload;
}
async markTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
return lastValueFrom(this.inputWriteService.MarkTableDeleted(data));
}
async unmarkTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
return lastValueFrom((this.inputWriteService as any).UnmarkTableDeleted(data));
}
}
+26 -24
View File
@@ -53,6 +53,9 @@ import { TableColumns } from '../inputs/dtos/input.model';
import { UpdateInputRequest } from '../inputs/dtos/old_interfaces';
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
type PipelineTable = { name: string; job_id?: string; is_deleted?: boolean; [key: string]: any };
type PipelineTablesConfig = { input_id?: string; tables: PipelineTable[] };
@ApiTags('PipelinesV2')
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
@UseFilters(new GrpcToHttpExceptionFilter())
@@ -62,6 +65,7 @@ export class PipelinesController {
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
private pipelinesClientService: PipelinesService,
private oldPipelinesService: OldPipelineService,
) {
@@ -222,30 +226,27 @@ export class PipelinesController {
language,
});
const result = await this.pipelinesClientService
.findOne({ id }, metadata)
.then((res) => {
//{pipeline:{tables: {tables: [], input_id: ''}}}
let tables = JSON.parse(res.pipeline.config.tables);
const input_id = tables?.input_id;
if (tables?.tables) tables = tables.tables;
Object.assign(res.pipeline, {
transformations: res.pipeline.transformations
? JSON.parse(res.pipeline.transformations)
: [],
config: {
cron: res.pipeline.config.cron,
tables,
input_id
},
properties: res.pipeline.properties
? JSON.parse(res.pipeline.properties)
: {},
});
return res;
});
const pipelineRes = await this.pipelinesClientService.findOne({ id }, metadata);
return result;
const parsed: PipelineTablesConfig = JSON.parse(pipelineRes.pipeline.config.tables);
const input_id = parsed.input_id;
const tables: PipelineTable[] = parsed.tables ?? [];
Object.assign(pipelineRes.pipeline, {
transformations: pipelineRes.pipeline.transformations
? JSON.parse(pipelineRes.pipeline.transformations)
: [],
config: {
cron: pipelineRes.pipeline.config.cron,
tables,
input_id,
},
properties: pipelineRes.pipeline.properties
? JSON.parse(pipelineRes.pipeline.properties)
: {},
});
return pipelineRes;
}
@Patch('/:id')
@@ -294,13 +295,14 @@ export class PipelinesController {
) {
this.logger.info('PipelinesController - update', { user });
const { customer_id, customer_name, user_id, username } = user;
const info: Info = {
user_id: user.user_id,
customer: user.customer_name,
customer_id: user.customer_id,
pipeline_id: pipelineId
};
const metadata = PackTheMetadata({
customer_id,
customer_name,
+5 -2
View File
@@ -12,6 +12,8 @@ import { ConnectorModule } from '../connector/connector.module';
import { InputsModule } from '../inputs/inputs.module';
import { TransformationsModule } from '../transformations/transformations.module';
import { PlatformApiModule } from '../platform-api/platform-api.module';
import { NimbusServicesModule } from 'src/services/nimbus/nimbus.module';
import { NimbusService } from 'src/services/nimbus/nimbus.service';
const client = new PipelinesClientConfiguration();
@@ -22,10 +24,11 @@ const client = new PipelinesClientConfiguration();
ConnectorModule,
InputsModule,
TransformationsModule,
PlatformApiModule
PlatformApiModule,
NimbusServicesModule
],
controllers: [PipelinesController],
providers: [PipelinesService, DadosferaLogger],
providers: [PipelinesService, DadosferaLogger, NimbusService],
exports: [PipelinesService],
})
export class PipelinesV2Module {}
+145 -12
View File
@@ -29,6 +29,11 @@ import ErrorCodes from 'src/utils/errorCodes';
import ErrorBuilder from 'src/utils/ErrorBuilder';
import { PlatformApiService } from '../platform-api/platform-api.service';
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
import { TableUpdate } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
import { AxiosError } from 'axios';
import { NimbusService } from 'src/services/nimbus/nimbus.service';
type RollbackPromise = () => Promise<any>;
export class PipelinesService implements OnModuleInit {
logger: DadosferaLogger;
@@ -42,7 +47,8 @@ export class PipelinesService implements OnModuleInit {
private readonly connectorService: ConnectorClientService,
private readonly inputsService: InputsService,
private readonly transformationsService: TransformationsService,
private readonly platformAPI: PlatformApiService
private readonly platformAPI: PlatformApiService,
private readonly nimbusService: NimbusService
) {
this.logger = dadosferaLogger.logger;
}
@@ -348,20 +354,150 @@ export class PipelinesService implements OnModuleInit {
async updatePipelineInput(pipelineId: string, inputId: string, updateInputDTO: UpdatePlatformInputRequest, info: Info, user: RequestUser, metadata: Metadata) {
this.logger.info('InputClientService - Update');
this.logger.info('Update Dynamo Reference');
const {
input: oldInput
} = await this.inputsService.findOne({
id: inputId,
info: info
});
this.logger.info('Update Dynamo Reference :' + JSON.stringify(oldInput));
const pipelineIdFormat = pipelineId.split('-').join('_');
const rollback: RollbackPromise[] = [];
const updateInputResponse = await this.inputsService.update(
inputId,
updateInputDTO,
info
)
);
this.logger.info('Dynamo Response', updateInputResponse);
const inputRollback = () => {
this.logger.info("exec rollback to input: " + JSON.stringify(oldInput));
return this.inputsService.rollbackUpdate(
{
id: inputId,
dataAssetUpdate: updateInputResponse.dataAssetUpdate,
tables: oldInput.tables,
info
}
) as Promise<any>;
}
rollback.push(inputRollback);
this.logger.info("Input Update Response: " + JSON.stringify(updateInputResponse))
const nimbusUpdates = updateInputResponse?.tablesUpdate || [];
nimbusUpdates.forEach(update => {
const nimbusRollback = () => {
return this.nimbusService.renameTable(
info.customer,
update.database,
{
table_name: update.table_name,
table_schema: update.table_schema
},
{
table_name: update.old_table_name,
table_schema: update.old_table_schema
}
);
}
rollback.push(nimbusRollback);
});
try {
await this.updateNimbus(info.customer, nimbusUpdates);
} catch (error) {
this.logger.error(error);
if (error instanceof AxiosError) {
this.logger.error(JSON.stringify(error.response.data));
}
await this.executeRenameRollback(rollback);
throw new Error("Error Nimbus updating tables");
}
try {
await this.updatePlatformJobs(
pipelineIdFormat,
updateInputResponse.input.type,
updateInputDTO,
user
);
} catch (error) {
this.logger.error(error);
await this.executeRenameRollback(rollback)
throw new Error("Error Platform API updating jobs");
}
return updateInputResponse;
}
private async executeRenameRollback(request: RollbackPromise[]) {
this.logger.info('rollback steps: ' + request.length)
const result = await Promise.allSettled(request.map(func => func()));
result.forEach(promise => {
this.logger.info("Promise finish with status: " + promise.status)
if (promise.status === "rejected") {
this.logger.error("reject with: " + JSON.stringify(promise.reason || {}))
}
if (promise.status === "fulfilled") {
this.logger.info("success with: " + JSON.stringify(promise.value || {}))
}
});
}
private async updateNimbus(customer: string, changes: TableUpdate[]) {
// throw new Error("teste error nimbus");
this.logger.info('Nimbus Changes: ' + JSON.stringify(changes));
if(!changes || changes.length === 0) return;
const requests = changes.map(change => {
return this.nimbusService.renameTable(customer, change.database, {
table_name: change.old_table_name,
table_schema: change.old_table_schema
}, {
table_name: change.table_name,
table_schema: change.table_schema
});
})
const values = await Promise.allSettled(requests);
const success = values.map(request => request.status === "fulfilled")
this.logger.info("Updates with succes: " + success.length);
values.forEach(promise => {
this.logger.info("Promise finish with status: " + promise.status)
if (promise.status === "rejected") {
this.logger.error("reject with: " + JSON.stringify(promise.reason || {}));
throw new Error(promise.reason );
}
if (promise.status === "fulfilled") {
this.logger.info("success with: " + JSON.stringify(promise.value || {}));
}
});
}
async updatePlatformJobs(pipelineId: string, pipelineType: string, updateInputDTO: UpdatePlatformInputRequest, user: RequestUser) {
const jobsUpdated = [];
for (const [index, table] of updateInputDTO.tables.entries()) {
const jobUpdate = {
job_id: `${pipelineIdFormat}_${index}`,
job_id: `${pipelineId}_${index}`,
}
if (table.type !== "incremental_with_qualify") {
delete table.destinations?.qualify;
}
if (table.memory) {
@@ -370,7 +506,7 @@ export class PipelinesService implements OnModuleInit {
}
}
this.logger.info('Updating input reference for table', table.name);
this.logger.info('Updating input reference for table: ' + table.name);
let hasUpdateSyncMode = false;
const jobSyncMode = {}
@@ -393,8 +529,8 @@ export class PipelinesService implements OnModuleInit {
if (table.type) {
hasUpdateSyncMode = true;
jobSyncMode['target_load_type'] = table.type;
}
if(hasUpdateSyncMode) {
@@ -431,17 +567,14 @@ export class PipelinesService implements OnModuleInit {
}));
const response = await this.platformAPI.proxy(
'POST',
`/pipelines/${pipelineIdFormat}/batch-update`,
'PUT',
`/pipeline/${pipelineId}/jobs`,
user,
{
job_updates: jobsUpdated
}
)
this.logger.info('Platform api response: ' + JSON.stringify(response));
return updateInputResponse;
}
}
@@ -11,6 +11,7 @@ import {
Inject,
BadRequestException,
HttpException,
NotFoundException,
} from '@nestjs/common';
import { ApiTags, ApiOperation } from '@nestjs/swagger';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
@@ -29,6 +30,7 @@ import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
import { CatalogService } from '../catalog/catalog.service';
import { PackTheMetadata } from '../../utils/PackTheMetadata';
import { ValidationTableDTO } from './platform-api.dto';
import { InputsService } from '../inputs/inputs.service';
type ValidateTablesDTO = {
@@ -54,6 +56,7 @@ export class PlatformApiController {
private readonly dynamoDBService: DynamoDBService,
private readonly customersService: CustomersService,
private readonly catalogService: CatalogService,
private readonly inputsService: InputsService,
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
@@ -845,6 +848,24 @@ export class PlatformApiController {
);
}
@Post('pipeline/:pipelineId/pipeline_run/:runId/cancel')
@ApiOperation({ summary: 'Cancel a running pipeline run' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async cancelPipelineRun(
@Param('pipelineId') pipelineId: string,
@Param('runId') runId: string,
@User() user: RequestUser,
) {
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
const normalizedRunId = this.normalizePipelineId(runId);
return this.platformApiService.proxy(
'POST',
`/pipeline/${normalizedPipelineId}/pipeline_run/${normalizedRunId}/cancel`,
user,
);
}
// ==================== JOBS - COLUMN EDITING ROUTES ====================
@Put('jobs/:jobId/input')
@@ -944,226 +965,52 @@ export class PlatformApiController {
);
}
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
@Get('jobs/jdbc/:jobId')
@ApiOperation({ summary: 'Get JDBC job details' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getJdbcJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
// Normalize job ID for Platform API (replace - with _)
const normalizedJobId = this.normalizeJobId(jobId);
return this.platformApiService.proxy('GET', `/jobs/jdbc/${normalizedJobId}`, user);
}
@Post('jobs/jdbc/:jobId/sync-mode')
@ApiOperation({ summary: 'Update JDBC job sync mode' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updateJdbcSyncMode(
@Param('jobId') jobId: string,
@Body() body: any,
@Delete('pipelines/:pipelineId/inputs/:inputId')
@ApiOperation({ summary: 'Mark a table as deleted and delete its associated job via platform-api' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
async deleteTable(
@Param('pipelineId') pipelineId: string,
@Param('inputId') inputId: string,
@Body() body: { table_name: string },
@User() user: RequestUser,
) {
// Normalize job ID for Platform API (replace - with _)
const normalizedJobId = this.normalizeJobId(jobId);
const tableName = body.table_name;
const info = {
customer_id: user.customer_id,
customer: user.customer_name,
user_id: user.user_id,
};
const result = await this.platformApiService.proxy(
'POST',
`/jobs/jdbc/${normalizedJobId}/sync-mode`,
user,
body,
);
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
await this.syncJdbcSyncModeToDynamoDB(jobId, body, user);
return result;
}
@Post('jobs/:jobId/rename-tables')
@ApiOperation({ summary: 'Rename job output tables and sync to catalog' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async renameJobTables(
@Param('jobId') jobId: string,
@Body() body: RenameTablesBody,
@User() user: RequestUser,
) {
const normalizedJobId = this.normalizeJobId(jobId);
const currentJob = await this.getJobByAnyConnectorType(normalizedJobId, user);
const result = await this.platformApiService.proxy(
'POST',
`/jobs/${normalizedJobId}/rename-tables`,
user,
body,
);
this.logger.info('deleteTable: marking table as deleted', { inputId, tableName });
const updatedInput: any = await this.inputsService.markTableDeleted({ input_id: inputId, table_name: tableName, info });
this.logger.info('deleteTable: table marked as deleted', { inputId, tableName });
try {
await this.syncTableRenameToCatalog(jobId, body, currentJob, user);
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
this.logger.info('deleteTable: fetching pipeline from platform-api', { pipelineId, normalizedPipelineId });
const platformPipeline = await this.platformApiService.proxy('GET', `/pipeline/${normalizedPipelineId}`, user);
this.logger.info('deleteTable: pipeline fetched', { jobCount: platformPipeline?.jobs?.length });
const job = platformPipeline?.jobs?.find((j: any) => j.input?.table_name === tableName);
if (!job) throw new NotFoundException(`Job for table '${tableName}' not found in pipeline`);
this.logger.info('deleteTable: deleting job from platform-api', { jobId: job.job_id });
await this.platformApiService.proxy('DELETE', `/jobs/${job.job_id}`, user);
this.logger.info('deleteTable: job deleted', { jobId: job.job_id });
return { name: tableName, is_deleted: updatedInput.is_deleted ?? true, deleted_at: updatedInput.deleted_at };
} catch (error) {
this.logger.error('Catalog sync failed, rolling back Snowflake rename', { jobId, error: error.message });
const reverseBody = this.buildSnowflakeRollbackBody(body, currentJob.output_config || {});
if (reverseBody) {
try {
await this.platformApiService.proxy('POST', `/jobs/${normalizedJobId}/rename-tables`, user, reverseBody);
this.logger.info('Snowflake rename rolled back', { jobId });
} catch (rollbackError) {
this.logger.error('Snowflake rollback failed', { jobId, error: rollbackError.message });
}
}
throw new HttpException('Table rename failed: catalog sync error, Snowflake reverted', 500);
}
return result;
}
private buildSnowflakeRollbackBody(
body: RenameTablesBody,
outputConfig: any,
): RenameTablesBody | null {
const reverse: RenameTablesBody = {};
if (body.raw) {
const nested = outputConfig.raw;
const oldTableName = nested?.table_name || outputConfig.table_name;
const oldTableSchema = nested?.table_schema || 'PUBLIC';
if (oldTableName) reverse.raw = { table_name: oldTableName, table_schema: oldTableSchema };
}
if (body.qualify) {
const nested = outputConfig.qualify;
if (nested?.table_name) reverse.qualify = { table_name: nested.table_name, table_schema: nested.table_schema || 'STAGED' };
}
return Object.keys(reverse).length > 0 ? reverse : null;
}
/**
* Sync table rename to Elasticsearch and Nimbus.
*
* For each target (raw, qualify):
* 1. Resolve old table name from output_config
* 2. Find the ES data asset by pipeline + table + schema
* 3. Update ES, Nimbus table-metadata, column-metadata, and data-preview
* 4. If any step fails, rollback all completed steps for that target
*/
private async syncTableRenameToCatalog(
jobId: string,
body: RenameTablesBody,
currentJob: any,
user: RequestUser,
): Promise<void> {
const pipelineId = this.extractPipelineIdFromJobId(jobId);
const outputConfig = currentJob.output_config || {};
const nimbusUrl = this.catalogService._getNimbusUrl({ info: { customer: user.customer_name } });
const databaseName = `DADOSFERA_PRD_${user.customer_name.toUpperCase()}`;
const targets = this.buildRenameTargets(body, outputConfig);
for (const { key, oldTableName, oldTableSchema, newValues } of targets) {
const rollbackSteps: Array<() => Promise<void>> = [];
this.logger.error('deleteTable: platform-api delete failed, attempting rollback', { tableName, error: error.message });
try {
const dataAsset = await this.elasticsearchService.findDataAssetByTable(
user.customer_name, oldTableName, oldTableSchema,
);
if (!dataAsset) {
this.logger.warn(`No data asset found for ${key}`, { jobId, pipelineId, oldTableName, oldTableSchema });
continue;
}
const { _es_id: esAssetId, nimbus_id: nimbusId } = dataAsset;
const oldValues = { table_name: oldTableName, table_schema: oldTableSchema };
// ES update
const esFields = { name: newValues.table_name, table_name: newValues.table_name, table_schema: newValues.table_schema, display_name: newValues.table_name };
await this.elasticsearchService.updateDataAsset(user.customer_name, esAssetId, esFields);
rollbackSteps.push(() => this.elasticsearchService.updateDataAsset(
user.customer_name, esAssetId,
{ name: oldTableName, table_name: oldTableName, table_schema: oldTableSchema, display_name: oldTableName },
));
const newTableNameUpper = newValues.table_name.toUpperCase();
const newTableSchemaUpper = newValues.table_schema.toUpperCase();
const oldTableNameUpper = oldTableName.toUpperCase();
const oldTableSchemaUpper = oldTableSchema.toUpperCase();
// Nimbus table-metadata
if (nimbusId) {
await this.catalogService.renameTableOnNimbus(nimbusUrl, nimbusId, { table_name: newTableNameUpper, table_schema: newTableSchemaUpper });
rollbackSteps.push(() => this.catalogService.renameTableOnNimbus(nimbusUrl, nimbusId, { table_name: oldTableNameUpper, table_schema: oldTableSchemaUpper }));
}
// Nimbus column-metadata
await this.catalogService.renameColumnMetadataOnNimbus(
nimbusUrl, databaseName, oldTableNameUpper, oldTableSchemaUpper, newTableNameUpper, newTableSchemaUpper,
);
rollbackSteps.push(() => this.catalogService.renameColumnMetadataOnNimbus(
nimbusUrl, databaseName, newTableNameUpper, newTableSchemaUpper, oldTableNameUpper, oldTableSchemaUpper,
));
// Nimbus data-preview
await this.catalogService.renameDataPreviewOnNimbus(
nimbusUrl, databaseName, oldTableNameUpper, oldTableSchemaUpper, newTableNameUpper, newTableSchemaUpper,
);
rollbackSteps.push(() => this.catalogService.renameDataPreviewOnNimbus(
nimbusUrl, databaseName, newTableNameUpper, newTableSchemaUpper, oldTableNameUpper, oldTableSchemaUpper,
));
this.logger.info(`Synced catalog rename for ${key}`, { jobId, oldTableName, newTableName: newValues.table_name });
} catch (error) {
this.logger.error(`Catalog sync failed for ${key}, rolling back catalog`, { jobId, error: error.message });
await this.executeRollback(rollbackSteps, key, jobId);
throw error;
await this.inputsService.unmarkTableDeleted({ input_id: inputId, table_name: tableName, info });
} catch (rollbackError) {
this.logger.error('deleteTable: rollback failed', { tableName, error: rollbackError.message });
}
throw error;
}
}
private buildRenameTargets(
body: RenameTablesBody,
outputConfig: any,
): Array<{ key: string; oldTableName: string; oldTableSchema: string; newValues: { table_name: string; table_schema: string } }> {
const DEFAULT_SCHEMAS = { raw: 'PUBLIC', qualify: 'STAGED' };
const targets: Array<{ key: string; oldTableName: string; oldTableSchema: string; newValues: { table_name: string; table_schema: string } }> = [];
for (const key of ['raw', 'qualify'] as const) {
if (!body[key]) continue;
const nested = outputConfig[key];
// qualify: only sync if output_config.qualify already exists
if (key === 'qualify' && !nested?.table_name) continue;
const oldTableName = nested?.table_name || outputConfig.table_name;
if (!oldTableName) continue;
targets.push({
key,
oldTableName,
oldTableSchema: nested?.table_schema || DEFAULT_SCHEMAS[key],
newValues: body[key],
});
}
return targets;
}
private async executeRollback(
steps: Array<() => Promise<void>>,
targetKey: string,
jobId: string,
): Promise<void> {
for (const rollback of steps.reverse()) {
try {
await rollback();
} catch (error) {
this.logger.error(`Rollback failed for ${targetKey}`, { jobId, error: error.message });
}
}
}
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
@Get('jobs/jdbc/configs/allowed_datatypes')
@ApiOperation({ summary: 'Get allowed datatypes for JDBC' })
@@ -1176,52 +1023,6 @@ export class PlatformApiController {
);
}
// ==================== JOBS - SINGER REPLICATION ROUTES ====================
@Get('jobs/singer/:jobId')
@ApiOperation({ summary: 'Get Singer job details' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getSingerJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
// Normalize job ID for Platform API (replace - with _)
const normalizedJobId = this.normalizeJobId(jobId);
return this.platformApiService.proxy('GET', `/jobs/singer/${normalizedJobId}`, user);
}
@Post('jobs/singer/:jobId/sync-mode')
@ApiOperation({ summary: 'Update Singer job sync mode' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updateSingerSyncMode(
@Param('jobId') jobId: string,
@Body() body: any,
@User() user: RequestUser,
) {
// Normalize job ID for Platform API (replace - with _)
const normalizedJobId = this.normalizeJobId(jobId);
const result = await this.platformApiService.proxy(
'POST',
`/jobs/singer/${normalizedJobId}/sync-mode`,
user,
body,
);
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
await this.syncSingerSyncModeToDynamoDB(jobId, body, user);
return result;
}
// ==================== JOBS - S3 ROUTES ====================
@Get('jobs/s3/:jobId')
@ApiOperation({ summary: 'Get S3 job details' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getS3Job(@Param('jobId') jobId: string, @User() user: RequestUser) {
// Normalize job ID for Platform API (replace - with _)
const normalizedJobId = this.normalizeJobId(jobId);
return this.platformApiService.proxy('GET', `/jobs/s3/${normalizedJobId}`, user);
}
// ==================== HEALTH ROUTE ====================
@Get('health')
@@ -8,9 +8,10 @@ import { ElasticsearchModule } from '../../services/elasticsearch';
import { DynamoDBModule } from '../../services/dynamodb';
import { CustomersModule } from '../customers/customers.module';
import { CatalogModule } from '../catalog/catalog.module';
import { InputsModule } from '../inputs/inputs.module';
@Module({
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule],
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule],
controllers: [PlatformApiController],
providers: [PlatformApiService, DadosferaLogger],
exports: [PlatformApiService],
@@ -89,12 +89,12 @@ export class PlatformApiService {
// Propagate non-2xx responses as HttpExceptions
if (response.status >= 400) {
this.logger.error('Platform API upstream error', {
this.logger.error('Platform API upstream error' + JSON.stringify({
status: response.status,
data: response.data,
path,
method: method.toUpperCase(),
});
}));
throw new HttpException(response.data, response.status);
}
@@ -107,6 +107,8 @@ export class PlatformApiService {
method: method.toUpperCase(),
});
this.logger.error(error)
if (error instanceof HttpException) {
throw error;
}
+8 -4
View File
@@ -125,9 +125,13 @@ export class DynamoDBService {
},
});
const { Item } = await this.documentClient.send(getCommand);
return Item as InputDocument | null;
try {
const { Item } = await this.documentClient.send(getCommand);
return Item as InputDocument | null;
} catch (error) {
this.logger.error('DynamoDB: findInput failed', { inputId, clientId, error: error.message });
throw error;
}
}
async deleteInput(clientId: string, inputId: string): Promise<void> {
@@ -203,7 +207,6 @@ export class DynamoDBService {
updatedTable.reference_column = changes.reference_column;
}
}
tables[tableIndex] = updatedTable;
// Save updated document
@@ -231,4 +234,5 @@ export class DynamoDBService {
throw error;
}
}
}
+9
View File
@@ -0,0 +1,9 @@
import { Module } from "@nestjs/common";
import { NimbusService } from "./nimbus.service";
import DadosferaLogger from "@dadosfera/dadosfera-logs";
@Module({
providers: [NimbusService, DadosferaLogger],
exports: [NimbusService],
})
export class NimbusServicesModule {}
+56
View File
@@ -0,0 +1,56 @@
import DadosferaLogger from "@dadosfera/dadosfera-logs";
import { Inject, Injectable } from "@nestjs/common";
import axios from "axios";
type TableUpdate = {
table_schema: string;
table_name: string;
}
@Injectable()
export class NimbusService {
private logger: DadosferaLogger;
constructor(
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
private buildUrl(customerName: string) {
if (process.env.ENV === 'prd') {
return `https://nimbus-${customerName}.dadosfera.ai`;
}
return `https://nimbus-${customerName}.${process.env.ENV.replace(
'local',
'stg',
)}.dadosfera.ai`;
}
async renameTable(customerName: string, database: string, old: TableUpdate, update: TableUpdate) {
const nimbusUrl = this.buildUrl(customerName);
const path = `/api/catalog/rename-tables/?database_name=${encodeURIComponent(database)}&table_name=${encodeURIComponent(old.table_name)}&table_schema=${encodeURIComponent(old.table_schema)}`;
try {
this.logger.info("Request for PATCH " + nimbusUrl + path);
this.logger.info("Payload: " + JSON.stringify(update));
const { data } = await axios.patch(nimbusUrl + path, {
table_name: update.table_name,
table_schema: update.table_schema
})
return data;
} catch (error) {
this.logger.error(error);
return {
message: error.message,
database,
old,
update
}
}
}
}