mirror of
https://github.com/dadosfera/maestro.git
synced 2026-10-01 15:29:08 +00:00
Compare commits
57
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
901518d26c | ||
|
|
4ab8147149 | ||
|
|
0086528d51 | ||
|
|
bf12598915 | ||
|
|
efb0b58648 | ||
|
|
8a9b58b612 | ||
|
|
837a9d7265 | ||
|
|
d0122a9c20 | ||
|
|
968f75b688 | ||
|
|
877cb9d281 | ||
|
|
59efdb6272 | ||
|
|
7b1049224c | ||
|
|
0369f10b5c | ||
|
|
ecca9106f0 | ||
|
|
efc1f49d92 | ||
|
|
2d46ac3213 | ||
|
|
271174176b | ||
|
|
dcf7aed51c | ||
|
|
7fbce8b5be | ||
|
|
049eca9700 | ||
|
|
3bcd78bb32 | ||
|
|
f603b64679 | ||
|
|
b640465624 | ||
|
|
98282471a9 | ||
|
|
8983430889 | ||
|
|
d3aeca12cd | ||
|
|
ce619942a3 | ||
|
|
71b3f278a5 | ||
|
|
c92542ed90 | ||
|
|
69a9d78642 | ||
|
|
bf4f3cfd8c | ||
|
|
01afa69fcb | ||
|
|
80f3913fd2 | ||
|
|
10293ad6a9 | ||
|
|
e4d0c9c3e6 | ||
|
|
21c57e5620 | ||
|
|
959210e354 | ||
|
|
f10953e949 | ||
|
|
13903b9bb9 | ||
|
|
2e3a13d421 | ||
|
|
fa17fc3001 | ||
|
|
b7171556b8 | ||
|
|
254a638392 | ||
|
|
74b3bd6b46 | ||
|
|
9a29ef5401 | ||
|
|
3f21faaa66 | ||
|
|
bf19d29a1d | ||
|
|
6523f707e3 | ||
|
|
8b0bf84d34 | ||
|
|
9ef4c51ba1 | ||
|
|
19521489fa | ||
|
|
2ce9aad005 | ||
|
|
2125884c6c | ||
|
|
f4c9226ef9 | ||
|
|
ff3999a6aa | ||
|
|
209470482a | ||
|
|
c26194554c |
+2
-2
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
|||||||
# run aws cli without mounting secret, because CI already has AWS credentials
|
# run aws cli without mounting secret, because CI already has AWS credentials
|
||||||
FROM build_base AS ci_image
|
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 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 . .
|
COPY . .
|
||||||
|
|
||||||
|
|
||||||
@@ -37,7 +37,7 @@ FROM build_base AS dev
|
|||||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
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
|
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
|
# flag --build-from-source is required to force-build sqlite3
|
||||||
RUN npm ci
|
RUN npm ci --ignore-scripts
|
||||||
COPY . .
|
COPY . .
|
||||||
ENTRYPOINT npm run start:dev
|
ENTRYPOINT npm run start:dev
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
|||||||
FROM build_base AS build
|
FROM build_base AS build
|
||||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
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
|
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 . .
|
COPY . .
|
||||||
RUN npm run build
|
RUN npm run build
|
||||||
|
|
||||||
|
|||||||
+186
-34
@@ -3001,14 +3001,7 @@
|
|||||||
"parameters": [],
|
"parameters": [],
|
||||||
"responses": {
|
"responses": {
|
||||||
"200": {
|
"200": {
|
||||||
"description": "",
|
"description": ""
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"type": "object"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"tags": [
|
"tags": [
|
||||||
@@ -3583,14 +3576,7 @@
|
|||||||
],
|
],
|
||||||
"responses": {
|
"responses": {
|
||||||
"200": {
|
"200": {
|
||||||
"description": "",
|
"description": ""
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"type": "object"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"tags": [
|
"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": {
|
"/platform/jobs/{jobId}/input": {
|
||||||
"put": {
|
"put": {
|
||||||
"operationId": "PlatformApiController_updateJobInput",
|
"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}": {
|
"/platform/jobs/jdbc/{jobId}": {
|
||||||
"get": {
|
"get": {
|
||||||
"operationId": "PlatformApiController_getJdbcJob",
|
"operationId": "PlatformApiController_getJdbcJob",
|
||||||
@@ -4747,6 +4814,42 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/platform/jobs/{jobId}/rename-tables": {
|
||||||
|
"post": {
|
||||||
|
"operationId": "PlatformApiController_renameJobTables",
|
||||||
|
"summary": "Rename job output tables and sync to catalog",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"name": "jobId",
|
||||||
|
"required": true,
|
||||||
|
"in": "path",
|
||||||
|
"schema": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"201": {
|
||||||
|
"description": "",
|
||||||
|
"content": {
|
||||||
|
"application/json": {
|
||||||
|
"schema": {
|
||||||
|
"type": "object"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"tags": [
|
||||||
|
"Platform API"
|
||||||
|
],
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
"/platform/jobs/jdbc/configs/allowed_datatypes": {
|
"/platform/jobs/jdbc/configs/allowed_datatypes": {
|
||||||
"get": {
|
"get": {
|
||||||
"operationId": "PlatformApiController_getJdbcAllowedDatatypes",
|
"operationId": "PlatformApiController_getJdbcAllowedDatatypes",
|
||||||
@@ -10125,6 +10228,21 @@
|
|||||||
"entities"
|
"entities"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
"Column": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"name": {
|
||||||
|
"type": "string"
|
||||||
|
},
|
||||||
|
"type": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"required": [
|
||||||
|
"name",
|
||||||
|
"type"
|
||||||
|
]
|
||||||
|
},
|
||||||
"TableColumns": {
|
"TableColumns": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
@@ -10140,13 +10258,7 @@
|
|||||||
"references": {
|
"references": {
|
||||||
"type": "array",
|
"type": "array",
|
||||||
"items": {
|
"items": {
|
||||||
"type": "string"
|
"$ref": "#/components/schemas/Column"
|
||||||
}
|
|
||||||
},
|
|
||||||
"identifier_columns": {
|
|
||||||
"type": "array",
|
|
||||||
"items": {
|
|
||||||
"type": "string"
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"destination": {
|
"destination": {
|
||||||
@@ -10154,13 +10266,20 @@
|
|||||||
},
|
},
|
||||||
"type": {
|
"type": {
|
||||||
"type": "string"
|
"type": "string"
|
||||||
|
},
|
||||||
|
"identifier_columns": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"reference_column": {
|
||||||
|
"$ref": "#/components/schemas/Column"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"required": [
|
"required": [
|
||||||
"name",
|
"name",
|
||||||
"columns",
|
"columns",
|
||||||
"references",
|
|
||||||
"identifier_columns",
|
|
||||||
"destination",
|
"destination",
|
||||||
"type"
|
"type"
|
||||||
]
|
]
|
||||||
@@ -10466,7 +10585,7 @@
|
|||||||
"enabled"
|
"enabled"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
"CustomerLink": {
|
"CustomerLinkItem": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
"href": {
|
"href": {
|
||||||
@@ -10488,28 +10607,61 @@
|
|||||||
"description"
|
"description"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
"CustomerLinksResponse": {
|
"CustomerSidebarSection": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
"links": {
|
"title": {
|
||||||
|
"type": "object"
|
||||||
|
},
|
||||||
|
"items": {
|
||||||
"type": "array",
|
"type": "array",
|
||||||
"items": {
|
"items": {
|
||||||
"$ref": "#/components/schemas/CustomerLink"
|
"oneOf": [
|
||||||
|
{
|
||||||
|
"$ref": "#/components/schemas/CustomerSidebarMenuItem"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"$ref": "#/components/schemas/CustomerSidebarLinkItem"
|
||||||
|
}
|
||||||
|
]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"required": [
|
"required": [
|
||||||
"links"
|
"title",
|
||||||
|
"items"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
"CustomerLinksConfig": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"home": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"$ref": "#/components/schemas/CustomerLinkItem"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"sidebar": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"$ref": "#/components/schemas/CustomerSidebarSection"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"CustomerLinksResponse": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"links": {
|
||||||
|
"$ref": "#/components/schemas/CustomerLinksConfig"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"CustomerLinkRequest": {
|
"CustomerLinkRequest": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
"links": {
|
"links": {
|
||||||
"type": "array",
|
"$ref": "#/components/schemas/CustomerLinksConfig"
|
||||||
"items": {
|
|
||||||
"$ref": "#/components/schemas/CustomerLink"
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"required": [
|
"required": [
|
||||||
|
|||||||
Generated
+83
-11
@@ -17,7 +17,7 @@
|
|||||||
"@aws-sdk/signature-v4": "^3.370.0",
|
"@aws-sdk/signature-v4": "^3.370.0",
|
||||||
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
|
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
|
||||||
"@dadosfera/protospack": "2.5.3",
|
"@dadosfera/protospack": "2.5.3",
|
||||||
"@dadosfera/protospack-v2": "3.38.0-beta.28",
|
"@dadosfera/protospack-v2": "^3.40.0-beta.5",
|
||||||
"@grpc/grpc-js": "^1.9.3",
|
"@grpc/grpc-js": "^1.9.3",
|
||||||
"@grpc/proto-loader": "^0.7.9",
|
"@grpc/proto-loader": "^0.7.9",
|
||||||
"@nestjs/cli": "^9.5.0",
|
"@nestjs/cli": "^9.5.0",
|
||||||
@@ -31,7 +31,7 @@
|
|||||||
"@nestjs/schematics": "^9.2.0",
|
"@nestjs/schematics": "^9.2.0",
|
||||||
"@nestjs/swagger": "^6.3.0",
|
"@nestjs/swagger": "^6.3.0",
|
||||||
"@nestjs/testing": "^9.4.3",
|
"@nestjs/testing": "^9.4.3",
|
||||||
"axios": "^0.30.2",
|
"axios": "0.30.3",
|
||||||
"cache-manager": "^5.1.4",
|
"cache-manager": "^5.1.4",
|
||||||
"cache-manager-ioredis-yet": "^1.1.0",
|
"cache-manager-ioredis-yet": "^1.1.0",
|
||||||
"class-transformer": "^0.5.1",
|
"class-transformer": "^0.5.1",
|
||||||
@@ -47,6 +47,7 @@
|
|||||||
"jwk-to-pem": "^2.0.5",
|
"jwk-to-pem": "^2.0.5",
|
||||||
"mixpanel": "^0.17.0",
|
"mixpanel": "^0.17.0",
|
||||||
"ms": "^3.0.0-canary.1",
|
"ms": "^3.0.0-canary.1",
|
||||||
|
"multer": "^2.0.2",
|
||||||
"openid-client": "^5.7.1",
|
"openid-client": "^5.7.1",
|
||||||
"passport": "^0.6.0",
|
"passport": "^0.6.0",
|
||||||
"passport-facebook": "^3.0.0",
|
"passport-facebook": "^3.0.0",
|
||||||
@@ -1744,10 +1745,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/@dadosfera/protospack-v2": {
|
"node_modules/@dadosfera/protospack-v2": {
|
||||||
"version": "3.38.0-beta.28",
|
"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.38.0-beta.28.tgz",
|
"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-w3Au0qschqZJ6OSHVDKaR2KeeCF2ncuJ4k7iZQceOLo1zxUU3TpDyLYyp/rp1DoSG763uIgqm3kJUXv7RlmfNw==",
|
"integrity": "sha512-iocKv/XXp2jKAasO5ONgm31cKfLgNsU4pEKZMr6YnR7nQaH11WcW7rnuagNxWoik++wLUqbYyf0bZWRDzMlCPA==",
|
||||||
"license": "ISC",
|
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@grpc/grpc-js": "^1.9.3",
|
"@grpc/grpc-js": "^1.9.3",
|
||||||
"rxjs": "^7.5.5"
|
"rxjs": "^7.5.5"
|
||||||
@@ -2928,6 +2928,20 @@
|
|||||||
"node": ">= 0.6"
|
"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": {
|
"node_modules/@nestjs/platform-express/node_modules/content-disposition": {
|
||||||
"version": "0.5.4",
|
"version": "0.5.4",
|
||||||
"resolved": "https://registry.npmjs.org/content-disposition/-/content-disposition-0.5.4.tgz",
|
"resolved": "https://registry.npmjs.org/content-disposition/-/content-disposition-0.5.4.tgz",
|
||||||
@@ -3097,6 +3111,24 @@
|
|||||||
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
|
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
|
||||||
"license": "MIT"
|
"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": {
|
"node_modules/@nestjs/platform-express/node_modules/negotiator": {
|
||||||
"version": "0.6.3",
|
"version": "0.6.3",
|
||||||
"resolved": "https://registry.npmjs.org/negotiator/-/negotiator-0.6.3.tgz",
|
"resolved": "https://registry.npmjs.org/negotiator/-/negotiator-0.6.3.tgz",
|
||||||
@@ -3121,6 +3153,25 @@
|
|||||||
"url": "https://github.com/sponsors/ljharb"
|
"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": {
|
"node_modules/@nestjs/platform-express/node_modules/safe-buffer": {
|
||||||
"version": "5.2.1",
|
"version": "5.2.1",
|
||||||
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.2.1.tgz",
|
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.2.1.tgz",
|
||||||
@@ -3195,6 +3246,19 @@
|
|||||||
"node": ">= 0.8"
|
"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": {
|
"node_modules/@nestjs/platform-express/node_modules/tslib": {
|
||||||
"version": "2.5.3",
|
"version": "2.5.3",
|
||||||
"resolved": "https://registry.npmjs.org/tslib/-/tslib-2.5.3.tgz",
|
"resolved": "https://registry.npmjs.org/tslib/-/tslib-2.5.3.tgz",
|
||||||
@@ -5534,9 +5598,9 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/axios": {
|
"node_modules/axios": {
|
||||||
"version": "0.30.2",
|
"version": "0.30.3",
|
||||||
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.2.tgz",
|
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.3.tgz",
|
||||||
"integrity": "sha512-0pE4RQ4UQi1jKY6p7u6i1Tkzqmu+d+/tHS7Q7rKunWLB9WyilBTpHHpXzPNMDj5hTbK0B0PTLSz07yqMBiF6xg==",
|
"integrity": "sha512-5/tmEb6TmE/ax3mdXBc/Mi6YdPGxQsv+0p5YlciXWt3PHIn0VamqCXhRMtScnwY3lbgSXLneOuXAKUhgmSRpwg==",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"follow-redirects": "^1.15.4",
|
"follow-redirects": "^1.15.4",
|
||||||
@@ -6509,7 +6573,6 @@
|
|||||||
"engines": [
|
"engines": [
|
||||||
"node >= 6.0"
|
"node >= 6.0"
|
||||||
],
|
],
|
||||||
"license": "MIT",
|
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"buffer-from": "^1.0.0",
|
"buffer-from": "^1.0.0",
|
||||||
"inherits": "^2.0.3",
|
"inherits": "^2.0.3",
|
||||||
@@ -9128,6 +9191,11 @@
|
|||||||
"url": "https://github.com/sponsors/sindresorhus"
|
"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": {
|
"node_modules/isexe": {
|
||||||
"version": "2.0.0",
|
"version": "2.0.0",
|
||||||
"resolved": "https://registry.npmjs.org/isexe/-/isexe-2.0.0.tgz",
|
"resolved": "https://registry.npmjs.org/isexe/-/isexe-2.0.0.tgz",
|
||||||
@@ -10725,7 +10793,6 @@
|
|||||||
"version": "2.0.2",
|
"version": "2.0.2",
|
||||||
"resolved": "https://registry.npmjs.org/multer/-/multer-2.0.2.tgz",
|
"resolved": "https://registry.npmjs.org/multer/-/multer-2.0.2.tgz",
|
||||||
"integrity": "sha512-u7f2xaZ/UG8oLXHvtF/oWTRvT44p9ecwBBqTwgJVq0+4BW1g8OW01TyMEGWBHbyMOYVHXslaut7qEQ1meATXgw==",
|
"integrity": "sha512-u7f2xaZ/UG8oLXHvtF/oWTRvT44p9ecwBBqTwgJVq0+4BW1g8OW01TyMEGWBHbyMOYVHXslaut7qEQ1meATXgw==",
|
||||||
"license": "MIT",
|
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"append-field": "^1.0.0",
|
"append-field": "^1.0.0",
|
||||||
"busboy": "^1.6.0",
|
"busboy": "^1.6.0",
|
||||||
@@ -11728,6 +11795,11 @@
|
|||||||
"url": "https://github.com/chalk/ansi-styles?sponsor=1"
|
"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": {
|
"node_modules/process-warning": {
|
||||||
"version": "1.0.0",
|
"version": "1.0.0",
|
||||||
"resolved": "https://registry.npmjs.org/process-warning/-/process-warning-1.0.0.tgz",
|
"resolved": "https://registry.npmjs.org/process-warning/-/process-warning-1.0.0.tgz",
|
||||||
|
|||||||
+8
-4
@@ -10,7 +10,7 @@
|
|||||||
},
|
},
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"co:login": "aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1",
|
"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",
|
"prebuild": "rimraf dist",
|
||||||
"build": "nest build",
|
"build": "nest build",
|
||||||
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
|
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
|
||||||
@@ -35,7 +35,7 @@
|
|||||||
"@aws-sdk/signature-v4": "^3.370.0",
|
"@aws-sdk/signature-v4": "^3.370.0",
|
||||||
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
|
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
|
||||||
"@dadosfera/protospack": "2.5.3",
|
"@dadosfera/protospack": "2.5.3",
|
||||||
"@dadosfera/protospack-v2": "3.38.0-beta.28",
|
"@dadosfera/protospack-v2": "^3.40.0-beta.5",
|
||||||
"@grpc/grpc-js": "^1.9.3",
|
"@grpc/grpc-js": "^1.9.3",
|
||||||
"@grpc/proto-loader": "^0.7.9",
|
"@grpc/proto-loader": "^0.7.9",
|
||||||
"@nestjs/cli": "^9.5.0",
|
"@nestjs/cli": "^9.5.0",
|
||||||
@@ -49,7 +49,7 @@
|
|||||||
"@nestjs/schematics": "^9.2.0",
|
"@nestjs/schematics": "^9.2.0",
|
||||||
"@nestjs/swagger": "^6.3.0",
|
"@nestjs/swagger": "^6.3.0",
|
||||||
"@nestjs/testing": "^9.4.3",
|
"@nestjs/testing": "^9.4.3",
|
||||||
"axios": "^0.30.2",
|
"axios": "0.30.3",
|
||||||
"cache-manager": "^5.1.4",
|
"cache-manager": "^5.1.4",
|
||||||
"cache-manager-ioredis-yet": "^1.1.0",
|
"cache-manager-ioredis-yet": "^1.1.0",
|
||||||
"class-transformer": "^0.5.1",
|
"class-transformer": "^0.5.1",
|
||||||
@@ -65,6 +65,7 @@
|
|||||||
"jwk-to-pem": "^2.0.5",
|
"jwk-to-pem": "^2.0.5",
|
||||||
"mixpanel": "^0.17.0",
|
"mixpanel": "^0.17.0",
|
||||||
"ms": "^3.0.0-canary.1",
|
"ms": "^3.0.0-canary.1",
|
||||||
|
"multer": "^2.0.2",
|
||||||
"openid-client": "^5.7.1",
|
"openid-client": "^5.7.1",
|
||||||
"passport": "^0.6.0",
|
"passport": "^0.6.0",
|
||||||
"passport-facebook": "^3.0.0",
|
"passport-facebook": "^3.0.0",
|
||||||
@@ -80,7 +81,7 @@
|
|||||||
"swagger-ui-express": "^4.6.3"
|
"swagger-ui-express": "^4.6.3"
|
||||||
},
|
},
|
||||||
"overrides": {
|
"overrides": {
|
||||||
"multer": "2.0.2",
|
"axios": "0.30.3",
|
||||||
"form-data": "^4.0.4",
|
"form-data": "^4.0.4",
|
||||||
"body-parser": "^1.20.3",
|
"body-parser": "^1.20.3",
|
||||||
"cross-spawn": "^7.0.5",
|
"cross-spawn": "^7.0.5",
|
||||||
@@ -117,5 +118,8 @@
|
|||||||
"ts-node": "^10.9.1",
|
"ts-node": "^10.9.1",
|
||||||
"tsconfig-paths": "^3.14.2",
|
"tsconfig-paths": "^3.14.2",
|
||||||
"typescript": "^4.9.5"
|
"typescript": "^4.9.5"
|
||||||
|
},
|
||||||
|
"resolutions": {
|
||||||
|
"axios": "0.30.3"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -489,6 +489,7 @@ export class AuthController {
|
|||||||
const userDto = {
|
const userDto = {
|
||||||
id: api_key.user_id,
|
id: api_key.user_id,
|
||||||
name: api_key.username,
|
name: api_key.username,
|
||||||
|
email: api_key.username,
|
||||||
customer: {
|
customer: {
|
||||||
id: api_key.customer_id,
|
id: api_key.customer_id,
|
||||||
name: api_key.customer_name,
|
name: api_key.customer_name,
|
||||||
|
|||||||
@@ -437,7 +437,8 @@ export class AuthClientService implements OnModuleInit {
|
|||||||
|
|
||||||
const userDto: UserDTO = {
|
const userDto: UserDTO = {
|
||||||
id: user.id,
|
id: user.id,
|
||||||
name: user.username,
|
name: user.name,
|
||||||
|
email: user.email,
|
||||||
jobTitle: user?.jobTitle || null,
|
jobTitle: user?.jobTitle || null,
|
||||||
department: user?.department || null,
|
department: user?.department || null,
|
||||||
hierarchy: user?.hierarchy || null,
|
hierarchy: user?.hierarchy || null,
|
||||||
|
|||||||
@@ -144,6 +144,7 @@ export interface BulkEditResponse {
|
|||||||
export type UserDTO = {
|
export type UserDTO = {
|
||||||
id: string,
|
id: string,
|
||||||
name: string,
|
name: string,
|
||||||
|
email: string,
|
||||||
jobTitle?: string,
|
jobTitle?: string,
|
||||||
department?: string,
|
department?: string,
|
||||||
hierarchy?: string,
|
hierarchy?: string,
|
||||||
|
|||||||
@@ -9,12 +9,12 @@ import {
|
|||||||
} from '@nestjs/common';
|
} from '@nestjs/common';
|
||||||
|
|
||||||
import { firstValueFrom, lastValueFrom } from 'rxjs';
|
import { firstValueFrom, lastValueFrom } from 'rxjs';
|
||||||
import { Link } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/entities';
|
|
||||||
import { DucClient } from '../duc/client.config';
|
import { DucClient } from '../duc/client.config';
|
||||||
import { ClientGrpc } from '@nestjs/microservices';
|
import { ClientGrpc } from '@nestjs/microservices';
|
||||||
import { ProtoServices } from '@dadosfera/protospack-v2/dist/lib/Duc';
|
import { ProtoServices } from '@dadosfera/protospack-v2/dist/lib/Duc';
|
||||||
import { CustomerUpdateRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
import { CustomerSetLinksRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||||
import { CustomersProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
import { CustomersProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||||
|
import { CustomerLinksConfig } from './dtos/customers';
|
||||||
import ErrorCodes from 'src/utils/errorCodes';
|
import ErrorCodes from 'src/utils/errorCodes';
|
||||||
import jwt from 'jsonwebtoken';
|
import jwt from 'jsonwebtoken';
|
||||||
import {
|
import {
|
||||||
@@ -67,12 +67,12 @@ export class CustomersService implements OnModuleInit {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
async getLinks(customerId: string) {
|
async getLinks(customerId: string): Promise<CustomerLinksConfig | null> {
|
||||||
try {
|
try {
|
||||||
const result = await lastValueFrom(
|
const result = await lastValueFrom(
|
||||||
this.customerService.CustomerFindOneById({ id: customerId }),
|
this.customerService.CustomerGetLinks({ customerId }),
|
||||||
);
|
);
|
||||||
return result.customer?.links || [];
|
return (result.links as CustomerLinksConfig) || null;
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||||
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
||||||
@@ -80,17 +80,17 @@ export class CustomersService implements OnModuleInit {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async setLinks(customerId: string, links: Link[]) {
|
async setLinks(customerId: string, links: CustomerLinksConfig) {
|
||||||
if (!customerId || !links) {
|
if (!customerId || !links) {
|
||||||
throw new HttpException(null, HttpStatus.BAD_REQUEST);
|
throw new HttpException(null, HttpStatus.BAD_REQUEST);
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
return await firstValueFrom(
|
return await firstValueFrom(
|
||||||
this.customerService.CustomerUpdate({
|
this.customerService.CustomerSetLinks({
|
||||||
id: customerId,
|
customerId,
|
||||||
links,
|
links: links as CustomerSetLinksRequest['links'],
|
||||||
} as CustomerUpdateRequest),
|
}),
|
||||||
);
|
);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
import { Link } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/entities';
|
|
||||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||||
|
|
||||||
export class CustomerLink implements Link {
|
export class CustomerLinkItem {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
href: string;
|
href: string;
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
@@ -9,15 +8,59 @@ export class CustomerLink implements Link {
|
|||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
description: string;
|
description: string;
|
||||||
@ApiPropertyOptional()
|
@ApiPropertyOptional()
|
||||||
iconSrc: string;
|
iconSrc?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export class CustomerSidebarLinkItem {
|
||||||
|
@ApiProperty()
|
||||||
|
type: 'link';
|
||||||
|
@ApiProperty({ type: Object })
|
||||||
|
title: Record<string, string>;
|
||||||
|
@ApiProperty()
|
||||||
|
link: string;
|
||||||
|
@ApiPropertyOptional()
|
||||||
|
icon?: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class CustomerSidebarMenuItem {
|
||||||
|
@ApiProperty()
|
||||||
|
type: 'menu';
|
||||||
|
@ApiProperty({ type: Object })
|
||||||
|
title: Record<string, string>;
|
||||||
|
@ApiPropertyOptional()
|
||||||
|
icon?: string;
|
||||||
|
@ApiProperty({ type: [CustomerSidebarLinkItem] })
|
||||||
|
items: CustomerSidebarLinkItem[];
|
||||||
|
}
|
||||||
|
|
||||||
|
export class CustomerSidebarSection {
|
||||||
|
@ApiProperty({ type: Object })
|
||||||
|
title: Record<string, string>;
|
||||||
|
@ApiProperty({
|
||||||
|
type: 'array',
|
||||||
|
items: {
|
||||||
|
oneOf: [
|
||||||
|
{ $ref: '#/components/schemas/CustomerSidebarMenuItem' },
|
||||||
|
{ $ref: '#/components/schemas/CustomerSidebarLinkItem' },
|
||||||
|
],
|
||||||
|
},
|
||||||
|
})
|
||||||
|
items: (CustomerSidebarMenuItem | CustomerSidebarLinkItem)[];
|
||||||
|
}
|
||||||
|
|
||||||
|
export class CustomerLinksConfig {
|
||||||
|
@ApiPropertyOptional({ type: [CustomerLinkItem] })
|
||||||
|
home?: CustomerLinkItem[];
|
||||||
|
@ApiPropertyOptional({ type: [CustomerSidebarSection] })
|
||||||
|
sidebar?: CustomerSidebarSection[];
|
||||||
|
}
|
||||||
|
|
||||||
export class CustomerLinkRequest {
|
export class CustomerLinkRequest {
|
||||||
@ApiProperty({ type: [CustomerLink] })
|
@ApiProperty({ type: CustomerLinksConfig })
|
||||||
links: CustomerLink[];
|
links: CustomerLinksConfig;
|
||||||
}
|
}
|
||||||
|
|
||||||
export class CustomerLinksResponse {
|
export class CustomerLinksResponse {
|
||||||
@ApiProperty({ type: [CustomerLink] })
|
@ApiPropertyOptional({ type: CustomerLinksConfig })
|
||||||
links: CustomerLink[];
|
links?: CustomerLinksConfig;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -11,17 +11,19 @@ export class TableColumns {
|
|||||||
name: string;
|
name: string;
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
columns: string[];
|
columns: string[];
|
||||||
@ApiProperty()
|
@ApiPropertyOptional({ type: [Column] })
|
||||||
references: Column[];
|
references: Column[];
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
identifier_columns: string[];
|
|
||||||
@ApiProperty()
|
|
||||||
destination: Record<'raw' | 'qualify', {
|
destination: Record<'raw' | 'qualify', {
|
||||||
table_name: string;
|
table_name: string;
|
||||||
table_schema: string;
|
table_schema: string;
|
||||||
}> | null;
|
}> | null;
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
type: string;
|
type: string;
|
||||||
|
@ApiPropertyOptional({ type: [String] })
|
||||||
|
identifier_columns?: string[];
|
||||||
|
@ApiPropertyOptional({ type: Column })
|
||||||
|
reference_column?: Column;
|
||||||
}
|
}
|
||||||
export class AvailableEntity {
|
export class AvailableEntity {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
|
|||||||
@@ -99,6 +99,7 @@ export class InputsController {
|
|||||||
customer: info.customer,
|
customer: info.customer,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
this.logger.info(JSON.stringify(body))
|
||||||
const response = await this.inputService.create({ body, info });
|
const response = await this.inputService.create({ body, info });
|
||||||
|
|
||||||
return response;
|
return response;
|
||||||
|
|||||||
@@ -17,6 +17,8 @@ import {
|
|||||||
InputCreateGenericRequest,
|
InputCreateGenericRequest,
|
||||||
InputCreateS3Request,
|
InputCreateS3Request,
|
||||||
InputNewCreateRequest,
|
InputNewCreateRequest,
|
||||||
|
InputUpdateResponse,
|
||||||
|
RollbackInputRequest,
|
||||||
TestConnectionRequest,
|
TestConnectionRequest,
|
||||||
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
|
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
|
||||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
||||||
@@ -71,8 +73,8 @@ export class InputsService {
|
|||||||
objectCamelToSnake(createInputResponse);
|
objectCamelToSnake(createInputResponse);
|
||||||
return createInputResponse;
|
return createInputResponse;
|
||||||
},
|
},
|
||||||
update: async (updateInputDTO: UpdateInputRequest) => {
|
update: async (updateInputDTO: UpdateInputRequest): Promise<InputUpdateResponse> => {
|
||||||
this.logger.info('InputClientService - Update');
|
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
||||||
const updateInputResponse = await lastValueFrom(
|
const updateInputResponse = await lastValueFrom(
|
||||||
this.inputWriteService.InputUpdate(updateInputDTO),
|
this.inputWriteService.InputUpdate(updateInputDTO),
|
||||||
);
|
);
|
||||||
@@ -164,6 +166,11 @@ export class InputsService {
|
|||||||
const inputCreateGenericRequest: InputCreateGenericRequest = {
|
const inputCreateGenericRequest: InputCreateGenericRequest = {
|
||||||
input: {
|
input: {
|
||||||
...body,
|
...body,
|
||||||
|
tables: (body.tables || []).map((table) => ({
|
||||||
|
...table,
|
||||||
|
identifier_columns: table.identifier_columns || [],
|
||||||
|
reference_column: table.reference_column || table.references?.[0],
|
||||||
|
})),
|
||||||
},
|
},
|
||||||
info,
|
info,
|
||||||
};
|
};
|
||||||
@@ -202,21 +209,44 @@ export class InputsService {
|
|||||||
async update(id: string, data, info: Info) {
|
async update(id: string, data, info: Info) {
|
||||||
// this.validateCron({ ...data, info });
|
// this.validateCron({ ...data, info });
|
||||||
try {
|
try {
|
||||||
const updateInputResponse: any = await this.OLD_inputClient.update({
|
const {
|
||||||
|
tablesUpdate,
|
||||||
|
dataAssetUpdate,
|
||||||
|
input
|
||||||
|
} = await this.OLD_inputClient.update({
|
||||||
id,
|
id,
|
||||||
info,
|
|
||||||
...data,
|
...data,
|
||||||
|
info,
|
||||||
});
|
});
|
||||||
|
|
||||||
updateInputResponse.input = this.adjustInputPayload(
|
const updateInputResponse = this.adjustInputPayload(
|
||||||
updateInputResponse?.input,
|
input,
|
||||||
);
|
);
|
||||||
return updateInputResponse;
|
return {
|
||||||
|
input: updateInputResponse,
|
||||||
|
tablesUpdate,
|
||||||
|
dataAssetUpdate
|
||||||
|
};
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
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) {
|
async remove(idRequest: IIdRequest) {
|
||||||
return lastValueFrom(this.inputWriteService.InputRemove(idRequest));
|
return lastValueFrom(this.inputWriteService.InputRemove(idRequest));
|
||||||
}
|
}
|
||||||
@@ -258,4 +288,12 @@ export class InputsService {
|
|||||||
};
|
};
|
||||||
return formatedPayload;
|
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));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -53,6 +53,9 @@ import { TableColumns } from '../inputs/dtos/input.model';
|
|||||||
import { UpdateInputRequest } from '../inputs/dtos/old_interfaces';
|
import { UpdateInputRequest } from '../inputs/dtos/old_interfaces';
|
||||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
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')
|
@ApiTags('PipelinesV2')
|
||||||
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
|
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
|
||||||
@UseFilters(new GrpcToHttpExceptionFilter())
|
@UseFilters(new GrpcToHttpExceptionFilter())
|
||||||
@@ -62,6 +65,7 @@ export class PipelinesController {
|
|||||||
constructor(
|
constructor(
|
||||||
@Inject(DadosferaLogger)
|
@Inject(DadosferaLogger)
|
||||||
dadosferaLogger: DadosferaLogger,
|
dadosferaLogger: DadosferaLogger,
|
||||||
|
|
||||||
private pipelinesClientService: PipelinesService,
|
private pipelinesClientService: PipelinesService,
|
||||||
private oldPipelinesService: OldPipelineService,
|
private oldPipelinesService: OldPipelineService,
|
||||||
) {
|
) {
|
||||||
@@ -222,30 +226,27 @@ export class PipelinesController {
|
|||||||
language,
|
language,
|
||||||
});
|
});
|
||||||
|
|
||||||
const result = await this.pipelinesClientService
|
const pipelineRes = await this.pipelinesClientService.findOne({ id }, metadata);
|
||||||
.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;
|
|
||||||
});
|
|
||||||
|
|
||||||
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')
|
@Patch('/:id')
|
||||||
@@ -294,13 +295,14 @@ export class PipelinesController {
|
|||||||
) {
|
) {
|
||||||
this.logger.info('PipelinesController - update', { user });
|
this.logger.info('PipelinesController - update', { user });
|
||||||
|
|
||||||
|
|
||||||
const { customer_id, customer_name, user_id, username } = user;
|
const { customer_id, customer_name, user_id, username } = user;
|
||||||
const info: Info = {
|
const info: Info = {
|
||||||
user_id: user.user_id,
|
user_id: user.user_id,
|
||||||
customer: user.customer_name,
|
customer: user.customer_name,
|
||||||
customer_id: user.customer_id,
|
customer_id: user.customer_id,
|
||||||
|
pipeline_id: pipelineId
|
||||||
};
|
};
|
||||||
|
|
||||||
const metadata = PackTheMetadata({
|
const metadata = PackTheMetadata({
|
||||||
customer_id,
|
customer_id,
|
||||||
customer_name,
|
customer_name,
|
||||||
|
|||||||
@@ -12,6 +12,8 @@ import { ConnectorModule } from '../connector/connector.module';
|
|||||||
import { InputsModule } from '../inputs/inputs.module';
|
import { InputsModule } from '../inputs/inputs.module';
|
||||||
import { TransformationsModule } from '../transformations/transformations.module';
|
import { TransformationsModule } from '../transformations/transformations.module';
|
||||||
import { PlatformApiModule } from '../platform-api/platform-api.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();
|
const client = new PipelinesClientConfiguration();
|
||||||
|
|
||||||
@@ -22,10 +24,11 @@ const client = new PipelinesClientConfiguration();
|
|||||||
ConnectorModule,
|
ConnectorModule,
|
||||||
InputsModule,
|
InputsModule,
|
||||||
TransformationsModule,
|
TransformationsModule,
|
||||||
PlatformApiModule
|
PlatformApiModule,
|
||||||
|
NimbusServicesModule
|
||||||
],
|
],
|
||||||
controllers: [PipelinesController],
|
controllers: [PipelinesController],
|
||||||
providers: [PipelinesService, DadosferaLogger],
|
providers: [PipelinesService, DadosferaLogger, NimbusService],
|
||||||
exports: [PipelinesService],
|
exports: [PipelinesService],
|
||||||
})
|
})
|
||||||
export class PipelinesV2Module {}
|
export class PipelinesV2Module {}
|
||||||
|
|||||||
@@ -29,6 +29,11 @@ import ErrorCodes from 'src/utils/errorCodes';
|
|||||||
import ErrorBuilder from 'src/utils/ErrorBuilder';
|
import ErrorBuilder from 'src/utils/ErrorBuilder';
|
||||||
import { PlatformApiService } from '../platform-api/platform-api.service';
|
import { PlatformApiService } from '../platform-api/platform-api.service';
|
||||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
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 {
|
export class PipelinesService implements OnModuleInit {
|
||||||
logger: DadosferaLogger;
|
logger: DadosferaLogger;
|
||||||
@@ -42,7 +47,8 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
private readonly connectorService: ConnectorClientService,
|
private readonly connectorService: ConnectorClientService,
|
||||||
private readonly inputsService: InputsService,
|
private readonly inputsService: InputsService,
|
||||||
private readonly transformationsService: TransformationsService,
|
private readonly transformationsService: TransformationsService,
|
||||||
private readonly platformAPI: PlatformApiService
|
private readonly platformAPI: PlatformApiService,
|
||||||
|
private readonly nimbusService: NimbusService
|
||||||
) {
|
) {
|
||||||
this.logger = dadosferaLogger.logger;
|
this.logger = dadosferaLogger.logger;
|
||||||
}
|
}
|
||||||
@@ -348,133 +354,227 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
async updatePipelineInput(pipelineId: string, inputId: string, updateInputDTO: UpdatePlatformInputRequest, info: Info, user: RequestUser, metadata: Metadata) {
|
async updatePipelineInput(pipelineId: string, inputId: string, updateInputDTO: UpdatePlatformInputRequest, info: Info, user: RequestUser, metadata: Metadata) {
|
||||||
this.logger.info('InputClientService - Update');
|
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 pipelineIdFormat = pipelineId.split('-').join('_');
|
||||||
|
const rollback: RollbackPromise[] = [];
|
||||||
|
|
||||||
const updateInputResponse = await this.inputsService.update(
|
const updateInputResponse = await this.inputsService.update(
|
||||||
inputId,
|
inputId,
|
||||||
updateInputDTO,
|
updateInputDTO,
|
||||||
info
|
info
|
||||||
)
|
);
|
||||||
|
|
||||||
const requests = [];
|
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>;
|
||||||
|
}
|
||||||
|
|
||||||
this.logger.info('Dynamo Response', updateInputResponse);
|
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()) {
|
for (const [index, table] of updateInputDTO.tables.entries()) {
|
||||||
const id = `${pipelineIdFormat}_${index}`;
|
const jobUpdate = {
|
||||||
this.logger.info('Updating input reference for table', table.name);
|
job_id: `${pipelineId}_${index}`,
|
||||||
const body = {}
|
}
|
||||||
|
|
||||||
|
if (table.type !== "incremental_with_qualify") {
|
||||||
|
delete table.destinations?.qualify;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (table.memory) {
|
||||||
|
jobUpdate["memory"] = {
|
||||||
|
amount: table.memory * 1000
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
this.logger.info('Updating input reference for table: ' + table.name);
|
||||||
|
let hasUpdateSyncMode = false;
|
||||||
|
|
||||||
|
const jobSyncMode = {}
|
||||||
|
|
||||||
if (table.columns) {
|
if (table.columns) {
|
||||||
body['column_include_list'] = table.columns;
|
hasUpdateSyncMode = true;
|
||||||
|
jobSyncMode['column_include_list'] = table.columns;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (table.reference_column) {
|
if (table.reference_column) {
|
||||||
body['incremental_column_name'] = table.reference_column.name;
|
hasUpdateSyncMode = true;
|
||||||
body['incremental_column_type'] = table.reference_column.type;
|
jobSyncMode['incremental_column_name'] = table.reference_column.name;
|
||||||
|
jobSyncMode['incremental_column_type'] = table.reference_column.type;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (table.identifier_columns) {
|
if (table.identifier_columns) {
|
||||||
body['primary_keys'] = table.identifier_columns;
|
hasUpdateSyncMode = true;
|
||||||
}
|
jobSyncMode['primary_keys'] = table.identifier_columns;
|
||||||
|
|
||||||
this.logger.info('Request body', body);
|
|
||||||
const updateCollumns = this.platformAPI.proxy(
|
|
||||||
'PATCH',
|
|
||||||
`/jobs/${id}/input`,
|
|
||||||
user,
|
|
||||||
body
|
|
||||||
)
|
|
||||||
requests.push(updateCollumns);
|
|
||||||
|
|
||||||
if (table.memory) {
|
|
||||||
this.logger.info('Updating memory allocation for table', table.name);
|
|
||||||
const updateMemory = this.platformAPI.proxy(
|
|
||||||
'PUT',
|
|
||||||
`/jobs/${id}/memory`,
|
|
||||||
user,
|
|
||||||
{
|
|
||||||
amount: table.memory
|
|
||||||
}
|
|
||||||
)
|
|
||||||
requests.push(updateMemory);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if (table.type) {
|
if (table.type) {
|
||||||
const updateSyncMode = this.updatePipelineSyncMode(table, id, user);
|
hasUpdateSyncMode = true;
|
||||||
requests.push(updateSyncMode);
|
|
||||||
|
jobSyncMode['target_load_type'] = table.type;
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
this.logger.info('Create Platform Request for each JOB');
|
if(hasUpdateSyncMode) {
|
||||||
|
jobUpdate["sync_mode"] = jobSyncMode;
|
||||||
|
}
|
||||||
|
|
||||||
if (updateInputDTO.cron) {
|
if (Object.keys(table.destinations).length > 1) {
|
||||||
const crnUpdatedRequest = new Promise(async (resolve, reject) => {
|
let hasChanges = false
|
||||||
const response = await this.updatePipelineCron(updateInputDTO.cron, pipelineIdFormat, user);
|
const jobRenameTables = {
|
||||||
|
raw: {},
|
||||||
if (response.error) {
|
qualify: {}
|
||||||
this.logger.error('Error updating pipeline cron', response.error);
|
|
||||||
return reject(new ErrorBuilder(response.error));
|
|
||||||
}
|
}
|
||||||
this.logger.error('Pipeline cron updated successfully', response);
|
|
||||||
return resolve(response);
|
if (Object.keys(table.destinations.raw).length > 1) {
|
||||||
});
|
hasChanges = true;
|
||||||
requests.push(crnUpdatedRequest);
|
jobRenameTables.raw = table.destinations.raw;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (Object.keys(table.destinations.qualify).length > 1) {
|
||||||
|
hasChanges = true;
|
||||||
|
jobRenameTables.qualify = table.destinations.qualify;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (hasChanges) {
|
||||||
|
jobUpdate['rename_tables'] = jobRenameTables;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
jobsUpdated.push(jobUpdate);
|
||||||
}
|
}
|
||||||
|
|
||||||
this.logger.info('Executing all request for the platform api');
|
this.logger.info('Request body:' + JSON.stringify({
|
||||||
|
jobs_updated: jobsUpdated
|
||||||
|
}));
|
||||||
|
|
||||||
const results = await Promise.allSettled(requests);
|
const response = await this.platformAPI.proxy(
|
||||||
this.logger.info('Platform api response', results);
|
'PUT',
|
||||||
|
`/pipeline/${pipelineId}/jobs`,
|
||||||
return updateInputResponse;
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
private async updatePipelineSyncMode(table: UpdateTableDTO, pipelineId: string, user: RequestUser) {
|
|
||||||
const body = {
|
|
||||||
target_load_type: table.type
|
|
||||||
}
|
|
||||||
|
|
||||||
if (table.type === 'incremental_with_qualify') {
|
|
||||||
body['incremental_column_name'] = table.reference_column.name;
|
|
||||||
body['incremental_column_type'] = table.reference_column.type;
|
|
||||||
body['primary_keys'] = table.identifier_columns;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (table.type === 'incremental') {
|
|
||||||
body['incremental_column_name'] = table.reference_column.name;
|
|
||||||
body['incremental_column_type'] = table.reference_column.type;
|
|
||||||
}
|
|
||||||
|
|
||||||
this.logger.info('Updating pipeline sync mode', {
|
|
||||||
pipelineId,
|
|
||||||
body
|
|
||||||
});
|
|
||||||
|
|
||||||
return this.platformAPI.proxy(
|
|
||||||
"POST",
|
|
||||||
`/jobs/jdbc/${pipelineId}/sync-mode`,
|
|
||||||
user,
|
user,
|
||||||
body
|
{
|
||||||
|
job_updates: jobsUpdated
|
||||||
|
}
|
||||||
)
|
)
|
||||||
|
this.logger.info('Platform api response: ' + JSON.stringify(response));
|
||||||
}
|
}
|
||||||
|
|
||||||
private async updatePipelineCron(cron: string, pipelineId: string, user: RequestUser) {
|
|
||||||
try {
|
|
||||||
const response = await this.platformAPI.proxy(
|
|
||||||
'PATCH',
|
|
||||||
`/pipeline/${pipelineId}`,
|
|
||||||
user,
|
|
||||||
{
|
|
||||||
cron
|
|
||||||
}
|
|
||||||
);
|
|
||||||
|
|
||||||
return response
|
|
||||||
} catch (error) {
|
|
||||||
return {
|
|
||||||
error: error.message
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import {
|
|||||||
Inject,
|
Inject,
|
||||||
BadRequestException,
|
BadRequestException,
|
||||||
HttpException,
|
HttpException,
|
||||||
|
NotFoundException,
|
||||||
} from '@nestjs/common';
|
} from '@nestjs/common';
|
||||||
import { ApiTags, ApiOperation } from '@nestjs/swagger';
|
import { ApiTags, ApiOperation } from '@nestjs/swagger';
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
@@ -29,6 +30,7 @@ import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
|
|||||||
import { CatalogService } from '../catalog/catalog.service';
|
import { CatalogService } from '../catalog/catalog.service';
|
||||||
import { PackTheMetadata } from '../../utils/PackTheMetadata';
|
import { PackTheMetadata } from '../../utils/PackTheMetadata';
|
||||||
import { ValidationTableDTO } from './platform-api.dto';
|
import { ValidationTableDTO } from './platform-api.dto';
|
||||||
|
import { InputsService } from '../inputs/inputs.service';
|
||||||
|
|
||||||
|
|
||||||
type ValidateTablesDTO = {
|
type ValidateTablesDTO = {
|
||||||
@@ -54,6 +56,7 @@ export class PlatformApiController {
|
|||||||
private readonly dynamoDBService: DynamoDBService,
|
private readonly dynamoDBService: DynamoDBService,
|
||||||
private readonly customersService: CustomersService,
|
private readonly customersService: CustomersService,
|
||||||
private readonly catalogService: CatalogService,
|
private readonly catalogService: CatalogService,
|
||||||
|
private readonly inputsService: InputsService,
|
||||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||||
) {
|
) {
|
||||||
this.logger = dadosferaLogger.logger;
|
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 ====================
|
// ==================== JOBS - COLUMN EDITING ROUTES ====================
|
||||||
|
|
||||||
@Put('jobs/:jobId/input')
|
@Put('jobs/:jobId/input')
|
||||||
@@ -944,226 +965,52 @@ export class PlatformApiController {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
|
@Delete('pipelines/:pipelineId/inputs/:inputId')
|
||||||
|
@ApiOperation({ summary: 'Mark a table as deleted and delete its associated job via platform-api' })
|
||||||
@Get('jobs/jdbc/:jobId')
|
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
|
||||||
@ApiOperation({ summary: 'Get JDBC job details' })
|
async deleteTable(
|
||||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
@Param('pipelineId') pipelineId: string,
|
||||||
async getJdbcJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
@Param('inputId') inputId: string,
|
||||||
// Normalize job ID for Platform API (replace - with _)
|
@Body() body: { table_name: string },
|
||||||
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,
|
|
||||||
@User() user: RequestUser,
|
@User() user: RequestUser,
|
||||||
) {
|
) {
|
||||||
// Normalize job ID for Platform API (replace - with _)
|
const tableName = body.table_name;
|
||||||
const normalizedJobId = this.normalizeJobId(jobId);
|
const info = {
|
||||||
|
customer_id: user.customer_id,
|
||||||
|
customer: user.customer_name,
|
||||||
|
user_id: user.user_id,
|
||||||
|
};
|
||||||
|
|
||||||
const result = await this.platformApiService.proxy(
|
this.logger.info('deleteTable: marking table as deleted', { inputId, tableName });
|
||||||
'POST',
|
const updatedInput: any = await this.inputsService.markTableDeleted({ input_id: inputId, table_name: tableName, info });
|
||||||
`/jobs/jdbc/${normalizedJobId}/sync-mode`,
|
this.logger.info('deleteTable: table marked as deleted', { inputId, tableName });
|
||||||
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,
|
|
||||||
);
|
|
||||||
|
|
||||||
try {
|
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) {
|
} catch (error) {
|
||||||
this.logger.error('Catalog sync failed, rolling back Snowflake rename', { jobId, error: error.message });
|
this.logger.error('deleteTable: platform-api delete failed, attempting rollback', { tableName, 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>> = [];
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const dataAsset = await this.elasticsearchService.findDataAssetByTable(
|
await this.inputsService.unmarkTableDeleted({ input_id: inputId, table_name: tableName, info });
|
||||||
user.customer_name, oldTableName, oldTableSchema,
|
} catch (rollbackError) {
|
||||||
);
|
this.logger.error('deleteTable: rollback failed', { tableName, error: rollbackError.message });
|
||||||
|
|
||||||
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;
|
|
||||||
}
|
}
|
||||||
|
throw error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private buildRenameTargets(
|
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
|
||||||
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 });
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
@Get('jobs/jdbc/configs/allowed_datatypes')
|
@Get('jobs/jdbc/configs/allowed_datatypes')
|
||||||
@ApiOperation({ summary: 'Get allowed datatypes for JDBC' })
|
@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 ====================
|
// ==================== HEALTH ROUTE ====================
|
||||||
|
|
||||||
@Get('health')
|
@Get('health')
|
||||||
|
|||||||
@@ -8,9 +8,10 @@ import { ElasticsearchModule } from '../../services/elasticsearch';
|
|||||||
import { DynamoDBModule } from '../../services/dynamodb';
|
import { DynamoDBModule } from '../../services/dynamodb';
|
||||||
import { CustomersModule } from '../customers/customers.module';
|
import { CustomersModule } from '../customers/customers.module';
|
||||||
import { CatalogModule } from '../catalog/catalog.module';
|
import { CatalogModule } from '../catalog/catalog.module';
|
||||||
|
import { InputsModule } from '../inputs/inputs.module';
|
||||||
|
|
||||||
@Module({
|
@Module({
|
||||||
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule],
|
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule],
|
||||||
controllers: [PlatformApiController],
|
controllers: [PlatformApiController],
|
||||||
providers: [PlatformApiService, DadosferaLogger],
|
providers: [PlatformApiService, DadosferaLogger],
|
||||||
exports: [PlatformApiService],
|
exports: [PlatformApiService],
|
||||||
|
|||||||
@@ -89,12 +89,12 @@ export class PlatformApiService {
|
|||||||
|
|
||||||
// Propagate non-2xx responses as HttpExceptions
|
// Propagate non-2xx responses as HttpExceptions
|
||||||
if (response.status >= 400) {
|
if (response.status >= 400) {
|
||||||
this.logger.error('Platform API upstream error', {
|
this.logger.error('Platform API upstream error' + JSON.stringify({
|
||||||
status: response.status,
|
status: response.status,
|
||||||
data: response.data,
|
data: response.data,
|
||||||
path,
|
path,
|
||||||
method: method.toUpperCase(),
|
method: method.toUpperCase(),
|
||||||
});
|
}));
|
||||||
throw new HttpException(response.data, response.status);
|
throw new HttpException(response.data, response.status);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -107,6 +107,8 @@ export class PlatformApiService {
|
|||||||
method: method.toUpperCase(),
|
method: method.toUpperCase(),
|
||||||
});
|
});
|
||||||
|
|
||||||
|
this.logger.error(error)
|
||||||
|
|
||||||
if (error instanceof HttpException) {
|
if (error instanceof HttpException) {
|
||||||
throw error;
|
throw error;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -125,9 +125,13 @@ export class DynamoDBService {
|
|||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|
||||||
const { Item } = await this.documentClient.send(getCommand);
|
try {
|
||||||
|
const { Item } = await this.documentClient.send(getCommand);
|
||||||
return Item as InputDocument | null;
|
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> {
|
async deleteInput(clientId: string, inputId: string): Promise<void> {
|
||||||
@@ -203,7 +207,6 @@ export class DynamoDBService {
|
|||||||
updatedTable.reference_column = changes.reference_column;
|
updatedTable.reference_column = changes.reference_column;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
tables[tableIndex] = updatedTable;
|
tables[tableIndex] = updatedTable;
|
||||||
|
|
||||||
// Save updated document
|
// Save updated document
|
||||||
@@ -231,4 +234,5 @@ export class DynamoDBService {
|
|||||||
throw error;
|
throw error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 {}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user