mirror of
https://github.com/dadosfera/maestro.git
synced 2026-09-20 15:24:49 +00:00
Compare commits
53
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
00cbadbb45 | ||
|
|
47ad527d38 | ||
|
|
51044a23b3 | ||
|
|
bb29d126c1 | ||
|
|
9d0f449eeb | ||
|
|
0eafa67e6f | ||
|
|
c3937472ec | ||
|
|
84d64424ca | ||
|
|
f0bfc5c94b | ||
|
|
ad86a6a698 | ||
|
|
af3b11ad54 | ||
|
|
d54f381998 | ||
|
|
3fd586753e | ||
|
|
9b3894f9e4 | ||
|
|
4a5f8f679a | ||
|
|
a444f5e5ec | ||
|
|
0d30c1cf83 | ||
|
|
15860519e1 | ||
|
|
cffda86eab | ||
|
|
a93fbfbbd9 | ||
|
|
e1cfc6e1a8 | ||
|
|
1e63df6536 | ||
|
|
8eb7fd0169 | ||
|
|
f8c6a8b747 | ||
|
|
8a9d6c2f9c | ||
|
|
bcb2a0b7cb | ||
|
|
2e181e70af | ||
|
|
ef615adb8e | ||
|
|
e03900811b | ||
|
|
5bc5fb0977 | ||
|
|
011032e3d4 | ||
|
|
17363e74f4 | ||
|
|
63efff6adf | ||
|
|
b5f569e522 | ||
|
|
5c77577992 | ||
|
|
4b9e113185 | ||
|
|
06d505c50a | ||
|
|
c70abfc826 | ||
|
|
db2c7d6c02 | ||
|
|
0a99ce1aa4 | ||
|
|
39c66f030a | ||
|
|
1223e21ac4 | ||
|
|
2177f6725c | ||
|
|
e8b982998f | ||
|
|
ccd4159c59 | ||
|
|
9e49abb40d | ||
|
|
21a82b64f3 | ||
|
|
2414fcf21e | ||
|
|
94fbdb2226 | ||
|
|
985170d7ae | ||
|
|
b38b8f51e2 | ||
|
|
ea16e62d6b | ||
|
|
6182705410 |
@@ -71,6 +71,11 @@ jobs:
|
|||||||
sudo mv helmfile /usr/local/bin/
|
sudo mv helmfile /usr/local/bin/
|
||||||
helmfile --version
|
helmfile --version
|
||||||
|
|
||||||
|
- name: Install Helm Diff plugin
|
||||||
|
run: |
|
||||||
|
helm plugin install https://github.com/databus23/helm-diff --version v3.9.3
|
||||||
|
helm diff version
|
||||||
|
|
||||||
- name: Debug Helm env
|
- name: Debug Helm env
|
||||||
run: |
|
run: |
|
||||||
helm env
|
helm env
|
||||||
@@ -102,4 +107,5 @@ jobs:
|
|||||||
- name: Run Helmfile Diff
|
- name: Run Helmfile Diff
|
||||||
env:
|
env:
|
||||||
ENV: ${{ needs.extract_environment.outputs.environment }}
|
ENV: ${{ needs.extract_environment.outputs.environment }}
|
||||||
|
HELM_PLUGINS: /home/runner/.local/share/helm/plugins
|
||||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml diff
|
run: helmfile -f deploy/helmfiles/${ENV}.yaml diff
|
||||||
|
|||||||
+1
-1
@@ -1,5 +1,5 @@
|
|||||||
FROM node:20-alpine AS base_image
|
FROM node:20-alpine AS base_image
|
||||||
RUN npm install -g npm@latest
|
RUN npm install -g npm@10.8.2
|
||||||
|
|
||||||
FROM base_image AS build_base
|
FROM base_image AS build_base
|
||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
|
|||||||
@@ -111,6 +111,8 @@ spec:
|
|||||||
value: "{{ .Values.maestro.redis_tls }}"
|
value: "{{ .Values.maestro.redis_tls }}"
|
||||||
- name: PLATFORM_API_URL
|
- name: PLATFORM_API_URL
|
||||||
value: {{ .Values.maestro.platform_api_url }}
|
value: {{ .Values.maestro.platform_api_url }}
|
||||||
|
- name: CONNECTIONS_API_URL
|
||||||
|
value: {{ .Values.maestro.connections_api_url | default "" | quote }}
|
||||||
- name: STORAGE_EXPLORER_API_URL
|
- name: STORAGE_EXPLORER_API_URL
|
||||||
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
|
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
|
||||||
- name: FIREBASE_BASE_URL
|
- name: FIREBASE_BASE_URL
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ maestro:
|
|||||||
cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36"
|
cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36"
|
||||||
redis_database: "1"
|
redis_database: "1"
|
||||||
platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com
|
platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com
|
||||||
|
connections_api_url: https://iy40eans64.execute-api.us-east-1.amazonaws.com
|
||||||
storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local:8000/api"
|
storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local:8000/api"
|
||||||
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/stg
|
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/stg
|
||||||
|
|
||||||
|
|||||||
+253
-90
@@ -3343,14 +3343,7 @@
|
|||||||
],
|
],
|
||||||
"responses": {
|
"responses": {
|
||||||
"200": {
|
"200": {
|
||||||
"description": "",
|
"description": ""
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"type": "object"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"tags": [
|
"tags": [
|
||||||
@@ -3647,6 +3640,46 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/pipelinesV2/{id}/upgrade": {
|
||||||
|
"patch": {
|
||||||
|
"operationId": "PipelinesController_upgradeConnector",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"name": "dadosfera-lang",
|
||||||
|
"in": "header",
|
||||||
|
"required": false,
|
||||||
|
"schema": {
|
||||||
|
"enum": [
|
||||||
|
"pt-br",
|
||||||
|
"en-us"
|
||||||
|
],
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "id",
|
||||||
|
"required": true,
|
||||||
|
"in": "path",
|
||||||
|
"schema": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"204": {
|
||||||
|
"description": ""
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"tags": [
|
||||||
|
"PipelinesV2"
|
||||||
|
],
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
"/pipelinesV2/init-upload": {
|
"/pipelinesV2/init-upload": {
|
||||||
"post": {
|
"post": {
|
||||||
"operationId": "PipelinesController_initUploadFile",
|
"operationId": "PipelinesController_initUploadFile",
|
||||||
@@ -3834,88 +3867,6 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"/pipelines/start/{id}": {
|
|
||||||
"post": {
|
|
||||||
"operationId": "PipelinesController_activate",
|
|
||||||
"summary": "",
|
|
||||||
"deprecated": true,
|
|
||||||
"description": "This method is deprecated. Please use route /pipelinesV2/start/:id instead",
|
|
||||||
"parameters": [
|
|
||||||
{
|
|
||||||
"name": "id",
|
|
||||||
"required": true,
|
|
||||||
"in": "path",
|
|
||||||
"schema": {
|
|
||||||
"type": "string"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"responses": {
|
|
||||||
"201": {
|
|
||||||
"description": "",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"type": "object"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"tags": [
|
|
||||||
"Pipelines"
|
|
||||||
],
|
|
||||||
"security": [
|
|
||||||
{
|
|
||||||
"access-token": []
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"access-token": []
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"/pipelines/{id}/status": {
|
|
||||||
"get": {
|
|
||||||
"operationId": "PipelinesController_getPipelineStatus",
|
|
||||||
"summary": "",
|
|
||||||
"deprecated": true,
|
|
||||||
"description": "This method is deprecated. Please use route /pipelinesV2/:id/status instead",
|
|
||||||
"parameters": [
|
|
||||||
{
|
|
||||||
"name": "id",
|
|
||||||
"required": true,
|
|
||||||
"in": "path",
|
|
||||||
"schema": {
|
|
||||||
"type": "string"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"responses": {
|
|
||||||
"200": {
|
|
||||||
"description": "",
|
|
||||||
"content": {
|
|
||||||
"application/json": {
|
|
||||||
"schema": {
|
|
||||||
"type": "object"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"tags": [
|
|
||||||
"Pipelines"
|
|
||||||
],
|
|
||||||
"security": [
|
|
||||||
{
|
|
||||||
"access-token": []
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"access-token": []
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"/transformations": {
|
"/transformations": {
|
||||||
"post": {
|
"post": {
|
||||||
"operationId": "TransformationsController_create",
|
"operationId": "TransformationsController_create",
|
||||||
@@ -4619,6 +4570,62 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/platform/pipelines/{pipelineId}/pipeline_run/{runId}/jobs": {
|
||||||
|
"get": {
|
||||||
|
"operationId": "PlatformApiController_getPipelineRunJobs",
|
||||||
|
"summary": "Get pipeline run jobs",
|
||||||
|
"description": "Proxies platform-api DB-backed job runs and returns `{ jobs: [...] }`.",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"name": "pipelineId",
|
||||||
|
"required": true,
|
||||||
|
"in": "path",
|
||||||
|
"schema": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "runId",
|
||||||
|
"required": true,
|
||||||
|
"in": "path",
|
||||||
|
"schema": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"200": {
|
||||||
|
"description": "DB-backed job runs for the selected pipeline run.",
|
||||||
|
"content": {
|
||||||
|
"application/json": {
|
||||||
|
"schema": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"jobs": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"type": "object"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"required": [
|
||||||
|
"jobs"
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"tags": [
|
||||||
|
"Platform API"
|
||||||
|
],
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
"/platform/jobs/{jobId}/input": {
|
"/platform/jobs/{jobId}/input": {
|
||||||
"put": {
|
"put": {
|
||||||
"operationId": "PlatformApiController_updateJobInput",
|
"operationId": "PlatformApiController_updateJobInput",
|
||||||
@@ -5531,6 +5538,48 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/catalog/custom-properties": {
|
||||||
|
"get": {
|
||||||
|
"operationId": "CatalogController_getCustomPropertyDefinitions",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"name": "dadosfera-lang",
|
||||||
|
"in": "header",
|
||||||
|
"required": false,
|
||||||
|
"schema": {
|
||||||
|
"enum": [
|
||||||
|
"pt-br",
|
||||||
|
"en-us"
|
||||||
|
],
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"responses": {
|
||||||
|
"200": {
|
||||||
|
"description": "",
|
||||||
|
"content": {
|
||||||
|
"application/json": {
|
||||||
|
"schema": {
|
||||||
|
"type": "object"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"tags": [
|
||||||
|
"Catalog"
|
||||||
|
],
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
"/catalog/data-asset/{id}": {
|
"/catalog/data-asset/{id}": {
|
||||||
"get": {
|
"get": {
|
||||||
"operationId": "CatalogController_getDataAsset",
|
"operationId": "CatalogController_getDataAsset",
|
||||||
@@ -5945,6 +5994,66 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
"/catalog/data-asset/{id}/certification-status": {
|
||||||
|
"put": {
|
||||||
|
"operationId": "CatalogController_updateDataAssetCertificationStatus",
|
||||||
|
"parameters": [
|
||||||
|
{
|
||||||
|
"name": "dadosfera-lang",
|
||||||
|
"in": "header",
|
||||||
|
"required": false,
|
||||||
|
"schema": {
|
||||||
|
"enum": [
|
||||||
|
"pt-br",
|
||||||
|
"en-us"
|
||||||
|
],
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "id",
|
||||||
|
"required": true,
|
||||||
|
"in": "path",
|
||||||
|
"schema": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"requestBody": {
|
||||||
|
"required": true,
|
||||||
|
"content": {
|
||||||
|
"application/json": {
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/components/schemas/IUpdateCertificationStatusRequest"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"responses": {
|
||||||
|
"200": {
|
||||||
|
"description": "",
|
||||||
|
"content": {
|
||||||
|
"application/json": {
|
||||||
|
"schema": {
|
||||||
|
"$ref": "#/components/schemas/IUpdateCertificationStatusRequest"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"tags": [
|
||||||
|
"Catalog"
|
||||||
|
],
|
||||||
|
"security": [
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"access-token": []
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
"/catalog/data-asset/{id}/manage-permissions": {
|
"/catalog/data-asset/{id}/manage-permissions": {
|
||||||
"put": {
|
"put": {
|
||||||
"operationId": "CatalogController_manageDataAssetPermissions",
|
"operationId": "CatalogController_manageDataAssetPermissions",
|
||||||
@@ -10978,6 +11087,37 @@
|
|||||||
"docs"
|
"docs"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
"CustomPropertyDto": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"key": {
|
||||||
|
"type": "string"
|
||||||
|
},
|
||||||
|
"value": {
|
||||||
|
"type": "string"
|
||||||
|
},
|
||||||
|
"type": {
|
||||||
|
"type": "string",
|
||||||
|
"enum": [
|
||||||
|
"text",
|
||||||
|
"number",
|
||||||
|
"date",
|
||||||
|
"boolean"
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"color": {
|
||||||
|
"type": "string"
|
||||||
|
},
|
||||||
|
"emoji": {
|
||||||
|
"type": "string"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"required": [
|
||||||
|
"key",
|
||||||
|
"value",
|
||||||
|
"type"
|
||||||
|
]
|
||||||
|
},
|
||||||
"IUpdateDataRequest": {
|
"IUpdateDataRequest": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
@@ -11006,6 +11146,12 @@
|
|||||||
},
|
},
|
||||||
"docs": {
|
"docs": {
|
||||||
"type": "string"
|
"type": "string"
|
||||||
|
},
|
||||||
|
"custom_properties": {
|
||||||
|
"type": "array",
|
||||||
|
"items": {
|
||||||
|
"$ref": "#/components/schemas/CustomPropertyDto"
|
||||||
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"required": [
|
"required": [
|
||||||
@@ -11025,6 +11171,23 @@
|
|||||||
"data_asset"
|
"data_asset"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
"IUpdateCertificationStatusRequest": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"certification_status": {
|
||||||
|
"type": "string",
|
||||||
|
"enum": [
|
||||||
|
"draft",
|
||||||
|
"in_review",
|
||||||
|
"approved",
|
||||||
|
"deprecated"
|
||||||
|
]
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"required": [
|
||||||
|
"certification_status"
|
||||||
|
]
|
||||||
|
},
|
||||||
"ICreateDataAsset": {
|
"ICreateDataAsset": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": {
|
"properties": {
|
||||||
|
|||||||
Generated
+4
-15
@@ -16,8 +16,7 @@
|
|||||||
"@aws-sdk/lib-dynamodb": "^3.414.0",
|
"@aws-sdk/lib-dynamodb": "^3.414.0",
|
||||||
"@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-v2": "^3.40.0-beta.14",
|
||||||
"@dadosfera/protospack-v2": "3.40.0-beta.8",
|
|
||||||
"@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",
|
||||||
@@ -1735,20 +1734,10 @@
|
|||||||
"winston-log2gelf": "^2.4.0"
|
"winston-log2gelf": "^2.4.0"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/@dadosfera/protospack": {
|
|
||||||
"version": "2.5.3",
|
|
||||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack/-/protospack-2.5.3.tgz",
|
|
||||||
"integrity": "sha512-yOLnd+s6n9VkPpZXO8HnUY27CQPHj/qs+ecddviA4Ldn0Gx4KGRgbVdsSSP45nPm0GHhCd2bHg4ap+la7xtRmA==",
|
|
||||||
"license": "ISC",
|
|
||||||
"dependencies": {
|
|
||||||
"rxjs": "^7.5.5"
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"node_modules/@dadosfera/protospack-v2": {
|
"node_modules/@dadosfera/protospack-v2": {
|
||||||
"version": "3.40.0-beta.8",
|
"version": "3.40.0-beta.14",
|
||||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.8.tgz",
|
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.14.tgz",
|
||||||
"integrity": "sha512-JE5qMjqB3UOM+tCUxB1EwYLQW0PecsaQIa1KDpKEaG3lrzG/H13z8iJi3WH/DuVav2EI94i9VcJWJ1Y0F7ribw==",
|
"integrity": "sha512-pv3pxq0x1XcBgf3ajD6QOFRLOduh8iEozKFA3AKlIW4gid+gT4iL0GcU2M+O7h0QFeO4JIzRZe/nEMN82nqk7A==",
|
||||||
"license": "ISC",
|
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@grpc/grpc-js": "^1.9.3",
|
"@grpc/grpc-js": "^1.9.3",
|
||||||
"rxjs": "^7.5.5"
|
"rxjs": "^7.5.5"
|
||||||
|
|||||||
+1
-2
@@ -34,8 +34,7 @@
|
|||||||
"@aws-sdk/lib-dynamodb": "^3.414.0",
|
"@aws-sdk/lib-dynamodb": "^3.414.0",
|
||||||
"@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-v2": "^3.40.0-beta.14",
|
||||||
"@dadosfera/protospack-v2": "3.40.0-beta.8",
|
|
||||||
"@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",
|
||||||
|
|||||||
@@ -17,7 +17,6 @@ import { ConnectionTestModule } from './modules/connection-test/connection-test.
|
|||||||
import { NetworkConfigModule } from './modules/network-config/network-config.module';
|
import { NetworkConfigModule } from './modules/network-config/network-config.module';
|
||||||
import { InputsModule } from './modules/inputs/inputs.module';
|
import { InputsModule } from './modules/inputs/inputs.module';
|
||||||
import { OauthModule } from './modules/oauth/oauth.module';
|
import { OauthModule } from './modules/oauth/oauth.module';
|
||||||
import { PipelinesModule } from './modules/pipelines/pipelines.module';
|
|
||||||
import { TransformationsModule } from './modules/transformations/transformations.module';
|
import { TransformationsModule } from './modules/transformations/transformations.module';
|
||||||
import { HealthModule } from './modules/health/health.module';
|
import { HealthModule } from './modules/health/health.module';
|
||||||
import { CatalogModule } from './modules/catalog/catalog.module';
|
import { CatalogModule } from './modules/catalog/catalog.module';
|
||||||
@@ -60,7 +59,6 @@ import { ReleaseNoteModule } from './modules/release_note/release_note.module';
|
|||||||
PermissionsModule,
|
PermissionsModule,
|
||||||
TermsOfUseModule,
|
TermsOfUseModule,
|
||||||
ConnectionTestModule,
|
ConnectionTestModule,
|
||||||
PipelinesModule,
|
|
||||||
TransformationsModule,
|
TransformationsModule,
|
||||||
UsersModule,
|
UsersModule,
|
||||||
RolesModule,
|
RolesModule,
|
||||||
|
|||||||
@@ -357,6 +357,16 @@ export const PERMISSIONS_GROUPS = {
|
|||||||
'es-es': 'Crear y editar atributos en el catálogo',
|
'es-es': 'Crear y editar atributos en el catálogo',
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
CERTIFY: {
|
||||||
|
seqid: 53,
|
||||||
|
claim: 'catalog:certify',
|
||||||
|
usage: PermissionUsages.PUBLIC,
|
||||||
|
name: {
|
||||||
|
'pt-br': 'Alterar o status de certificação dos Ativos',
|
||||||
|
'en-us': "Change Assets' certification status",
|
||||||
|
'es-es': 'Cambiar el estado de certificación de los Activos',
|
||||||
|
},
|
||||||
|
},
|
||||||
DELETE: {
|
DELETE: {
|
||||||
seqid: 1,
|
seqid: 1,
|
||||||
claim: 'catalog:delete',
|
claim: 'catalog:delete',
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ import {
|
|||||||
HttpStatus,
|
HttpStatus,
|
||||||
Res,
|
Res,
|
||||||
} from '@nestjs/common';
|
} from '@nestjs/common';
|
||||||
|
import { ValidationPipe } from '../../pipes/object-validation.pipe';
|
||||||
import {
|
import {
|
||||||
ApiCreatedResponse,
|
ApiCreatedResponse,
|
||||||
ApiHeaders,
|
ApiHeaders,
|
||||||
@@ -46,6 +47,7 @@ import {
|
|||||||
IMakeAComment,
|
IMakeAComment,
|
||||||
IOneDataAsset,
|
IOneDataAsset,
|
||||||
IPreviewResponse,
|
IPreviewResponse,
|
||||||
|
IUpdateCertificationStatusRequest,
|
||||||
IUpdateDataRequest,
|
IUpdateDataRequest,
|
||||||
TriggerCatalogReq,
|
TriggerCatalogReq,
|
||||||
TriggerCatalogRes,
|
TriggerCatalogRes,
|
||||||
@@ -269,6 +271,23 @@ export class CatalogController {
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Get('custom-properties')
|
||||||
|
@RequireSomePermission(
|
||||||
|
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||||
|
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||||
|
)
|
||||||
|
async getCustomPropertyDefinitions(@User() user: RequestUser) {
|
||||||
|
const { customer_id, customer_name, user_id, username } = user;
|
||||||
|
const metadata = PackTheMetadata({
|
||||||
|
customer_id,
|
||||||
|
customer_name,
|
||||||
|
user_id,
|
||||||
|
username,
|
||||||
|
});
|
||||||
|
|
||||||
|
return this.catalogService.getCustomPropertyDefinitions(metadata);
|
||||||
|
}
|
||||||
|
|
||||||
@Get('data-asset/:id')
|
@Get('data-asset/:id')
|
||||||
@RequireSomePermission(
|
@RequireSomePermission(
|
||||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||||
@@ -493,6 +512,8 @@ export class CatalogController {
|
|||||||
language,
|
language,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
delete (body as any).certification_status;
|
||||||
|
|
||||||
const result = await this.catalogService.updateOneDataAsset({
|
const result = await this.catalogService.updateOneDataAsset({
|
||||||
body,
|
body,
|
||||||
data_asset_id,
|
data_asset_id,
|
||||||
@@ -506,6 +527,33 @@ export class CatalogController {
|
|||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Put('data-asset/:id/certification-status')
|
||||||
|
@RequireSomePermission(
|
||||||
|
PERMISSIONS_GROUPS.CATALOG.permissions.CERTIFY,
|
||||||
|
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||||
|
)
|
||||||
|
async updateDataAssetCertificationStatus(
|
||||||
|
@User() user: RequestUser,
|
||||||
|
@Language() language: LanguageEnum,
|
||||||
|
@Param('id') data_asset_id: string,
|
||||||
|
@Body(new ValidationPipe()) body: IUpdateCertificationStatusRequest,
|
||||||
|
): Promise<IUpdateCertificationStatusRequest> {
|
||||||
|
const { customer_id, customer_name, user_id, username } = user;
|
||||||
|
const metadata = PackTheMetadata({
|
||||||
|
customer_id,
|
||||||
|
customer_name,
|
||||||
|
user_id,
|
||||||
|
username,
|
||||||
|
language,
|
||||||
|
});
|
||||||
|
|
||||||
|
return this.catalogService.updateCertificationStatus({
|
||||||
|
body,
|
||||||
|
data_asset_id,
|
||||||
|
metadata,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
@Post('data-asset/:id/docs')
|
@Post('data-asset/:id/docs')
|
||||||
@RequireSomePermission(
|
@RequireSomePermission(
|
||||||
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
||||||
|
|||||||
@@ -5,20 +5,17 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|||||||
import { CatalogController } from './catalog.controller';
|
import { CatalogController } from './catalog.controller';
|
||||||
import { CatalogClientConfiguration } from './catalog-client';
|
import { CatalogClientConfiguration } from './catalog-client';
|
||||||
import { ClientsModule } from '@nestjs/microservices';
|
import { ClientsModule } from '@nestjs/microservices';
|
||||||
import { PipelinesModule as OldPipelineModule } from 'src/modules/pipelines/pipelines.module';
|
|
||||||
import { UsersModule } from '../users/users.module';
|
import { UsersModule } from '../users/users.module';
|
||||||
import { RolesModule } from '../roles/roles.module';
|
import { RolesModule } from '../roles/roles.module';
|
||||||
import { CustomersModule } from '../customers/customers.module';
|
import { CustomersModule } from '../customers/customers.module';
|
||||||
import { ShareModule } from './share/share.module';
|
import { ShareModule } from './share/share.module';
|
||||||
import { CatalogService } from './catalog.service';
|
import { CatalogService } from './catalog.service';
|
||||||
import { MixpanelModule } from '../mixpanel/mixpanel.module';
|
|
||||||
|
|
||||||
const client = new CatalogClientConfiguration();
|
const client = new CatalogClientConfiguration();
|
||||||
|
|
||||||
@Module({
|
@Module({
|
||||||
imports: [
|
imports: [
|
||||||
ClientsModule.register([client.providerOptions]),
|
ClientsModule.register([client.providerOptions]),
|
||||||
OldPipelineModule,
|
|
||||||
UsersModule,
|
UsersModule,
|
||||||
RolesModule,
|
RolesModule,
|
||||||
CustomersModule,
|
CustomersModule,
|
||||||
|
|||||||
@@ -29,6 +29,7 @@ import {
|
|||||||
AssetReporter,
|
AssetReporter,
|
||||||
BatchRemoveRlsRulesRequest,
|
BatchRemoveRlsRulesRequest,
|
||||||
CreateDataDocsDTO,
|
CreateDataDocsDTO,
|
||||||
|
IUpdateCertificationStatusRequest,
|
||||||
IUpdateDataRequest,
|
IUpdateDataRequest,
|
||||||
TriggerCatalogReq,
|
TriggerCatalogReq,
|
||||||
} from './dtos';
|
} from './dtos';
|
||||||
@@ -119,6 +120,10 @@ class CatalogService implements OnModuleInit {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async getCustomPropertyDefinitions(metadata: Metadata) {
|
||||||
|
return lastValueFrom(this.catalogReadService.GetCustomPropertyDefinitions({}, metadata));
|
||||||
|
}
|
||||||
|
|
||||||
async createDataAsset(data: Messages.CreateDataAssetRequest, metadata) {
|
async createDataAsset(data: Messages.CreateDataAssetRequest, metadata) {
|
||||||
this.logger.info('CatalogService - Manage Data assets permissions');
|
this.logger.info('CatalogService - Manage Data assets permissions');
|
||||||
if (!data.embed) data.embed = undefined;
|
if (!data.embed) data.embed = undefined;
|
||||||
@@ -384,6 +389,28 @@ class CatalogService implements OnModuleInit {
|
|||||||
return { data_asset: asset[0] };
|
return { data_asset: asset[0] };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async updateCertificationStatus(data: {
|
||||||
|
data_asset_id: string;
|
||||||
|
body: IUpdateCertificationStatusRequest;
|
||||||
|
metadata: Metadata;
|
||||||
|
}) {
|
||||||
|
const { body, data_asset_id, metadata } = data;
|
||||||
|
|
||||||
|
await lastValueFrom(
|
||||||
|
this.catalogWriteService.UpdateDataAsset(
|
||||||
|
{
|
||||||
|
id: data_asset_id,
|
||||||
|
changes: JSON.stringify({
|
||||||
|
certification_status: body.certification_status,
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
metadata,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
|
||||||
|
return { certification_status: body.certification_status };
|
||||||
|
}
|
||||||
|
|
||||||
async updateOneDataAsset(data: {
|
async updateOneDataAsset(data: {
|
||||||
data_asset_id: string;
|
data_asset_id: string;
|
||||||
customer_id: string;
|
customer_id: string;
|
||||||
|
|||||||
@@ -1,4 +1,10 @@
|
|||||||
import { ApiProperty, ApiPropertyOptional, PickType } from '@nestjs/swagger';
|
import { ApiProperty, ApiPropertyOptional, PickType } from '@nestjs/swagger';
|
||||||
|
import {
|
||||||
|
IsEnum,
|
||||||
|
IsNotEmpty,
|
||||||
|
IsOptional,
|
||||||
|
IsString,
|
||||||
|
} from 'class-validator';
|
||||||
import { CreateDataAssetRequest } from '@dadosfera/protospack-v2/dist/lib/Catalog/interfaces/messages';
|
import { CreateDataAssetRequest } from '@dadosfera/protospack-v2/dist/lib/Catalog/interfaces/messages';
|
||||||
|
|
||||||
export enum DataAssetShareType {
|
export enum DataAssetShareType {
|
||||||
@@ -6,6 +12,12 @@ export enum DataAssetShareType {
|
|||||||
public = 'public',
|
public = 'public',
|
||||||
private = 'private',
|
private = 'private',
|
||||||
}
|
}
|
||||||
|
export enum CertificationStatus {
|
||||||
|
draft = 'draft',
|
||||||
|
in_review = 'in_review',
|
||||||
|
approved = 'approved',
|
||||||
|
deprecated = 'deprecated',
|
||||||
|
}
|
||||||
export enum OrderEnum {
|
export enum OrderEnum {
|
||||||
asc = 'asc',
|
asc = 'asc',
|
||||||
desc = 'desc',
|
desc = 'desc',
|
||||||
@@ -191,6 +203,27 @@ export class IData {
|
|||||||
day_opening: number;
|
day_opening: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
export enum CustomPropertyType {
|
||||||
|
TEXT = 'text',
|
||||||
|
NUMBER = 'number',
|
||||||
|
DATE = 'date',
|
||||||
|
BOOLEAN = 'boolean',
|
||||||
|
}
|
||||||
|
|
||||||
|
export class CustomPropertyDto {
|
||||||
|
@ApiProperty()
|
||||||
|
key: string;
|
||||||
|
@ApiProperty()
|
||||||
|
value: string;
|
||||||
|
@ApiProperty({ enum: CustomPropertyType })
|
||||||
|
type: CustomPropertyType;
|
||||||
|
@ApiPropertyOptional()
|
||||||
|
color?: string;
|
||||||
|
@ApiPropertyOptional()
|
||||||
|
emoji?: string;
|
||||||
|
}
|
||||||
|
|
||||||
export class IUpdateDataRequest {
|
export class IUpdateDataRequest {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
name: string;
|
name: string;
|
||||||
@@ -204,7 +237,16 @@ export class IUpdateDataRequest {
|
|||||||
share_type?: DataAssetShareType;
|
share_type?: DataAssetShareType;
|
||||||
@ApiPropertyOptional()
|
@ApiPropertyOptional()
|
||||||
docs?: string;
|
docs?: string;
|
||||||
|
@ApiPropertyOptional({ type: [CustomPropertyDto] })
|
||||||
|
custom_properties?: CustomPropertyDto[];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export class IUpdateCertificationStatusRequest {
|
||||||
|
@ApiProperty({ enum: CertificationStatus })
|
||||||
|
@IsEnum(CertificationStatus)
|
||||||
|
certification_status: CertificationStatus;
|
||||||
|
}
|
||||||
|
|
||||||
export class ICreateDataAsset implements CreateDataAssetRequest {
|
export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
display_name: string;
|
display_name: string;
|
||||||
|
|||||||
@@ -22,6 +22,9 @@ import {
|
|||||||
ConnectionTestListTablesRes,
|
ConnectionTestListTablesRes,
|
||||||
GetTableMetadataRes,
|
GetTableMetadataRes,
|
||||||
GetTableMetadataReq,
|
GetTableMetadataReq,
|
||||||
|
RefreshCatalogReq,
|
||||||
|
RefreshCatalogRes,
|
||||||
|
RefreshCatalogStatusReq,
|
||||||
} from './dto/connection-test';
|
} from './dto/connection-test';
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
import { Authenticated } from 'src/decorators/authentication.decorator';
|
import { Authenticated } from 'src/decorators/authentication.decorator';
|
||||||
@@ -84,7 +87,7 @@ export class ConnectionTestController {
|
|||||||
});
|
});
|
||||||
return this.connectionTestService.connectionTestListSchemas(
|
return this.connectionTestService.connectionTestListSchemas(
|
||||||
body,
|
body,
|
||||||
user.customer_name,
|
user,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -101,7 +104,7 @@ export class ConnectionTestController {
|
|||||||
});
|
});
|
||||||
return this.connectionTestService.connectionTestListTables(
|
return this.connectionTestService.connectionTestListTables(
|
||||||
body,
|
body,
|
||||||
user.customer_name,
|
user,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -118,7 +121,38 @@ export class ConnectionTestController {
|
|||||||
});
|
});
|
||||||
return this.connectionTestService.getTableMetadata(
|
return this.connectionTestService.getTableMetadata(
|
||||||
body,
|
body,
|
||||||
user.customer_name,
|
user,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Post('refresh-catalog')
|
||||||
|
@ApiOkResponse({ type: RefreshCatalogRes })
|
||||||
|
@HttpCode(HttpStatus.ACCEPTED)
|
||||||
|
async refreshCatalog(
|
||||||
|
@User() user: RequestUser,
|
||||||
|
@Body(new ValidationPipe()) body: RefreshCatalogReq,
|
||||||
|
) {
|
||||||
|
this.logger.info('/connection-test/refresh-catalog', {
|
||||||
|
user: user.user_id,
|
||||||
|
customer: user.customer_name,
|
||||||
|
connection: body.connection_id,
|
||||||
|
});
|
||||||
|
return this.connectionTestService.refreshCatalog(body, user);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Post('refresh-catalog/status')
|
||||||
|
@ApiOkResponse({ type: RefreshCatalogRes })
|
||||||
|
@HttpCode(HttpStatus.OK)
|
||||||
|
async refreshCatalogStatus(
|
||||||
|
@User() user: RequestUser,
|
||||||
|
@Body(new ValidationPipe()) body: RefreshCatalogStatusReq,
|
||||||
|
) {
|
||||||
|
this.logger.info('/connection-test/refresh-catalog/status', {
|
||||||
|
user: user.user_id,
|
||||||
|
customer: user.customer_name,
|
||||||
|
connection: body.connection_id,
|
||||||
|
session: body.session_id,
|
||||||
|
});
|
||||||
|
return this.connectionTestService.refreshCatalogStatus(body, user);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,10 +5,17 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|||||||
import { ClientsModule } from '@nestjs/microservices';
|
import { ClientsModule } from '@nestjs/microservices';
|
||||||
import { ConnectionTestClientConfiguration } from './connection-test-client.config';
|
import { ConnectionTestClientConfiguration } from './connection-test-client.config';
|
||||||
import { ConnectionModule } from '../connection/connection.module';
|
import { ConnectionModule } from '../connection/connection.module';
|
||||||
|
import { ConnectionsApiModule } from '../connections-api/connections-api.module';
|
||||||
|
import { PlatformApiModule } from '../platform-api/platform-api.module';
|
||||||
const client = new ConnectionTestClientConfiguration();
|
const client = new ConnectionTestClientConfiguration();
|
||||||
@Module({
|
@Module({
|
||||||
controllers: [ConnectionTestController],
|
controllers: [ConnectionTestController],
|
||||||
providers: [ConnectionTestService, DadosferaLogger],
|
providers: [ConnectionTestService, DadosferaLogger],
|
||||||
imports: [ClientsModule.register([client.providerOptions]), ConnectionModule],
|
imports: [
|
||||||
|
ClientsModule.register([client.providerOptions]),
|
||||||
|
ConnectionModule,
|
||||||
|
ConnectionsApiModule,
|
||||||
|
PlatformApiModule,
|
||||||
|
],
|
||||||
})
|
})
|
||||||
export class ConnectionTestModule {}
|
export class ConnectionTestModule {}
|
||||||
|
|||||||
@@ -0,0 +1,205 @@
|
|||||||
|
import { ConnectionTestService } from './connection-test.service';
|
||||||
|
import { RequestUser } from 'src/decorators/user.decorator';
|
||||||
|
|
||||||
|
describe('ConnectionTestService catalog cache', () => {
|
||||||
|
const user: RequestUser = {
|
||||||
|
user_id: 'user-id',
|
||||||
|
username: 'user@example.com',
|
||||||
|
permissions: [],
|
||||||
|
customer_id: 'customer-id',
|
||||||
|
customer_name: 'customer-name',
|
||||||
|
customer_tier: 'standard',
|
||||||
|
access_token: 'token',
|
||||||
|
customer_modules: [],
|
||||||
|
roles: [],
|
||||||
|
};
|
||||||
|
const grpcClient = { getService: jest.fn().mockReturnValue({}) };
|
||||||
|
const connectionsService = {};
|
||||||
|
const connectionsApiService = { proxy: jest.fn() };
|
||||||
|
const platformApiService = { proxy: jest.fn() };
|
||||||
|
let service: ConnectionTestService;
|
||||||
|
|
||||||
|
beforeEach(() => {
|
||||||
|
jest.clearAllMocks();
|
||||||
|
service = new ConnectionTestService(
|
||||||
|
grpcClient as any,
|
||||||
|
connectionsService as any,
|
||||||
|
connectionsApiService as any,
|
||||||
|
platformApiService as any,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('keeps the existing schemas response contract', async () => {
|
||||||
|
connectionsApiService.proxy.mockResolvedValue({
|
||||||
|
schemas: [{ schema_name: 'analytics' }, { schema_name: 'public' }],
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.connectionTestListSchemas(
|
||||||
|
{ connection_id: 'config-id', plugin: 'postgresql' },
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: true,
|
||||||
|
schema_list: ['analytics', 'public'],
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it('keeps the existing tables response contract', async () => {
|
||||||
|
connectionsApiService.proxy.mockResolvedValue({
|
||||||
|
tables: [{ table_name: 'customers' }, { table_name: 'orders' }],
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.connectionTestListTables(
|
||||||
|
{
|
||||||
|
connection_id: 'config-id',
|
||||||
|
plugin: 'postgresql',
|
||||||
|
schema: 'public',
|
||||||
|
},
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: true,
|
||||||
|
table_list: ['customers', 'orders'],
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it('maps cached columns to the existing table metadata contract', async () => {
|
||||||
|
connectionsApiService.proxy.mockResolvedValue({
|
||||||
|
columns: [
|
||||||
|
{
|
||||||
|
column_name: 'id',
|
||||||
|
data_type: 'bigint',
|
||||||
|
is_primary_key: true,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.getTableMetadata(
|
||||||
|
{
|
||||||
|
connection_id: 'config-id',
|
||||||
|
plugin: 'postgresql',
|
||||||
|
schema: 'public',
|
||||||
|
table_list: ['customers'],
|
||||||
|
},
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: true,
|
||||||
|
tables_metadata: [
|
||||||
|
{
|
||||||
|
table_name: 'customers',
|
||||||
|
columns: [
|
||||||
|
{
|
||||||
|
name: 'id',
|
||||||
|
type: 'bigint',
|
||||||
|
is_primary_key: true,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
references: [],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
});
|
||||||
|
expect(connectionsApiService.proxy).toHaveBeenCalledWith(
|
||||||
|
'GET',
|
||||||
|
'/connection_catalog/config-id/schemas/public/tables/customers/columns',
|
||||||
|
user,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('submits a catalog refresh without holding the request open', async () => {
|
||||||
|
platformApiService.proxy.mockResolvedValue({
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.refreshCatalog(
|
||||||
|
{ connection_id: 'config-id', plugin: 'postgresql' },
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: true,
|
||||||
|
status: 'PENDING',
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(platformApiService.proxy).toHaveBeenCalledWith(
|
||||||
|
'POST',
|
||||||
|
'/connection_test',
|
||||||
|
user,
|
||||||
|
{
|
||||||
|
customer_id: user.customer_name,
|
||||||
|
plugin: 'postgresql',
|
||||||
|
task: {
|
||||||
|
task_type: 'refresh_catalog',
|
||||||
|
connection: {
|
||||||
|
provider: 'connection_manager',
|
||||||
|
config_id: 'config-id',
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('keeps polling without changing the catalog pointer while pending', async () => {
|
||||||
|
platformApiService.proxy.mockResolvedValue({ status: 'PENDING' });
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.refreshCatalogStatus(
|
||||||
|
{
|
||||||
|
connection_id: 'config-id',
|
||||||
|
plugin: 'postgresql',
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
},
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: false,
|
||||||
|
status: 'PENDING',
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(connectionsApiService.proxy).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('publishes the catalog pointer after the refresh finishes', async () => {
|
||||||
|
platformApiService.proxy.mockResolvedValue({ status: 'DONE' });
|
||||||
|
connectionsApiService.proxy.mockResolvedValue({
|
||||||
|
last_catalog_refresh_status: 'SUCCESS',
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(
|
||||||
|
service.refreshCatalogStatus(
|
||||||
|
{
|
||||||
|
connection_id: 'config/id',
|
||||||
|
plugin: 'postgresql',
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
},
|
||||||
|
user,
|
||||||
|
),
|
||||||
|
).resolves.toEqual({
|
||||||
|
operation_result: true,
|
||||||
|
status: 'DONE',
|
||||||
|
session_id: 'session-id',
|
||||||
|
date: '20260731',
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(connectionsApiService.proxy).toHaveBeenCalledWith(
|
||||||
|
'PUT',
|
||||||
|
'/connection_config/config%2Fid/catalog_metadata',
|
||||||
|
user,
|
||||||
|
{
|
||||||
|
last_catalog_refresh_status: 'SUCCESS',
|
||||||
|
last_catalog_connection_test_date: '20260731',
|
||||||
|
last_catalog_connection_test_session_id: 'session-id',
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
import { Inject, Injectable } from '@nestjs/common';
|
import { HttpException, HttpStatus, Inject, Injectable } from '@nestjs/common';
|
||||||
import { ClientGrpc } from '@nestjs/microservices';
|
import { ClientGrpc } from '@nestjs/microservices';
|
||||||
import { ConnectionTest } from '@dadosfera/protospack-v2';
|
import { ConnectionTest } from '@dadosfera/protospack-v2';
|
||||||
import { lastValueFrom } from 'rxjs';
|
import { lastValueFrom } from 'rxjs';
|
||||||
@@ -13,6 +13,9 @@ import {
|
|||||||
ConnectionTestPingRes,
|
ConnectionTestPingRes,
|
||||||
GetTableMetadataReq,
|
GetTableMetadataReq,
|
||||||
GetTableMetadataRes,
|
GetTableMetadataRes,
|
||||||
|
RefreshCatalogReq,
|
||||||
|
RefreshCatalogRes,
|
||||||
|
RefreshCatalogStatusReq,
|
||||||
} from './dto/connection-test';
|
} from './dto/connection-test';
|
||||||
import { ConnectionClientService } from '../connection/client.service';
|
import { ConnectionClientService } from '../connection/client.service';
|
||||||
import {
|
import {
|
||||||
@@ -21,6 +24,8 @@ import {
|
|||||||
} from '../connection/dtos/connection';
|
} from '../connection/dtos/connection';
|
||||||
import { RequestUser } from 'src/decorators/user.decorator';
|
import { RequestUser } from 'src/decorators/user.decorator';
|
||||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||||
|
import { ConnectionsApiService } from '../connections-api/connections-api.service';
|
||||||
|
import { PlatformApiService } from '../platform-api/platform-api.service';
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class ConnectionTestService {
|
export class ConnectionTestService {
|
||||||
@@ -28,6 +33,8 @@ export class ConnectionTestService {
|
|||||||
constructor(
|
constructor(
|
||||||
@Inject('ConnectionTestGrpcClient') private readonly grpcClient: ClientGrpc,
|
@Inject('ConnectionTestGrpcClient') private readonly grpcClient: ClientGrpc,
|
||||||
private connectionsService: ConnectionClientService,
|
private connectionsService: ConnectionClientService,
|
||||||
|
private connectionsApiService: ConnectionsApiService,
|
||||||
|
private platformApiService: PlatformApiService,
|
||||||
) {
|
) {
|
||||||
this.connectionTestReadClient =
|
this.connectionTestReadClient =
|
||||||
grpcClient.getService<ConnectionTest.ReadService.ConnectionTestReadServices>(
|
grpcClient.getService<ConnectionTest.ReadService.ConnectionTestReadServices>(
|
||||||
@@ -147,45 +154,137 @@ export class ConnectionTestService {
|
|||||||
}
|
}
|
||||||
async connectionTestListSchemas(
|
async connectionTestListSchemas(
|
||||||
body: ConnectionTestListSchemasReq,
|
body: ConnectionTestListSchemasReq,
|
||||||
customer_name: string,
|
user: RequestUser,
|
||||||
): Promise<ConnectionTestListSchemasRes> {
|
): Promise<ConnectionTestListSchemasRes> {
|
||||||
const { connection_id, plugin } = body;
|
const result = await this.connectionsApiService.proxy(
|
||||||
return lastValueFrom(
|
'GET',
|
||||||
this.connectionTestReadClient.ListSchemas({
|
`/connection_catalog/${encodeURIComponent(body.connection_id)}/schemas`,
|
||||||
connection_id,
|
user,
|
||||||
customer_name,
|
|
||||||
plugin,
|
|
||||||
}),
|
|
||||||
);
|
);
|
||||||
|
return {
|
||||||
|
operation_result: true,
|
||||||
|
schema_list: result.schemas.map((schema) => schema.schema_name),
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
async connectionTestListTables(
|
async connectionTestListTables(
|
||||||
body: ConnectionTestListTablesReq,
|
body: ConnectionTestListTablesReq,
|
||||||
customer_name: string,
|
user: RequestUser,
|
||||||
): Promise<ConnectionTestListTablesRes> {
|
): Promise<ConnectionTestListTablesRes> {
|
||||||
const { connection_id, plugin, schema } = body;
|
const result = await this.connectionsApiService.proxy(
|
||||||
return lastValueFrom(
|
'GET',
|
||||||
this.connectionTestReadClient.ListTables({
|
`/connection_catalog/${encodeURIComponent(body.connection_id)}` +
|
||||||
connection_id,
|
`/schemas/${encodeURIComponent(body.schema)}/tables`,
|
||||||
customer_name,
|
user,
|
||||||
plugin,
|
|
||||||
schema,
|
|
||||||
}),
|
|
||||||
);
|
);
|
||||||
|
return {
|
||||||
|
operation_result: true,
|
||||||
|
table_list: result.tables.map((table) => table.table_name),
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
async getTableMetadata(
|
async getTableMetadata(
|
||||||
body: GetTableMetadataReq,
|
body: GetTableMetadataReq,
|
||||||
customer_name: string,
|
user: RequestUser,
|
||||||
): Promise<GetTableMetadataRes> {
|
): Promise<GetTableMetadataRes> {
|
||||||
const { schema, plugin, table_list, connection_id } = body;
|
const tables_metadata = await Promise.all(
|
||||||
return lastValueFrom(
|
body.table_list.map(async (table_name) => {
|
||||||
this.connectionTestReadClient.GetTableMetadata({
|
const result = await this.connectionsApiService.proxy(
|
||||||
connection_id,
|
'GET',
|
||||||
customer_name,
|
`/connection_catalog/${encodeURIComponent(body.connection_id)}` +
|
||||||
plugin,
|
`/schemas/${encodeURIComponent(body.schema)}` +
|
||||||
schema,
|
`/tables/${encodeURIComponent(table_name)}/columns`,
|
||||||
table_list,
|
user,
|
||||||
|
);
|
||||||
|
return {
|
||||||
|
table_name,
|
||||||
|
columns: result.columns.map((column) => ({
|
||||||
|
name: column.column_name,
|
||||||
|
type: column.data_type,
|
||||||
|
is_primary_key: column.is_primary_key,
|
||||||
|
})),
|
||||||
|
references: [],
|
||||||
|
};
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
|
return { operation_result: true, tables_metadata };
|
||||||
|
}
|
||||||
|
|
||||||
|
async refreshCatalog(
|
||||||
|
body: RefreshCatalogReq,
|
||||||
|
user: RequestUser,
|
||||||
|
): Promise<RefreshCatalogRes> {
|
||||||
|
const task = await this.platformApiService.proxy(
|
||||||
|
'POST',
|
||||||
|
'/connection_test',
|
||||||
|
user,
|
||||||
|
{
|
||||||
|
customer_id: user.customer_name,
|
||||||
|
plugin: body.plugin,
|
||||||
|
task: {
|
||||||
|
task_type: 'refresh_catalog',
|
||||||
|
connection: {
|
||||||
|
provider: 'connection_manager',
|
||||||
|
config_id: body.connection_id,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
if (!task.session_id || !task.date) {
|
||||||
|
throw new HttpException(
|
||||||
|
'Platform API did not return a catalog refresh task identifier',
|
||||||
|
HttpStatus.BAD_GATEWAY,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
operation_result: true,
|
||||||
|
status: 'PENDING',
|
||||||
|
session_id: task.session_id,
|
||||||
|
date: task.date,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
async refreshCatalogStatus(
|
||||||
|
body: RefreshCatalogStatusReq,
|
||||||
|
user: RequestUser,
|
||||||
|
): Promise<RefreshCatalogRes> {
|
||||||
|
const result = await this.platformApiService.proxy(
|
||||||
|
'POST',
|
||||||
|
'/connection_test/status',
|
||||||
|
user,
|
||||||
|
{
|
||||||
|
session_id: body.session_id,
|
||||||
|
date: body.date,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
if (result.status === 'DONE') {
|
||||||
|
await this.connectionsApiService.proxy(
|
||||||
|
'PUT',
|
||||||
|
`/connection_config/${encodeURIComponent(
|
||||||
|
body.connection_id,
|
||||||
|
)}/catalog_metadata`,
|
||||||
|
user,
|
||||||
|
{
|
||||||
|
last_catalog_refresh_status: 'SUCCESS',
|
||||||
|
last_catalog_connection_test_date: body.date,
|
||||||
|
last_catalog_connection_test_session_id: body.session_id,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
} else if (result.status === 'ERROR' || result.status === 'EXPIRED') {
|
||||||
|
throw new HttpException(
|
||||||
|
`Catalog refresh finished with status ${result.status}`,
|
||||||
|
HttpStatus.BAD_GATEWAY,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
operation_result: result.status === 'DONE',
|
||||||
|
status: result.status,
|
||||||
|
session_id: body.session_id,
|
||||||
|
date: body.date,
|
||||||
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||||
import { IsString, IsOptional } from 'class-validator';
|
import { IsIn, IsString, IsOptional } from 'class-validator';
|
||||||
import { DatabaseConnectionPropertiesDto } from 'src/modules/connection/dtos/connection';
|
import { DatabaseConnectionPropertiesDto } from 'src/modules/connection/dtos/connection';
|
||||||
import { CreateConnectionDto } from 'src/modules/connection/dtos/connection';
|
import { CreateConnectionDto } from 'src/modules/connection/dtos/connection';
|
||||||
export class ColumnDto {
|
export class ColumnDto {
|
||||||
@@ -7,6 +7,8 @@ export class ColumnDto {
|
|||||||
name: string;
|
name: string;
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
type: string;
|
type: string;
|
||||||
|
@ApiProperty()
|
||||||
|
is_primary_key: boolean;
|
||||||
}
|
}
|
||||||
export class TableMetadataDto {
|
export class TableMetadataDto {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
@@ -131,3 +133,37 @@ export class GetTableMetadataRes {
|
|||||||
@ApiProperty({ type: [TableMetadataDto] })
|
@ApiProperty({ type: [TableMetadataDto] })
|
||||||
tables_metadata: TableMetadataDto[];
|
tables_metadata: TableMetadataDto[];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export class RefreshCatalogReq {
|
||||||
|
@ApiProperty()
|
||||||
|
@IsString()
|
||||||
|
connection_id: string;
|
||||||
|
|
||||||
|
@ApiProperty({ enum: ['oracle', 'mysql', 'postgresql', 'sqlserver'] })
|
||||||
|
@IsIn(['oracle', 'mysql', 'postgresql', 'sqlserver'])
|
||||||
|
plugin: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class RefreshCatalogStatusReq extends RefreshCatalogReq {
|
||||||
|
@ApiProperty()
|
||||||
|
@IsString()
|
||||||
|
session_id: string;
|
||||||
|
|
||||||
|
@ApiProperty()
|
||||||
|
@IsString()
|
||||||
|
date: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class RefreshCatalogRes {
|
||||||
|
@ApiProperty()
|
||||||
|
operation_result: boolean;
|
||||||
|
|
||||||
|
@ApiProperty()
|
||||||
|
status: string;
|
||||||
|
|
||||||
|
@ApiProperty()
|
||||||
|
session_id: string;
|
||||||
|
|
||||||
|
@ApiProperty()
|
||||||
|
date: string;
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,11 @@
|
|||||||
|
export const CONNECTIONS_API_CONFIG = {
|
||||||
|
getUrl: (): string => {
|
||||||
|
const url = process.env.CONNECTIONS_API_URL;
|
||||||
|
if (!url) {
|
||||||
|
throw new Error('CONNECTIONS_API_URL environment variable is not set');
|
||||||
|
}
|
||||||
|
return url;
|
||||||
|
},
|
||||||
|
region: process.env.AWS_REGION || 'us-east-1',
|
||||||
|
timeout: parseInt(process.env.CONNECTIONS_API_TIMEOUT || '30000', 10),
|
||||||
|
};
|
||||||
@@ -0,0 +1,10 @@
|
|||||||
|
import { Module } from '@nestjs/common';
|
||||||
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
|
|
||||||
|
import { ConnectionsApiService } from './connections-api.service';
|
||||||
|
|
||||||
|
@Module({
|
||||||
|
providers: [ConnectionsApiService, DadosferaLogger],
|
||||||
|
exports: [ConnectionsApiService],
|
||||||
|
})
|
||||||
|
export class ConnectionsApiModule {}
|
||||||
@@ -0,0 +1,99 @@
|
|||||||
|
import { Injectable, Inject, HttpException } from '@nestjs/common';
|
||||||
|
import { SignatureV4 } from '@aws-sdk/signature-v4';
|
||||||
|
import { Sha256 } from '@aws-crypto/sha256-js';
|
||||||
|
import { defaultProvider } from '@aws-sdk/credential-provider-node';
|
||||||
|
import axios, { AxiosResponse, Method } from 'axios';
|
||||||
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
|
|
||||||
|
import { RequestUser } from '../../decorators/user.decorator';
|
||||||
|
import { CONNECTIONS_API_CONFIG } from './connections-api.config';
|
||||||
|
|
||||||
|
@Injectable()
|
||||||
|
export class ConnectionsApiService {
|
||||||
|
private signer: SignatureV4;
|
||||||
|
private logger: any;
|
||||||
|
|
||||||
|
constructor(@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger) {
|
||||||
|
this.logger = dadosferaLogger.logger;
|
||||||
|
this.signer = new SignatureV4({
|
||||||
|
service: 'execute-api',
|
||||||
|
region: CONNECTIONS_API_CONFIG.region,
|
||||||
|
credentials: defaultProvider(),
|
||||||
|
sha256: Sha256,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async proxy(
|
||||||
|
method: string,
|
||||||
|
path: string,
|
||||||
|
user: RequestUser,
|
||||||
|
body?: any,
|
||||||
|
query?: Record<string, string>,
|
||||||
|
): Promise<any> {
|
||||||
|
const baseUrl = CONNECTIONS_API_CONFIG.getUrl();
|
||||||
|
const url = new URL(`${baseUrl}${path}`);
|
||||||
|
|
||||||
|
if (query) {
|
||||||
|
Object.entries(query).forEach(([key, value]) => {
|
||||||
|
if (value !== undefined && value !== null) {
|
||||||
|
url.searchParams.set(key, String(value));
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
const headers: Record<string, string> = {
|
||||||
|
host: url.hostname,
|
||||||
|
'content-type': 'application/json',
|
||||||
|
customer_name: user.customer_name || '',
|
||||||
|
customer_id: user.customer_id || '',
|
||||||
|
'x-user-id': user.user_id || '',
|
||||||
|
'x-username': user.username || '',
|
||||||
|
'x-customer-tier': user.customer_tier || '',
|
||||||
|
'x-customer-id': user.customer_id || '',
|
||||||
|
};
|
||||||
|
const requestToSign = {
|
||||||
|
method: method.toUpperCase(),
|
||||||
|
protocol: url.protocol,
|
||||||
|
hostname: url.hostname,
|
||||||
|
port: url.port ? parseInt(url.port, 10) : undefined,
|
||||||
|
path: url.pathname + url.search,
|
||||||
|
headers,
|
||||||
|
body: body ? JSON.stringify(body) : undefined,
|
||||||
|
};
|
||||||
|
|
||||||
|
try {
|
||||||
|
const signedRequest = await this.signer.sign(requestToSign);
|
||||||
|
const response: AxiosResponse = await axios({
|
||||||
|
method: method as Method,
|
||||||
|
url: url.href,
|
||||||
|
headers: signedRequest.headers as Record<string, string>,
|
||||||
|
data: body,
|
||||||
|
timeout: CONNECTIONS_API_CONFIG.timeout,
|
||||||
|
validateStatus: () => true,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (response.status >= 400) {
|
||||||
|
throw new HttpException(response.data, response.status);
|
||||||
|
}
|
||||||
|
return response.data;
|
||||||
|
} catch (error) {
|
||||||
|
this.logger.error('Connections API proxy error', {
|
||||||
|
error: error.message,
|
||||||
|
path,
|
||||||
|
method: method.toUpperCase(),
|
||||||
|
});
|
||||||
|
if (error instanceof HttpException) {
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
if (error.response) {
|
||||||
|
throw new HttpException(error.response.data, error.response.status);
|
||||||
|
}
|
||||||
|
if (error.code === 'ECONNREFUSED') {
|
||||||
|
throw new HttpException('Connections API service unavailable', 503);
|
||||||
|
}
|
||||||
|
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
|
||||||
|
throw new HttpException('Connections API request timeout', 504);
|
||||||
|
}
|
||||||
|
throw new HttpException('Internal server error', 500);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,4 +1,8 @@
|
|||||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
export interface Info {
|
||||||
|
user_id: string;
|
||||||
|
customer_id: string;
|
||||||
|
customer: string;
|
||||||
|
}
|
||||||
|
|
||||||
interface Values {
|
interface Values {
|
||||||
jdbc_user: string;
|
jdbc_user: string;
|
||||||
|
|||||||
@@ -23,6 +23,7 @@ import {
|
|||||||
} 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';
|
||||||
import { CreateInputReq } from './dtos/input.model';
|
import { CreateInputReq } from './dtos/input.model';
|
||||||
|
import { Metadata } from '@grpc/grpc-js';
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
|
|
||||||
@@ -73,10 +74,10 @@ export class InputsService {
|
|||||||
objectCamelToSnake(createInputResponse);
|
objectCamelToSnake(createInputResponse);
|
||||||
return createInputResponse;
|
return createInputResponse;
|
||||||
},
|
},
|
||||||
update: async (updateInputDTO: UpdateInputRequest): Promise<InputUpdateResponse> => {
|
update: async (updateInputDTO: UpdateInputRequest, metadata: Metadata): Promise<InputUpdateResponse> => {
|
||||||
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
||||||
const updateInputResponse = await lastValueFrom(
|
const updateInputResponse = await lastValueFrom(
|
||||||
this.inputWriteService.InputUpdate(updateInputDTO),
|
this.inputWriteService.InputUpdate(updateInputDTO, metadata),
|
||||||
);
|
);
|
||||||
|
|
||||||
return updateInputResponse;
|
return updateInputResponse;
|
||||||
@@ -206,7 +207,7 @@ export class InputsService {
|
|||||||
return findOneInputResponse;
|
return findOneInputResponse;
|
||||||
}
|
}
|
||||||
|
|
||||||
async update(id: string, data, info: Info) {
|
async update(id: string, data, info: Info, metadata?: Metadata) {
|
||||||
// this.validateCron({ ...data, info });
|
// this.validateCron({ ...data, info });
|
||||||
try {
|
try {
|
||||||
const {
|
const {
|
||||||
@@ -217,7 +218,7 @@ export class InputsService {
|
|||||||
id,
|
id,
|
||||||
...data,
|
...data,
|
||||||
info,
|
info,
|
||||||
});
|
}, metadata);
|
||||||
|
|
||||||
const updateInputResponse = this.adjustInputPayload(
|
const updateInputResponse = this.adjustInputPayload(
|
||||||
input,
|
input,
|
||||||
|
|||||||
@@ -1,78 +0,0 @@
|
|||||||
import { ConflictException, Inject, OnModuleInit } from '@nestjs/common';
|
|
||||||
import { ClientGrpc } from '@nestjs/microservices';
|
|
||||||
import {
|
|
||||||
PipelineServicesNames,
|
|
||||||
PipelinesServiceInterface,
|
|
||||||
} from '@dadosfera/protospack';
|
|
||||||
import { lastValueFrom } from 'rxjs';
|
|
||||||
|
|
||||||
import { IIdRequest } from './interfaces';
|
|
||||||
|
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|
||||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
|
||||||
|
|
||||||
export class PipelinesClientService implements OnModuleInit {
|
|
||||||
private pipelineService: PipelinesServiceInterface;
|
|
||||||
logger: DadosferaLogger;
|
|
||||||
|
|
||||||
constructor(
|
|
||||||
@Inject(DadosferaLogger)
|
|
||||||
dadosferaLogger: DadosferaLogger,
|
|
||||||
@Inject(PipelinesClientConfiguration.name)
|
|
||||||
private readonly grpcClient: ClientGrpc,
|
|
||||||
) {
|
|
||||||
this.logger = dadosferaLogger.logger;
|
|
||||||
}
|
|
||||||
|
|
||||||
onModuleInit() {
|
|
||||||
this.pipelineService =
|
|
||||||
this.grpcClient.getService<PipelinesServiceInterface>(
|
|
||||||
PipelineServicesNames.PipelineService,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
async getPipelineStatus(data) {
|
|
||||||
this.logger.info('PipelinesClientService - GetPipelineStatus');
|
|
||||||
|
|
||||||
const statusPipelineResponse = await lastValueFrom(
|
|
||||||
this.pipelineService.getPipelineStatus(data),
|
|
||||||
)
|
|
||||||
.then((res) => {
|
|
||||||
const statusArray =
|
|
||||||
res.status?.sort((a, b) => {
|
|
||||||
if (a.id < b.id) {
|
|
||||||
return 1;
|
|
||||||
} else {
|
|
||||||
return -1;
|
|
||||||
}
|
|
||||||
}) || [];
|
|
||||||
return { status: statusArray };
|
|
||||||
})
|
|
||||||
.catch((err) => {
|
|
||||||
this.logger.error(err.message);
|
|
||||||
throw new Error(err);
|
|
||||||
});
|
|
||||||
this.logger.info('Done');
|
|
||||||
|
|
||||||
return statusPipelineResponse;
|
|
||||||
}
|
|
||||||
|
|
||||||
async runPipeline({ id, info }: IIdRequest) {
|
|
||||||
this.logger.info('PipelinesClientService - RunPipeline');
|
|
||||||
const statusPipelineResponse = await lastValueFrom(
|
|
||||||
this.pipelineService.triggerPipeline({ id, info }),
|
|
||||||
).catch((err) => {
|
|
||||||
this.logger.error(err.message);
|
|
||||||
throw new Error(err);
|
|
||||||
});
|
|
||||||
|
|
||||||
if (statusPipelineResponse.status == false) {
|
|
||||||
throw new ConflictException(
|
|
||||||
'This pipeline is not ready yet to execute, Try again later!',
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
this.logger.info('Done');
|
|
||||||
return statusPipelineResponse;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-36
@@ -1,36 +0,0 @@
|
|||||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
|
||||||
|
|
||||||
export interface ICreatePipelineDto {
|
|
||||||
input: IdRequest;
|
|
||||||
transformations: IdRequest[];
|
|
||||||
output: IdRequest;
|
|
||||||
tags: string[];
|
|
||||||
name: string;
|
|
||||||
description: string;
|
|
||||||
info: Info;
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface IdRequest {
|
|
||||||
id: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface IIdRequest {
|
|
||||||
id: string;
|
|
||||||
info: Info;
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface IUpdatePipelineRequest {
|
|
||||||
input: IdRequest;
|
|
||||||
transformations: IdRequest[];
|
|
||||||
output: IdRequest;
|
|
||||||
tags: string[];
|
|
||||||
name: string;
|
|
||||||
description: string;
|
|
||||||
id: string;
|
|
||||||
info: Info;
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface IGetPipelineLogsRequest {
|
|
||||||
id: string;
|
|
||||||
details: string;
|
|
||||||
}
|
|
||||||
@@ -1,33 +0,0 @@
|
|||||||
import {
|
|
||||||
ClientsProviderAsyncOptions,
|
|
||||||
GrpcOptions,
|
|
||||||
Transport,
|
|
||||||
} from '@nestjs/microservices';
|
|
||||||
import { PipelinePackages, PipelineProtoFilePath } from '@dadosfera/protospack';
|
|
||||||
import { credentials } from '@grpc/grpc-js';
|
|
||||||
|
|
||||||
const isLocalConnection =
|
|
||||||
process.env.PIFACTORY_URL.startsWith('pi-factory:') ||
|
|
||||||
process.env.PIFACTORY_URL.includes('0.0.0.0');
|
|
||||||
|
|
||||||
export class PipelinesClientConfiguration {
|
|
||||||
public name = 'PipelinesClientConfiguration';
|
|
||||||
private config: GrpcOptions = {
|
|
||||||
transport: Transport.GRPC,
|
|
||||||
options: {
|
|
||||||
url: process.env.PIFACTORY_URL,
|
|
||||||
package: PipelinePackages,
|
|
||||||
credentials: isLocalConnection ? undefined : credentials.createSsl(),
|
|
||||||
protoPath: PipelineProtoFilePath,
|
|
||||||
loader: {
|
|
||||||
keepCase: true,
|
|
||||||
enums: String,
|
|
||||||
defaults: false,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
};
|
|
||||||
providerOptions: ClientsProviderAsyncOptions = {
|
|
||||||
name: this.name,
|
|
||||||
...this.config,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
@@ -1,72 +0,0 @@
|
|||||||
import { Body, Controller, Get, Inject, Param, Post } from '@nestjs/common';
|
|
||||||
import { ApiOperation, ApiTags } from '@nestjs/swagger';
|
|
||||||
import {
|
|
||||||
AuthenticateCondition,
|
|
||||||
Authenticated,
|
|
||||||
RequireSomePermission,
|
|
||||||
} from 'src/decorators/authentication.decorator';
|
|
||||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
|
||||||
import { PipelinesService } from './pipelines.service';
|
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|
||||||
import { ApiInternalOnlyController } from 'src/decorators/swagger.decorator';
|
|
||||||
|
|
||||||
@ApiInternalOnlyController()
|
|
||||||
@ApiTags('Pipelines')
|
|
||||||
@Controller('pipelines')
|
|
||||||
@Authenticated()
|
|
||||||
export class PipelinesController {
|
|
||||||
logger: DadosferaLogger;
|
|
||||||
constructor(
|
|
||||||
@Inject(DadosferaLogger)
|
|
||||||
dadosferaLogger: DadosferaLogger,
|
|
||||||
private pipelineService: PipelinesService,
|
|
||||||
) {
|
|
||||||
this.logger = dadosferaLogger.logger;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post('start/:id')
|
|
||||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
|
||||||
@ApiOperation({
|
|
||||||
deprecated: true,
|
|
||||||
description:
|
|
||||||
'This method is deprecated. Please use route /pipelinesV2/start/:id instead',
|
|
||||||
})
|
|
||||||
async activate(@Param('id') id: string, @Body() body) {
|
|
||||||
const { info } = body;
|
|
||||||
|
|
||||||
this.logger.info(
|
|
||||||
process.env.DEV_URL + `/pipeline/start/${id} - ON START PIPELINE ROUTE`,
|
|
||||||
{
|
|
||||||
user: body.info.user_id,
|
|
||||||
customer: body.info.customer,
|
|
||||||
},
|
|
||||||
);
|
|
||||||
|
|
||||||
const response = await this.pipelineService.runPipeline({ id, info });
|
|
||||||
|
|
||||||
return response;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Get(':id/status')
|
|
||||||
@ApiOperation({
|
|
||||||
deprecated: true,
|
|
||||||
description:
|
|
||||||
'This method is deprecated. Please use route /pipelinesV2/:id/status instead',
|
|
||||||
})
|
|
||||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
|
||||||
async getPipelineStatus(@Body() body, @Param('id') id: string) {
|
|
||||||
body.id = id;
|
|
||||||
|
|
||||||
this.logger.info(
|
|
||||||
process.env.DEV_URL + `/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`,
|
|
||||||
{
|
|
||||||
user: body.info.user_id,
|
|
||||||
customer: body.info.customer,
|
|
||||||
},
|
|
||||||
);
|
|
||||||
|
|
||||||
const response = await this.pipelineService.getPipelineStatus(body);
|
|
||||||
|
|
||||||
return response;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,19 +0,0 @@
|
|||||||
import { Module } from '@nestjs/common';
|
|
||||||
import { ClientsModule } from '@nestjs/microservices';
|
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|
||||||
|
|
||||||
import { PipelinesController } from './pipelines.controller';
|
|
||||||
import { PipelinesService } from './pipelines.service';
|
|
||||||
|
|
||||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
|
||||||
import { PipelinesClientService } from './client.service';
|
|
||||||
|
|
||||||
const client = new PipelinesClientConfiguration();
|
|
||||||
|
|
||||||
@Module({
|
|
||||||
imports: [ClientsModule.register([client.providerOptions])],
|
|
||||||
controllers: [PipelinesController],
|
|
||||||
providers: [PipelinesService, PipelinesClientService, DadosferaLogger],
|
|
||||||
exports: [PipelinesService],
|
|
||||||
})
|
|
||||||
export class PipelinesModule {}
|
|
||||||
@@ -1,33 +0,0 @@
|
|||||||
import { HttpException, HttpStatus, Injectable } from '@nestjs/common';
|
|
||||||
import { PipelinesClientService } from './client.service';
|
|
||||||
import { IIdRequest } from './interfaces';
|
|
||||||
import { objectCamelToSnake } from 'src/utils/CaseConverter';
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class PipelinesService {
|
|
||||||
constructor(private pipelineClient: PipelinesClientService) {}
|
|
||||||
|
|
||||||
async getPipelineStatus(data: IIdRequest) {
|
|
||||||
try {
|
|
||||||
const pipelineStatusResponse =
|
|
||||||
await this.pipelineClient.getPipelineStatus(data);
|
|
||||||
|
|
||||||
return objectCamelToSnake(pipelineStatusResponse);
|
|
||||||
} catch (err) {
|
|
||||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
async runPipeline({ id, info }: IIdRequest) {
|
|
||||||
try {
|
|
||||||
const triggerPipelineResponse = await this.pipelineClient.runPipeline({
|
|
||||||
id,
|
|
||||||
info,
|
|
||||||
});
|
|
||||||
|
|
||||||
return objectCamelToSnake(triggerPipelineResponse);
|
|
||||||
} catch (err) {
|
|
||||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,5 +1,4 @@
|
|||||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
|
||||||
|
|
||||||
export class PipelineInputsDTO {
|
export class PipelineInputsDTO {
|
||||||
@ApiProperty()
|
@ApiProperty()
|
||||||
@@ -62,6 +61,12 @@ export interface IIdRequest {
|
|||||||
info: Info;
|
info: Info;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface Info {
|
||||||
|
user_id: string;
|
||||||
|
customer_id: string;
|
||||||
|
customer: string;
|
||||||
|
}
|
||||||
|
|
||||||
export interface IUpdatePipelineRequest {
|
export interface IUpdatePipelineRequest {
|
||||||
input: IdRequest;
|
input: IdRequest;
|
||||||
transformations: IdRequest[];
|
transformations: IdRequest[];
|
||||||
|
|||||||
@@ -34,7 +34,6 @@ import { Messages } from '@dadosfera/protospack-v2/dist/lib/PipelineV2';
|
|||||||
import { RequestUser, User } from 'src/decorators/user.decorator';
|
import { RequestUser, User } from 'src/decorators/user.decorator';
|
||||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||||
|
|
||||||
import { PipelinesService as OldPipelineService } from 'src/modules/pipelines/pipelines.service';
|
|
||||||
import {
|
import {
|
||||||
ICompleteUploadCSVFile,
|
ICompleteUploadCSVFile,
|
||||||
ICreatePipelineCSVFile,
|
ICreatePipelineCSVFile,
|
||||||
@@ -63,9 +62,7 @@ export class PipelinesController {
|
|||||||
constructor(
|
constructor(
|
||||||
@Inject(DadosferaLogger)
|
@Inject(DadosferaLogger)
|
||||||
dadosferaLogger: DadosferaLogger,
|
dadosferaLogger: DadosferaLogger,
|
||||||
|
|
||||||
private pipelinesClientService: PipelinesService,
|
private pipelinesClientService: PipelinesService,
|
||||||
private oldPipelinesService: OldPipelineService,
|
|
||||||
) {
|
) {
|
||||||
this.logger = dadosferaLogger.logger;
|
this.logger = dadosferaLogger.logger;
|
||||||
}
|
}
|
||||||
@@ -201,7 +198,7 @@ export class PipelinesController {
|
|||||||
customer: body.info.customer,
|
customer: body.info.customer,
|
||||||
});
|
});
|
||||||
|
|
||||||
const response = await this.oldPipelinesService.getPipelineStatus(body);
|
const response = await this.pipelinesClientService.getPipelineStatus(body);
|
||||||
|
|
||||||
return response;
|
return response;
|
||||||
}
|
}
|
||||||
@@ -320,7 +317,6 @@ 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 info: Info = {
|
const info: Info = {
|
||||||
user_id: user.user_id,
|
user_id: user.user_id,
|
||||||
customer: user.customer_name,
|
customer: user.customer_name,
|
||||||
@@ -328,13 +324,7 @@ export class PipelinesController {
|
|||||||
pipeline_id: pipelineId
|
pipeline_id: pipelineId
|
||||||
};
|
};
|
||||||
|
|
||||||
const metadata = PackTheMetadata({
|
const metadata = PackTheMetadata(user);
|
||||||
customer_id,
|
|
||||||
customer_name,
|
|
||||||
user_id,
|
|
||||||
username,
|
|
||||||
language,
|
|
||||||
});
|
|
||||||
|
|
||||||
const response = await this.pipelinesClientService.updatePipelineInput(
|
const response = await this.pipelinesClientService.updatePipelineInput(
|
||||||
pipelineId,
|
pipelineId,
|
||||||
@@ -368,6 +358,21 @@ export class PipelinesController {
|
|||||||
return response;
|
return response;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Patch('/:id/upgrade')
|
||||||
|
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||||
|
@HttpCode(HttpStatus.NO_CONTENT)
|
||||||
|
async upgradeConnector(
|
||||||
|
@Language() language: LanguageEnum,
|
||||||
|
@Param('id') id: string,
|
||||||
|
@User() user: RequestUser
|
||||||
|
) {
|
||||||
|
this.logger.info('PipelinesController - upgrade connector');
|
||||||
|
|
||||||
|
const metadata = PackTheMetadata(user);
|
||||||
|
|
||||||
|
await this.pipelinesClientService.upgrade(id, metadata);
|
||||||
|
}
|
||||||
|
|
||||||
@Delete(':id')
|
@Delete(':id')
|
||||||
@ApiNoContentResponse()
|
@ApiNoContentResponse()
|
||||||
@HttpCode(HttpStatus.NO_CONTENT)
|
@HttpCode(HttpStatus.NO_CONTENT)
|
||||||
@@ -497,7 +502,7 @@ export class PipelinesController {
|
|||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
|
||||||
const response = await this.oldPipelinesService.runPipeline({ id, info });
|
const response = await this.pipelinesClientService.runPipeline({ id, info });
|
||||||
|
|
||||||
return response;
|
return response;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,7 +7,6 @@ import { PipelinesService } from './pipelines.service';
|
|||||||
|
|
||||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||||
|
|
||||||
import { PipelinesModule as OldPipelineModule } from 'src/modules/pipelines/pipelines.module';
|
|
||||||
import { ConnectorModule } from '../connector/connector.module';
|
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';
|
||||||
@@ -21,7 +20,6 @@ const client = new PipelinesClientConfiguration();
|
|||||||
@Module({
|
@Module({
|
||||||
imports: [
|
imports: [
|
||||||
ClientsModule.register([client.providerOptions]),
|
ClientsModule.register([client.providerOptions]),
|
||||||
OldPipelineModule,
|
|
||||||
ConnectorModule,
|
ConnectorModule,
|
||||||
InputsModule,
|
InputsModule,
|
||||||
TransformationsModule,
|
TransformationsModule,
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
/* eslint-disable no-async-promise-executor */
|
/* eslint-disable no-async-promise-executor */
|
||||||
import {
|
import {
|
||||||
BadRequestException,
|
BadRequestException,
|
||||||
|
ConflictException,
|
||||||
HttpException,
|
HttpException,
|
||||||
HttpStatus,
|
HttpStatus,
|
||||||
Inject,
|
Inject,
|
||||||
@@ -17,7 +18,7 @@ import { lastValueFrom } from 'rxjs';
|
|||||||
|
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||||
import { ICreatePipelineV2Req, UpdatePlatformInputRequest, UpdateTableDTO } from './interfaces';
|
import { ICreatePipelineV2Req, IIdRequest, UpdatePlatformInputRequest, UpdateTableDTO } from './interfaces';
|
||||||
import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages';
|
import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages';
|
||||||
import { Metadata } from '@grpc/grpc-js';
|
import { Metadata } from '@grpc/grpc-js';
|
||||||
import { ConnectorClientService } from '../connector/client.service';
|
import { ConnectorClientService } from '../connector/client.service';
|
||||||
@@ -176,6 +177,15 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
return updatePipelineResponse;
|
return updatePipelineResponse;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async upgrade(id: string, metadata: Metadata) {
|
||||||
|
await lastValueFrom(
|
||||||
|
this.pipelineWriteService.Upgrade(
|
||||||
|
{ id },
|
||||||
|
metadata,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
async remove(data: { id: string; metadata: Metadata; user: RequestUser }) {
|
async remove(data: { id: string; metadata: Metadata; user: RequestUser }) {
|
||||||
const { id, metadata, user } = data;
|
const { id, metadata, user } = data;
|
||||||
const info = {
|
const info = {
|
||||||
@@ -373,7 +383,8 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
const updateInputResponse = await this.inputsService.update(
|
const updateInputResponse = await this.inputsService.update(
|
||||||
inputId,
|
inputId,
|
||||||
updateInputDTO,
|
updateInputDTO,
|
||||||
info
|
info,
|
||||||
|
metadata
|
||||||
);
|
);
|
||||||
|
|
||||||
const inputRollback = () => {
|
const inputRollback = () => {
|
||||||
@@ -394,34 +405,36 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
|
|
||||||
const nimbusUpdates = updateInputResponse?.tablesUpdate || [];
|
const nimbusUpdates = updateInputResponse?.tablesUpdate || [];
|
||||||
|
|
||||||
nimbusUpdates.forEach(update => {
|
if (user.customer_modules.includes('catalog')) {
|
||||||
const nimbusRollback = () => {
|
nimbusUpdates.forEach(update => {
|
||||||
return this.nimbusService.renameTable(
|
const nimbusRollback = () => {
|
||||||
info.customer,
|
return this.nimbusService.renameTable(
|
||||||
update.database,
|
info.customer,
|
||||||
{
|
update.database,
|
||||||
table_name: update.table_name,
|
{
|
||||||
table_schema: update.table_schema
|
table_name: update.table_name,
|
||||||
},
|
table_schema: update.table_schema
|
||||||
{
|
},
|
||||||
table_name: update.old_table_name,
|
{
|
||||||
table_schema: update.old_table_schema
|
table_name: update.old_table_name,
|
||||||
}
|
table_schema: update.old_table_schema
|
||||||
);
|
}
|
||||||
}
|
);
|
||||||
rollback.push(nimbusRollback);
|
}
|
||||||
});
|
rollback.push(nimbusRollback);
|
||||||
|
});
|
||||||
|
|
||||||
try {
|
try {
|
||||||
await this.updateNimbus(info.customer, nimbusUpdates);
|
await this.updateNimbus(info.customer, nimbusUpdates);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
this.logger.error(error);
|
this.logger.error(error);
|
||||||
if (error instanceof AxiosError) {
|
if (error instanceof AxiosError) {
|
||||||
this.logger.error(JSON.stringify(error.response.data));
|
this.logger.error(JSON.stringify(error.response.data));
|
||||||
}
|
}
|
||||||
await this.executeRenameRollback(rollback);
|
await this.executeRenameRollback(rollback);
|
||||||
|
|
||||||
throw new Error("Error Nimbus updating tables");
|
throw new Error("Error Nimbus updating tables");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -623,4 +636,49 @@ export class PipelinesService implements OnModuleInit {
|
|||||||
return assets;
|
return assets;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async getPipelineStatus(data) {
|
||||||
|
this.logger.info('PipelinesClientService - GetPipelineStatus');
|
||||||
|
|
||||||
|
const statusPipelineResponse = await lastValueFrom(
|
||||||
|
this.pipelineReadService.PipelineV2GetPipelineV2Status(data),
|
||||||
|
)
|
||||||
|
.then((res) => {
|
||||||
|
const statusArray =
|
||||||
|
res.status?.sort((a, b) => {
|
||||||
|
if (a.id < b.id) {
|
||||||
|
return 1;
|
||||||
|
} else {
|
||||||
|
return -1;
|
||||||
|
}
|
||||||
|
}) || [];
|
||||||
|
return { status: statusArray };
|
||||||
|
})
|
||||||
|
.catch((err) => {
|
||||||
|
this.logger.error(err.message);
|
||||||
|
throw new Error(err);
|
||||||
|
});
|
||||||
|
this.logger.info('Done');
|
||||||
|
|
||||||
|
return statusPipelineResponse;
|
||||||
|
}
|
||||||
|
|
||||||
|
async runPipeline({ id, info }: IIdRequest) {
|
||||||
|
this.logger.info('PipelinesClientService - RunPipeline');
|
||||||
|
const statusPipelineResponse = await lastValueFrom(
|
||||||
|
this.pipelineWriteService.PipelineV2TriggerPipelineV2({ id, info }),
|
||||||
|
).catch((err) => {
|
||||||
|
this.logger.error(err.message);
|
||||||
|
throw new Error(err);
|
||||||
|
});
|
||||||
|
|
||||||
|
if (statusPipelineResponse.status == false) {
|
||||||
|
throw new ConflictException(
|
||||||
|
'This pipeline is not ready yet to execute, Try again later!',
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
this.logger.info('Done');
|
||||||
|
return statusPipelineResponse;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ import {
|
|||||||
NotFoundException,
|
NotFoundException,
|
||||||
UseGuards,
|
UseGuards,
|
||||||
} from '@nestjs/common';
|
} from '@nestjs/common';
|
||||||
import { ApiTags, ApiOperation } from '@nestjs/swagger';
|
import { ApiTags, ApiOperation, ApiOkResponse } from '@nestjs/swagger';
|
||||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||||
|
|
||||||
import {
|
import {
|
||||||
@@ -72,6 +72,10 @@ export class PlatformApiController {
|
|||||||
return id?.replace(/-/g, '_') || '';
|
return id?.replace(/-/g, '_') || '';
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private decodePathParam(value: string): string {
|
||||||
|
return value ? decodeURIComponent(value) : '';
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Denormalize ID back to UUID format (replace _ with -).
|
* Denormalize ID back to UUID format (replace _ with -).
|
||||||
* Used when we receive a normalized ID but need the original UUID.
|
* Used when we receive a normalized ID but need the original UUID.
|
||||||
@@ -832,6 +836,40 @@ export class PlatformApiController {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Get('pipelines/:pipelineId/pipeline_run/:runId/jobs')
|
||||||
|
@ApiOperation({
|
||||||
|
summary: 'Get pipeline run jobs',
|
||||||
|
description: 'Proxies platform-api DB-backed job runs and returns `{ jobs: [...] }`.',
|
||||||
|
})
|
||||||
|
@ApiOkResponse({
|
||||||
|
description: 'DB-backed job runs for the selected pipeline run.',
|
||||||
|
schema: {
|
||||||
|
type: 'object',
|
||||||
|
properties: {
|
||||||
|
jobs: {
|
||||||
|
type: 'array',
|
||||||
|
items: { type: 'object' },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
required: ['jobs'],
|
||||||
|
},
|
||||||
|
})
|
||||||
|
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||||
|
async getPipelineRunJobs(
|
||||||
|
@Param('pipelineId') pipelineId: string,
|
||||||
|
@Param('runId') runId: string,
|
||||||
|
@User() user: RequestUser,
|
||||||
|
) {
|
||||||
|
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
|
||||||
|
const decodedRunId = this.decodePathParam(runId);
|
||||||
|
|
||||||
|
return this.platformApiService.proxy(
|
||||||
|
'GET',
|
||||||
|
`/pipeline/${normalizedPipelineId}/pipeline_run/${decodedRunId}/jobs`,
|
||||||
|
user,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
// ==================== JOBS - COLUMN EDITING ROUTES ====================
|
// ==================== JOBS - COLUMN EDITING ROUTES ====================
|
||||||
|
|
||||||
@Put('jobs/:jobId/input')
|
@Put('jobs/:jobId/input')
|
||||||
|
|||||||
+5
-2
@@ -1,5 +1,8 @@
|
|||||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
export interface Info {
|
||||||
|
user_id: string;
|
||||||
|
customer_id: string;
|
||||||
|
customer: string;
|
||||||
|
}
|
||||||
export interface ICreateTransformationsRequest {
|
export interface ICreateTransformationsRequest {
|
||||||
transformations: Transformation[];
|
transformations: Transformation[];
|
||||||
info: Info;
|
info: Info;
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ interface IMetadata {
|
|||||||
details?: string;
|
details?: string;
|
||||||
sensitive?: string;
|
sensitive?: string;
|
||||||
roles?: string[];
|
roles?: string[];
|
||||||
|
customer_modules?: string[];
|
||||||
is_data_manager?: boolean;
|
is_data_manager?: boolean;
|
||||||
access_token?: string;
|
access_token?: string;
|
||||||
host?: string;
|
host?: string;
|
||||||
|
|||||||
Reference in New Issue
Block a user