mirror of
https://github.com/dadosfera/maestro.git
synced 2026-09-01 04:08:16 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8d04c1d495 |
@@ -1,12 +0,0 @@
|
||||
node_modules
|
||||
dist
|
||||
.git
|
||||
*.log
|
||||
npm-debug.log*
|
||||
.DS_Store
|
||||
.env
|
||||
.env.*
|
||||
coverage
|
||||
.nyc_output
|
||||
*.tgz
|
||||
!protospack.tgz
|
||||
@@ -56,9 +56,9 @@ jobs:
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
curl -fsSLO https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
wget https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
tar -xzf helmfile_0.148.0_linux_amd64.tar.gz
|
||||
sudo mv helmfile /usr/local/bin/
|
||||
mv helmfile /usr/local/bin/
|
||||
helmfile --version
|
||||
|
||||
- name: Install Helm Diff Plugin
|
||||
@@ -105,7 +105,7 @@ jobs:
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
curl -fsSLO https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
wget https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
tar -xzf helmfile_0.148.0_linux_amd64.tar.gz
|
||||
sudo mv helmfile /usr/local/bin/
|
||||
helmfile --version
|
||||
|
||||
@@ -66,21 +66,13 @@ jobs:
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
curl -fsSLO https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
wget https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
tar -xzf helmfile_0.148.0_linux_amd64.tar.gz
|
||||
sudo mv helmfile /usr/local/bin/
|
||||
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
|
||||
run: |
|
||||
helm env
|
||||
echo "HOME=$HOME"
|
||||
ls -R $HOME/.local/share/helm || true
|
||||
- name: Install Helm Diff Plugin
|
||||
run: helm plugin install https://github.com/databus23/helm-diff || true
|
||||
|
||||
- name: Authenticate with OKE cluster
|
||||
env:
|
||||
@@ -107,5 +99,4 @@ jobs:
|
||||
- name: Run Helmfile Diff
|
||||
env:
|
||||
ENV: ${{ needs.extract_environment.outputs.environment }}
|
||||
HELM_PLUGINS: /home/runner/.local/share/helm/plugins
|
||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml diff
|
||||
|
||||
+3
-4
@@ -1,5 +1,4 @@
|
||||
FROM node:20-alpine AS base_image
|
||||
RUN npm install -g npm@10.8.2
|
||||
FROM node:18.17-alpine AS base_image
|
||||
|
||||
FROM base_image AS build_base
|
||||
WORKDIR /app
|
||||
@@ -22,7 +21,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
# run aws cli without mounting secret, because CI already has AWS credentials
|
||||
FROM build_base AS ci_image
|
||||
RUN aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
RUN npm ci --ignore-scripts
|
||||
RUN npm ci
|
||||
COPY . .
|
||||
|
||||
|
||||
@@ -37,7 +36,7 @@ FROM build_base AS dev
|
||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
||||
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
# flag --build-from-source is required to force-build sqlite3
|
||||
RUN npm ci --ignore-scripts
|
||||
RUN npm ci
|
||||
COPY . .
|
||||
ENTRYPOINT npm run start:dev
|
||||
|
||||
|
||||
@@ -1,47 +0,0 @@
|
||||
FROM node:22-alpine AS base_image
|
||||
RUN npm install -g npm@latest
|
||||
|
||||
FROM base_image AS build_base
|
||||
WORKDIR /app
|
||||
RUN apk update
|
||||
RUN apk add --no-cache \
|
||||
aws-cli \
|
||||
chromium \
|
||||
nss \
|
||||
freetype \
|
||||
harfbuzz \
|
||||
ca-certificates \
|
||||
ttf-freefont
|
||||
COPY package*.json ./
|
||||
|
||||
ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
PUPPETEER_EXECUTABLE_PATH=/usr/bin/chromium-browser
|
||||
|
||||
|
||||
# Local build with secrets
|
||||
FROM build_base AS build
|
||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
||||
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
RUN npm run build
|
||||
|
||||
|
||||
FROM base_image
|
||||
WORKDIR /app
|
||||
COPY --from=build /app/dist ./dist
|
||||
COPY --from=build /app/node_modules ./node_modules
|
||||
COPY --from=build /app/package*.json ./
|
||||
RUN apk update
|
||||
RUN apk add --no-cache \
|
||||
chromium \
|
||||
nss \
|
||||
freetype \
|
||||
harfbuzz \
|
||||
ca-certificates \
|
||||
ttf-freefont
|
||||
|
||||
ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
PUPPETEER_EXECUTABLE_PATH=/usr/bin/chromium-browser
|
||||
|
||||
ENTRYPOINT ["npm", "run", "start:prod"]
|
||||
@@ -2,8 +2,8 @@
|
||||
<image src="./assets/maestro.svg" style="width:10rem">
|
||||
</p>
|
||||
|
||||
# Maestro
|
||||
|
||||
# Maestro
|
||||
|
||||
Maestro é a API principal da Dadosfera. É responsável pela comunicação do Frontend com nossos microsserviços.
|
||||
|
||||
|
||||
Binary file not shown.
@@ -48,9 +48,6 @@ spec:
|
||||
{{- toYaml .Values.resources | nindent 12 }}
|
||||
{{- end }}
|
||||
env:
|
||||
# Auth Provider Configuration (cognito or keycloak)
|
||||
- name: AUTH_PROVIDER
|
||||
value: {{ .Values.maestro.auth_provider | default "cognito" | quote }}
|
||||
- name: AWS_IDENTITY_POOL_ID
|
||||
value: {{ .Values.maestro.aws_identity_pool_id }}
|
||||
- name: AWS_REGION
|
||||
@@ -107,14 +104,6 @@ spec:
|
||||
value: {{ .Values.maestro.redis_host }}
|
||||
- name: REDIS_PORT
|
||||
value: "{{ .Values.maestro.redis_port }}"
|
||||
- name: REDIS_TLS
|
||||
value: "{{ .Values.maestro.redis_tls }}"
|
||||
- name: PLATFORM_API_URL
|
||||
value: {{ .Values.maestro.platform_api_url }}
|
||||
- name: STORAGE_EXPLORER_API_URL
|
||||
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
|
||||
- name: FIREBASE_BASE_URL
|
||||
value: {{ .Values.maestro.firebase_base_url }}
|
||||
- name: JWT_PRIVATE_KEY
|
||||
valueFrom:
|
||||
secretKeyRef:
|
||||
@@ -135,14 +124,3 @@ spec:
|
||||
secretKeyRef:
|
||||
name: prd-{{ .Values.app_name }}
|
||||
key: AWS_DEFAULT_REGION
|
||||
# Elasticsearch
|
||||
- name: ELASTICSEARCH_URL
|
||||
valueFrom:
|
||||
secretKeyRef:
|
||||
name: prd-{{ .Values.app_name }}
|
||||
key: ELASTICSEARCH_URL
|
||||
- name: ELASTICSEARCH_API_KEY
|
||||
valueFrom:
|
||||
secretKeyRef:
|
||||
name: prd-{{ .Values.app_name }}
|
||||
key: ELASTICSEARCH_API_KEY
|
||||
|
||||
@@ -4,18 +4,9 @@ metadata:
|
||||
annotations:
|
||||
nginx.ingress.kubernetes.io/whitelist-source-range: "69.49.241.121/32" # hostgator ip
|
||||
nginx.ingress.kubernetes.io/proxy-body-size: "0"
|
||||
nginx.ingress.kubernetes.io/proxy-read-timeout: "300"
|
||||
nginx.ingress.kubernetes.io/proxy-connect-timeout: "300"
|
||||
nginx.ingress.kubernetes.io/proxy-send-timeout: "300"
|
||||
nginx.ingress.kubernetes.io/server-snippet: |
|
||||
underscores_in_headers on;
|
||||
ignore_invalid_headers on;
|
||||
nginx.ingress.kubernetes.io/proxy-buffer-size: "16k"
|
||||
nginx.ingress.kubernetes.io/proxy-buffers-number: "8"
|
||||
nginx.ingress.kubernetes.io/proxy-busy-buffers-size: "64k"
|
||||
{{- if .Values.maestro.restricted_ip}}
|
||||
nginx.ingress.kubernetes.io/whitelist-source-range: {{ .Values.maestro.restricted_ip }}
|
||||
{{- end }}
|
||||
|
||||
generation: 1
|
||||
labels:
|
||||
|
||||
@@ -38,15 +38,3 @@ spec:
|
||||
version: "AWSCURRENT"
|
||||
property: token
|
||||
|
||||
- secretKey: ELASTICSEARCH_URL
|
||||
remoteRef:
|
||||
key: {{ .Values.maestro.env }}/microservices/elasticsearch
|
||||
version: "AWSCURRENT"
|
||||
property: ELASTICSEARCH_URL
|
||||
|
||||
- secretKey: ELASTICSEARCH_API_KEY
|
||||
remoteRef:
|
||||
key: {{ .Values.maestro.env }}/microservices/elasticsearch
|
||||
version: "AWSCURRENT"
|
||||
property: ELASTICSEARCH_API_KEY
|
||||
|
||||
|
||||
@@ -8,9 +8,6 @@ maestro:
|
||||
open_group_id: e3f98a2f-7748-4981-8505-7695c8ca8218
|
||||
cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36"
|
||||
redis_database: "1"
|
||||
platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com
|
||||
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
|
||||
|
||||
hostname: maestro.stg.dadosfera.ai
|
||||
|
||||
|
||||
@@ -27,9 +27,6 @@ resources:
|
||||
cpu: 2000m
|
||||
memory: 2Gi
|
||||
maestro:
|
||||
# Auth provider: "cognito" (default) or "keycloak"
|
||||
# Note: maestro doesn't connect to Keycloak directly, only duc does
|
||||
auth_provider: "cognito"
|
||||
aws_identity_pool_id: "us-east-1_Mrezsw9Sn"
|
||||
duc_url: duc.dadosfera.ai
|
||||
in_factory_url: in-factory.dadosfera.ai
|
||||
@@ -46,16 +43,12 @@ maestro:
|
||||
upload_file_agent_connection: cbc2f881-58c4-4d60-8003-0979b0b5b911
|
||||
open_customer_id: f239718a-a271-4ef9-ae7e-02a2f0f3aa6e
|
||||
open_group_id: 401573bb-334f-44b2-b30e-88d4cea31ae9
|
||||
platform_api_url: https://oz8v2zid1e.execute-api.us-east-1.amazonaws.com
|
||||
storage_explorer_api_url: "https://storage-explorer-{customer}.dadosfera.ai/api"
|
||||
dedicated_proxy: ""
|
||||
restricted_ip: ""
|
||||
redis_host: "aaapzppmlyamkocqwstpo7zvopczyyiyuy6xzm2g6c5k4mq3a66be4a-0.redis.sa-saopaulo-1.oci.oraclecloud.com"
|
||||
redis_port: "6379"
|
||||
redis_database: "0"
|
||||
redis_tls: "true"
|
||||
cookie_secret: "13cc5e136d3074bcc05bec8697092ec1f5f376bf"
|
||||
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/prd
|
||||
autoscaling:
|
||||
enabled: false
|
||||
minReplicas: 1
|
||||
|
||||
+793
-2780
File diff suppressed because it is too large
Load Diff
Vendored
+8
-1
@@ -16,7 +16,14 @@ declare global {
|
||||
OPEN_CUSTOMER_ID: string;
|
||||
DEDICATED_PROXY: string;
|
||||
COOKIE_SECRET: string;
|
||||
REDIS_TLS?: string;
|
||||
|
||||
// Autodrive Configuration
|
||||
AUTODRIVE_USERNAME?: string;
|
||||
AUTODRIVE_PASSWORD?: string;
|
||||
AUTODRIVE_BASE_URL?: string;
|
||||
AUTODRIVE_MODEL?: string;
|
||||
AUTODRIVE_KEY?: string;
|
||||
AUTO_DRIVE_KEY?: string;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Generated
+1773
-2581
File diff suppressed because it is too large
Load Diff
+6
-20
@@ -10,7 +10,7 @@
|
||||
},
|
||||
"scripts": {
|
||||
"co:login": "aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1",
|
||||
"proto-update": "npm i @dadosfera/protospack-v2@v3.40.0-beta.1 --save-exact",
|
||||
"proto-update": "npm i @dadosfera/protospack-v2@latest --save-exact",
|
||||
"prebuild": "rimraf dist",
|
||||
"build": "nest build",
|
||||
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
|
||||
@@ -27,14 +27,10 @@
|
||||
"test:e2e": "jest --config ./test/jest-e2e.json"
|
||||
},
|
||||
"dependencies": {
|
||||
"@aws-crypto/sha256-js": "^5.2.0",
|
||||
"@aws-sdk/client-dynamodb": "^3.414.0",
|
||||
"@aws-sdk/client-secrets-manager": "^3.414.0",
|
||||
"@aws-sdk/credential-provider-node": "^3.940.0",
|
||||
"@aws-sdk/lib-dynamodb": "^3.414.0",
|
||||
"@aws-sdk/signature-v4": "^3.370.0",
|
||||
"@dadosfera/dadosfera-logs": "^1.0.0-beta.4",
|
||||
"@dadosfera/protospack-v2": "^3.40.0-beta.14",
|
||||
"@dadosfera/protospack": "2.5.3",
|
||||
"@dadosfera/protospack-v2": "3.38.0-beta.10",
|
||||
"@grpc/grpc-js": "^1.9.3",
|
||||
"@grpc/proto-loader": "^0.7.9",
|
||||
"@nestjs/cli": "^9.5.0",
|
||||
@@ -48,7 +44,7 @@
|
||||
"@nestjs/schematics": "^9.2.0",
|
||||
"@nestjs/swagger": "^6.3.0",
|
||||
"@nestjs/testing": "^9.4.3",
|
||||
"axios": "0.30.3",
|
||||
"axios": "^0.27.2",
|
||||
"cache-manager": "^5.1.4",
|
||||
"cache-manager-ioredis-yet": "^1.1.0",
|
||||
"class-transformer": "^0.5.1",
|
||||
@@ -64,7 +60,6 @@
|
||||
"jwk-to-pem": "^2.0.5",
|
||||
"mixpanel": "^0.17.0",
|
||||
"ms": "^3.0.0-canary.1",
|
||||
"multer": "^2.0.2",
|
||||
"openid-client": "^5.7.1",
|
||||
"passport": "^0.6.0",
|
||||
"passport-facebook": "^3.0.0",
|
||||
@@ -80,17 +75,11 @@
|
||||
"swagger-ui-express": "^4.6.3"
|
||||
},
|
||||
"overrides": {
|
||||
"axios": "0.30.3",
|
||||
"form-data": "^4.0.4",
|
||||
"body-parser": "^1.20.3",
|
||||
"cross-spawn": "^7.0.5",
|
||||
"glob": "^10.5.0",
|
||||
"path-to-regexp": "^3.3.0",
|
||||
"semver": "^7.5.2"
|
||||
"multer": "1.4.5-lts.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/cache-manager": "^4.0.6",
|
||||
"@types/cookie-parser": "^1.4.9",
|
||||
"@types/cache-manager": "^4.0.6",
|
||||
"@types/express": "^4.17.17",
|
||||
"@types/express-session": "^1.18.1",
|
||||
"@types/jest": "27.0.2",
|
||||
@@ -117,8 +106,5 @@
|
||||
"ts-node": "^10.9.1",
|
||||
"tsconfig-paths": "^3.14.2",
|
||||
"typescript": "^4.9.5"
|
||||
},
|
||||
"resolutions": {
|
||||
"axios": "0.30.3"
|
||||
}
|
||||
}
|
||||
|
||||
+2
-7
@@ -17,6 +17,7 @@ import { ConnectionTestModule } from './modules/connection-test/connection-test.
|
||||
import { NetworkConfigModule } from './modules/network-config/network-config.module';
|
||||
import { InputsModule } from './modules/inputs/inputs.module';
|
||||
import { OauthModule } from './modules/oauth/oauth.module';
|
||||
import { PipelinesModule } from './modules/pipelines/pipelines.module';
|
||||
import { TransformationsModule } from './modules/transformations/transformations.module';
|
||||
import { HealthModule } from './modules/health/health.module';
|
||||
import { CatalogModule } from './modules/catalog/catalog.module';
|
||||
@@ -32,10 +33,6 @@ import { NetworkPolicyModule } from './modules/network-policy/network-policy.mod
|
||||
import { AssignModule } from './modules/assign/assign.module';
|
||||
import { ShareMetadataModule } from './modules/share-metadata/share-metadata.module';
|
||||
import { ApiKeyModule } from './modules/api-key/api-key.module';
|
||||
import { PlatformApiModule } from './modules/platform-api/platform-api.module';
|
||||
import { StorageExplorerModule } from './modules/storage-explorer/storage-explorer.module';
|
||||
import { ReleaseNoteModule } from './modules/release_note/release_note.module';
|
||||
|
||||
|
||||
@Module({
|
||||
providers: [
|
||||
@@ -59,6 +56,7 @@ import { ReleaseNoteModule } from './modules/release_note/release_note.module';
|
||||
PermissionsModule,
|
||||
TermsOfUseModule,
|
||||
ConnectionTestModule,
|
||||
PipelinesModule,
|
||||
TransformationsModule,
|
||||
UsersModule,
|
||||
RolesModule,
|
||||
@@ -75,11 +73,8 @@ import { ReleaseNoteModule } from './modules/release_note/release_note.module';
|
||||
ApiKeyModule,
|
||||
IdentityProviderModule,
|
||||
NetworkPolicyModule,
|
||||
PlatformApiModule,
|
||||
StorageExplorerModule,
|
||||
//Always leave HealthModule last, so it is on the bottom of swagger
|
||||
HealthModule,
|
||||
ReleaseNoteModule,
|
||||
],
|
||||
})
|
||||
export class AppModule {}
|
||||
|
||||
@@ -153,7 +153,6 @@ export class AuthenticationGuard
|
||||
user_id: accessTokenPayload.user_id,
|
||||
username: accessTokenPayload.username,
|
||||
permissions: accessTokenPayload.permissions,
|
||||
roles: accessTokenPayload.roles,
|
||||
customer_id: accessTokenPayload.customer_id,
|
||||
customer_name: accessTokenPayload.customer_name,
|
||||
customer_tier: accessTokenPayload.customer_tier,
|
||||
|
||||
@@ -1,21 +0,0 @@
|
||||
import jwt, { JwtPayload } from 'jsonwebtoken';
|
||||
|
||||
export function extractUserFrom(aRawJwt: string) {
|
||||
const decodedToken = jwt.decode(aRawJwt, {
|
||||
complete: true,
|
||||
});
|
||||
|
||||
const payload = decodedToken.payload as JwtPayload;
|
||||
|
||||
return {
|
||||
user_id: payload.user_id,
|
||||
username: payload.username,
|
||||
permissions: payload.permissions,
|
||||
roles: payload.roles,
|
||||
customer_id: payload.customer_id,
|
||||
customer_name: payload.customer_name,
|
||||
customer_tier: payload.customer_tier,
|
||||
customer_modules: payload.customer_modules,
|
||||
access_token: aRawJwt,
|
||||
}
|
||||
}
|
||||
@@ -116,44 +116,6 @@ export const PERMISSIONS_GROUPS = {
|
||||
},
|
||||
},
|
||||
},
|
||||
IMPORT_FILES: {
|
||||
title: {
|
||||
'pt-br': 'Coletar | Importar arquivos',
|
||||
'en-us': 'Collect | Import files',
|
||||
'es-es': 'Colecta | Importar archivos',
|
||||
},
|
||||
permissions: {
|
||||
VIEW: {
|
||||
seqid: 48,
|
||||
claim: 'import-file:view',
|
||||
usage: PermissionUsages.PUBLIC,
|
||||
name: {
|
||||
'pt-br': 'Importar arquivos',
|
||||
'en-us': 'Import files',
|
||||
'es-es': 'Importar archivos',
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
AI_CHAT: {
|
||||
title: {
|
||||
'pt-br': 'AutodriveDDF',
|
||||
'en-us': 'AutodriveDDF',
|
||||
'es-es': 'AutodriveDDF',
|
||||
},
|
||||
permissions: {
|
||||
VIEW: {
|
||||
seqid: 49,
|
||||
claim: 'ai-chat:view',
|
||||
usage: PermissionUsages.PUBLIC,
|
||||
name: {
|
||||
'pt-br': 'AutodriveDDF',
|
||||
'en-us': 'AutodriveDDF',
|
||||
'es-es': 'AutodriveDDF',
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
CONNECTION: {
|
||||
title: {
|
||||
'pt-br': 'Coletar | Fontes de dados',
|
||||
@@ -357,16 +319,6 @@ export const PERMISSIONS_GROUPS = {
|
||||
'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: {
|
||||
seqid: 1,
|
||||
claim: 'catalog:delete',
|
||||
@@ -400,25 +352,6 @@ export const PERMISSIONS_GROUPS = {
|
||||
},
|
||||
},
|
||||
},
|
||||
LINEAGE: {
|
||||
title: {
|
||||
'pt-br': 'Explorar | Linhagem',
|
||||
'en-us': 'Explore | Lineage',
|
||||
'es-es': 'Explorar | Linaje',
|
||||
},
|
||||
permissions: {
|
||||
VIEW: {
|
||||
seqid: 50,
|
||||
claim: 'lineage:view',
|
||||
usage: PermissionUsages.PUBLIC,
|
||||
name: {
|
||||
'pt-br': 'Acessar ao módulo de Linhagem',
|
||||
'en-us': 'Access to Lineage module',
|
||||
'es-es': 'Acceda al módulo de Linaje',
|
||||
},
|
||||
}
|
||||
},
|
||||
},
|
||||
EMBED: {
|
||||
title: {
|
||||
'pt-br': 'Analisar | Incorporação',
|
||||
@@ -678,35 +611,6 @@ export const PERMISSIONS_GROUPS = {
|
||||
},
|
||||
},
|
||||
},
|
||||
STORAGE_EXPLORER: {
|
||||
title: {
|
||||
'pt-br': 'Storage Explorer',
|
||||
'en-us': 'Storage Explorer',
|
||||
'es-es': 'Storage Explorer',
|
||||
},
|
||||
permissions: {
|
||||
READ: {
|
||||
seqid: 51,
|
||||
claim: 'storage-explorer:read',
|
||||
usage: PermissionUsages.PUBLIC,
|
||||
name: {
|
||||
'pt-br': 'Ler dados do Storage Explorer',
|
||||
'en-us': 'Read Storage Explorer data',
|
||||
'es-es': 'Leer datos del Storage Explorer',
|
||||
},
|
||||
},
|
||||
WRITE: {
|
||||
seqid: 52,
|
||||
claim: 'storage-explorer:write',
|
||||
usage: PermissionUsages.PUBLIC,
|
||||
name: {
|
||||
'pt-br': 'Escrever dados no Storage Explorer',
|
||||
'en-us': 'Write Storage Explorer data',
|
||||
'es-es': 'Escribir datos en Storage Explorer',
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
};
|
||||
export interface DadosferaModule {
|
||||
name: string;
|
||||
|
||||
@@ -12,7 +12,6 @@ export interface RequestUser {
|
||||
customer_tier: string;
|
||||
access_token: string;
|
||||
customer_modules: string[];
|
||||
roles: string[];
|
||||
}
|
||||
|
||||
export const User: (options?: { required?: boolean }) => ParameterDecorator =
|
||||
|
||||
@@ -1,65 +0,0 @@
|
||||
import {
|
||||
BadRequestException,
|
||||
CanActivate,
|
||||
ExecutionContext,
|
||||
Inject,
|
||||
Injectable,
|
||||
OnModuleInit,
|
||||
} from '@nestjs/common';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import { map, Observable } from 'rxjs';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import {
|
||||
ReadService,
|
||||
ProtoServices,
|
||||
} from '@dadosfera/protospack-v2/dist/lib/PipelineV2';
|
||||
import { PipelinesClientConfiguration } from 'src/modules/pipelinesV2/pipelines-client';
|
||||
import { PlatformApiService } from 'src/modules/platform-api/platform-api.service';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@Injectable()
|
||||
export class PipelineExecutionGuard implements CanActivate {
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private readonly platformApiService: PlatformApiService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
async canActivate(context: ExecutionContext): Promise<boolean> {
|
||||
try {
|
||||
this.logger.info(
|
||||
'PipelineExecutionGuard: Checking if pipeline can be executed...',
|
||||
);
|
||||
const request = context.switchToHttp().getRequest();
|
||||
const pipelineId = request.params.pipelineId;
|
||||
const user = request.user;
|
||||
const idRegex = /[^0-9a-zA-Z_$]+/g;
|
||||
const convertedId = pipelineId.replace(idRegex, '_');
|
||||
|
||||
const status = await this.platformApiService.proxy(
|
||||
'GET',
|
||||
`/pipeline/${convertedId}/pipeline_run`,
|
||||
user,
|
||||
);
|
||||
|
||||
const currentStatus = status[status.length - 1]
|
||||
|
||||
this.logger.info('Pipeline current status response:' + JSON.stringify(currentStatus));
|
||||
|
||||
if (currentStatus.last_status.toLowerCase() === 'running') {
|
||||
this.logger.error('Pipeline is running, cannot update input now');
|
||||
throw new BadRequestException('Pipeline is running, cannot update input now');
|
||||
} else {
|
||||
return true;
|
||||
}
|
||||
} catch (error) {
|
||||
this.logger.error('Error in PipelineExecutionGuard: ' + error.message);
|
||||
throw new BadRequestException('Error checking pipeline status: ' + error.message);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
+4
-25
@@ -18,41 +18,21 @@ async function bootstrap() {
|
||||
});
|
||||
const logger = new DadosferaLogger();
|
||||
|
||||
const corsOrigins = [];
|
||||
|
||||
if (process.env.ENV === 'local') {
|
||||
corsOrigins.push('http://localhost:4200');
|
||||
} else {
|
||||
corsOrigins.push(
|
||||
'https://app.stg.dadosfera.ai',
|
||||
'https://app.dadosfera.ai',
|
||||
'https://private-frontend.stg.dadosfera.ai',
|
||||
'https://unimed.dadosfera.ai',
|
||||
'https://boston-scientific.dadosfera.ai',
|
||||
'https://plataforma.dadosfera.ai'
|
||||
);
|
||||
}
|
||||
|
||||
const app = await NestFactory.create(AppModule, {
|
||||
logger,
|
||||
cors: {
|
||||
origin: corsOrigins,
|
||||
origin: '*',
|
||||
methods: 'GET,HEAD,PUT,PATCH,POST,DELETE',
|
||||
preflightContinue: false,
|
||||
optionsSuccessStatus: 204,
|
||||
credentials: true,
|
||||
credentials: true
|
||||
},
|
||||
});
|
||||
|
||||
app.use(helmet());
|
||||
app.use(cookieParser(process.env.COOKIE_SECRET));
|
||||
|
||||
if (process.env.ENV !== 'local') {
|
||||
if (process.env.ENV === 'prd') {
|
||||
app.use('/catalog/register-dataset', json({ limit: '10mb' }));
|
||||
app.use(
|
||||
'/catalog/register-dataset',
|
||||
urlencoded({ extended: true, limit: '10mb' }),
|
||||
);
|
||||
app.use('/catalog/register-dataset', urlencoded({ extended: true, limit: '10mb' }));
|
||||
}
|
||||
|
||||
configureSwagger(app);
|
||||
@@ -111,4 +91,3 @@ function configureSwagger(app: INestApplication) {
|
||||
);
|
||||
}
|
||||
bootstrap();
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Controller, Body, Put, Get, NotFoundException} from '@nestjs/common';
|
||||
import { Controller, Post, Body, Put, Get} from '@nestjs/common';
|
||||
import { AssignService } from './assign.service';
|
||||
import { CreateAssignDto } from './dto/create-assign.dto';
|
||||
import { Authenticated, RequireModule, RequireSomePermission } from 'src/decorators/authentication.decorator';
|
||||
@@ -28,10 +28,6 @@ export class AssignController {
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.EMBED_ASSIGNED)
|
||||
async get(@User() user: RequestUser) {
|
||||
const metadata = PackTheMetadata(user);
|
||||
try {
|
||||
return await this.assignService.get(metadata);
|
||||
} catch (error) {
|
||||
throw new NotFoundException(error.message)
|
||||
}
|
||||
return await this.assignService.get(metadata);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,6 @@ import {
|
||||
Req,
|
||||
Param,
|
||||
Res,
|
||||
UnauthorizedException,
|
||||
} from '@nestjs/common';
|
||||
import {
|
||||
ApiHeaders,
|
||||
@@ -55,7 +54,7 @@ import jwt, { JwtPayload } from 'jsonwebtoken';
|
||||
import { LanguageEnum } from 'src/utils/languages.enum';
|
||||
import { Language } from 'src/decorators/language.decorator';
|
||||
import { ApiInternalOnlyEndpoint } from 'src/decorators/swagger.decorator';
|
||||
import { ApiKeyService } from 'src/modules/api-key/api-key.service';
|
||||
import { Cookie } from 'express-session';
|
||||
|
||||
type CookiesValues = {
|
||||
accessToken?: string;
|
||||
@@ -75,7 +74,6 @@ export class AuthController {
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private authClient: AuthClientService,
|
||||
private apiKeyService: ApiKeyService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
|
||||
@@ -105,13 +103,14 @@ export class AuthController {
|
||||
const data = await this.authClient.signIn({ username, password, totp }, metadata);
|
||||
|
||||
if (data.tokens) {
|
||||
this.authClient.writeAuthSession(res, {
|
||||
this.addTokenInCookie(res, {
|
||||
accessToken: data.tokens.accessToken,
|
||||
refreshToken: data.tokens.refreshToken,
|
||||
userId: data.user.id
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
return res.send(data);
|
||||
} catch (error) {
|
||||
this.logger.error('/auth - SignIn - ERROR', error);
|
||||
@@ -128,8 +127,27 @@ export class AuthController {
|
||||
) {
|
||||
try {
|
||||
this.logger.info('/auth - SignOut');
|
||||
|
||||
this.authClient.cleanUpAuthSession(res);
|
||||
const exp = 1000 * 60 * 3;
|
||||
|
||||
res.cookie('ddf-auth', '', {
|
||||
domain: 'dadosfera.local',
|
||||
maxAge: Date.now() - exp,
|
||||
expires: new Date(),
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
|
||||
res.cookie('ddf-refresh-auth', '', {
|
||||
domain: 'dadosfera.local',
|
||||
maxAge: Date.now() - exp,
|
||||
expires: new Date(),
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
|
||||
this.logger.info('Clean cookie sessions');
|
||||
|
||||
return res.send();
|
||||
} catch (error) {
|
||||
@@ -158,9 +176,8 @@ export class AuthController {
|
||||
|
||||
const data = await this.authClient.refreshAccessToken({ refreshToken, userId }, metadata);
|
||||
|
||||
this.authClient.writeAuthSession(res, {
|
||||
this.addTokenInCookie(res, {
|
||||
accessToken: data.accessToken,
|
||||
refreshToken: data.refreshToken,
|
||||
userId
|
||||
});
|
||||
|
||||
@@ -192,14 +209,13 @@ export class AuthController {
|
||||
) {
|
||||
this.logger.info('/auth - change-password');
|
||||
|
||||
const { oldPassword, newPassword, totpCode } = body;
|
||||
const { oldPassword, newPassword } = body;
|
||||
const { authorization: accessToken } = headers;
|
||||
|
||||
return this.authClient.changePassword({
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -217,8 +233,7 @@ export class AuthController {
|
||||
|
||||
const { username } = body;
|
||||
|
||||
await this.authClient.resetPassword({ username }, metadata);
|
||||
return { authProvider: process.env.AUTH_PROVIDER || 'cognito' };
|
||||
return this.authClient.resetPassword({ username }, metadata);
|
||||
}
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@@ -478,57 +493,114 @@ export class AuthController {
|
||||
@Get('me')
|
||||
async getMe(@Req() req: Request, @Res() res: Response) {
|
||||
this.logger.info('GET /auth/me ')
|
||||
this.logger.info(JSON.stringify(req.headers));
|
||||
// Lê cookies
|
||||
const accessToken = req.cookies['ddf-auth'];
|
||||
const userId = req.cookies['ddf-user-id'];
|
||||
|
||||
// Check for API key header first
|
||||
const apiKey = req.get('X-Api-key');
|
||||
if (apiKey) {
|
||||
this.logger.info('Authenticating via X-Api-key header');
|
||||
const { api_key } = await this.apiKeyService.get(apiKey);
|
||||
|
||||
const userDto = {
|
||||
id: api_key.user_id,
|
||||
name: api_key.username,
|
||||
email: api_key.username,
|
||||
this.logger.info('Has cookie: ' + Boolean(accessToken))
|
||||
let payload: any;
|
||||
let userInfo: any = {};
|
||||
try {
|
||||
// Decodifica e valida o JWT de acesso
|
||||
const decoded: any = accessToken && jwt.decode(accessToken, { complete: true });
|
||||
if (!decoded) throw new Error('Invalid token')
|
||||
const { kid } = decoded.header;
|
||||
// Busca a chave pública
|
||||
const { keys } = await this.authClient.getPublicKeys();
|
||||
const pemValue = keys.find((k) => k.kid === kid)?.pem;
|
||||
if (!pemValue) throw new Error('Public key not found');
|
||||
jwt.verify(accessToken, pemValue);
|
||||
payload = decoded.payload;
|
||||
userInfo = {
|
||||
id: payload.user_id,
|
||||
name: payload.username,
|
||||
customer: {
|
||||
id: api_key.customer_id,
|
||||
name: api_key.customer_name,
|
||||
tier: api_key.customer_tier,
|
||||
id: payload.customer_id,
|
||||
name: payload.customer_name,
|
||||
tier: payload.customer_tier,
|
||||
}
|
||||
};
|
||||
return res.status(200).json(userInfo);
|
||||
} catch (err) {
|
||||
this.logger.error(err.message);
|
||||
const refreshToken = req.cookies['ddf-refresh-auth'];
|
||||
|
||||
return res.status(200).json(userDto);
|
||||
}
|
||||
|
||||
// Get token and headers
|
||||
const accessToken = req.cookies['ddf-auth'];
|
||||
const refreshToken = req.cookies['ddf-refresh-auth'];
|
||||
const userId = req.cookies['ddf-user-id'];
|
||||
const resourceHost = req.headers["x-original-url"] as string || "" ;
|
||||
|
||||
const hasUserSession = Boolean(accessToken) && Boolean(userId);
|
||||
this.logger.info('Has User Session: ' + hasUserSession);
|
||||
|
||||
if (!hasUserSession) {
|
||||
throw new UnauthorizedException()
|
||||
}
|
||||
|
||||
try {
|
||||
const userDto = await this.authClient.validateUserSession(accessToken, resourceHost);
|
||||
return res.status(200).json(userDto);
|
||||
} catch (error) {
|
||||
|
||||
if (!refreshToken) {
|
||||
this.logger.info('Token is invalid')
|
||||
this.logger.info('Has Refresh Token: '+ Boolean(refreshToken))
|
||||
// Se access token inválido, tenta refresh
|
||||
if (!refreshToken || !userId) {
|
||||
this.logger.error('Invalid refresh token or customer name');
|
||||
throw new UnauthorizedException("Invalid refresh token or customer name");
|
||||
};
|
||||
return res.status(401).json({ error: 'Not authenticated' });
|
||||
}
|
||||
try {
|
||||
// Chama refreshAccessToken
|
||||
const metadata = PackTheMetadata({
|
||||
});
|
||||
this.logger.info('Call Refresh Token')
|
||||
const data = await this.authClient.refreshAccessToken({ refreshToken, userId }, metadata);
|
||||
this.logger.info('Finish Refresh Token')
|
||||
// Retorna novo access token e dados mínimos
|
||||
this.addTokenInCookie(res, {
|
||||
accessToken: data.accessToken,
|
||||
userId
|
||||
});
|
||||
// Decodifica novo token
|
||||
const decoded: any = jwt.decode(data.accessToken, { complete: true });
|
||||
const payload = decoded.payload;
|
||||
userInfo = {
|
||||
id: payload.user_id,
|
||||
name: payload.username,
|
||||
customer: {
|
||||
id: payload.customer_id,
|
||||
name: payload.customer_name,
|
||||
tier: payload.customer_tier,
|
||||
}
|
||||
};
|
||||
return res.status(200).json(userInfo);
|
||||
} catch (refreshErr) {
|
||||
this.logger.error(refreshErr)
|
||||
return res.status(401).json({ error: 'Not authenticated' });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const {
|
||||
authSession,
|
||||
user
|
||||
} = await this.authClient.refreshUserSession(refreshToken, userId, resourceHost);
|
||||
this.authClient.writeAuthSession(res, authSession);
|
||||
return res.status(200).json(user);
|
||||
private addTokenInCookie(res: Response, data: CookiesValues) {
|
||||
let exp = 1000 * 60 * 5; // 5 minutes
|
||||
|
||||
if (data.accessToken) {
|
||||
const { exp: expiration } = jwt.decode(data.accessToken) as JwtPayload;
|
||||
exp = (expiration - 30) * 1000; // exp em segundos, maxAge em ms
|
||||
|
||||
this.logger.info('Set Cookie ddf-auth')
|
||||
res.cookie('ddf-auth', data.accessToken, {
|
||||
domain: 'stg.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
|
||||
if (data.refreshToken) {
|
||||
this.logger.info('Set Cookie ddf-refresh-auth')
|
||||
res.cookie('ddf-refresh-auth', data.refreshToken, {
|
||||
domain: 'stg.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
|
||||
if (data.userId) {
|
||||
this.logger.info('Set Cookie ddf-refresh-auth')
|
||||
res.cookie('ddf-user-id', data.userId, {
|
||||
domain: 'stg.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8,11 +8,10 @@ import { AuthClientService } from './auth.service';
|
||||
import { DucClient } from '../duc/client.config';
|
||||
import { GoogleLoginStrategy } from './passport-strategies/google-strategy';
|
||||
import { getOauthSecrets } from 'src/utils/OauthSecrets';
|
||||
import { ApiKeyModule } from '../api-key/api-key.module';
|
||||
const client = new DucClient();
|
||||
|
||||
@Module({
|
||||
imports: [ClientsModule.register([client.providerOptions]), ApiKeyModule],
|
||||
imports: [ClientsModule.register([client.providerOptions])],
|
||||
controllers: [AuthController],
|
||||
providers: [
|
||||
AuthClientService,
|
||||
|
||||
@@ -1,20 +1,10 @@
|
||||
import {
|
||||
OnModuleInit,
|
||||
Inject,
|
||||
Injectable,
|
||||
ForbiddenException,
|
||||
HttpException,
|
||||
HttpStatus,
|
||||
} from '@nestjs/common';
|
||||
import { OnModuleInit, Inject, Injectable, ForbiddenException } from '@nestjs/common';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { ProtoServices } from '@dadosfera/protospack-v2/dist/lib/Duc';
|
||||
import {
|
||||
AuthProtoService as AuthServiceInterface,
|
||||
UsersProtoService,
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import { AuthProtoService as AuthServiceInterface, IdentityProviderProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import {
|
||||
AuthSnowflakeSignInRequest,
|
||||
AuthSignInRequest,
|
||||
@@ -31,24 +21,18 @@ import {
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { DucClient } from '../duc/client.config';
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
import { BulkEditResponse, UserDTO } from './dtos/login';
|
||||
import jwt, { JwtPayload } from 'jsonwebtoken';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import { Request, Response } from 'express';
|
||||
import { BulkEditResponse } from './dtos/login';
|
||||
|
||||
type AuthSession = {
|
||||
accessToken?: string;
|
||||
refreshToken?: string;
|
||||
userId?: string;
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AuthClientService implements OnModuleInit {
|
||||
|
||||
|
||||
logger: DadosferaLogger;
|
||||
|
||||
private authService: AuthServiceInterface;
|
||||
private userService: UsersProtoService;
|
||||
|
||||
private authService: AuthServiceInterface;
|
||||
private identityProviderService: IdentityProviderProtoService;
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
@@ -62,8 +46,8 @@ export class AuthClientService implements OnModuleInit {
|
||||
ProtoServices.AuthProtoService,
|
||||
);
|
||||
|
||||
this.userService = this.grpcClient.getService<UsersProtoService>(
|
||||
ProtoServices.UsersProtoService,
|
||||
this.identityProviderService = this.grpcClient.getService<IdentityProviderProtoService>(
|
||||
ProtoServices.IdentityProviderProtoService,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -79,11 +63,11 @@ export class AuthClientService implements OnModuleInit {
|
||||
return lastValueFrom(this.authService.AuthSnowflakeSignIn(input));
|
||||
}
|
||||
|
||||
checkDedicatedProxy({ customer }: AuthSignInResponse) {
|
||||
checkDedicatedProxy({
|
||||
customer
|
||||
}: AuthSignInResponse) {
|
||||
const DEDICATED_PROXY = process.env.DEDICATED_PROXY || '';
|
||||
this.logger.info(
|
||||
'SignIn - Setting customer ID for dedicated proxy: ' + DEDICATED_PROXY,
|
||||
);
|
||||
this.logger.info('SignIn - Setting customer ID for dedicated proxy: ' + DEDICATED_PROXY);
|
||||
this.logger.info('Customer ID: ' + customer.id);
|
||||
|
||||
if (DEDICATED_PROXY !== '' && DEDICATED_PROXY !== customer.id) {
|
||||
@@ -91,9 +75,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
}
|
||||
|
||||
// Bloquear o customer de acesso o maestro publico
|
||||
this.logger.info(
|
||||
'Check if customer have network policy: ' + customer.modules,
|
||||
);
|
||||
this.logger.info('Check if customer have network policy: ' + customer.modules);
|
||||
const hasNetworkPolicyModule = customer.modules.includes('network-policy');
|
||||
if (hasNetworkPolicyModule && DEDICATED_PROXY === '') {
|
||||
throw new ForbiddenException();
|
||||
@@ -112,6 +94,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
result = await lastValueFrom(
|
||||
this.authService.AuthSignIn({ username, password, totp }, metadata),
|
||||
);
|
||||
|
||||
} catch (error) {
|
||||
this.logger.error('SignIn - Error during sign-in');
|
||||
this.logger.error(error);
|
||||
@@ -122,7 +105,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
this.checkDedicatedProxy(result);
|
||||
}
|
||||
|
||||
return result;
|
||||
return result
|
||||
}
|
||||
|
||||
async refreshAccessToken(
|
||||
@@ -132,10 +115,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
this.logger.info('RefreshAccessToken');
|
||||
|
||||
return lastValueFrom(
|
||||
this.authService.AuthRefreshAccessToken(
|
||||
{ refreshToken, userId },
|
||||
metadata,
|
||||
),
|
||||
this.authService.AuthRefreshAccessToken({ refreshToken, userId }, metadata),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -143,7 +123,6 @@ export class AuthClientService implements OnModuleInit {
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
}: AuthChangePasswordRequest) {
|
||||
this.logger.info('ChangePassword');
|
||||
|
||||
@@ -152,7 +131,6 @@ export class AuthClientService implements OnModuleInit {
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
}),
|
||||
);
|
||||
}
|
||||
@@ -305,191 +283,4 @@ export class AuthClientService implements OnModuleInit {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
public async validateUserSession(accessToken: any, resourceHost: string) {
|
||||
const payload = await this.validateJwtToken(accessToken);
|
||||
|
||||
const userDto = await this.getUserfromPayload(payload);
|
||||
|
||||
this.validateResourceAccess(resourceHost, userDto);
|
||||
return userDto;
|
||||
}
|
||||
|
||||
public async refreshUserSession(
|
||||
refreshToken: string,
|
||||
userId: string,
|
||||
originHeader: string,
|
||||
): Promise<{
|
||||
user: UserDTO;
|
||||
authSession: AuthSession;
|
||||
}> {
|
||||
const metadata = PackTheMetadata({});
|
||||
|
||||
this.logger.info('Call Refresh Token');
|
||||
const refreshCredentials = await this.refreshAccessToken(
|
||||
{ refreshToken, userId },
|
||||
metadata,
|
||||
);
|
||||
this.logger.info('Finish Refresh Token');
|
||||
|
||||
const userDto = await this.validateUserSession(
|
||||
refreshCredentials.accessToken,
|
||||
originHeader,
|
||||
);
|
||||
return {
|
||||
user: userDto,
|
||||
authSession: {
|
||||
accessToken: refreshCredentials.accessToken,
|
||||
refreshToken: refreshCredentials.refreshToken,
|
||||
userId,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
public writeAuthSession(res: Response, data: AuthSession) {
|
||||
let exp = 1000 * 60 * 5; // 5 minutes
|
||||
|
||||
if (data.accessToken) {
|
||||
const { exp: expiration } = jwt.decode(data.accessToken) as JwtPayload;
|
||||
exp = (expiration - 30) * 1000; // exp em segundos, maxAge em ms
|
||||
|
||||
this.logger.info('Set Cookie ddf-auth');
|
||||
res.cookie('ddf-auth', data.accessToken, {
|
||||
domain: '.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
|
||||
if (data.refreshToken) {
|
||||
this.logger.info('Set Cookie ddf-refresh-auth');
|
||||
res.cookie('ddf-refresh-auth', data.refreshToken, {
|
||||
domain: '.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
|
||||
if (data.userId) {
|
||||
this.logger.info('Set Cookie ddf-refresh-auth');
|
||||
res.cookie('ddf-user-id', data.userId, {
|
||||
domain: '.dadosfera.ai',
|
||||
maxAge: exp,
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
public cleanUpAuthSession(res: Response) {
|
||||
const exp = 1000 * 60 * 3;
|
||||
|
||||
res.cookie('ddf-auth', '', {
|
||||
domain: 'dadosfera.ai',
|
||||
maxAge: Date.now() - exp,
|
||||
expires: new Date(),
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
|
||||
res.cookie('ddf-refresh-auth', '', {
|
||||
domain: 'dadosfera.ai',
|
||||
maxAge: Date.now() - exp,
|
||||
expires: new Date(),
|
||||
httpOnly: true,
|
||||
secure: true,
|
||||
sameSite: 'none', // Necessário para cookies em requisições cross-site
|
||||
});
|
||||
|
||||
this.logger.info('Clean cookie sessions');
|
||||
}
|
||||
|
||||
private async validateJwtToken(token: string) {
|
||||
const decoded: any = token && jwt.decode(token, { complete: true });
|
||||
if (!decoded) throw new Error('Invalid token');
|
||||
|
||||
const { kid } = decoded.header;
|
||||
// Busca a chave pública
|
||||
const { keys } = await this.getPublicKeys();
|
||||
const pemValue = keys.find((k) => k.kid === kid)?.pem;
|
||||
if (!pemValue) throw new Error('Public key not found');
|
||||
jwt.verify(token, pemValue);
|
||||
|
||||
return decoded.payload;
|
||||
}
|
||||
|
||||
private async getUserfromPayload(payload: JwtPayload): Promise<UserDTO> {
|
||||
this.logger.info('getUser');
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_id: payload.customer_id,
|
||||
});
|
||||
|
||||
const { user } = await lastValueFrom(
|
||||
this.userService.UserFindOneById({ id: payload.user_id }, metadata),
|
||||
);
|
||||
|
||||
const userDto: UserDTO = {
|
||||
id: user.id,
|
||||
name: user.name,
|
||||
email: user.email,
|
||||
jobTitle: user?.jobTitle || null,
|
||||
department: user?.department || null,
|
||||
hierarchy: user?.hierarchy || null,
|
||||
customer: {
|
||||
id: payload.customer_id,
|
||||
name: payload.customer_name,
|
||||
tier: payload.customer_tier,
|
||||
},
|
||||
};
|
||||
|
||||
return userDto;
|
||||
}
|
||||
|
||||
private validateResourceAccess(host: string, user: UserDTO) {
|
||||
this.logger.info(
|
||||
"Validate whether the source URL is a resource belonging to the user's client",
|
||||
);
|
||||
this.logger.info('Host: ' + host);
|
||||
this.logger.info('Customer: ' + user.customer.name);
|
||||
|
||||
const hostParts = host.split('.');
|
||||
const domain = hostParts[0];
|
||||
const isResouceStg = hostParts[1] === 'stg';
|
||||
|
||||
const notFoundCustomerInDomain = !domain.includes('-')
|
||||
|
||||
if (notFoundCustomerInDomain) {
|
||||
this.logger.info(`Not found Customer Name in domain`);
|
||||
return;
|
||||
}
|
||||
|
||||
const domainParts = domain.split('-');
|
||||
|
||||
const customerInDomain = domainParts[domainParts.length - 1];
|
||||
|
||||
if (isResouceStg && process.env.ENV !== 'stg') {
|
||||
this.logger.error(`Customer ${user.customer.name} cannot access ${host}`);
|
||||
throw new HttpException(
|
||||
`Customer ${user.customer.name} cannot access ${host}`,
|
||||
HttpStatus.FORBIDDEN
|
||||
);
|
||||
}
|
||||
|
||||
if (customerInDomain != user.customer.name) {
|
||||
this.logger.error(`Customer ${user.customer.name} cannot access ${host}`);
|
||||
throw new HttpException(
|
||||
`Customer ${user.customer.name} cannot access ${host}`,
|
||||
HttpStatus.FORBIDDEN
|
||||
);
|
||||
}
|
||||
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -140,17 +140,3 @@ export interface BulkEditResponse {
|
||||
successfulUsers: string[];
|
||||
failedUsers: string[];
|
||||
}
|
||||
|
||||
export type UserDTO = {
|
||||
id: string,
|
||||
name: string,
|
||||
email: string,
|
||||
jobTitle?: string,
|
||||
department?: string,
|
||||
hierarchy?: string,
|
||||
customer: {
|
||||
id: string,
|
||||
name: string,
|
||||
tier: string,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,12 +17,14 @@ import {
|
||||
HttpStatus,
|
||||
Res,
|
||||
} from '@nestjs/common';
|
||||
import { ValidationPipe } from '../../pipes/object-validation.pipe';
|
||||
import {
|
||||
ApiCreatedResponse,
|
||||
ApiHeaders,
|
||||
ApiOkResponse,
|
||||
ApiTags,
|
||||
ApiOperation,
|
||||
ApiParam,
|
||||
ApiResponse,
|
||||
} from '@nestjs/swagger';
|
||||
import {
|
||||
Authenticated,
|
||||
@@ -47,7 +49,6 @@ import {
|
||||
IMakeAComment,
|
||||
IOneDataAsset,
|
||||
IPreviewResponse,
|
||||
IUpdateCertificationStatusRequest,
|
||||
IUpdateDataRequest,
|
||||
TriggerCatalogReq,
|
||||
TriggerCatalogRes,
|
||||
@@ -90,10 +91,6 @@ export class CatalogController {
|
||||
@Query() query: ICatalogAllRequest,
|
||||
): Promise<ICatalogAllResponse> {
|
||||
const { user_id, customer_name, customer_id, username, permissions } = user;
|
||||
this.logger.info(`/catalog - searchCatalog`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const is_data_manager = permissions.includes(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
|
||||
@@ -130,10 +127,6 @@ export class CatalogController {
|
||||
@Res() res: Response
|
||||
) {
|
||||
const { user_id, customer_name, customer_id, username, permissions } = user;
|
||||
this.logger.info(`/catalog/download - searchCatalog`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const is_data_manager = permissions.includes(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
|
||||
@@ -170,10 +163,6 @@ export class CatalogController {
|
||||
async findByPipelineAndObject(@User() user: RequestUser, @Query() query) {
|
||||
const { username, user_id, customer_id, customer_name, permissions } = user;
|
||||
const { pipeline, object } = query;
|
||||
this.logger.info(`/catalog - ON GET DATA ASSET BY PIPELINE AND OBJECT`, {
|
||||
username,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
if (!pipeline || !object) {
|
||||
throw new BadRequestException('Query params not provided');
|
||||
@@ -226,10 +215,6 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
async findAllTags(@Body() body) {
|
||||
this.logger.info(`/catalog - ON FIND ALL TAGS ROUTE`, {
|
||||
user: body.info.user_id,
|
||||
customer: body.info.customer,
|
||||
});
|
||||
|
||||
const { user_id, customer, customer_id } = body.info;
|
||||
const metadata = PackTheMetadata({
|
||||
@@ -243,51 +228,6 @@ export class CatalogController {
|
||||
return res;
|
||||
}
|
||||
|
||||
@Get('schemas')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
async findSchemas(@User() user: RequestUser) {
|
||||
const { username, user_id, customer_id, customer_name } = user;
|
||||
this.logger.info(`/catalog - ON FIND SCHEMAS ROUTE`, {
|
||||
username,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
username,
|
||||
user_id,
|
||||
customer_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
try {
|
||||
const res = await this.catalogService.findSchemas(metadata);
|
||||
return res;
|
||||
} catch (error) {
|
||||
throw new HttpException(error.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@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')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
@@ -300,10 +240,6 @@ export class CatalogController {
|
||||
) {
|
||||
const { username, user_id, customer_id, customer_name, permissions } = user;
|
||||
|
||||
this.logger.info(`GET /data-asset/${id}`, {
|
||||
username,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const is_data_manager = permissions.includes(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
|
||||
@@ -356,10 +292,6 @@ export class CatalogController {
|
||||
) {
|
||||
const { username, user_id, customer_id, customer_name, permissions } = user;
|
||||
|
||||
this.logger.info(`/catalog - ON GET ONE DASHBOARD METABASE ROUTE`, {
|
||||
username,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const is_data_manager = permissions.includes(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
|
||||
@@ -411,10 +343,6 @@ export class CatalogController {
|
||||
): Promise<IColumnsMetadataResponse> {
|
||||
const { customer_name, customer_id, user_id, username } = user;
|
||||
|
||||
this.logger.info(`/catalog - columns-metadata`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_name,
|
||||
@@ -442,10 +370,6 @@ export class CatalogController {
|
||||
): Promise<IPreviewResponse> {
|
||||
const { customer_name, customer_id, user_id, username, customer_modules } = user;
|
||||
|
||||
this.logger.info(`/catalog - ON GET DATA DOCS ROUTE`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_name,
|
||||
@@ -470,26 +394,32 @@ export class CatalogController {
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@Param('id') id: string,
|
||||
@Query('asset_type') asset_type: string,
|
||||
): Promise<IDocsResponse> {
|
||||
const { customer_name, customer_id, user_id, username } = user;
|
||||
try {
|
||||
const { customer_name, customer_id, user_id, username } = user;
|
||||
|
||||
this.logger.info(`/catalog - ON GET DATA DOCS ROUTE`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
this.logger.info(`/catalog - ON GET DATA DOCS ROUTE`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
id,
|
||||
});
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_name,
|
||||
customer_id,
|
||||
user_id,
|
||||
username,
|
||||
language,
|
||||
});
|
||||
const metadata = PackTheMetadata({
|
||||
customer_name,
|
||||
customer_id,
|
||||
user_id,
|
||||
username,
|
||||
language,
|
||||
});
|
||||
|
||||
const docs = await this.catalogService.getDataDocs(id, asset_type, metadata);
|
||||
const docs = await this.catalogService.getDataDocs(id, metadata);
|
||||
|
||||
return { docs };
|
||||
return { docs };
|
||||
} catch (error) {
|
||||
this.logger.error(`Error in getDataAssetDocs for id ${id}: ${error.message}`);
|
||||
this.logger.error(`Error details: ${JSON.stringify(error)}`);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
@Put('data-asset/:id')
|
||||
@@ -512,8 +442,6 @@ export class CatalogController {
|
||||
language,
|
||||
});
|
||||
|
||||
delete (body as any).certification_status;
|
||||
|
||||
const result = await this.catalogService.updateOneDataAsset({
|
||||
body,
|
||||
data_asset_id,
|
||||
@@ -527,33 +455,6 @@ export class CatalogController {
|
||||
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')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
||||
@@ -564,36 +465,22 @@ export class CatalogController {
|
||||
@Headers() headers,
|
||||
@Param('id') table_id: string,
|
||||
@Body('docs') docs: string,
|
||||
@Query('asset_type') asset_type: string,
|
||||
) {
|
||||
const { user_id, customer_name, customer_id, username } = user;
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_id,
|
||||
customer_name,
|
||||
user_id,
|
||||
username,
|
||||
});
|
||||
const { user_id, customer_name } = user;
|
||||
|
||||
this.logger.info(`/catalog - ON POST DATA DOCS ROUTE`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const body = {
|
||||
const res = await this.catalogService.createDataDocs({
|
||||
table_id,
|
||||
docs,
|
||||
asset_type,
|
||||
info: {
|
||||
customer: customer_name,
|
||||
},
|
||||
}
|
||||
|
||||
const res = await this.catalogService.createDataDocs(body, metadata);
|
||||
});
|
||||
|
||||
return res;
|
||||
}
|
||||
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Put('data-asset/:id/manage-permissions')
|
||||
async manageDataAssetPermissions(
|
||||
@@ -1019,4 +906,88 @@ export class CatalogController {
|
||||
this.logger.error(error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Post('/data-asset/:nimbus_id/docs/ai')
|
||||
@ApiOperation({
|
||||
summary: 'Save documentation for data asset',
|
||||
description: 'Saves documentation content for a data asset',
|
||||
})
|
||||
@ApiParam({
|
||||
name: 'nimbus_id',
|
||||
description: 'Nimbus ID of the data asset',
|
||||
type: 'string',
|
||||
})
|
||||
@ApiResponse({
|
||||
status: 201,
|
||||
description: 'Documentation saved successfully',
|
||||
})
|
||||
async saveDocumentation(
|
||||
@Param('nimbus_id') nimbusId: string,
|
||||
@Body() body: { docs: string },
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
|
||||
const metadata = PackTheMetadata(user);
|
||||
|
||||
try {
|
||||
await this.catalogService.updateDataAssetDocumentation(nimbusId, body.docs, metadata);
|
||||
return {
|
||||
message: 'Documentation saved successfully',
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error(`Error saving documentation for ${nimbusId}: ${error.message}`);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
@Post('/data-asset/:nimbus_id/docs/generate-ai')
|
||||
@ApiOperation({
|
||||
summary: 'Generate AI documentation for data asset',
|
||||
description: 'Generates comprehensive documentation for a data asset using AI (Autodrive)',
|
||||
})
|
||||
@ApiParam({
|
||||
name: 'nimbus_id',
|
||||
description: 'Nimbus ID of the data asset',
|
||||
type: 'string',
|
||||
})
|
||||
@ApiResponse({
|
||||
status: 201,
|
||||
description: 'AI documentation generated successfully',
|
||||
schema: {
|
||||
type: 'object',
|
||||
properties: {
|
||||
message: { type: 'string' },
|
||||
documentation: { type: 'string' },
|
||||
},
|
||||
},
|
||||
})
|
||||
@ApiResponse({
|
||||
status: 400,
|
||||
description: 'Bad request - invalid nimbus_id or missing data',
|
||||
})
|
||||
@ApiResponse({
|
||||
status: 404,
|
||||
description: 'Data asset not found',
|
||||
})
|
||||
@ApiResponse({
|
||||
status: 500,
|
||||
description: 'Internal server error during AI generation',
|
||||
})
|
||||
async generateAiDocumentation(
|
||||
@Param('nimbus_id') dataAssetId: string,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
const metadata = PackTheMetadata(user);
|
||||
|
||||
try {
|
||||
const result = await this.catalogService.generateAiDocumentation(dataAssetId, metadata, user);
|
||||
return {
|
||||
message: 'AI documentation generated successfully',
|
||||
documentation: result,
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error(`Error generating AI documentation for ${dataAssetId}: ${error.message}`);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5,17 +5,20 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { CatalogController } from './catalog.controller';
|
||||
import { CatalogClientConfiguration } from './catalog-client';
|
||||
import { ClientsModule } from '@nestjs/microservices';
|
||||
import { PipelinesModule as OldPipelineModule } from 'src/modules/pipelines/pipelines.module';
|
||||
import { UsersModule } from '../users/users.module';
|
||||
import { RolesModule } from '../roles/roles.module';
|
||||
import { CustomersModule } from '../customers/customers.module';
|
||||
import { ShareModule } from './share/share.module';
|
||||
import { CatalogService } from './catalog.service';
|
||||
import { MixpanelModule } from '../mixpanel/mixpanel.module';
|
||||
|
||||
const client = new CatalogClientConfiguration();
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
ClientsModule.register([client.providerOptions]),
|
||||
OldPipelineModule,
|
||||
UsersModule,
|
||||
RolesModule,
|
||||
CustomersModule,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,42 @@
|
||||
/**
|
||||
* Constantes relacionadas ao Autodrive
|
||||
*
|
||||
*/
|
||||
export const AUTODRIVE_CONSTANTS = {
|
||||
// URLs e endpoints (apenas do ENV)
|
||||
BASE_URL: process.env.BASE_URL_AUTODRIVE || process.env.AUTODRIVE_BASE_URL,
|
||||
|
||||
// Credenciais (apenas do ENV, sem fallback para segurança)
|
||||
USERNAME: process.env.AUTODRIVE_USERNAME,
|
||||
PASSWORD: process.env.AUTODRIVE_PASSWORD,
|
||||
|
||||
// Modelo padrão
|
||||
DEFAULT_MODEL: process.env.AUTODRIVE_MODEL || "gpt-4o",
|
||||
|
||||
// Timeouts (em milissegundos)
|
||||
ASK_TIMEOUT: 120000,
|
||||
ANSWER_TIMEOUT: 180000,
|
||||
|
||||
// Headers
|
||||
HEADERS: {
|
||||
'Content-Type': 'application/json',
|
||||
},
|
||||
} as const;
|
||||
|
||||
/**
|
||||
* Chaves geográficas para detecção de dados de localização
|
||||
*/
|
||||
export const GEOGRAPHIC_KEYS = [
|
||||
'country', 'countries', 'city', 'cities',
|
||||
'region', 'regions', 'location', 'state',
|
||||
'states', 'address'
|
||||
] as const;
|
||||
|
||||
/**
|
||||
* Países comuns para detecção automática
|
||||
*/
|
||||
export const COMMON_COUNTRIES = [
|
||||
'brazil', 'brasil', 'usa', 'united states',
|
||||
'canada', 'mexico', 'argentina', 'chile',
|
||||
'colombia'
|
||||
] as const;
|
||||
@@ -1,10 +1,4 @@
|
||||
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';
|
||||
|
||||
export enum DataAssetShareType {
|
||||
@@ -12,12 +6,6 @@ export enum DataAssetShareType {
|
||||
public = 'public',
|
||||
private = 'private',
|
||||
}
|
||||
export enum CertificationStatus {
|
||||
draft = 'draft',
|
||||
in_review = 'in_review',
|
||||
approved = 'approved',
|
||||
deprecated = 'deprecated',
|
||||
}
|
||||
export enum OrderEnum {
|
||||
asc = 'asc',
|
||||
desc = 'desc',
|
||||
@@ -110,8 +98,6 @@ export class IDataAsset {
|
||||
embed?: EmbedObject;
|
||||
@ApiPropertyOptional({ enum: DataAssetShareType })
|
||||
share_type?: DataAssetShareType;
|
||||
@ApiPropertyOptional()
|
||||
docs?: string;
|
||||
}
|
||||
|
||||
export class IOneDataAsset {
|
||||
@@ -161,24 +147,6 @@ export class ICatalogAllRequest {
|
||||
description: 'Tipo de ordenação - `asc`: crescente; `desc`: decrescente ',
|
||||
})
|
||||
order?: OrderEnum;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description: 'ID do usuário owner para filtrar data assets',
|
||||
example: 'user-id-1,user-id-2',
|
||||
})
|
||||
owner?: string;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description: 'Data inicial para filtro de catálogo (formato: YYYY-MM-DD)',
|
||||
example: '2025-01-01',
|
||||
})
|
||||
catalog_date_from?: string;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description: 'Data final para filtro de catálogo (formato: YYYY-MM-DD)',
|
||||
example: '2025-12-31',
|
||||
})
|
||||
catalog_date_to?: string;
|
||||
}
|
||||
|
||||
export class ICatalogAllResponse {
|
||||
@@ -203,27 +171,6 @@ export class IData {
|
||||
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 {
|
||||
@ApiProperty()
|
||||
name: string;
|
||||
@@ -235,18 +182,7 @@ export class IUpdateDataRequest {
|
||||
embed: EmbedObject;
|
||||
@ApiPropertyOptional({ enum: DataAssetShareType })
|
||||
share_type?: DataAssetShareType;
|
||||
@ApiPropertyOptional()
|
||||
docs?: string;
|
||||
@ApiPropertyOptional({ type: [CustomPropertyDto] })
|
||||
custom_properties?: CustomPropertyDto[];
|
||||
}
|
||||
|
||||
export class IUpdateCertificationStatusRequest {
|
||||
@ApiProperty({ enum: CertificationStatus })
|
||||
@IsEnum(CertificationStatus)
|
||||
certification_status: CertificationStatus;
|
||||
}
|
||||
|
||||
export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||
@ApiProperty()
|
||||
display_name: string;
|
||||
@@ -260,8 +196,6 @@ export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||
location: string;
|
||||
@ApiPropertyOptional()
|
||||
embed: EmbedObject;
|
||||
@ApiPropertyOptional()
|
||||
docs: string;
|
||||
}
|
||||
|
||||
export class IPreview {
|
||||
@@ -394,10 +328,3 @@ export type AssetReporter = {
|
||||
created_at: string;
|
||||
tags: string;
|
||||
}
|
||||
|
||||
export type CreateDataDocsDTO = {
|
||||
table_id: string;
|
||||
docs: string;
|
||||
asset_type: string;
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
/**
|
||||
|
||||
*/
|
||||
export const AI_DOCUMENTATION_PROMPT = `crie uma documentação em Portugues, Ingles e Espanhol seguindo essas instruções
|
||||
1. Persona: como profissional de governança e engenharia de dados
|
||||
2. Tarefa: ao receber as informações da tabela criar uma documentação com o seguinte escopo
|
||||
**A primeira linha do documento tem que conter a seguinte informação: ## Document languages: EN / BR / ES
|
||||
**A segunda linha tem que obrigatoriamente conter a escrita Table: nome da tabela
|
||||
**A terceira linha tem que obrigatoriamente conter a escrita Table Schema: nome do table schema
|
||||
**DIRETRIZ CRUCIAL DE CONSISTÊNCIA E COMPLETUDE DE SCHEMA:**
|
||||
**1. Fonte Exclusiva de Metadados:** O 'Table Schema' definido na linha acima é a ÚNICA fonte de verdade para o schema dos dados a serem documentados. TODAS as informações subsequentes, especialmente na seção 'Estrutura da Tabela' (incluindo a lista de colunas, seus nomes, tipos de dados, descrições e exemplos) DEVEM ser extraídas EXCLUSIVAMENTE de metadados que correspondem a ESTE 'Table Schema'. Se os dados de entrada que você recebeu contiverem informações para a mesma tabela ou colunas mas de schemas diferentes (ex: um schema 'bronze' e um 'silver'), você DEVE IGNORAR TOTALMENTE as informações dos schemas divergentes para esta tarefa de documentação e utilizar APENAS as do 'Table Schema' aqui especificado.
|
||||
**2. Listagem Completa de Colunas:** Sua principal tarefa na seção 'Estrutura da Tabela' é identificar e listar TODAS as colunas que pertencem ao 'Table Schema' especificado. Verifique nos dados de entrada fornecidos se há uma indicação explícita do número total de colunas para esta tabela neste schema (por exemplo, um campo como 'Num_columns' ou similar nos metadados da tabela). Você deve se esforçar para listar exatamente essa quantidade de colunas. Se essa contagem não estiver disponível, liste todas as colunas que você puder identificar como pertencentes exclusivamente a este 'Table Schema'. A completude em relação ao schema especificado é essencial.
|
||||
|
||||
**Depois de "Estrutura da tablea", incluir a mensagem "Este documento foi gerado por IA", traduzida corretamente para cada idioma.**
|
||||
**Obrigatoriamente:Após finalizar a versão em Inglês, começar a versão em Português** **Após finalizar a versão em Português, começar a versão em Espanhol** **Antes de começar cada versão, colocar um título como:** - \`## English Version\` (para inglês)
|
||||
- \`## Versão em Português\` (para português)
|
||||
- \`## Versión en Español\` (para espanhol)
|
||||
- Descrição: fornece uma visão geral do ativo de dados,
|
||||
destacando seu propósito e principal funcionalidade.
|
||||
Esta sessão resume o conteúdo e o objetivo do ativo, ajudando os usuários a entender rapidamente o que o ativo representa
|
||||
e como pode ser utilizado em suas análises e decisões.
|
||||
- Sugestão de Domínio de Dados:
|
||||
Analise cuidadosamente os dados da tabela e sugira o domínio mais apropriado. Inclua:
|
||||
- Domínio Sugerido: [Nome do domínio]
|
||||
- Motivo: [Explicação breve sobre porque a tabela pertence a este domínio]
|
||||
- Observações: [Qualquer observação adicional relevante]
|
||||
|
||||
Exemplos de Domínios de Dados para referência:
|
||||
- Financeiro: Dados sobre transações, receitas, despesas, etc.
|
||||
- Recursos Humanos: Dados sobre funcionários, cargos, salários, etc.
|
||||
- Produtos: Dados sobre produtos, categorias, preços, etc.
|
||||
- Fornecedores: Dados sobre fornecedores, produtos fornecidos, localizações, etc.
|
||||
- Marketing: Dados sobre campanhas, leads, conversões, etc.
|
||||
- Vendas: Dados sobre vendas, clientes, produtos vendidos, etc.
|
||||
- Operações: Dados sobre processos, logística, produção, etc.
|
||||
- Clientes: Dados sobre clientes, interações, histórico, etc.
|
||||
-Tags Sugeridas:
|
||||
A IA deve gerar tags relevantes **com base nos dados da tabela**.
|
||||
- **IMPORTANTE: Analise cuidadosamente os dados de preview da tabela (PREVIEW DATA) para encontrar países. Procure em todas as colunas por nomes de países, cidades ou regiões.**
|
||||
- **Garanta que as tags estejam separadas por espaços vazios, todas na mesma linha, exemplo: #marketing #sales #australia #canada, limitar até 3 países que mais aparecem** - **Os países DEVEM ser extraídos dos dados de preview da tabela. Procure em colunas como City, Country, Region, Location, etc.** - Por que esta tabela é interessante:
|
||||
Nesta sessão, é destacada a importância do ativo, explicando como ele pode ser útil para os usuários.
|
||||
São abordadas as formas como o ativo pode melhorar a tomada de decisões, identificar padrões relevantes ou fornecer insights valiosos.
|
||||
O objetivo é ressaltar a utilidade prática e o impacto positivo que o ativo pode ter em suas atividades.
|
||||
- Análises potencialmente úteis feitas com esses dados:
|
||||
Aqui são listadas algumas das análises que podem ser realizadas com o ativo de dados. Inclui sugestões de dashboards,
|
||||
relatórios ou outros tipos de análises que aproveitam as informações fornecidas pelo ativo.
|
||||
O objetivo é oferecer maneiras de utilizar os dados para obter insights valiosos e apoiar a tomada de decisões informadas.
|
||||
- Links Úteis:
|
||||
Os Links Úteis oferecem recursos adicionais relacionados ao ativo de dados, incluindo guias,
|
||||
artigos ou outras fontes de informação que podem ajudar os usuários a compreender melhor o ativo e suas aplicações. Além disso,
|
||||
inclui um link rápido dentro da Dadosfera para ativos relacionados diretamente com o ativo em questão, facilitando a navegação entre os ativos.
|
||||
- Estrutura da Tabela:
|
||||
A Estrutura da Tabela detalha TODAS as colunas e os dados disponíveis no ativo, conforme pertencentes ao 'Table Schema' principal definido no início deste documento.
|
||||
**Instrução Detalhada para Estrutura da Tabela:**
|
||||
Siga rigorosamente estes passos:
|
||||
1. Identifique nos dados de entrada (metadados da tabela e das colunas) todas as colunas que pertencem EXCLUSIVAMENTE ao 'Table Schema' especificado no cabeçalho deste documento. Se houver uma contagem de colunas (ex: 'Num_columns') para este schema específico, assegure-se de listar essa quantidade.
|
||||
2. Para CADA uma dessas colunas identificadas, formate a saída da seguinte maneira, **SEM utilizar NENHUM marcador de lista (como traços ou asteriscos) no início de cada entrada de coluna**. Cada coluna deve ser apresentada como um bloco de texto. Inclua uma linha em branco entre a documentação de cada coluna para separação visual.
|
||||
- Apresente o NOME_DA_COLUNA em maiúsculas, seguido pelo (TIPO_DE_DADO_EXTRAÍDO_DOS_METADADOS_DO_SCHEMA_CORRETO) entre parênteses.
|
||||
- O **NOME_DA_COLUNA (TIPO_DE_DADO_EXTRAÍDO_DOS_METADADOS_DO_SCHEMA_CORRETO)** deve estar na primeira linha do bloco da coluna e **inteiramente em negrito**.
|
||||
- Na linha seguinte, a etiqueta "**Descrição:**" deve estar **em negrito**, seguida pelo texto da descrição da coluna.
|
||||
- Na linha seguinte à descrição, a etiqueta "**Exemplo:**" deve estar **em negrito**, seguida pelo valor do exemplo. Se o exemplo for um valor literal ou código, formate-o entre crases (\`) se apropriado.
|
||||
- Se houver informações adicionais relevantes (como "Valores Possíveis:", "Observações:", etc.), coloque a etiqueta correspondente **em negrito** em uma nova linha, seguida pelo seu texto.
|
||||
|
||||
Este documento foi gerado por IA.
|
||||
|
||||
NOME_COLUNA_1 (TIPO_DADO_SCHEMA_CORRETO_1):
|
||||
Descrição: [Descrição da coluna 1, do schema correto]
|
||||
Exemplo: \`[Exemplo de valor para coluna 1, do schema correto]\`
|
||||
|
||||
NOME_COLUNA_2 (TIPO_DADO_SCHEMA_CORRETO_2):
|
||||
Descrição: [Descrição da coluna 2, do schema correto]
|
||||
Exemplo: \`[Exemplo de valor para coluna 2, do schema correto]\`
|
||||
|
||||
(continue este formato com início de cada coluna para TODAS as colunas do 'Table Schema' especificado, garanta com que NUNCA tenha TRAÇO OU PONTO no inicio)
|
||||
|
||||
3. Contexto : O usuário ira cadastrar um ativo de dados na nossa plataforma e para ter um bom catalogo ele ira querer gerar a documentação padronizada mas explicativa e
|
||||
automática
|
||||
4. Restrições : A documentação deve seguir obrigatoriamente o mesmo padrão principalmente na parte de estrutura de dados
|
||||
5. Objetivo: O principal objetivo é gerar uma documentação acessível, clara,
|
||||
automática e padronizada para os usuários que desejem cadastrar um ativo de dados na plataforma`;
|
||||
|
||||
/**
|
||||
* Configurações para a geração de documentação com IA
|
||||
*/
|
||||
export const AI_DOCUMENTATION_CONFIG = {
|
||||
FETCH_K: 250,
|
||||
K: 100,
|
||||
} as const;
|
||||
@@ -160,7 +160,7 @@ export class ShareService implements OnModuleInit {
|
||||
});
|
||||
|
||||
const { documentation } = await lastValueFrom(
|
||||
this.catalogReadService.GetDatasetDoc({ id }, metadata),
|
||||
this.catalogReadService.GetDatasetDoc({ id, type: undefined }, metadata),
|
||||
);
|
||||
console.log(documentation);
|
||||
const docs = JSON.parse(documentation);
|
||||
@@ -180,7 +180,7 @@ export class ShareService implements OnModuleInit {
|
||||
return data_assets.map((data_asset) => {
|
||||
const owner = customer_users.find(
|
||||
(u) => u.id === data_asset.owner,
|
||||
)?.email;
|
||||
)?.username;
|
||||
|
||||
const roles = [];
|
||||
const users = [];
|
||||
@@ -190,7 +190,7 @@ export class ShareService implements OnModuleInit {
|
||||
}
|
||||
for (const user_id of data_asset.users) {
|
||||
const user = customer_users.find((r) => r.id === user_id);
|
||||
if (user) users.push({ id: user.id, email: user.email });
|
||||
if (user) users.push({ id: user.id, username: user.username });
|
||||
}
|
||||
return {
|
||||
...data_asset,
|
||||
@@ -262,7 +262,6 @@ export class ShareService implements OnModuleInit {
|
||||
user_id: accessTokenPayload.user_id,
|
||||
username: accessTokenPayload.username,
|
||||
permissions: accessTokenPayload.permissions,
|
||||
roles: accessTokenPayload.roles,
|
||||
customer_id: accessTokenPayload.customer_id,
|
||||
customer_name: accessTokenPayload.customer_name,
|
||||
customer_tier: accessTokenPayload.customer_tier,
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
/**
|
||||
* Tipos relacionados à geração de documentação com IA
|
||||
*/
|
||||
|
||||
export interface AutodriveCredentials {
|
||||
username: string;
|
||||
password: string;
|
||||
baseUrl: string;
|
||||
model: string;
|
||||
authHeader?: string;
|
||||
}
|
||||
|
||||
export interface AutodriveAskPayload {
|
||||
question: string;
|
||||
fetch_k: number;
|
||||
k: number;
|
||||
model: string;
|
||||
}
|
||||
|
||||
export interface AutodriveAskResponse {
|
||||
answer?: string;
|
||||
question_id?: string;
|
||||
dataset_id?: string;
|
||||
}
|
||||
|
||||
export interface AutodriveAnswerResponse {
|
||||
status: 'started' | 'success' | 'failed';
|
||||
answer?: string;
|
||||
status_reason?: string;
|
||||
}
|
||||
|
||||
export interface AutodriveUploadResponse {
|
||||
dataset_id: string;
|
||||
}
|
||||
|
||||
export interface DatasetStatusResponse {
|
||||
status: 'processing' | 'success' | 'failed';
|
||||
status_reason?: string;
|
||||
}
|
||||
|
||||
export interface ColumnData {
|
||||
name: string;
|
||||
type: string;
|
||||
description?: string;
|
||||
nullable?: string;
|
||||
}
|
||||
|
||||
export interface ColumnsMetadata {
|
||||
columns: ColumnData[];
|
||||
}
|
||||
|
||||
export interface DataPreview {
|
||||
preview: any[];
|
||||
}
|
||||
|
||||
export interface FormattedDataForAI {
|
||||
dataAsset: any;
|
||||
dataPreview: any[];
|
||||
columnsData?: ColumnsMetadata;
|
||||
}
|
||||
@@ -136,32 +136,4 @@ export class CustomersController {
|
||||
const result = await this.customersService.getAccessDashboardUrl(user.customer_name, metadata);
|
||||
return result;
|
||||
}
|
||||
|
||||
@Get(':id/organization-info')
|
||||
@Authenticated()
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.USERS.permissions.ADMIN)
|
||||
@ApiOkResponse({ description: 'Organization information' })
|
||||
async getOrganizationInfo(@Param('id') id: string) {
|
||||
this.logger.info('getOrganizationInfo', { id });
|
||||
return this.customersService.getOrganizationInfo(id);
|
||||
}
|
||||
|
||||
@Put(':id/organization-info')
|
||||
@Authenticated()
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.USERS.permissions.ADMIN)
|
||||
@HttpCode(HttpStatus.OK)
|
||||
@ApiOkResponse({ description: 'Organization information updated' })
|
||||
async updateOrganizationInfo(
|
||||
@Param('id') id: string,
|
||||
@Body() body: {
|
||||
companyName: string;
|
||||
companySite: string;
|
||||
domain: string;
|
||||
cnpj: string;
|
||||
description: string;
|
||||
},
|
||||
) {
|
||||
return this.customersService.updateOrganizationInfo(id, body);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -9,12 +9,12 @@ import {
|
||||
} from '@nestjs/common';
|
||||
|
||||
import { firstValueFrom, lastValueFrom } from 'rxjs';
|
||||
import { Link } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/entities';
|
||||
import { DucClient } from '../duc/client.config';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import { ProtoServices } from '@dadosfera/protospack-v2/dist/lib/Duc';
|
||||
import { CustomerSetLinksRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { CustomerUpdateRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { CustomersProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import { CustomerLinksConfig } from './dtos/customers';
|
||||
import ErrorCodes from 'src/utils/errorCodes';
|
||||
import jwt from 'jsonwebtoken';
|
||||
import {
|
||||
@@ -67,12 +67,12 @@ export class CustomersService implements OnModuleInit {
|
||||
)
|
||||
}
|
||||
|
||||
async getLinks(customerId: string): Promise<CustomerLinksConfig | null> {
|
||||
async getLinks(customerId: string) {
|
||||
try {
|
||||
const result = await lastValueFrom(
|
||||
this.customerService.CustomerGetLinks({ customerId }),
|
||||
this.customerService.CustomerFindOneById({ id: customerId }),
|
||||
);
|
||||
return (result.links as CustomerLinksConfig) || null;
|
||||
return result.customer?.links || [];
|
||||
} catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
||||
@@ -80,17 +80,17 @@ export class CustomersService implements OnModuleInit {
|
||||
}
|
||||
}
|
||||
|
||||
async setLinks(customerId: string, links: CustomerLinksConfig) {
|
||||
async setLinks(customerId: string, links: Link[]) {
|
||||
if (!customerId || !links) {
|
||||
throw new HttpException(null, HttpStatus.BAD_REQUEST);
|
||||
}
|
||||
|
||||
try {
|
||||
return await firstValueFrom(
|
||||
this.customerService.CustomerSetLinks({
|
||||
customerId,
|
||||
links: links as CustomerSetLinksRequest['links'],
|
||||
}),
|
||||
this.customerService.CustomerUpdate({
|
||||
id: customerId,
|
||||
links,
|
||||
} as CustomerUpdateRequest),
|
||||
);
|
||||
} catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||
@@ -223,57 +223,4 @@ export class CustomersService implements OnModuleInit {
|
||||
})
|
||||
)
|
||||
}
|
||||
|
||||
async updateOrganizationInfo(
|
||||
customerId: string,
|
||||
data: {
|
||||
companyName: string;
|
||||
companySite: string;
|
||||
domain: string;
|
||||
cnpj: string;
|
||||
description: string;
|
||||
},
|
||||
) {
|
||||
try {
|
||||
const result = await lastValueFrom(
|
||||
this.customerService.OrganizationUpdate({
|
||||
customerId,
|
||||
companyName: data.companyName || '',
|
||||
companySite: data.companySite || '',
|
||||
domain: data.domain || '',
|
||||
cnpj: data.cnpj || '',
|
||||
description: data.description || '',
|
||||
}),
|
||||
);
|
||||
|
||||
return result;
|
||||
} catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
||||
else throw err;
|
||||
}
|
||||
}
|
||||
|
||||
async getOrganizationInfo(customerId: string) {
|
||||
try {
|
||||
const customerResponse = await lastValueFrom(
|
||||
this.customerService.CustomerFindOneById({ id: customerId })
|
||||
);
|
||||
|
||||
const customer = customerResponse.customer;
|
||||
|
||||
return {
|
||||
companyName: customer.companyName || '',
|
||||
companySite: customer.companySite || '',
|
||||
domain: customer.domain || '',
|
||||
cnpj: customer.cnpj || '',
|
||||
description: customer.description || ''
|
||||
};
|
||||
} catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
||||
else throw err;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { Link } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/entities';
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
|
||||
export class CustomerLinkItem {
|
||||
export class CustomerLink implements Link {
|
||||
@ApiProperty()
|
||||
href: string;
|
||||
@ApiProperty()
|
||||
@@ -8,59 +9,15 @@ export class CustomerLinkItem {
|
||||
@ApiProperty()
|
||||
description: string;
|
||||
@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 {
|
||||
@ApiProperty({ type: CustomerLinksConfig })
|
||||
links: CustomerLinksConfig;
|
||||
@ApiProperty({ type: [CustomerLink] })
|
||||
links: CustomerLink[];
|
||||
}
|
||||
|
||||
export class CustomerLinksResponse {
|
||||
@ApiPropertyOptional({ type: CustomerLinksConfig })
|
||||
links?: CustomerLinksConfig;
|
||||
}
|
||||
@ApiProperty({ type: [CustomerLink] })
|
||||
links: CustomerLink[];
|
||||
}
|
||||
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
|
||||
export class OrganizationUpdateRequest {
|
||||
@ApiProperty()
|
||||
name: string;
|
||||
@ApiPropertyOptional()
|
||||
companySite: string;
|
||||
@ApiProperty()
|
||||
domain: string;
|
||||
@ApiPropertyOptional()
|
||||
info: string;
|
||||
@ApiPropertyOptional()
|
||||
cnpj: string;
|
||||
}
|
||||
|
||||
export class OrganizationResponse {
|
||||
@ApiProperty()
|
||||
name: string;
|
||||
@ApiPropertyOptional()
|
||||
companySite: string;
|
||||
@ApiProperty()
|
||||
domain: string;
|
||||
@ApiPropertyOptional()
|
||||
info: string;
|
||||
@ApiPropertyOptional()
|
||||
cnpj: string;
|
||||
}
|
||||
@@ -11,19 +11,10 @@ export class TableColumns {
|
||||
name: string;
|
||||
@ApiProperty()
|
||||
columns: string[];
|
||||
@ApiPropertyOptional({ type: [Column] })
|
||||
@ApiProperty()
|
||||
references: Column[];
|
||||
@ApiProperty()
|
||||
destination: Record<'raw' | 'qualify', {
|
||||
table_name: string;
|
||||
table_schema: string;
|
||||
}> | null;
|
||||
@ApiProperty()
|
||||
type: string;
|
||||
@ApiPropertyOptional({ type: [String] })
|
||||
identifier_columns?: string[];
|
||||
@ApiPropertyOptional({ type: Column })
|
||||
reference_column?: Column;
|
||||
}
|
||||
export class AvailableEntity {
|
||||
@ApiProperty()
|
||||
|
||||
@@ -1,8 +1,4 @@
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
interface Values {
|
||||
jdbc_user: string;
|
||||
|
||||
@@ -99,7 +99,6 @@ export class InputsController {
|
||||
customer: info.customer,
|
||||
});
|
||||
|
||||
this.logger.info(JSON.stringify(body))
|
||||
const response = await this.inputService.create({ body, info });
|
||||
|
||||
return response;
|
||||
|
||||
@@ -17,8 +17,6 @@ import {
|
||||
InputCreateGenericRequest,
|
||||
InputCreateS3Request,
|
||||
InputNewCreateRequest,
|
||||
InputUpdateResponse,
|
||||
RollbackInputRequest,
|
||||
TestConnectionRequest,
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
|
||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
||||
@@ -73,8 +71,8 @@ export class InputsService {
|
||||
objectCamelToSnake(createInputResponse);
|
||||
return createInputResponse;
|
||||
},
|
||||
update: async (updateInputDTO: UpdateInputRequest): Promise<InputUpdateResponse> => {
|
||||
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
||||
update: async (updateInputDTO: UpdateInputRequest) => {
|
||||
this.logger.info('InputClientService - Update');
|
||||
const updateInputResponse = await lastValueFrom(
|
||||
this.inputWriteService.InputUpdate(updateInputDTO),
|
||||
);
|
||||
@@ -166,11 +164,6 @@ export class InputsService {
|
||||
const inputCreateGenericRequest: InputCreateGenericRequest = {
|
||||
input: {
|
||||
...body,
|
||||
tables: (body.tables || []).map((table) => ({
|
||||
...table,
|
||||
identifier_columns: table.identifier_columns || [],
|
||||
reference_column: table.reference_column || table.references?.[0],
|
||||
})),
|
||||
},
|
||||
info,
|
||||
};
|
||||
@@ -207,46 +200,23 @@ export class InputsService {
|
||||
}
|
||||
|
||||
async update(id: string, data, info: Info) {
|
||||
// this.validateCron({ ...data, info });
|
||||
this.validateCron({ ...data, info });
|
||||
try {
|
||||
const {
|
||||
tablesUpdate,
|
||||
dataAssetUpdate,
|
||||
input
|
||||
} = await this.OLD_inputClient.update({
|
||||
const updateInputResponse: any = await this.OLD_inputClient.update({
|
||||
id,
|
||||
...data,
|
||||
info,
|
||||
...data,
|
||||
});
|
||||
|
||||
const updateInputResponse = this.adjustInputPayload(
|
||||
input,
|
||||
updateInputResponse.input = this.adjustInputPayload(
|
||||
updateInputResponse?.input,
|
||||
);
|
||||
return {
|
||||
input: updateInputResponse,
|
||||
tablesUpdate,
|
||||
dataAssetUpdate
|
||||
};
|
||||
return updateInputResponse;
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
|
||||
async rollbackUpdate(
|
||||
data: RollbackInputRequest
|
||||
) {
|
||||
this.logger.info('PipelinesClientService - rollbackUpdate');
|
||||
this.logger.info('Rolling back input update with data: ' + JSON.stringify(data));
|
||||
const updatePipelineResponse = await lastValueFrom(
|
||||
this.inputWriteService.RollbackInputUpdate(
|
||||
data
|
||||
),
|
||||
);
|
||||
this.logger.info('Done');
|
||||
|
||||
return updatePipelineResponse;
|
||||
}
|
||||
|
||||
async remove(idRequest: IIdRequest) {
|
||||
return lastValueFrom(this.inputWriteService.InputRemove(idRequest));
|
||||
}
|
||||
@@ -288,12 +258,4 @@ export class InputsService {
|
||||
};
|
||||
return formatedPayload;
|
||||
}
|
||||
|
||||
async markTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
|
||||
return lastValueFrom(this.inputWriteService.MarkTableDeleted(data));
|
||||
}
|
||||
|
||||
async unmarkTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
|
||||
return lastValueFrom((this.inputWriteService as any).UnmarkTableDeleted(data));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,46 +1,29 @@
|
||||
import { Body, Controller, Inject, Param, Post, Req } from '@nestjs/common';
|
||||
import { init } from 'mixpanel';
|
||||
import { Authenticated } from 'src/decorators/authentication.decorator';
|
||||
import { ApiInternalOnlyController } from 'src/decorators/swagger.decorator';
|
||||
import { RequestUser } from 'src/decorators/user.decorator';
|
||||
import { RequestUser, User } from 'src/decorators/user.decorator';
|
||||
import { MixpanelService } from './mixpanel.service';
|
||||
import { extractUserFrom } from 'src/authentication/extract-user';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@ApiInternalOnlyController()
|
||||
@Controller('trackEvent')
|
||||
export class MixpanelController {
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private mixpanelService: MixpanelService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
private mixpanelService: MixpanelService
|
||||
) {}
|
||||
@Post(':id')
|
||||
async trackEvent(
|
||||
@Param('id') id,
|
||||
@Body() body,
|
||||
@User() user: RequestUser,
|
||||
@Req() request
|
||||
) {
|
||||
this.logger.info(`POST Track Event: ${id}`)
|
||||
delete body.info;
|
||||
|
||||
const anonymousUser = {
|
||||
username: "anonymous",
|
||||
customer_name: "anonymous"
|
||||
} as RequestUser
|
||||
|
||||
const hasToken = request.headers['authorization'];
|
||||
|
||||
const user = hasToken ? extractUserFrom(hasToken) : anonymousUser;
|
||||
|
||||
this.logger.info(`Has user: ${typeof hasToken == "string"}`)
|
||||
|
||||
await this.mixpanelService.track(id, user, request, body)
|
||||
|
||||
this.logger.info(`Event successful`)
|
||||
return { id, body, user: user.username };
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
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
@@ -0,0 +1,36 @@
|
||||
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;
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
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,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
import { Body, Controller, Get, Inject, Param, Post } from '@nestjs/common';
|
||||
import { ApiOperation, ApiTags } from '@nestjs/swagger';
|
||||
import {
|
||||
AuthenticateCondition,
|
||||
Authenticated,
|
||||
} 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()
|
||||
@AuthenticateCondition((req, user) => {
|
||||
let action;
|
||||
|
||||
switch (req.method) {
|
||||
case 'POST':
|
||||
action = 'CREATE';
|
||||
break;
|
||||
|
||||
case 'PUT':
|
||||
action = 'UPDATE';
|
||||
break;
|
||||
|
||||
default:
|
||||
action = req.method;
|
||||
}
|
||||
|
||||
return user.permissions.includes(
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions[action].seqid,
|
||||
);
|
||||
})
|
||||
export class PipelinesController {
|
||||
logger: DadosferaLogger;
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private pipelineService: PipelinesService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
@Post('start/:id')
|
||||
@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',
|
||||
})
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,33 @@
|
||||
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,13 +1,5 @@
|
||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||
|
||||
export class PipelineInputsDTO {
|
||||
@ApiProperty()
|
||||
tables: Array<{
|
||||
name: string,
|
||||
type: string,
|
||||
|
||||
}>
|
||||
}
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export class IPipelineV2 {
|
||||
@ApiProperty()
|
||||
@@ -61,12 +53,6 @@ export interface IIdRequest {
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
@@ -137,30 +123,3 @@ export class PipelineFindAllReq {
|
||||
@ApiPropertyOptional()
|
||||
type?: string | undefined;
|
||||
}
|
||||
|
||||
export interface UpdateTableDTO {
|
||||
name: string;
|
||||
type: string;
|
||||
columns: string[];
|
||||
destinations: {
|
||||
raw: {
|
||||
table_schema: string;
|
||||
table_name: string;
|
||||
};
|
||||
qualify: {
|
||||
table_schema: string;
|
||||
table_name: string;
|
||||
};
|
||||
};
|
||||
identifier_columns: string[];
|
||||
reference_column: {
|
||||
name: string;
|
||||
type: string;
|
||||
};
|
||||
memory: number;
|
||||
}
|
||||
|
||||
export interface UpdatePlatformInputRequest {
|
||||
cron: string;
|
||||
tables: Array<UpdateTableDTO>;
|
||||
}
|
||||
|
||||
@@ -14,8 +14,7 @@ import {
|
||||
Patch,
|
||||
HttpException,
|
||||
BadRequestException,
|
||||
UseGuards,
|
||||
Res,
|
||||
CacheTTL,
|
||||
} from '@nestjs/common';
|
||||
import {
|
||||
ApiCreatedResponse,
|
||||
@@ -25,8 +24,8 @@ import {
|
||||
ApiTags,
|
||||
} from '@nestjs/swagger';
|
||||
import {
|
||||
AuthenticateCondition,
|
||||
RequireAllPermissions,
|
||||
RequireSomePermission,
|
||||
} from 'src/decorators/authentication.decorator';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
@@ -35,6 +34,7 @@ import { Messages } from '@dadosfera/protospack-v2/dist/lib/PipelineV2';
|
||||
import { RequestUser, User } from 'src/decorators/user.decorator';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
|
||||
import { PipelinesService as OldPipelineService } from 'src/modules/pipelines/pipelines.service';
|
||||
import {
|
||||
ICompleteUploadCSVFile,
|
||||
ICreatePipelineCSVFile,
|
||||
@@ -42,34 +42,52 @@ import {
|
||||
IPipelineV2,
|
||||
IInitUploadCSVFile,
|
||||
PipelineFindAllReq,
|
||||
UpdatePlatformInputRequest,
|
||||
} from './interfaces';
|
||||
import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filter';
|
||||
import { LanguageEnum } from 'src/utils/languages.enum';
|
||||
import { Language } from 'src/decorators/language.decorator';
|
||||
import { ApiInternalOnlyEndpoint } from 'src/decorators/swagger.decorator';
|
||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
||||
import { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard';
|
||||
|
||||
type PipelineTable = { name: string; job_id?: string; is_deleted?: boolean; [key: string]: any };
|
||||
type PipelineTablesConfig = { input_id?: string; tables: PipelineTable[] };
|
||||
|
||||
@ApiTags('PipelinesV2')
|
||||
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
|
||||
@UseFilters(new GrpcToHttpExceptionFilter())
|
||||
@Controller('pipelinesV2')
|
||||
@AuthenticateCondition((req, user) => {
|
||||
let action;
|
||||
|
||||
switch (req.method) {
|
||||
case 'POST':
|
||||
action = 'CREATE';
|
||||
break;
|
||||
|
||||
case 'PUT':
|
||||
action = 'UPDATE';
|
||||
break;
|
||||
|
||||
case 'PATCH':
|
||||
action = 'UPDATE';
|
||||
break;
|
||||
|
||||
default:
|
||||
action = req.method;
|
||||
}
|
||||
|
||||
return user.permissions.includes(
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions[action].seqid,
|
||||
);
|
||||
})
|
||||
export class PipelinesController {
|
||||
logger: DadosferaLogger;
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private pipelinesClientService: PipelinesService,
|
||||
private oldPipelinesService: OldPipelineService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
@Get('monitoring-dashboard')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getMonitoringDashboard(@User() user: RequestUser) {
|
||||
this.logger.info('PipelinesController - getMonitoringDashboard', { user });
|
||||
|
||||
@@ -82,7 +100,6 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Post()
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@ApiCreatedResponse({ type: IPipelineV2 })
|
||||
async create(
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -109,7 +126,6 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get()
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async findAll(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -134,7 +150,6 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get('/download-logs')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async downloadLogs(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -165,7 +180,6 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get(':id/config')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW,PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineproperties(
|
||||
@Language() language: LanguageEnum,
|
||||
@User() user: RequestUser,
|
||||
@@ -177,7 +191,6 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get(':id/objects')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineObjects(
|
||||
@Language() language: LanguageEnum,
|
||||
@User() user: RequestUser,
|
||||
@@ -189,9 +202,7 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get(':id/status')
|
||||
@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(`/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`, {
|
||||
@@ -199,13 +210,12 @@ export class PipelinesController {
|
||||
customer: body.info.customer,
|
||||
});
|
||||
|
||||
const response = await this.pipelinesClientService.getPipelineStatus(body);
|
||||
const response = await this.oldPipelinesService.getPipelineStatus(body);
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
@Get('/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async findOne(
|
||||
@Language() language: LanguageEnum,
|
||||
@User() user: RequestUser,
|
||||
@@ -223,56 +233,31 @@ export class PipelinesController {
|
||||
language,
|
||||
});
|
||||
|
||||
const pipelineRes = await this.pipelinesClientService.findOne({ id }, metadata);
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
@Get("/:id/data-assets")
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET
|
||||
)
|
||||
async findAllDataAssetByPipeline(
|
||||
@Language() language: LanguageEnum,
|
||||
@Param('id') id: string,
|
||||
@User() user: RequestUser,
|
||||
@Query('object') object: string
|
||||
) {
|
||||
const payload = {
|
||||
pipeline: id,
|
||||
object: object,
|
||||
};
|
||||
|
||||
this.logger.info(`GET pipelinesV2/:id/data-assets` + JSON.stringify(payload));
|
||||
|
||||
const result =
|
||||
await this.pipelinesClientService.findAllDataAssetByPipeline(payload, user);
|
||||
const result = await this.pipelinesClientService
|
||||
.findOne({ id }, metadata)
|
||||
.then((res) => {
|
||||
//{pipeline:{tables: {tables: [], input_id: ''}}}
|
||||
let tables = JSON.parse(res.pipeline.config.tables);
|
||||
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,
|
||||
},
|
||||
properties: res.pipeline.properties
|
||||
? JSON.parse(res.pipeline.properties)
|
||||
: {},
|
||||
});
|
||||
return res;
|
||||
});
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@Patch('/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async update(
|
||||
@Language() language: LanguageEnum,
|
||||
@Body() updatePipelineDto,
|
||||
@@ -306,54 +291,12 @@ export class PipelinesController {
|
||||
return response;
|
||||
}
|
||||
|
||||
@Patch('/:pipelineId/inputs/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
@UseGuards(PipelineExecutionGuard)
|
||||
async updatePipelineInput(
|
||||
@Language() language: LanguageEnum,
|
||||
@Body() pipelineInputDTO: UpdatePlatformInputRequest,
|
||||
@Param('id') inputId: string,
|
||||
@Param('pipelineId') pipelineId: string,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
this.logger.info('PipelinesController - update', { user });
|
||||
|
||||
const { customer_id, customer_name, user_id, username } = user;
|
||||
const info: Info = {
|
||||
user_id: user.user_id,
|
||||
customer: user.customer_name,
|
||||
customer_id: user.customer_id,
|
||||
pipeline_id: pipelineId
|
||||
};
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_id,
|
||||
customer_name,
|
||||
user_id,
|
||||
username,
|
||||
language,
|
||||
});
|
||||
|
||||
const response = await this.pipelinesClientService.updatePipelineInput(
|
||||
pipelineId,
|
||||
inputId,
|
||||
pipelineInputDTO,
|
||||
info,
|
||||
user,
|
||||
metadata,
|
||||
);
|
||||
|
||||
this.logger.info('PipelinesController - update: OK', { user });
|
||||
return response;
|
||||
}
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Put('/:id')
|
||||
@ApiOperation({
|
||||
deprecated: true,
|
||||
description: 'This method is deprecated. Please use PATCH instead',
|
||||
})
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updateDeprecated(
|
||||
@Language() language: LanguageEnum,
|
||||
@Body() updatePipelineDto,
|
||||
@@ -366,25 +309,9 @@ export class PipelinesController {
|
||||
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')
|
||||
@ApiNoContentResponse()
|
||||
@HttpCode(HttpStatus.NO_CONTENT)
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
|
||||
async delete(@Param('id') id: string, @User() user: RequestUser) {
|
||||
this.logger.info('PipelinesController - delete', { user });
|
||||
const metadata = PackTheMetadata({
|
||||
@@ -398,7 +325,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/init-upload')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async initUploadFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: IInitUploadCSVFile,
|
||||
@@ -430,7 +357,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/complete-upload')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async completeUploadFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: ICompleteUploadCSVFile,
|
||||
@@ -448,7 +375,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/file')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async uploadedFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: ICreatePipelineCSVFile,
|
||||
@@ -498,7 +425,6 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('start/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async activate(@Param('id') id: string, @Body() body) {
|
||||
const { info } = body;
|
||||
|
||||
@@ -510,7 +436,7 @@ export class PipelinesController {
|
||||
},
|
||||
);
|
||||
|
||||
const response = await this.pipelinesClientService.runPipeline({ id, info });
|
||||
const response = await this.oldPipelinesService.runPipeline({ id, info });
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
@@ -7,28 +7,23 @@ import { PipelinesService } from './pipelines.service';
|
||||
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
|
||||
import { PipelinesModule as OldPipelineModule } from 'src/modules/pipelines/pipelines.module';
|
||||
import { ConnectorModule } from '../connector/connector.module';
|
||||
import { InputsModule } from '../inputs/inputs.module';
|
||||
import { TransformationsModule } from '../transformations/transformations.module';
|
||||
import { PlatformApiModule } from '../platform-api/platform-api.module';
|
||||
import { NimbusServicesModule } from 'src/services/nimbus/nimbus.module';
|
||||
import { NimbusService } from 'src/services/nimbus/nimbus.service';
|
||||
import { CatalogModule } from '../catalog/catalog.module';
|
||||
|
||||
const client = new PipelinesClientConfiguration();
|
||||
|
||||
@Module({
|
||||
imports: [
|
||||
ClientsModule.register([client.providerOptions]),
|
||||
OldPipelineModule,
|
||||
ConnectorModule,
|
||||
InputsModule,
|
||||
TransformationsModule,
|
||||
PlatformApiModule,
|
||||
NimbusServicesModule,
|
||||
CatalogModule
|
||||
],
|
||||
controllers: [PipelinesController],
|
||||
providers: [PipelinesService, DadosferaLogger, NimbusService],
|
||||
providers: [PipelinesService, DadosferaLogger],
|
||||
exports: [PipelinesService],
|
||||
})
|
||||
export class PipelinesV2Module {}
|
||||
|
||||
@@ -1,7 +1,5 @@
|
||||
/* eslint-disable no-async-promise-executor */
|
||||
import {
|
||||
BadRequestException,
|
||||
ConflictException,
|
||||
HttpException,
|
||||
HttpStatus,
|
||||
Inject,
|
||||
@@ -18,7 +16,7 @@ import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
import { ICreatePipelineV2Req, IIdRequest, UpdatePlatformInputRequest, UpdateTableDTO } from './interfaces';
|
||||
import { ICreatePipelineV2Req } from './interfaces';
|
||||
import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages';
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
import { ConnectorClientService } from '../connector/client.service';
|
||||
@@ -28,17 +26,6 @@ import { TransformationsService } from '../transformations/transformations.servi
|
||||
import { getObjValueFromPath, objHasPath } from 'src/utils/ObjValueFromPath';
|
||||
import ErrorCodes from 'src/utils/errorCodes';
|
||||
import ErrorBuilder from 'src/utils/ErrorBuilder';
|
||||
import { PlatformApiService } from '../platform-api/platform-api.service';
|
||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
||||
import { TableUpdate } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
|
||||
import { AxiosError } from 'axios';
|
||||
import { NimbusService } from 'src/services/nimbus/nimbus.service';
|
||||
import { PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import { IDataAsset } from '../catalog/dtos';
|
||||
import { CatalogService } from '../catalog/catalog.service';
|
||||
|
||||
type RollbackPromise = () => Promise<any>;
|
||||
|
||||
export class PipelinesService implements OnModuleInit {
|
||||
logger: DadosferaLogger;
|
||||
@@ -52,9 +39,6 @@ export class PipelinesService implements OnModuleInit {
|
||||
private readonly connectorService: ConnectorClientService,
|
||||
private readonly inputsService: InputsService,
|
||||
private readonly transformationsService: TransformationsService,
|
||||
private readonly platformAPI: PlatformApiService,
|
||||
private readonly nimbusService: NimbusService,
|
||||
private readonly catalogService: CatalogService
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
@@ -154,7 +138,6 @@ export class PipelinesService implements OnModuleInit {
|
||||
const findOnePipelineResponse = await lastValueFrom(
|
||||
this.pipelineReadService.PipelineV2FindOne(data, metadata),
|
||||
);
|
||||
console.log('pipeline find one response', findOnePipelineResponse);
|
||||
this.logger.info('Done');
|
||||
|
||||
return findOnePipelineResponse;
|
||||
@@ -177,15 +160,6 @@ export class PipelinesService implements OnModuleInit {
|
||||
return updatePipelineResponse;
|
||||
}
|
||||
|
||||
async upgrade(id: string, metadata: Metadata) {
|
||||
await lastValueFrom(
|
||||
this.pipelineWriteService.Upgrade(
|
||||
{ id },
|
||||
metadata,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
async remove(data: { id: string; metadata: Metadata; user: RequestUser }) {
|
||||
const { id, metadata, user } = data;
|
||||
const info = {
|
||||
@@ -365,317 +339,4 @@ export class PipelinesService implements OnModuleInit {
|
||||
|
||||
return res;
|
||||
}
|
||||
|
||||
async updatePipelineInput(pipelineId: string, inputId: string, updateInputDTO: UpdatePlatformInputRequest, info: Info, user: RequestUser, metadata: Metadata) {
|
||||
this.logger.info('InputClientService - Update');
|
||||
|
||||
const {
|
||||
input: oldInput
|
||||
} = await this.inputsService.findOne({
|
||||
id: inputId,
|
||||
info: info
|
||||
});
|
||||
|
||||
this.logger.info('Update Dynamo Reference :' + JSON.stringify(oldInput));
|
||||
const pipelineIdFormat = pipelineId.split('-').join('_');
|
||||
const rollback: RollbackPromise[] = [];
|
||||
|
||||
const updateInputResponse = await this.inputsService.update(
|
||||
inputId,
|
||||
updateInputDTO,
|
||||
info
|
||||
);
|
||||
|
||||
const inputRollback = () => {
|
||||
this.logger.info("exec rollback to input: " + JSON.stringify(oldInput));
|
||||
return this.inputsService.rollbackUpdate(
|
||||
{
|
||||
id: inputId,
|
||||
dataAssetUpdate: updateInputResponse.dataAssetUpdate,
|
||||
tables: oldInput.tables,
|
||||
info
|
||||
}
|
||||
) as Promise<any>;
|
||||
}
|
||||
|
||||
rollback.push(inputRollback);
|
||||
|
||||
this.logger.info("Input Update Response: " + JSON.stringify(updateInputResponse))
|
||||
|
||||
const nimbusUpdates = updateInputResponse?.tablesUpdate || [];
|
||||
|
||||
nimbusUpdates.forEach(update => {
|
||||
const nimbusRollback = () => {
|
||||
return this.nimbusService.renameTable(
|
||||
info.customer,
|
||||
update.database,
|
||||
{
|
||||
table_name: update.table_name,
|
||||
table_schema: update.table_schema
|
||||
},
|
||||
{
|
||||
table_name: update.old_table_name,
|
||||
table_schema: update.old_table_schema
|
||||
}
|
||||
);
|
||||
}
|
||||
rollback.push(nimbusRollback);
|
||||
});
|
||||
|
||||
try {
|
||||
await this.updateNimbus(info.customer, nimbusUpdates);
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
if (error instanceof AxiosError) {
|
||||
this.logger.error(JSON.stringify(error.response.data));
|
||||
}
|
||||
await this.executeRenameRollback(rollback);
|
||||
|
||||
throw new Error("Error Nimbus updating tables");
|
||||
}
|
||||
|
||||
try {
|
||||
await this.updatePlatformJobs(
|
||||
pipelineIdFormat,
|
||||
updateInputResponse.input.type,
|
||||
updateInputDTO,
|
||||
user
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
await this.executeRenameRollback(rollback)
|
||||
throw new Error("Error Platform API updating jobs");
|
||||
}
|
||||
|
||||
return updateInputResponse;
|
||||
}
|
||||
|
||||
private async executeRenameRollback(request: RollbackPromise[]) {
|
||||
this.logger.info('rollback steps: ' + request.length)
|
||||
const result = await Promise.allSettled(request.map(func => func()));
|
||||
result.forEach(promise => {
|
||||
this.logger.info("Promise finish with status: " + promise.status)
|
||||
|
||||
if (promise.status === "rejected") {
|
||||
this.logger.error("reject with: " + JSON.stringify(promise.reason || {}))
|
||||
}
|
||||
|
||||
if (promise.status === "fulfilled") {
|
||||
this.logger.info("success with: " + JSON.stringify(promise.value || {}))
|
||||
}
|
||||
});
|
||||
|
||||
}
|
||||
|
||||
private async updateNimbus(customer: string, changes: TableUpdate[]) {
|
||||
// throw new Error("teste error nimbus");
|
||||
this.logger.info('Nimbus Changes: ' + JSON.stringify(changes));
|
||||
if(!changes || changes.length === 0) return;
|
||||
|
||||
const requests = changes.map(change => {
|
||||
return this.nimbusService.renameTable(customer, change.database, {
|
||||
table_name: change.old_table_name,
|
||||
table_schema: change.old_table_schema
|
||||
}, {
|
||||
table_name: change.table_name,
|
||||
table_schema: change.table_schema
|
||||
});
|
||||
})
|
||||
|
||||
const values = await Promise.allSettled(requests);
|
||||
|
||||
const success = values.map(request => request.status === "fulfilled")
|
||||
|
||||
this.logger.info("Updates with succes: " + success.length);
|
||||
|
||||
values.forEach(promise => {
|
||||
this.logger.info("Promise finish with status: " + promise.status)
|
||||
|
||||
if (promise.status === "rejected") {
|
||||
this.logger.error("reject with: " + JSON.stringify(promise.reason || {}));
|
||||
throw new Error(promise.reason );
|
||||
}
|
||||
|
||||
if (promise.status === "fulfilled") {
|
||||
this.logger.info("success with: " + JSON.stringify(promise.value || {}));
|
||||
}
|
||||
});
|
||||
|
||||
}
|
||||
|
||||
async updatePlatformJobs(pipelineId: string, pipelineType: string, updateInputDTO: UpdatePlatformInputRequest, user: RequestUser) {
|
||||
const jobsUpdated = [];
|
||||
|
||||
for (const [index, table] of updateInputDTO.tables.entries()) {
|
||||
const jobUpdate = {
|
||||
job_id: `${pipelineId}_${index}`,
|
||||
}
|
||||
|
||||
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) {
|
||||
hasUpdateSyncMode = true;
|
||||
jobSyncMode['column_include_list'] = table.columns;
|
||||
}
|
||||
|
||||
if (table.reference_column) {
|
||||
hasUpdateSyncMode = true;
|
||||
jobSyncMode['incremental_column_name'] = table.reference_column.name;
|
||||
jobSyncMode['incremental_column_type'] = table.reference_column.type;
|
||||
}
|
||||
|
||||
if (table.identifier_columns) {
|
||||
hasUpdateSyncMode = true;
|
||||
jobSyncMode['primary_keys'] = table.identifier_columns;
|
||||
}
|
||||
|
||||
if (table.type) {
|
||||
hasUpdateSyncMode = true;
|
||||
|
||||
jobSyncMode['target_load_type'] = table.type;
|
||||
}
|
||||
|
||||
if(hasUpdateSyncMode) {
|
||||
jobUpdate["sync_mode"] = jobSyncMode;
|
||||
}
|
||||
|
||||
if (Object.keys(table.destinations).length > 1) {
|
||||
let hasChanges = false
|
||||
const jobRenameTables = {
|
||||
raw: {},
|
||||
qualify: {}
|
||||
}
|
||||
|
||||
if (Object.keys(table.destinations.raw).length > 1) {
|
||||
hasChanges = true;
|
||||
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('Request body:' + JSON.stringify({
|
||||
jobs_updated: jobsUpdated
|
||||
}));
|
||||
|
||||
const response = await this.platformAPI.proxy(
|
||||
'PUT',
|
||||
`/pipeline/${pipelineId}/jobs`,
|
||||
user,
|
||||
{
|
||||
job_updates: jobsUpdated
|
||||
}
|
||||
)
|
||||
this.logger.info('Platform api response: ' + JSON.stringify(response));
|
||||
}
|
||||
|
||||
async findAllDataAssetByPipeline(data: {
|
||||
pipeline: string,
|
||||
object?: string
|
||||
}, user: RequestUser) {
|
||||
const metadata = PackTheMetadata(user);
|
||||
|
||||
const isDataAdmin = user.permissions.includes(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
|
||||
);
|
||||
|
||||
let has_permission = false;
|
||||
|
||||
const {
|
||||
data_assets: resultString
|
||||
} = await lastValueFrom(
|
||||
this.pipelineReadService.FindAllDataAssetByPipeline(data, metadata)
|
||||
);
|
||||
|
||||
const result = JSON.parse(resultString) as any;
|
||||
const data_assets: IDataAsset[] = []
|
||||
result.forEach(data_asset => {
|
||||
if (data_asset?.owner === user.username) has_permission = true;
|
||||
|
||||
for (const role of user.roles) {
|
||||
if (data_asset.roles.includes(role)) has_permission = true;
|
||||
}
|
||||
|
||||
if (data_asset.users.includes(user.user_id)) has_permission = true;
|
||||
|
||||
if (isDataAdmin || has_permission) {
|
||||
delete data_asset.p_roles;
|
||||
delete data_asset.p_users;
|
||||
data_assets.push(data_asset as IDataAsset);
|
||||
}
|
||||
});
|
||||
|
||||
const assets = await this.catalogService.getAssetsUsersAndRoles(data_assets, user.customer_id);
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,11 +0,0 @@
|
||||
export const PLATFORM_API_CONFIG = {
|
||||
getUrl: (): string => {
|
||||
const url = process.env.PLATFORM_API_URL;
|
||||
if (!url) {
|
||||
throw new Error('PLATFORM_API_URL environment variable is not set');
|
||||
}
|
||||
return url;
|
||||
},
|
||||
region: process.env.AWS_REGION || 'us-east-1',
|
||||
timeout: parseInt(process.env.PLATFORM_API_TIMEOUT || '30000', 10),
|
||||
};
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,9 +0,0 @@
|
||||
import { ApiProperty } from "@nestjs/swagger";
|
||||
|
||||
export class ValidationTableDTO {
|
||||
@ApiProperty()
|
||||
tables: Array<{
|
||||
table_name: string;
|
||||
table_schema: string;
|
||||
}>
|
||||
}
|
||||
@@ -1,19 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { PlatformApiController } from './platform-api.controller';
|
||||
import { PlatformApiService } from './platform-api.service';
|
||||
import { ElasticsearchModule } from '../../services/elasticsearch';
|
||||
import { DynamoDBModule } from '../../services/dynamodb';
|
||||
import { CustomersModule } from '../customers/customers.module';
|
||||
import { CatalogModule } from '../catalog/catalog.module';
|
||||
import { InputsModule } from '../inputs/inputs.module';
|
||||
|
||||
@Module({
|
||||
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule],
|
||||
controllers: [PlatformApiController],
|
||||
providers: [PlatformApiService, DadosferaLogger],
|
||||
exports: [PlatformApiService],
|
||||
})
|
||||
export class PlatformApiModule {}
|
||||
@@ -1,131 +0,0 @@
|
||||
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 { PLATFORM_API_CONFIG } from './platform-api.config';
|
||||
|
||||
@Injectable()
|
||||
export class PlatformApiService {
|
||||
private signer: SignatureV4;
|
||||
private logger: any;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
this.signer = new SignatureV4({
|
||||
service: 'execute-api',
|
||||
region: PLATFORM_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 = PLATFORM_API_CONFIG.getUrl();
|
||||
const url = new URL(`${baseUrl}${path}`);
|
||||
|
||||
// Add query params
|
||||
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',
|
||||
// Forward user context headers
|
||||
// Note: platform-api expects customer_name in the 'customer_id' header (contract inconsistency)
|
||||
'customer_id': user.customer_name || '',
|
||||
'customer_name': user.customer_name || '',
|
||||
'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,
|
||||
};
|
||||
|
||||
this.logger.info('Proxying request to platform-api', {
|
||||
method: method.toUpperCase(),
|
||||
path,
|
||||
customer_id: user.customer_id,
|
||||
user_id: user.user_id,
|
||||
});
|
||||
|
||||
try {
|
||||
// Sign with IAM v4
|
||||
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: PLATFORM_API_CONFIG.timeout,
|
||||
validateStatus: () => true, // Don't throw on non-2xx
|
||||
});
|
||||
|
||||
// Propagate non-2xx responses as HttpExceptions
|
||||
if (response.status >= 400) {
|
||||
this.logger.error('Platform API upstream error' + JSON.stringify({
|
||||
status: response.status,
|
||||
data: response.data,
|
||||
path,
|
||||
method: method.toUpperCase(),
|
||||
}));
|
||||
throw new HttpException(response.data, response.status);
|
||||
}
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.logger.error('Platform API proxy error', {
|
||||
error: error.message,
|
||||
status: error.response?.status,
|
||||
path,
|
||||
method: method.toUpperCase(),
|
||||
});
|
||||
|
||||
this.logger.error(error)
|
||||
|
||||
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('Platform API service unavailable', 503);
|
||||
}
|
||||
|
||||
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
|
||||
throw new HttpException('Platform API request timeout', 504);
|
||||
}
|
||||
|
||||
throw new HttpException('Internal server error', 500);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,14 +0,0 @@
|
||||
export type ReleaseNoteDTO = {
|
||||
id: string;
|
||||
date: string;
|
||||
tag: string;
|
||||
title: string;
|
||||
visible: boolean;
|
||||
expiryDate: string;
|
||||
content: string;
|
||||
showEmojis: boolean;
|
||||
image?: string;
|
||||
link?: string;
|
||||
linkText?: string;
|
||||
};
|
||||
|
||||
@@ -1,20 +0,0 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteController } from './release_note.controller';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
|
||||
describe('ReleaseNoteController', () => {
|
||||
let controller: ReleaseNoteController;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
controllers: [ReleaseNoteController],
|
||||
providers: [ReleaseNoteService],
|
||||
}).compile();
|
||||
|
||||
controller = module.get<ReleaseNoteController>(ReleaseNoteController);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(controller).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -1,26 +0,0 @@
|
||||
import { Controller, Get, Inject } from '@nestjs/common';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
import { Authenticated } from 'src/decorators/authentication.decorator';
|
||||
import { Language } from 'src/decorators/language.decorator';
|
||||
import { LanguageEnum } from 'src/utils/languages.enum';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@Controller('release_note')
|
||||
@Authenticated()
|
||||
export class ReleaseNoteController {
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private readonly releaseNoteService: ReleaseNoteService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
@Get()
|
||||
async getLatestReleaseNote(@Language() language: LanguageEnum) {
|
||||
this.logger.info(`Fetching latest release note for language: ${language}`);
|
||||
return await this.releaseNoteService.getLatestReleaseNote(language);
|
||||
}
|
||||
}
|
||||
@@ -1,10 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
import { ReleaseNoteController } from './release_note.controller';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@Module({
|
||||
controllers: [ReleaseNoteController],
|
||||
providers: [ReleaseNoteService, DadosferaLogger]
|
||||
})
|
||||
export class ReleaseNoteModule {}
|
||||
@@ -1,18 +0,0 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
|
||||
describe('ReleaseNoteService', () => {
|
||||
let service: ReleaseNoteService;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
providers: [ReleaseNoteService],
|
||||
}).compile();
|
||||
|
||||
service = module.get<ReleaseNoteService>(ReleaseNoteService);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(service).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -1,46 +0,0 @@
|
||||
import { Inject, Injectable } from '@nestjs/common';
|
||||
import axios, { AxiosInstance } from 'axios';
|
||||
import { LanguageEnum } from 'src/utils/languages.enum';
|
||||
import { ReleaseNoteDTO } from './dto/release_note.dto';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@Injectable()
|
||||
export class ReleaseNoteService {
|
||||
client: AxiosInstance;
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
this.client = axios.create({
|
||||
baseURL: process.env.FIREBASE_BASE_URL,
|
||||
});
|
||||
}
|
||||
|
||||
async getLatestReleaseNote(lang: LanguageEnum) {
|
||||
try {
|
||||
const lng = lang.split('-');
|
||||
const language = lng[0] + '-' + lng[1].toUpperCase();
|
||||
|
||||
const endpoint = `/release_note/${language}.json`;
|
||||
const {
|
||||
data,
|
||||
status,
|
||||
config
|
||||
} = await this.client.get<ReleaseNoteDTO>(endpoint)
|
||||
this.logger.info(`Fetched release note for language: ${lang} with status: ${status}`);
|
||||
this.logger.info(`Request URL: ${config.baseURL}/${config.url}`);
|
||||
|
||||
return data;
|
||||
} catch (error) {
|
||||
this.logger.error(`Error fetching release note: ${error.message}`);
|
||||
|
||||
if (axios.isAxiosError(error)) {
|
||||
this.logger.error(`Axios error details: ${error.toJSON()}`);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -229,33 +229,27 @@ export class RolesService {
|
||||
const [roleTreated] = this.getRolesPermissionsName([role.role]);
|
||||
return { role: roleTreated };
|
||||
}
|
||||
|
||||
getRolesPermissionsName(roles: GetRolesPermissionsName[]): RoleDto[] {
|
||||
const newRoles: RoleDto[] = [];
|
||||
for (const role of roles) {
|
||||
const newRole: RoleDto = this.formatRole(role);
|
||||
const allPermissions = this.permissionsService.getAllPermissions(
|
||||
this.language,
|
||||
);
|
||||
const newPermissions = role.permissions.map((p) => {
|
||||
const permission = allPermissions.find((per) => per.seqid === p.seqid);
|
||||
return {
|
||||
...p,
|
||||
name: permission.name,
|
||||
id: p.seqid,
|
||||
};
|
||||
});
|
||||
const newRole: RoleDto = {
|
||||
...role,
|
||||
permissions: newPermissions,
|
||||
isPublic: role.isPublic,
|
||||
};
|
||||
newRoles.push(newRole);
|
||||
}
|
||||
return newRoles;
|
||||
}
|
||||
|
||||
private formatRole(role: GetRolesPermissionsName): RoleDto {
|
||||
const allPermissions = this.permissionsService.getAllPermissions(
|
||||
this.language
|
||||
);
|
||||
const newPermissions = role.permissions.map((p) => {
|
||||
const permission = allPermissions.find((per) => per.seqid === p.seqid);
|
||||
return {
|
||||
...p,
|
||||
name: permission.name,
|
||||
id: p.seqid,
|
||||
};
|
||||
});
|
||||
|
||||
return {
|
||||
...role,
|
||||
permissions: newPermissions,
|
||||
isPublic: role.isPublic,
|
||||
};;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
export const STORAGE_EXPLORER_CONFIG = {
|
||||
getUrl: (customerName: string): string => {
|
||||
const urlTemplate = process.env.STORAGE_EXPLORER_API_URL;
|
||||
if (!urlTemplate) {
|
||||
throw new Error('STORAGE_EXPLORER_API_URL environment variable is not set');
|
||||
}
|
||||
// Replace {customer_id} placeholder with actual customer ID
|
||||
// For local: http://172.17.0.1:8000/api (no placeholder)
|
||||
// For prod: https://storage-explorer-{customer_id}.dadosfera.ai/api
|
||||
return urlTemplate.replace('{customer}', customerName);
|
||||
},
|
||||
timeout: parseInt(process.env.STORAGE_EXPLORER_TIMEOUT || '30000', 10),
|
||||
};
|
||||
@@ -1,383 +0,0 @@
|
||||
import {
|
||||
Controller,
|
||||
Get,
|
||||
Post,
|
||||
Put,
|
||||
Param,
|
||||
Body,
|
||||
Query,
|
||||
Inject,
|
||||
UseInterceptors,
|
||||
UploadedFiles,
|
||||
Headers,
|
||||
} from '@nestjs/common';
|
||||
import { ApiTags, ApiOperation, ApiConsumes } from '@nestjs/swagger';
|
||||
import { FilesInterceptor } from '@nestjs/platform-express';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import FormData from 'form-data';
|
||||
|
||||
import {
|
||||
Authenticated,
|
||||
RequireAllPermissions,
|
||||
} from '../../decorators/authentication.decorator';
|
||||
import { User, RequestUser } from '../../decorators/user.decorator';
|
||||
import { StorageExplorerService } from './storage-explorer.service';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
|
||||
@ApiTags('Storage Explorer')
|
||||
@Controller('storage-explorer')
|
||||
@Authenticated()
|
||||
export class StorageExplorerController {
|
||||
private logger: any;
|
||||
|
||||
constructor(
|
||||
private readonly storageExplorerService: StorageExplorerService,
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// TABLE OPERATIONS
|
||||
// ============================================
|
||||
|
||||
@ApiOperation({ summary: 'Validate table name in PostgreSQL and Snowflake' })
|
||||
@Post('tables/validate-name')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async validateTableName(
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'POST',
|
||||
'/tables/validate-name',
|
||||
user,
|
||||
|
||||
body,
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Create a new table' })
|
||||
@Post('tables')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async createTable(
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'POST',
|
||||
'/tables/',
|
||||
user,
|
||||
|
||||
body,
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'List all tables with pagination' })
|
||||
@Get('tables')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async listTables(
|
||||
@Query('page') page: number,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
'/tables/',
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ page },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Get table details by ID' })
|
||||
@Get('tables/:tableId')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getTable(
|
||||
@Param('tableId') tableId: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/tables/${tableId}`,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Link a dataset to a table' })
|
||||
@Post('tables/:tableId/datasets/:datasetId')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async linkDatasetToTable(
|
||||
@Param('tableId') tableId: string,
|
||||
@Param('datasetId') datasetId: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'POST',
|
||||
`/tables/${tableId}/datasets/${datasetId}`,
|
||||
user,
|
||||
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Get all datasets linked to a table' })
|
||||
@Get('tables/:tableId/datasets')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getTableDatasets(
|
||||
@Param('tableId') tableId: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/tables/${tableId}/datasets`,
|
||||
user,
|
||||
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Get table schema' })
|
||||
@Get('tables/:tableId/schema')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getTableSchema(
|
||||
@Param('tableId') tableId: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/tables/${tableId}/schema`,
|
||||
user,
|
||||
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Validate schema compatibility between table and dataset' })
|
||||
@Post('tables/:tableId/validate-compatibility/:datasetId')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async validateSchemaCompatibility(
|
||||
@Param('tableId') tableId: string,
|
||||
@Param('datasetId') datasetId: string,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'POST',
|
||||
`/tables/${tableId}/validate-compatibility/${datasetId}`,
|
||||
user,
|
||||
|
||||
);
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// DATASET OPERATIONS
|
||||
// ============================================
|
||||
|
||||
@ApiOperation({ summary: 'Get dataset preview data' })
|
||||
@Get('datasets/:datasetId/preview')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getDatasetPreview(
|
||||
@Param('datasetId') datasetId: string,
|
||||
@Query('limit') limit: number,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/datasets/${datasetId}/preview`,
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ limit },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Get dataset schema information' })
|
||||
@Get('datasets/:datasetId/schema')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getDatasetSchema(
|
||||
@Param('datasetId') datasetId: string,
|
||||
@Query('force_refresh') forceRefresh: boolean,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/datasets/${datasetId}/schema`,
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ force_refresh: forceRefresh },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'List all datasets for a specific upload' })
|
||||
@Get('datasets/upload/:uploadId')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async listDatasetsByUpload(
|
||||
@Param('uploadId') uploadId: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
`/datasets/upload/${uploadId}`,
|
||||
user,
|
||||
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Refresh dataset schema with new parsing options (Excel)' })
|
||||
@Put('datasets/:datasetId/refresh-schema')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async refreshDatasetSchema(
|
||||
@Param('datasetId') datasetId: string,
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'PUT',
|
||||
`/datasets/${datasetId}/refresh-schema`,
|
||||
user,
|
||||
|
||||
body,
|
||||
);
|
||||
}
|
||||
|
||||
// ============================================
|
||||
// STORAGE OPERATIONS
|
||||
// ============================================
|
||||
|
||||
@ApiOperation({ summary: 'List file explorer uploads with pagination' })
|
||||
@Get('storage/uploads/history')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async listFileExplorerUploads(
|
||||
@Query('page') page: number,
|
||||
@Query('limit') limit: number,
|
||||
@Query('folder_path') folderPath: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
'/storage/uploads/history',
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ page, limit, folder_path: folderPath },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Browse folders and files in storage' })
|
||||
@Get('storage/browse')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async browseStorage(
|
||||
@Query('path') path: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
'/storage/browse',
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ path },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Upload multiple files to storage' })
|
||||
@Post('storage/upload/batch')
|
||||
@ApiConsumes('multipart/form-data')
|
||||
@UseInterceptors(FilesInterceptor('files'))
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async batchUpload(
|
||||
@UploadedFiles() files: Array<Express.Multer.File>,
|
||||
@Body('folder_path') folderPath: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
// Create FormData to forward files to storage-explorer API
|
||||
const formData = new FormData();
|
||||
|
||||
// Add files
|
||||
if (files && files.length > 0) {
|
||||
files.forEach((file) => {
|
||||
formData.append('files', file.buffer, {
|
||||
filename: file.originalname,
|
||||
contentType: file.mimetype,
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
// Add folder_path
|
||||
if (folderPath) {
|
||||
formData.append('folder_path', folderPath);
|
||||
}
|
||||
|
||||
return this.storageExplorerService.proxyFormData(
|
||||
'POST',
|
||||
'/storage/upload/batch',
|
||||
user,
|
||||
|
||||
formData,
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Create a new folder in storage' })
|
||||
@Post('storage/folder/create')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.WRITE)
|
||||
async createFolder(
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'POST',
|
||||
'/storage/folder/create',
|
||||
user,
|
||||
|
||||
body,
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Download a file from storage' })
|
||||
@Get('storage/download')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async downloadFile(
|
||||
@Query('file_path') filePath: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
'/storage/download',
|
||||
user,
|
||||
|
||||
undefined,
|
||||
{ file_path: filePath },
|
||||
);
|
||||
}
|
||||
|
||||
@ApiOperation({ summary: 'Get detailed file metadata' })
|
||||
@Get('storage/metadata')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.STORAGE_EXPLORER.permissions.READ)
|
||||
async getFileMetadata(
|
||||
@Query('file_path') filePath: string,
|
||||
@User() user: RequestUser,
|
||||
|
||||
) {
|
||||
return this.storageExplorerService.proxy(
|
||||
'GET',
|
||||
'/storage/metadata',
|
||||
user,
|
||||
undefined,
|
||||
{ file_path: filePath },
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -1,13 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { StorageExplorerController } from './storage-explorer.controller';
|
||||
import { StorageExplorerService } from './storage-explorer.service';
|
||||
|
||||
@Module({
|
||||
imports: [],
|
||||
controllers: [StorageExplorerController],
|
||||
providers: [StorageExplorerService, DadosferaLogger],
|
||||
exports: [StorageExplorerService],
|
||||
})
|
||||
export class StorageExplorerModule {}
|
||||
@@ -1,178 +0,0 @@
|
||||
import { Injectable, Inject, HttpException } from '@nestjs/common';
|
||||
import axios, { AxiosResponse, Method } from 'axios';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { RequestUser } from '../../decorators/user.decorator';
|
||||
import { STORAGE_EXPLORER_CONFIG } from './storage-explorer.config';
|
||||
|
||||
@Injectable()
|
||||
export class StorageExplorerService {
|
||||
private logger: any;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
async proxy(
|
||||
method: string,
|
||||
path: string,
|
||||
user: RequestUser,
|
||||
body?: any,
|
||||
query?: Record<string, any>
|
||||
): Promise<any> {
|
||||
// Validate customer_id is present for multi-tenant isolation
|
||||
if (!user.customer_id) {
|
||||
throw new HttpException('Customer ID is required for storage operations', 400);
|
||||
}
|
||||
|
||||
// Get customer-specific storage-explorer URL
|
||||
const baseUrl = STORAGE_EXPLORER_CONFIG.getUrl(user.customer_name);
|
||||
const url = new URL(`${baseUrl}${path}`);
|
||||
|
||||
// Add query params
|
||||
if (query) {
|
||||
Object.entries(query).forEach(([key, value]) => {
|
||||
if (value !== undefined && value !== null) {
|
||||
url.searchParams.set(key, String(value));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
const headers: Record<string, string> = {
|
||||
'content-type': 'application/json',
|
||||
};
|
||||
|
||||
this.logger.info('Proxying request to storage-explorer', {
|
||||
method: method.toUpperCase(),
|
||||
path,
|
||||
customer_id: user.customer_id,
|
||||
storage_url: baseUrl,
|
||||
user_id: user.user_id,
|
||||
});
|
||||
|
||||
try {
|
||||
const response: AxiosResponse = await axios({
|
||||
method: method as Method,
|
||||
url: url.href,
|
||||
headers,
|
||||
data: body,
|
||||
timeout: STORAGE_EXPLORER_CONFIG.timeout,
|
||||
validateStatus: () => true, // Don't throw on non-2xx
|
||||
});
|
||||
|
||||
// Propagate non-2xx responses as HttpExceptions
|
||||
if (response.status >= 400) {
|
||||
throw new HttpException(response.data, response.status);
|
||||
}
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.logger.error('Storage Explorer API proxy error', {
|
||||
error: error.message,
|
||||
status: error.response?.status,
|
||||
path,
|
||||
storage_url: baseUrl,
|
||||
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('Storage Explorer API service unavailable', 503);
|
||||
}
|
||||
|
||||
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
|
||||
throw new HttpException('Storage Explorer API request timeout', 504);
|
||||
}
|
||||
|
||||
throw new HttpException('Internal server error', 500);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Proxy with file upload support (multipart/form-data)
|
||||
*/
|
||||
async proxyFormData(
|
||||
method: string,
|
||||
path: string,
|
||||
user: RequestUser,
|
||||
formData: any,
|
||||
query?: Record<string, any>,
|
||||
): Promise<any> {
|
||||
// Validate customer_id is present for multi-tenant isolation
|
||||
if (!user.customer_id) {
|
||||
throw new HttpException('Customer ID is required for storage operations', 400);
|
||||
}
|
||||
|
||||
// Get customer-specific storage-explorer URL
|
||||
const baseUrl = STORAGE_EXPLORER_CONFIG.getUrl(user.customer_name);
|
||||
const url = new URL(`${baseUrl}${path}`);
|
||||
|
||||
// Add query params
|
||||
if (query) {
|
||||
Object.entries(query).forEach(([key, value]) => {
|
||||
if (value !== undefined && value !== null) {
|
||||
url.searchParams.set(key, String(value));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
const headers: Record<string, string> = {
|
||||
// Let axios set Content-Type for multipart/form-data with boundary
|
||||
...formData.getHeaders?.(),
|
||||
};
|
||||
|
||||
this.logger.info('Proxying form data request to storage-explorer', {
|
||||
method: method.toUpperCase(),
|
||||
path,
|
||||
customer_id: user.customer_id,
|
||||
storage_url: baseUrl,
|
||||
user_id: user.user_id,
|
||||
});
|
||||
|
||||
try {
|
||||
const response: AxiosResponse = await axios({
|
||||
method: method as Method,
|
||||
url: url.href,
|
||||
headers,
|
||||
data: formData,
|
||||
timeout: STORAGE_EXPLORER_CONFIG.timeout,
|
||||
maxContentLength: Infinity,
|
||||
maxBodyLength: Infinity,
|
||||
validateStatus: () => true,
|
||||
});
|
||||
|
||||
if (response.status >= 400) {
|
||||
throw new HttpException(response.data, response.status);
|
||||
}
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.logger.error('Storage Explorer API form data proxy error', {
|
||||
error: error.message,
|
||||
status: error.response?.status,
|
||||
path,
|
||||
storage_url: baseUrl,
|
||||
method: method.toUpperCase(),
|
||||
});
|
||||
|
||||
if (error instanceof HttpException) {
|
||||
throw error;
|
||||
}
|
||||
|
||||
if (error.response) {
|
||||
throw new HttpException(error.response.data, error.response.status);
|
||||
}
|
||||
|
||||
throw new HttpException('Internal server error', 500);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -124,24 +124,4 @@ export class ThemeController {
|
||||
}
|
||||
}
|
||||
|
||||
@Post('/:id/theme/reset')
|
||||
@ApiOkResponse({ type: CustomerThemeResponse })
|
||||
async resetTheme(@Param('id') id: string) {
|
||||
this.logger.info('getCustomerTheme with id' + id);
|
||||
|
||||
try {
|
||||
await this.themeService.resetTheme(id);
|
||||
|
||||
return { theme: null };
|
||||
}catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND) {
|
||||
this.logger.error('Error - getCustomerTheme - Expect CUSTOMER.NOT_FOUND');
|
||||
throw new HttpException(err.details, HttpStatus.NOT_FOUND);
|
||||
} else {
|
||||
this.logger.error('Error - getCustomerTheme Unknown Error:' + err?.message);
|
||||
return { theme: null };
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -44,18 +44,6 @@ export class ThemeService implements OnModuleInit {
|
||||
);
|
||||
}
|
||||
|
||||
async resetTheme(id: string) {
|
||||
const { theme } = await firstValueFrom(
|
||||
this.themeService.ResetCustomerTheme({
|
||||
id
|
||||
}),
|
||||
);
|
||||
|
||||
return {
|
||||
theme
|
||||
}
|
||||
}
|
||||
|
||||
async createThemeByCustomer(id: string, theme: CustomerThemeRequest & Files) {
|
||||
if (!id) {
|
||||
this.logger.error('Error - saveCustomertheme - not found id:' + id);
|
||||
|
||||
+2
-5
@@ -1,8 +1,5 @@
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export interface ICreateTransformationsRequest {
|
||||
transformations: Transformation[];
|
||||
info: Info;
|
||||
|
||||
@@ -38,12 +38,6 @@ export class User {
|
||||
department?: string;
|
||||
@ApiProperty()
|
||||
hierarchy?: string;
|
||||
@ApiProperty()
|
||||
bio?: string;
|
||||
@ApiProperty()
|
||||
companyName?: string;
|
||||
@ApiProperty()
|
||||
personalSite?: string;
|
||||
@ApiPropertyOptional()
|
||||
customer?: Customer;
|
||||
@ApiProperty()
|
||||
@@ -64,8 +58,6 @@ export class UserNoRolesAndCustomer extends OmitType(UserNoRoles, [
|
||||
export class IUserByCustomer extends OmitType(User, ['customer']) {
|
||||
@ApiPropertyOptional()
|
||||
permissions?: string[];
|
||||
@ApiPropertyOptional()
|
||||
authProvider?: string;
|
||||
}
|
||||
|
||||
export class CreateUserReq {
|
||||
@@ -117,12 +109,6 @@ export class UpdateUserReq {
|
||||
@ApiPropertyOptional()
|
||||
hierarchy?: string;
|
||||
@ApiPropertyOptional()
|
||||
bio?: string;
|
||||
@ApiPropertyOptional()
|
||||
personalSite?: string;
|
||||
@ApiPropertyOptional()
|
||||
companyName?: string;
|
||||
@ApiPropertyOptional()
|
||||
roleNames?: string[];
|
||||
}
|
||||
|
||||
|
||||
@@ -266,6 +266,7 @@ export class UsersController {
|
||||
}
|
||||
|
||||
@Patch(':id')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.USERS.permissions.ADMIN)
|
||||
@ApiOkResponse({ type: UpdateUserRes })
|
||||
async updateUser(
|
||||
@User() user: RequestUser,
|
||||
@@ -273,18 +274,6 @@ export class UsersController {
|
||||
@Param('id') id: string,
|
||||
@Language() language: LanguageEnum,
|
||||
) {
|
||||
|
||||
const isSameUser = user.user_id === id;
|
||||
const isSuperAdmin = user.permissions.includes(PERMISSIONS_GROUPS.USERS.permissions.ADMIN.seqid)
|
||||
if (!isSameUser && !isSuperAdmin) {
|
||||
throw new ErrorBuilder(ErrorCodes.AUTH.FORBIDDEN);
|
||||
}
|
||||
|
||||
if (isSameUser && !isSuperAdmin && body.roleNames) {
|
||||
// Prevent users from updating their own roles
|
||||
delete body.roleNames;
|
||||
}
|
||||
|
||||
this.logger.info('updateUser', { user });
|
||||
this.userService.setLanguage(language);
|
||||
return await this.userService.updateUser(body, id, user.customer_id);
|
||||
|
||||
@@ -125,7 +125,6 @@ export class UsersService implements OnModuleInit {
|
||||
return { permissions };
|
||||
});
|
||||
res.user.permissions = permissions;
|
||||
res.user.authProvider = process.env.AUTH_PROVIDER || 'cognito';
|
||||
return res;
|
||||
}
|
||||
|
||||
@@ -149,23 +148,20 @@ export class UsersService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async updateUser(req: UpdateUserReq, id: string, customerId: string) {
|
||||
const { roleNames, ...updateUserDTO } = req;
|
||||
if (roleNames && roleNames.length > 0) {
|
||||
const { department, hierarchy, jobTitle, name, roleNames, email } = req;
|
||||
if (roleNames) {
|
||||
await this.setRoles({ roleNames, userId: id }, customerId);
|
||||
}
|
||||
|
||||
const { user } = await lastValueFrom(
|
||||
this.usersClientService.UserUpdate({
|
||||
department: updateUserDTO.department,
|
||||
email: updateUserDTO.email,
|
||||
hierarchy: updateUserDTO.hierarchy,
|
||||
jobTitle: updateUserDTO.jobTitle,
|
||||
name: updateUserDTO.name,
|
||||
bio: updateUserDTO.bio,
|
||||
companyName: updateUserDTO.companyName,
|
||||
personalSite: updateUserDTO.personalSite,
|
||||
name,
|
||||
customerId,
|
||||
id,
|
||||
department,
|
||||
hierarchy,
|
||||
jobTitle,
|
||||
email,
|
||||
metabaseUserId: undefined,
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -14,10 +14,7 @@ export class ValidationPipe implements PipeTransform<any> {
|
||||
return value;
|
||||
}
|
||||
const object = plainToInstance(metatype, value);
|
||||
const errors = await validate(object, {
|
||||
forbidUnknownValues: false,
|
||||
whitelist: true,
|
||||
});
|
||||
const errors = await validate(object);
|
||||
if (errors.length > 0) {
|
||||
const errorMessages = errors.map((err) => err.constraints);
|
||||
throw new BadRequestException(errorMessages);
|
||||
|
||||
@@ -1,4 +0,0 @@
|
||||
export const DYNAMODB_CONFIG = {
|
||||
region: () => process.env.AWS_REGION || 'us-east-1',
|
||||
inputsTable: () => process.env.INPUTS_DB || 'dadosfera-inputs-prd',
|
||||
};
|
||||
@@ -1,9 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { DynamoDBService } from './dynamodb.service';
|
||||
|
||||
@Module({
|
||||
providers: [DynamoDBService, DadosferaLogger],
|
||||
exports: [DynamoDBService],
|
||||
})
|
||||
export class DynamoDBModule {}
|
||||
@@ -1,238 +0,0 @@
|
||||
import { Injectable, Inject } from '@nestjs/common';
|
||||
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
|
||||
import {
|
||||
DynamoDBDocumentClient,
|
||||
GetCommand,
|
||||
PutCommand,
|
||||
DeleteCommand,
|
||||
TranslateConfig,
|
||||
} from '@aws-sdk/lib-dynamodb';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { v4 as uuid } from 'uuid';
|
||||
import { DYNAMODB_CONFIG } from './dynamodb.config';
|
||||
|
||||
export interface ReferenceColumn {
|
||||
name: string;
|
||||
type: string;
|
||||
}
|
||||
|
||||
export interface InputDocument {
|
||||
id: string;
|
||||
client_id: string;
|
||||
user_id: string;
|
||||
created_at: string;
|
||||
name: string;
|
||||
description?: string;
|
||||
plugin: string;
|
||||
type: string;
|
||||
tables?: Array<{
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: ReferenceColumn;
|
||||
}>;
|
||||
credentials?: Record<string, any>;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class DynamoDBService {
|
||||
private documentClient: DynamoDBDocumentClient;
|
||||
private logger: any;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
|
||||
const dynamoConfig = { region: DYNAMODB_CONFIG.region() };
|
||||
const marshallOptions: TranslateConfig = {
|
||||
marshallOptions: {
|
||||
removeUndefinedValues: true,
|
||||
},
|
||||
};
|
||||
|
||||
const dynamoDb = new DynamoDBClient(dynamoConfig);
|
||||
this.documentClient = DynamoDBDocumentClient.from(dynamoDb, marshallOptions);
|
||||
}
|
||||
|
||||
async createInput(
|
||||
clientId: string,
|
||||
userId: string,
|
||||
data: {
|
||||
name: string;
|
||||
description?: string;
|
||||
plugin: string;
|
||||
type: string;
|
||||
tables?: Array<{
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: ReferenceColumn;
|
||||
}>;
|
||||
},
|
||||
): Promise<InputDocument> {
|
||||
const tableName = DYNAMODB_CONFIG.inputsTable();
|
||||
const id = uuid();
|
||||
const created_at = new Date().toISOString();
|
||||
|
||||
const item: InputDocument = {
|
||||
id,
|
||||
client_id: clientId,
|
||||
user_id: userId,
|
||||
created_at,
|
||||
name: data.name,
|
||||
description: data.description,
|
||||
plugin: data.plugin,
|
||||
type: data.type,
|
||||
tables: data.tables,
|
||||
};
|
||||
|
||||
this.logger.info('DynamoDB: Creating input', {
|
||||
tableName,
|
||||
inputId: id,
|
||||
plugin: data.plugin,
|
||||
});
|
||||
|
||||
const putCommand = new PutCommand({
|
||||
TableName: tableName,
|
||||
Item: item,
|
||||
});
|
||||
|
||||
try {
|
||||
await this.documentClient.send(putCommand);
|
||||
this.logger.info('DynamoDB: Input created successfully', { inputId: id });
|
||||
return item;
|
||||
} catch (error) {
|
||||
this.logger.error('DynamoDB: Failed to create input', {
|
||||
tableName,
|
||||
inputId: id,
|
||||
region: DYNAMODB_CONFIG.region(),
|
||||
error: error.message,
|
||||
errorName: error.name,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async findInput(clientId: string, inputId: string): Promise<InputDocument | null> {
|
||||
const tableName = DYNAMODB_CONFIG.inputsTable();
|
||||
|
||||
const getCommand = new GetCommand({
|
||||
TableName: tableName,
|
||||
Key: {
|
||||
id: inputId,
|
||||
client_id: clientId,
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
const { Item } = await this.documentClient.send(getCommand);
|
||||
return Item as InputDocument | null;
|
||||
} catch (error) {
|
||||
this.logger.error('DynamoDB: findInput failed', { inputId, clientId, error: error.message });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async deleteInput(clientId: string, inputId: string): Promise<void> {
|
||||
const tableName = DYNAMODB_CONFIG.inputsTable();
|
||||
|
||||
this.logger.info('DynamoDB: Deleting input', {
|
||||
tableName,
|
||||
inputId,
|
||||
});
|
||||
|
||||
const deleteCommand = new DeleteCommand({
|
||||
TableName: tableName,
|
||||
Key: {
|
||||
id: inputId,
|
||||
client_id: clientId,
|
||||
},
|
||||
});
|
||||
|
||||
await this.documentClient.send(deleteCommand);
|
||||
|
||||
this.logger.info('DynamoDB: Input deleted successfully', { inputId });
|
||||
}
|
||||
|
||||
/**
|
||||
* Update a specific table entry in the input document.
|
||||
* Fetches the current document, updates the matching table, and saves.
|
||||
*/
|
||||
async updateInputTable(
|
||||
clientId: string,
|
||||
inputId: string,
|
||||
tableName: string,
|
||||
changes: {
|
||||
type?: string;
|
||||
columns?: string[];
|
||||
reference_column?: ReferenceColumn | null;
|
||||
},
|
||||
): Promise<void> {
|
||||
const dynamoTableName = DYNAMODB_CONFIG.inputsTable();
|
||||
|
||||
this.logger.info('DynamoDB: Updating input table', {
|
||||
inputId,
|
||||
tableName,
|
||||
changes: Object.keys(changes),
|
||||
});
|
||||
|
||||
// Get current document
|
||||
const current = await this.findInput(clientId, inputId);
|
||||
if (!current) {
|
||||
this.logger.warn('DynamoDB: Input not found for update', { inputId });
|
||||
return;
|
||||
}
|
||||
|
||||
// Find and update the matching table
|
||||
const tables = current.tables || [];
|
||||
const tableIndex = tables.findIndex((t) => t.name === tableName);
|
||||
|
||||
if (tableIndex === -1) {
|
||||
this.logger.warn('DynamoDB: Table not found in input', {
|
||||
inputId,
|
||||
tableName,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// Merge changes into the table entry
|
||||
const updatedTable = { ...tables[tableIndex] };
|
||||
if ('type' in changes) updatedTable.type = changes.type;
|
||||
if ('columns' in changes) updatedTable.columns = changes.columns;
|
||||
if ('reference_column' in changes) {
|
||||
if (changes.reference_column === null) {
|
||||
delete updatedTable.reference_column;
|
||||
} else {
|
||||
updatedTable.reference_column = changes.reference_column;
|
||||
}
|
||||
}
|
||||
tables[tableIndex] = updatedTable;
|
||||
|
||||
// Save updated document
|
||||
const putCommand = new PutCommand({
|
||||
TableName: dynamoTableName,
|
||||
Item: {
|
||||
...current,
|
||||
tables,
|
||||
updated_at: new Date().toISOString(),
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
await this.documentClient.send(putCommand);
|
||||
this.logger.info('DynamoDB: Input table updated successfully', {
|
||||
inputId,
|
||||
tableName,
|
||||
});
|
||||
} catch (error) {
|
||||
this.logger.error('DynamoDB: Failed to update input table', {
|
||||
inputId,
|
||||
tableName,
|
||||
error: error.message,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,3 +0,0 @@
|
||||
export * from './dynamodb.service';
|
||||
export * from './dynamodb.module';
|
||||
export * from './dynamodb.config';
|
||||
@@ -1,5 +0,0 @@
|
||||
export const ELASTICSEARCH_CONFIG = {
|
||||
getUrl: () => process.env.ELASTICSEARCH_URL || 'http://localhost:9200',
|
||||
getApiKey: () => process.env.ELASTICSEARCH_API_KEY || '',
|
||||
timeout: 10000, // 10 seconds
|
||||
};
|
||||
@@ -1,9 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { ElasticsearchService } from './elasticsearch.service';
|
||||
|
||||
@Module({
|
||||
providers: [ElasticsearchService, DadosferaLogger],
|
||||
exports: [ElasticsearchService],
|
||||
})
|
||||
export class ElasticsearchModule {}
|
||||
@@ -1,457 +0,0 @@
|
||||
import { Injectable, Inject } from '@nestjs/common';
|
||||
import axios, { AxiosInstance, AxiosError } from 'axios';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { ELASTICSEARCH_CONFIG } from './elasticsearch.config';
|
||||
|
||||
interface MultiLang {
|
||||
'en-us': string;
|
||||
'pt-br': string;
|
||||
'es-es': string;
|
||||
}
|
||||
|
||||
interface MultiLangArray {
|
||||
'en-us': string[];
|
||||
'pt-br': string[];
|
||||
'es-es': string[];
|
||||
}
|
||||
|
||||
interface ConnectorInfo {
|
||||
plugin: string;
|
||||
name: MultiLang;
|
||||
image: string;
|
||||
version: string;
|
||||
tags: string[];
|
||||
}
|
||||
|
||||
interface PipelineDocument {
|
||||
id: string;
|
||||
name: MultiLang;
|
||||
description: MultiLang;
|
||||
customer_id: string;
|
||||
user_id: string;
|
||||
username: string;
|
||||
status: string;
|
||||
created_at: string;
|
||||
updated_at: string;
|
||||
last_status_updated: string;
|
||||
tags: string[];
|
||||
// Connector metadata
|
||||
connection_id: string;
|
||||
connector_name: string;
|
||||
connector_plugin: string;
|
||||
connector_version: string;
|
||||
image_url: string;
|
||||
// Config
|
||||
config: {
|
||||
cron: string;
|
||||
tables?: string;
|
||||
};
|
||||
properties: string;
|
||||
type: string;
|
||||
in_use: number;
|
||||
keywords: MultiLangArray;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class ElasticsearchService {
|
||||
private client: AxiosInstance;
|
||||
private logger: any;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
this.client = axios.create({
|
||||
baseURL: ELASTICSEARCH_CONFIG.getUrl(),
|
||||
headers: {
|
||||
Authorization: `ApiKey ${ELASTICSEARCH_CONFIG.getApiKey()}`,
|
||||
'Content-Type': 'application/json',
|
||||
},
|
||||
timeout: ELASTICSEARCH_CONFIG.timeout,
|
||||
});
|
||||
}
|
||||
|
||||
private getIndex(customerName: string): string {
|
||||
return `${customerName}_pipelines`;
|
||||
}
|
||||
|
||||
private formatMultiLang(value: string): MultiLang {
|
||||
return {
|
||||
'en-us': value,
|
||||
'pt-br': value,
|
||||
'es-es': value,
|
||||
};
|
||||
}
|
||||
|
||||
private formatMultiLangArray(value: string[] = []): MultiLangArray {
|
||||
return {
|
||||
'en-us': value,
|
||||
'pt-br': value,
|
||||
'es-es': value,
|
||||
};
|
||||
}
|
||||
|
||||
buildPipelineDocument(
|
||||
pipelineId: string,
|
||||
data: {
|
||||
name: string;
|
||||
description?: string;
|
||||
user_id: string;
|
||||
username: string;
|
||||
customer_id: string;
|
||||
plugin: string;
|
||||
connection_id: string;
|
||||
cron?: string;
|
||||
tables?: string;
|
||||
properties?: Record<string, any>;
|
||||
type?: string;
|
||||
status?: string;
|
||||
created_at?: string;
|
||||
keywords?: string[];
|
||||
},
|
||||
connector: ConnectorInfo | null,
|
||||
): PipelineDocument {
|
||||
const now = new Date().toISOString();
|
||||
|
||||
return {
|
||||
id: pipelineId,
|
||||
connection_id: data.connection_id,
|
||||
connector_name: connector?.name?.['en-us'] || '',
|
||||
connector_plugin: connector?.plugin || data.plugin,
|
||||
connector_version: connector?.version || '1.0.0',
|
||||
created_at: data.created_at || now,
|
||||
updated_at: now,
|
||||
customer_id: data.customer_id,
|
||||
description: this.formatMultiLang(data.description || ''),
|
||||
image_url: connector?.image || '',
|
||||
keywords: this.formatMultiLangArray(data.keywords),
|
||||
name: this.formatMultiLang(data.name || ''),
|
||||
status: data.status || 'CREATED',
|
||||
user_id: data.user_id,
|
||||
username: data.username,
|
||||
config: {
|
||||
cron: data.cron,
|
||||
tables: data.tables,
|
||||
},
|
||||
properties: data.properties ? JSON.stringify(data.properties) : '{}',
|
||||
type: data.type,
|
||||
in_use: 1,
|
||||
last_status_updated: now,
|
||||
tags: connector?.tags || [],
|
||||
};
|
||||
}
|
||||
|
||||
async getConnectorByPlugin(plugin: string): Promise<ConnectorInfo | null> {
|
||||
this.logger.info('Elasticsearch: Looking up connector', { plugin });
|
||||
|
||||
try {
|
||||
const response = await this.client.post('/connectors/_search', {
|
||||
query: {
|
||||
term: { plugin: plugin },
|
||||
},
|
||||
size: 1,
|
||||
});
|
||||
|
||||
const hits = response.data.hits?.hits || [];
|
||||
if (hits.length === 0) {
|
||||
this.logger.warn('Elasticsearch: Connector not found', { plugin });
|
||||
return null;
|
||||
}
|
||||
|
||||
const source = hits[0]._source;
|
||||
return {
|
||||
plugin: source.plugin,
|
||||
name: source.name,
|
||||
image: source.image,
|
||||
version: source.version,
|
||||
tags: source.tags || [],
|
||||
};
|
||||
} catch (error) {
|
||||
this.handleError('getConnectorByPlugin', error, { plugin });
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async createPipeline(
|
||||
customerName: string,
|
||||
pipelineId: string,
|
||||
data: {
|
||||
name: string;
|
||||
description?: string;
|
||||
user_id: string;
|
||||
username: string;
|
||||
customer_id: string;
|
||||
status?: string;
|
||||
created_at?: string;
|
||||
plugin: string;
|
||||
connection_id: string;
|
||||
cron?: string;
|
||||
tables?: string;
|
||||
properties?: Record<string, any>;
|
||||
type?: string;
|
||||
},
|
||||
connector: ConnectorInfo | null,
|
||||
): Promise<any> {
|
||||
const index = this.getIndex(customerName);
|
||||
const document = this.buildPipelineDocument(pipelineId, data, connector);
|
||||
|
||||
this.logger.info('Elasticsearch: Creating pipeline', {
|
||||
index,
|
||||
pipelineId,
|
||||
plugin: document.connector_plugin,
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await this.client.post(
|
||||
`${index}/_doc/${pipelineId}`,
|
||||
document,
|
||||
{ params: { refresh: 'wait_for' } },
|
||||
);
|
||||
|
||||
this.logger.info('Elasticsearch: Pipeline created successfully', {
|
||||
pipelineId,
|
||||
result: response.data.result,
|
||||
});
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.handleError('createPipeline', error, { pipelineId, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async updatePipeline(
|
||||
customerName: string,
|
||||
pipelineId: string,
|
||||
changes: {
|
||||
name?: string;
|
||||
description?: string;
|
||||
cron?: string;
|
||||
status?: string;
|
||||
tags?: string[];
|
||||
},
|
||||
): Promise<any> {
|
||||
const index = this.getIndex(customerName);
|
||||
const now = new Date().toISOString();
|
||||
|
||||
this.logger.info('Elasticsearch: Updating pipeline', {
|
||||
index,
|
||||
pipelineId,
|
||||
fields: Object.keys(changes),
|
||||
});
|
||||
|
||||
try {
|
||||
// Fetch current document
|
||||
const currentDoc = await this.client.get(`${index}/_doc/${pipelineId}`);
|
||||
const current = currentDoc.data._source;
|
||||
|
||||
// Build updated document, preserving existing values
|
||||
const updated: Record<string, any> = {
|
||||
...current,
|
||||
updated_at: now,
|
||||
};
|
||||
|
||||
if ('name' in changes) {
|
||||
updated.name = this.formatMultiLang(changes.name);
|
||||
}
|
||||
|
||||
if ('description' in changes) {
|
||||
updated.description = this.formatMultiLang(changes.description);
|
||||
}
|
||||
|
||||
if ('cron' in changes) {
|
||||
updated.config = {
|
||||
...current.config,
|
||||
cron: changes.cron,
|
||||
};
|
||||
}
|
||||
|
||||
if ('status' in changes) {
|
||||
updated.status = changes.status;
|
||||
updated.last_status_updated = now;
|
||||
}
|
||||
|
||||
if ('tags' in changes) {
|
||||
updated.tags = changes.tags;
|
||||
}
|
||||
|
||||
const response = await this.client.post(
|
||||
`${index}/_doc/${pipelineId}`,
|
||||
updated,
|
||||
{ params: { refresh: 'wait_for' } },
|
||||
);
|
||||
|
||||
this.logger.info('Elasticsearch: Pipeline updated successfully', {
|
||||
pipelineId,
|
||||
result: response.data.result,
|
||||
});
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.handleError('updatePipeline', error, { pipelineId, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async getPipeline(
|
||||
customerName: string,
|
||||
pipelineId: string,
|
||||
): Promise<PipelineDocument | null> {
|
||||
const index = this.getIndex(customerName);
|
||||
|
||||
this.logger.info('Elasticsearch: Getting pipeline', {
|
||||
index,
|
||||
pipelineId,
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await this.client.get(`${index}/_doc/${pipelineId}`);
|
||||
return response.data._source as PipelineDocument;
|
||||
} catch (error) {
|
||||
if (error instanceof AxiosError && error.response?.status === 404) {
|
||||
this.logger.warn('Elasticsearch: Pipeline not found', {
|
||||
pipelineId,
|
||||
index,
|
||||
});
|
||||
return null;
|
||||
}
|
||||
this.handleError('getPipeline', error, { pipelineId, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async deletePipeline(
|
||||
customerName: string,
|
||||
pipelineId: string,
|
||||
): Promise<any> {
|
||||
const index = this.getIndex(customerName);
|
||||
|
||||
this.logger.info('Elasticsearch: Deleting pipeline', {
|
||||
index,
|
||||
pipelineId,
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await this.client.delete(
|
||||
`${index}/_doc/${pipelineId}`,
|
||||
{ params: { refresh: 'wait_for' } },
|
||||
);
|
||||
|
||||
this.logger.info('Elasticsearch: Pipeline deleted successfully', {
|
||||
pipelineId,
|
||||
result: response.data.result,
|
||||
});
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
// If document not found, log warning but don't throw
|
||||
if (error instanceof AxiosError && error.response?.status === 404) {
|
||||
this.logger.warn('Elasticsearch: Pipeline not found for deletion', {
|
||||
pipelineId,
|
||||
index,
|
||||
});
|
||||
return { result: 'not_found' };
|
||||
}
|
||||
|
||||
this.handleError('deletePipeline', error, { pipelineId, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private getDataAssetIndex(customerName: string): string {
|
||||
return `${customerName}_data_assets_catalog`;
|
||||
}
|
||||
|
||||
async findDataAssetByTable(
|
||||
customerName: string,
|
||||
tableName: string,
|
||||
tableSchema: string,
|
||||
): Promise<{ id: string; nimbus_id: number | null; [key: string]: any } | null> {
|
||||
const index = this.getDataAssetIndex(customerName);
|
||||
|
||||
this.logger.info('Elasticsearch: Searching data asset', {
|
||||
index,
|
||||
tableName,
|
||||
tableSchema,
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await this.client.post(`/${index}/_search`, {
|
||||
query: {
|
||||
bool: {
|
||||
must: [
|
||||
{ term: { 'table_name.keyword': tableName.toUpperCase() } },
|
||||
{ term: { 'table_schema.keyword': tableSchema.toUpperCase() } },
|
||||
],
|
||||
},
|
||||
},
|
||||
size: 1,
|
||||
});
|
||||
|
||||
const hits = response.data.hits?.hits || [];
|
||||
if (hits.length === 0) {
|
||||
this.logger.warn('Elasticsearch: Data asset not found', { tableName, tableSchema, index });
|
||||
return null;
|
||||
}
|
||||
|
||||
return { ...hits[0]._source, _es_id: hits[0]._id };
|
||||
} catch (error) {
|
||||
this.handleError('findDataAssetByTable', error, { tableName, tableSchema, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async updateDataAsset(
|
||||
customerName: string,
|
||||
assetId: string,
|
||||
updates: Record<string, any>,
|
||||
): Promise<any> {
|
||||
const index = this.getDataAssetIndex(customerName);
|
||||
|
||||
this.logger.info('Elasticsearch: Updating data asset', {
|
||||
index,
|
||||
assetId,
|
||||
fields: Object.keys(updates),
|
||||
});
|
||||
|
||||
try {
|
||||
const response = await this.client.post(
|
||||
`/${index}/_update/${assetId}`,
|
||||
{ doc: updates },
|
||||
{ params: { refresh: 'wait_for' } },
|
||||
);
|
||||
|
||||
this.logger.info('Elasticsearch: Data asset updated', {
|
||||
assetId,
|
||||
result: response.data.result,
|
||||
});
|
||||
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.handleError('updateDataAsset', error, { assetId, index });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private handleError(
|
||||
operation: string,
|
||||
error: any,
|
||||
context: Record<string, any>,
|
||||
): void {
|
||||
if (error instanceof AxiosError) {
|
||||
this.logger.error(`Elasticsearch: ${operation} failed`, {
|
||||
...context,
|
||||
status: error.response?.status,
|
||||
statusText: error.response?.statusText,
|
||||
errorData: error.response?.data,
|
||||
message: error.message,
|
||||
});
|
||||
} else {
|
||||
this.logger.error(`Elasticsearch: ${operation} failed`, {
|
||||
...context,
|
||||
message: error.message,
|
||||
stack: error.stack,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,3 +0,0 @@
|
||||
export * from './elasticsearch.module';
|
||||
export * from './elasticsearch.service';
|
||||
export * from './elasticsearch.config';
|
||||
@@ -1,9 +0,0 @@
|
||||
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 {}
|
||||
@@ -1,56 +0,0 @@
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -14,7 +14,7 @@ import { redisStore } from 'cache-manager-ioredis-yet';
|
||||
keyPrefix: 'maestro:sso',
|
||||
}
|
||||
|
||||
if (process.env.REDIS_TLS === 'true') {
|
||||
if (process.env.ENV !== 'local') {
|
||||
baseRedisConfig['tls'] = {
|
||||
servername: process.env.REDIS_HOST,
|
||||
}
|
||||
|
||||
@@ -1,90 +0,0 @@
|
||||
import CronParser from 'cron-parser';
|
||||
|
||||
export enum ScheduleLimits {
|
||||
MINUTE = 'minute',
|
||||
HOUR = 'hour',
|
||||
DAY = 'day',
|
||||
UNLIMITED = 'unlimited',
|
||||
}
|
||||
|
||||
const SECONDS_IN_MINUTE = 60;
|
||||
const SECONDS_IN_HOUR = 3600;
|
||||
const SECONDS_IN_DAY = 86400;
|
||||
|
||||
/**
|
||||
* Airflow preset schedules mapped to cron expressions.
|
||||
* @once is special - it means run only once (no recurring schedule).
|
||||
*/
|
||||
const AIRFLOW_PRESETS: Record<string, string | null> = {
|
||||
'@once': null, // No recurring schedule - always valid
|
||||
'@hourly': '0 * * * *', // Every hour
|
||||
'@daily': '0 0 * * *', // Every day at midnight
|
||||
'@weekly': '0 0 * * 0', // Every week on Sunday
|
||||
'@monthly': '0 0 1 * *', // First day of every month
|
||||
'@yearly': '0 0 1 1 *', // First day of every year
|
||||
'@annually': '0 0 1 1 *', // Same as @yearly
|
||||
};
|
||||
|
||||
/**
|
||||
* Convert Airflow preset to cron expression.
|
||||
* Returns null for @once (no recurring schedule).
|
||||
* Returns original string if not an Airflow preset.
|
||||
*/
|
||||
export function convertAirflowPresetToCron(schedule: string): string | null {
|
||||
const preset = AIRFLOW_PRESETS[schedule.toLowerCase()];
|
||||
if (preset !== undefined) {
|
||||
return preset;
|
||||
}
|
||||
return schedule;
|
||||
}
|
||||
|
||||
export function getMinimumIntervalSeconds(scheduleLimit: string): number {
|
||||
switch (scheduleLimit) {
|
||||
case ScheduleLimits.MINUTE:
|
||||
return SECONDS_IN_MINUTE;
|
||||
case ScheduleLimits.HOUR:
|
||||
return SECONDS_IN_HOUR;
|
||||
case ScheduleLimits.DAY:
|
||||
return SECONDS_IN_DAY;
|
||||
case ScheduleLimits.UNLIMITED:
|
||||
default:
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
export function getCronIntervalSeconds(cron: string): number {
|
||||
const interval = CronParser.parseExpression(cron);
|
||||
const nextDate = interval.next().toDate();
|
||||
const afterNextDate = interval.next().toDate();
|
||||
return Math.floor((afterNextDate.getTime() - nextDate.getTime()) / 1000);
|
||||
}
|
||||
|
||||
export function validateCronAgainstScheduleLimit(
|
||||
cron: string,
|
||||
scheduleLimit: string,
|
||||
): { valid: boolean; message?: string } {
|
||||
if (!cron) return { valid: true };
|
||||
|
||||
// Convert Airflow presets to cron expressions
|
||||
const cronExpression = convertAirflowPresetToCron(cron);
|
||||
|
||||
// @once returns null - no recurring schedule, always valid
|
||||
if (cronExpression === null) {
|
||||
return { valid: true };
|
||||
}
|
||||
|
||||
try {
|
||||
const cronInterval = getCronIntervalSeconds(cronExpression);
|
||||
const minInterval = getMinimumIntervalSeconds(scheduleLimit);
|
||||
|
||||
if (cronInterval < minInterval) {
|
||||
return {
|
||||
valid: false,
|
||||
message: `Schedule interval (${cronInterval}s) is below customer limit (${scheduleLimit}: ${minInterval}s minimum)`,
|
||||
};
|
||||
}
|
||||
return { valid: true };
|
||||
} catch (error) {
|
||||
return { valid: false, message: `Invalid cron expression: ${error.message}` };
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user